From d9fd752a1afe3d5c693476941534afe857151d5f Mon Sep 17 00:00:00 2001 From: jaturapornchai <39358840+jaturapornchai@users.noreply.github.com> Date: Fri, 10 Apr 2026 04:40:22 +0700 Subject: [PATCH] U4: Worker refresh 15min + warmup worker (ping every 2min) Shorten the scan/health/exam cycle from 1h to 15min so dead/new models are discovered faster, and add a lightweight warmup pinger that hits passing models every 2min to keep upstream connections hot and surface failures between full cycles. The warmup uses its own Redis lock (worker:warmup-leader) so it can run alongside the main worker on any replica without fighting over the existing worker:leader key. Co-Authored-By: Claude Opus 4.6 (1M context) --- src/lib/worker/index.ts | 12 ++- src/lib/worker/warmup.ts | 178 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 186 insertions(+), 4 deletions(-) create mode 100644 src/lib/worker/warmup.ts diff --git a/src/lib/worker/index.ts b/src/lib/worker/index.ts index ba55226..e4e247c 100644 --- a/src/lib/worker/index.ts +++ b/src/lib/worker/index.ts @@ -3,6 +3,7 @@ import { scanModels } from "./scanner"; import { checkHealth } from "./health"; import { runExams } from "./exam"; import { acquireLeader, renewLeader, releaseLeader } from "./leader"; +import { startWarmup } from "./warmup"; export { scanModels } from "./scanner"; export { checkHealth } from "./health"; @@ -90,7 +91,7 @@ export async function runWorkerCycle(): Promise { await setState("status", "running"); await setState("last_run", new Date().toISOString()); - const next = new Date(Date.now() + 60 * 60 * 1000).toISOString(); + const next = new Date(Date.now() + 15 * 60 * 1000).toISOString(); await setState("next_run", next); await logWorker("worker", "Worker cycle started"); @@ -151,7 +152,7 @@ export async function runWorkerCycle(): Promise { export function startWorker(): void { if (workerTimer) return; // already started - logWorker("worker", "Worker starting — running immediately then every 1h"); + logWorker("worker", "Worker starting — running immediately then every 15min"); // Run once immediately (async, don't block) runWorkerCycle().catch((err) => { @@ -160,14 +161,17 @@ export function startWorker(): void { setState("status", "error"); }); - // Then every 1 hour + // Then every 15 minutes workerTimer = setInterval(() => { runWorkerCycle().catch((err) => { logWorker("worker", `Scheduled cycle error: ${err}`, "error"); isRunning = false; setState("status", "error"); }); - }, 60 * 60 * 1000); + }, 15 * 60 * 1000); + + // Start background warmup pinger (independent schedule) + startWarmup(); } export async function getWorkerStatus(): Promise { diff --git a/src/lib/worker/warmup.ts b/src/lib/worker/warmup.ts new file mode 100644 index 0000000..1912c7b --- /dev/null +++ b/src/lib/worker/warmup.ts @@ -0,0 +1,178 @@ +import { getSqlClient } from "@/lib/db/schema"; +import { getRedis } from "@/lib/redis"; +import { getNextApiKey } from "@/lib/api-keys"; +import { PROVIDER_URLS } from "@/lib/providers"; +import { recordOutcome } from "@/lib/live-score"; + +// ─── Warmup worker ─── +// Pings passing models every 2min to keep upstream connections warm and +// detect dead providers fast. Uses its own Redis lock so it can run +// concurrently with the main worker cycle (separate schedule). + +const WARMUP_INTERVAL_MS = 2 * 60 * 1000; // 2 minutes +const WARMUP_TIMEOUT_MS = 5_000; +const MAX_MODELS_PER_WARMUP = 30; +const WARMUP_LEADER_KEY = "worker:warmup-leader"; +const WARMUP_LEADER_TTL_SEC = 90; + +let warmupTimer: ReturnType | null = null; +let isWarming = false; + +function workerId(): string { + return process.env.HOSTNAME || process.env.COMPUTERNAME || "local"; +} + +async function acquireWarmupLeader(): Promise { + try { + const redis = getRedis(); + const me = workerId(); + const result = await redis.set(WARMUP_LEADER_KEY, me, "EX", WARMUP_LEADER_TTL_SEC, "NX"); + if (result === "OK") return true; + const holder = await redis.get(WARMUP_LEADER_KEY); + return holder === me; + } catch { + // Redis down → single-replica fallback + return true; + } +} + +async function renewWarmupLeader(): Promise { + try { + const redis = getRedis(); + await redis.expire(WARMUP_LEADER_KEY, WARMUP_LEADER_TTL_SEC); + } catch { + // silent + } +} + +async function releaseWarmupLeader(): Promise { + try { + const redis = getRedis(); + const me = workerId(); + const holder = await redis.get(WARMUP_LEADER_KEY); + if (holder === me) { + await redis.del(WARMUP_LEADER_KEY); + } + } catch { + // silent + } +} + +async function logWorker(step: string, message: string, level = "info"): Promise { + try { + const sql = getSqlClient(); + await sql`INSERT INTO worker_logs (step, message, level) VALUES (${step}, ${message}, ${level})`; + } catch { + // silent + } +} + +interface WarmupModel { + id: string; + provider: string; + model_id: string; +} + +async function pingModel(m: WarmupModel): Promise<{ success: boolean; latency: number }> { + const url = PROVIDER_URLS[m.provider]; + if (!url) return { success: false, latency: 0 }; + + const apiKey = getNextApiKey(m.provider); + if (!apiKey && m.provider !== "ollama" && m.provider !== "pollinations") { + return { success: false, latency: 0 }; + } + + const start = Date.now(); + try { + const res = await fetch(url, { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: `Bearer ${apiKey || "dummy"}`, + }, + body: JSON.stringify({ + model: m.model_id, + messages: [{ role: "user", content: "." }], + max_tokens: 1, + temperature: 0, + }), + signal: AbortSignal.timeout(WARMUP_TIMEOUT_MS), + }); + return { success: res.ok, latency: Date.now() - start }; + } catch { + return { success: false, latency: Date.now() - start }; + } +} + +export async function runWarmupCycle(): Promise { + if (isWarming) return; + + const isLeader = await acquireWarmupLeader(); + if (!isLeader) return; + + isWarming = true; + try { + const sql = getSqlClient(); + // Passing models that aren't in cooldown + const models = await sql` + SELECT m.id, m.provider, m.model_id + FROM models m + INNER JOIN ( + SELECT DISTINCT ON (model_id) model_id, passed + FROM exam_attempts WHERE finished_at IS NOT NULL + ORDER BY model_id, started_at DESC + ) ea ON m.id = ea.model_id AND ea.passed = true + LEFT JOIN ( + SELECT DISTINCT ON (model_id) model_id, cooldown_until + FROM health_logs ORDER BY model_id, id DESC + ) h ON h.model_id = m.id + WHERE h.cooldown_until IS NULL OR h.cooldown_until < now() + LIMIT ${MAX_MODELS_PER_WARMUP} + `; + + if (models.length === 0) { + await logWorker("warmup", "No models to warm up"); + return; + } + + let success = 0; + let failed = 0; + const CONCURRENCY = 5; + let idx = 0; + + async function worker() { + while (idx < models.length) { + const m = models[idx++]; + const { success: ok, latency } = await pingModel(m); + if (ok) { + success++; + recordOutcome(m.provider, m.model_id, true, latency); + } else { + failed++; + recordOutcome(m.provider, m.model_id, false, latency); + } + } + } + + await Promise.all(Array.from({ length: CONCURRENCY }, worker)); + await renewWarmupLeader(); + await logWorker("warmup", `🔥 Pinged ${models.length} models — ${success} ok, ${failed} failed`); + } catch (err) { + await logWorker("warmup", `Warmup cycle error: ${err}`, "error"); + } finally { + await releaseWarmupLeader(); + isWarming = false; + } +} + +export function startWarmup(): void { + if (warmupTimer) return; + logWorker("warmup", "Warmup worker starting — ping every 2min"); + // Delay first run so the main gateway has time to stabilize + setTimeout(() => { + runWarmupCycle().catch((err) => logWorker("warmup", `Initial warmup error: ${err}`, "error")); + }, 30_000); + warmupTimer = setInterval(() => { + runWarmupCycle().catch((err) => logWorker("warmup", `Scheduled warmup error: ${err}`, "error")); + }, WARMUP_INTERVAL_MS); +}