// How many slots a REMOTE worker is worth, probed from the remote itself. // // A remote runs its own worker pool, and its /api/worker/health has always // returned that pool's per-worker summary — which every caller promptly threw // away, reading only `ok`. So the number of transcriptions a remote can // actually take in parallel was sitting on the wire while the primary treated // every remote as one slot. This module keeps that answer: the count of the // remote's enabled, non-degraded workers. // // TTL-CACHED. reconfigure() is synchronous and runs on // every batch start and settings save; a network probe cannot live inside it. // So the cache is read synchronously (cachedRemoteSlots), staleness is asked // synchronously (remoteCapacityStale), and the probe itself is fired // fire-and-forget by the pool when the answer is stale — no timer, no // background loop. The staleness cost is bounded and cheap: a remote that // gained a worker in the last minute runs one slot short until the next probe. // // An explicit `remote.slots` in the worker config bypasses all of this — the // operator's number is not second-guessed by a probe. import type { Worker } from "../lib/workers"; const REMOTE_CAPACITY_TTL_MS = 60_000; type Cached = { at: number; // null = the probe failed (unreachable, bad body). Cached too, so a dead // remote is not re-probed on every reconfigure inside one TTL window. slots: number | null; }; const cache = new Map(); const inFlight = new Map>(); // Exported for tests, which need a clean slate between fixtures. export function clearRemoteCapacityCache(): void { cache.clear(); inFlight.clear(); } // The cached slot count, or null when unknown (never probed, or the last probe // failed). Deliberately returns a STALE value rather than null — bounded // staleness beats a pool that collapses to one slot for the duration of a // probe. export function cachedRemoteSlots(workerId: string): number | null { return cache.get(workerId)?.slots ?? null; } export function remoteCapacityStale(workerId: string): boolean { const hit = cache.get(workerId); return !hit || Date.now() - hit.at >= REMOTE_CAPACITY_TTL_MS; } type HealthBody = { ok?: boolean; workers?: Array<{ state?: string; degraded?: boolean }>; }; // Probe the remote's /api/worker/health and cache how many transcriptions it // can take at once: its enabled, non-degraded workers. In-flight probes are // shared, failures are cached (as null) for the TTL, and the count is floored // at 1 — a configured remote must never silently vanish from the pool because // its operator momentarily disabled everything; the health ping in // transcribeOne is what handles a remote that is genuinely down. export function probeRemoteCapacity( worker: Worker, timeoutMs = 3000, ): Promise { const running = inFlight.get(worker.id); if (running) return running; const probe = (async (): Promise => { const baseUrl = (worker.remote?.baseUrl ?? "").replace(/\/+$/, ""); let slots: number | null = null; if (baseUrl) { try { const token = worker.remote?.token?.trim(); const res = await fetch(`${baseUrl}/api/worker/health`, { headers: token ? { Authorization: `Bearer ${token}` } : {}, signal: AbortSignal.timeout(timeoutMs), }); if (res.ok) { const body = (await res.json()) as HealthBody; if (body.ok === true && Array.isArray(body.workers)) { slots = Math.max( 1, body.workers.filter( (w) => w.state === "enabled" && w.degraded !== true, ).length, ); } } } catch { // Unreachable/timeout — cached as unknown below. } } cache.set(worker.id, { at: Date.now(), slots }); return slots; })(); inFlight.set(worker.id, probe); void probe.finally(() => inFlight.delete(worker.id)); return probe; }