// The LLM fan-out's claim on an endpoint worker: lease a slot, make sure the // endpoint actually serves the primary's model tag, hand back the baseUrl. // // WHY VERIFICATION IS NOT OPTIONAL. Freshness identity pins the CONFIGURED // model string (and, for digest, numCtx via the prompt variant); only baseUrl // may vary per call. An endpoint that silently served a different model would // not fail — it would write records whose provenance looks current and whose // content came from the wrong model. The cheap guard is ollama's own /api/tags: // before an endpoint's first use (and every VERIFY_TTL_MS thereafter) the // dispatcher confirms the tag is present, and an endpoint that lacks it — or // cannot be reached — is DEGRADED through the pool's existing machinery, not // retried per item. // // This module is deliberately thin: no queue, no state beyond the verification // cache. Claiming is the pool's tryAcquire (non-parking — a parked acquire // inside a runPool slot would deadlock the batch), and capacity for a // dispatcher's limit() is the pool's freeSlots. import { getWorkerPool, type Lease, type WorkerPool } from "../jobs/workerPool"; export type LlmSlot = { lease: Lease; baseUrl: string; workerId: string; }; const VERIFY_TTL_MS = 10 * 60_000; // worker.id -> when it was last confirmed to serve which model. const verified = new Map(); // Exported for tests, which need a clean slate between fixtures. export function clearLlmVerificationCache(): void { verified.clear(); } // The remote term of a dispatcher's limit(): free llm slots serving any of the // given operations. Pure counting — no lease is taken. export function freeLlmSlots( ops: readonly string[], pool: WorkerPool = getWorkerPool(), ): number { return pool.freeSlots(ops, { kind: "llm" }); } async function servesModel( baseUrl: string, model: string, timeoutMs: number, ): Promise { try { const res = await fetch(`${baseUrl}/api/tags`, { signal: AbortSignal.timeout(timeoutMs), }); if (!res.ok) return false; const body = (await res.json()) as { models?: Array<{ name?: unknown }> }; const names = (body.models ?? []) .map((m) => m.name) .filter((n): n is string => typeof n === "string"); // ollama names a bare tag ":latest"; a config of "qwen2.5" matches // it. Anything more specific must match exactly. return names.some( (n) => n === model || (!model.includes(":") && n === `${model}:latest`), ); } catch { return false; } } // Claim a verified endpoint slot for one operation, or null when none is free. // An endpoint that fails verification (missing tag, unreachable) is degraded // and the next candidate tried, so one bad endpoint costs one probe rather // than poisoning the run. export async function acquireLlmSlot( op: string, model: string, log: (msg: string) => void, pool: WorkerPool = getWorkerPool(), timeoutMs = 3000, ): Promise { for (;;) { const lease = pool.tryAcquire([op], { kind: "llm" }); if (!lease) return null; const workerId = lease.worker.id; const baseUrl = (lease.worker.llm?.baseUrl ?? "").replace(/\/+$/, ""); if (!baseUrl) { // Sanitizers should make this unreachable; degrade rather than loop on it. pool.markDegraded(workerId); lease.release(); continue; } const hit = verified.get(workerId); const fresh = hit && hit.model === model && Date.now() - hit.at < VERIFY_TTL_MS; if (fresh || (await servesModel(baseUrl, model, timeoutMs))) { if (!fresh) verified.set(workerId, { at: Date.now(), model }); return { lease, baseUrl, workerId }; } verified.delete(workerId); log( `LLM endpoint ${workerId} (${baseUrl}) does not serve "${model}" (or is unreachable) — ` + `degraded for this run. Pull the exact model tag there, or re-enable it from the Workers page once fixed.`, ); pool.markDegraded(workerId); lease.release(); } }