File size: 9,130 Bytes
35743bd
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
/**
 * Local Provider Health Check
 *
 * Background polling of local provider_nodes (localhost) to detect
 * when they are up or down. Uses GET /models with a 5s timeout.
 *
 * Health status is stored in-memory (no DB migration needed).
 * Backoff schedule: 30s → 60s → 120s → 300s max on consecutive failures.
 * Resets to 30s on first success after failure.
 *
 * Uses Promise.allSettled so one slow/down node doesn't block others.
 */

import { getProviderNodes } from "@/lib/localDb";

// ── Types ────────────────────────────────────────────────────────────────

export interface HealthStatus {
  nodeId: string;
  prefix: string;
  isHealthy: boolean;
  lastCheck: Date;
  lastError?: string;
  consecutiveFailures: number;
  responseTimeMs?: number;
}

// ── Config ───────────────────────────────────────────────────────────────

const BACKOFF_SCHEDULE = [30_000, 60_000, 120_000, 300_000];
const CHECK_TIMEOUT_MS = 5_000;
const INITIAL_DELAY_MS = 15_000; // Wait for server boot before first sweep
const LOG_PREFIX = "[LocalHealthCheck]";
const TRUE_ENV_VALUES = new Set(["1", "true", "yes", "on"]);

function isBuildProcess(): boolean {
  return typeof process !== "undefined" && process.env.NEXT_PHASE === "phase-production-build";
}

function isAutomatedTestProcess(): boolean {
  return (
    typeof process !== "undefined" &&
    (process.env.NODE_ENV === "test" ||
      process.env.VITEST !== undefined ||
      process.argv.some((arg) => arg.includes("test")))
  );
}

// ── State (globalThis survives HMR re-evaluation) ───────────────────────

declare global {
  var __omnirouteLocalHC:
    | {
        initialized: boolean;
        sweepTimer: ReturnType<typeof setTimeout> | null;
        healthCache: Map<string, HealthStatus>;
        sweepInProgress: boolean;
      }
    | undefined;
}

function getLHCState() {
  if (!globalThis.__omnirouteLocalHC) {
    globalThis.__omnirouteLocalHC = {
      initialized: false,
      sweepTimer: null,
      healthCache: new Map(),
      sweepInProgress: false,
    };
  }
  return globalThis.__omnirouteLocalHC;
}

const healthCache = getLHCState().healthCache;

// ── Helpers ──────────────────────────────────────────────────────────────

function isEnvFlagEnabled(name: string): boolean {
  const value = process.env[name];
  if (!value) return false;
  return TRUE_ENV_VALUES.has(value.trim().toLowerCase());
}

function isLocalHealthCheckDisabled(): boolean {
  return (
    isEnvFlagEnabled("OMNIROUTE_DISABLE_LOCAL_HEALTHCHECK") ||
    isBuildProcess() ||
    isAutomatedTestProcess()
  );
}

function isLocalhostUrl(baseUrl: string): boolean {
  try {
    const u = new URL(baseUrl);
    // Block credentials in URL to prevent SSRF via user@host (e.g., http://localhost@evil.com)
    if (u.username || u.password) return false;
    // Note: URL.hostname returns "[::1]" WITH brackets for IPv6 — both forms checked.
    // Verified: node -e "new URL('http://[::1]:8080').hostname" → "[::1]"
    // Strictly matching 172.16.0.0/12 (Docker/local) and explicitly blocking ::1 per SSRF hardening
    return (
      u.hostname === "localhost" ||
      u.hostname === "127.0.0.1" ||
      /^172\.(1[6-9]|2[0-9]|3[0-1])\.\d{1,3}\.\d{1,3}$/.test(u.hostname)
    );
  } catch {
    return false;
  }
}

function getNextInterval(failures: number): number {
  return BACKOFF_SCHEDULE[Math.min(failures, BACKOFF_SCHEDULE.length - 1)];
}

// ── Core ─────────────────────────────────────────────────────────────────

async function checkNode(node: {
  id: string;
  prefix: string;
  baseUrl: string;
}): Promise<HealthStatus> {
  const url = `${node.baseUrl.replace(/\/+$/, "")}/models`;
  const start = Date.now();
  const prev = healthCache.get(node.id);

  try {
    const res = await fetch(url, { signal: AbortSignal.timeout(CHECK_TIMEOUT_MS) });
    // Consume/cancel response body to free resources
    res.body?.cancel().catch(() => {});
    const isHealthy = res.ok || res.status === 401; // 401 = server up but auth required
    return {
      nodeId: node.id,
      prefix: node.prefix,
      isHealthy,
      lastCheck: new Date(),
      consecutiveFailures: isHealthy ? 0 : (prev?.consecutiveFailures ?? 0) + 1,
      responseTimeMs: Date.now() - start,
      lastError: isHealthy ? undefined : `HTTP ${res.status}`,
    };
  } catch (err: unknown) {
    const message = err instanceof Error ? err.message : "Connection failed";
    return {
      nodeId: node.id,
      prefix: node.prefix,
      isHealthy: false,
      lastCheck: new Date(),
      consecutiveFailures: (prev?.consecutiveFailures ?? 0) + 1,
      responseTimeMs: Date.now() - start,
      lastError: message,
    };
  }
}

/** Single sweep: check all local provider_nodes in parallel. */
export async function sweep(): Promise<void> {
  const state = getLHCState();
  if (state.sweepInProgress) return;
  state.sweepInProgress = true;

  try {
    let nodes: Array<{ id: string; prefix: string; baseUrl: string }>;
    try {
      const raw = await getProviderNodes();
      nodes = (Array.isArray(raw) ? raw : []).filter(
        (n: Record<string, unknown>) =>
          typeof n.baseUrl === "string" && isLocalhostUrl(n.baseUrl as string)
      ) as Array<{ id: string; prefix: string; baseUrl: string }>;
    } catch (err) {
      console.error(LOG_PREFIX, "Failed to load provider_nodes:", err);
      return;
    }

    // Prune stale entries for deleted nodes
    const currentNodeIds = new Set(nodes.map((n) => n.id));
    for (const key of healthCache.keys()) {
      if (!currentNodeIds.has(key)) healthCache.delete(key);
    }

    if (nodes.length === 0) return;

    const results = await Promise.allSettled(nodes.map((node) => checkNode(node)));

    for (const result of results) {
      if (result.status === "fulfilled") {
        const status = result.value;
        const prev = healthCache.get(status.nodeId);

        // Log state transitions
        if (prev && prev.isHealthy !== status.isHealthy) {
          const emoji = status.isHealthy ? "✅" : "❌";
          console.log(
            LOG_PREFIX,
            `${emoji} ${status.prefix} is now ${status.isHealthy ? "healthy" : "unhealthy"}${status.lastError ? ` (${status.lastError})` : ""} [${status.responseTimeMs}ms]`
          );
        }

        healthCache.set(status.nodeId, status);
      }
    }
  } finally {
    state.sweepInProgress = false;
    scheduleSweep();
  }
}

function scheduleSweep(): void {
  const state = getLHCState();
  if (!state.initialized) return;
  if (state.sweepTimer) clearTimeout(state.sweepTimer);

  // Use the maximum consecutive failures across all nodes to determine interval
  let maxFailures = 0;
  for (const status of healthCache.values()) {
    if (status.consecutiveFailures > maxFailures) {
      maxFailures = status.consecutiveFailures;
    }
  }

  const interval = getNextInterval(maxFailures);
  state.sweepTimer = setTimeout(sweep, interval);
}

// ── Public API ───────────────────────────────────────────────────────────

/** Get health status for a specific provider_node. */
export function getHealthStatus(nodeId: string): HealthStatus | undefined {
  return healthCache.get(nodeId);
}

/** Check if a provider_node is healthy. Returns true if never checked (optimistic). */
export function isNodeHealthy(nodeId: string): boolean {
  const status = healthCache.get(nodeId);
  return status?.isHealthy ?? true;
}

/** Get all health statuses (for monitoring API). */
export function getAllHealthStatuses(): Record<string, HealthStatus> {
  return Object.fromEntries(healthCache);
}

/** Start the health check scheduler (idempotent). */
export function initLocalHealthCheck(): void {
  const state = getLHCState();
  if (state.initialized || isLocalHealthCheckDisabled()) return;
  state.initialized = true;

  console.log(
    LOG_PREFIX,
    `Starting local provider health check (initial delay ${INITIAL_DELAY_MS / 1000}s)`
  );

  state.sweepTimer = setTimeout(() => {
    sweep().catch((err) => console.error(LOG_PREFIX, "Initial sweep failed:", err));
  }, INITIAL_DELAY_MS);
}

/** Stop the scheduler (for tests / hot-reload). */
export function stopLocalHealthCheck(): void {
  const state = getLHCState();
  if (state.sweepTimer) {
    clearTimeout(state.sweepTimer);
    state.sweepTimer = null;
  }
  state.initialized = false;
}

// Auto-initialize on first import (same pattern as tokenHealthCheck.ts:272)
initLocalHealthCheck();