// Global transcription worker pool — the scheduler that distributes individual // transcriptions across configured workers. // // Like the JobRegistry (registry.ts), this is a process-wide singleton held on // globalThis so every request/server-action shares ONE view of which worker has // a free slot. That global-ness is the whole point: "the GPU worker has exactly // 1 slot across everything" can only be enforced in one shared place. A batch // (runWhisperBatch) no longer owns a pLimit; it acquires a lease per video from // this pool and releases it when the transcription finishes. // // Scheduling is a priority-ordered free-list (decided with the user): every // enabled worker runs concurrently, and when several have a free slot the // highest-priority one is handed out first. When NO worker is eligible (all // disabled/draining/full), acquire() parks on the waiter queue — so a batch // pauses, neither failing nor busy-looping, and resumes the instant a worker is // re-enabled. // // Waiters are ordered foreground-before-background, not plain FIFO: a manual // (foreground) acquire is served before any parked background one, while FIFO is // preserved WITHIN each class. This mirrors the registry queue's foreground/ // background insertion (see registry.ts enqueue) so the same "manual preempts // auto" priority axis governs BOTH schedulers. The auto-runner marks its // per-video transcription acquires `background: true`, so a manual transcribe // jumps ahead of queued auto work — the running auto unit isn't interrupted, it // drains, and the freed slot goes to the manual waiter. import { workerMatches, type Worker, type WorkerKind } from "../lib/workers"; import { cachedRemoteSlots, probeRemoteCapacity, remoteCapacityStale, } from "../controller/remoteCapacity"; import { getSettings } from "../lib/settings"; import { getPaths } from "../lib/paths"; import { readWorkerDefaults } from "./workerDefaults"; import { compareTier } from "./scheduler"; import type { SchedulerTier } from "./jobKinds"; // Consecutive transport/exec failures before a worker is auto-marked degraded // (skipped by the scheduler until re-enabled). See markFailure / Phase 5. export const WORKER_DEGRADE_THRESHOLD = 3; export type WorkerRuntimeState = "enabled" | "draining" | "disabled"; // A granted slot. release() is idempotent; the controller calls it in a finally. export type Lease = { worker: Worker; release: () => void; }; // Client-safe snapshot of one worker for the Workers page / active-jobs payload. export type WorkerSummary = { id: string; name: string; kind: Worker["kind"]; appId?: string; priority: number; // A worker is one slot: busy is true while its single transcription runs. busy: boolean; state: WorkerRuntimeState; degraded: boolean; enabled: boolean; // persisted intent (config.enabled) // Capability tags, verbatim from the config. What the sweep console reads to // say where an operation can run — through the same workerMatches rule the // scheduler grants by, so the display and the dispatch cannot disagree. tags?: string[]; // The engine's compute device, verbatim from the per-worker config ("cpu", // "cuda:0", …). UNDEFINED MEANS UNKNOWN, not CPU: it is the engine binary's own // default, which on parakeet-cli may well be the GPU. // // Exposed because digestYield() has to answer "is anything competing for GPU // shaders right now?", and `kind: "local"` does not answer it — this box runs a // GPU worker beside two explicitly CPU-pinned ones, and treating those as GPU // contention stopped the digest lane dead for work that used no shaders at all. device?: string; }; type PoolEntry = { config: Worker; // One slot per worker: busy while its single transcription runs. busy: boolean; state: WorkerRuntimeState; degraded: boolean; consecutiveFailures: number; // Set when a worker is removed from settings while still busy: drain then drop. retireWhenIdle: boolean; // Registered by the in-flight transcription (if its engine supports it) to // request a graceful "stop & keep partial". Cleared when the lease releases. activeStop?: () => void; }; type Waiter = { resolve: (lease: Lease) => void; reject: (err: Error) => void; onAbort?: () => void; // Priority tier (see enqueueWaiter). Background waiters (auto transcriptions) // queue behind foreground (manual) ones for free slots; an urgent waiter would // jump ahead of both. Same compareTier ordering the registry's scheduler uses. tier: SchedulerTier; // Capability requirement (see workerMatches). A freed slot is granted to the // FIRST waiter it matches, not blindly to the head — pump() scans past a // waiter whose requirement this slot cannot satisfy. requires?: readonly string[]; // A further restriction on WHICH worker may serve this waiter (acquire's // opts.only). Matched alongside `requires`, never instead of it. only?: WorkerFilter; }; // Narrows an acquire to particular workers — by id, by kind — on top of the // capability rule. The one-off file transcription uses it to stay on LOCAL // workers, or on the one worker the caller named. export type WorkerFilter = (worker: Worker) => boolean; export class WorkerPool { // Insertion order is the tiebreak for equal priority, so use a Map (ordered). private entries = new Map(); private waiters: Waiter[] = []; private initialized = false; // When set, the pool is in a temporary "pause all": every worker is forced // disabled and this records each worker's pre-pause state so resumeAll can // restore it. Runtime-only (not persisted) — a restart returns workers to // their configured enabled state. See pauseAll / resumeAll. private pausedSnapshot: Map | null = null; // Lazily seed from settings on first use. Idempotent. private ensureInit(): void { if (this.initialized) return; this.initialized = true; // Initial seed: apply the persisted enabled flags as the starting state. this.reconfigure(getSettings().workers, { applyEnabled: true }); // Then overlay a saved "default" arrangement, if one exists, so the operator // gets their preferred set of enabled workers back at launch. this.applyDefaults(); } // Apply the persisted "default" worker arrangement (Workers-page "Set as // default") over the settings-seeded state: each listed worker starts enabled, // every other worker starts disabled. A no-op when no default has been saved // (readWorkerDefaults returns null) — the settings.enabled seed stands. This // governs the launch-time seed only; a later settings save still re-applies // settings.enabled via reconfigure(applyEnabled: true). private applyDefaults(): void { const defaults = readWorkerDefaults(getPaths()); if (!defaults) return; const enabled = new Set(defaults.enabledWorkerIds); for (const entry of this.entries.values()) { if (enabled.has(entry.config.id)) { entry.state = "enabled"; entry.degraded = false; entry.consecutiveFailures = 0; } else { entry.state = "disabled"; } } this.pump(); } // Re-sync the pool with a worker list (defaults to current settings). Updates // each worker's config in place, adds new workers, and retires removed ones — // but NEVER drops a worker that still has in-flight leases (that would // oversubscribe hardware the moment it's re-added). A busy removed worker is // marked retireWhenIdle and dropped when its last lease releases. // // `applyEnabled` controls whether an EXISTING worker's runtime state is reset // from its persisted `enabled` flag. The settings-save path passes true (a new // enabled intent should take effect). The batch-start path passes false (the // default) so an operator's runtime enable/disable/drain from the Workers page // is preserved across the reconfigure a batch does on start — otherwise every // batch would silently re-enable a worker the operator just turned off. reconfigure(workers?: Worker[], opts?: { applyEnabled?: boolean }): void { this.initialized = true; const applyEnabled = opts?.applyEnabled === true; const next = workers ?? getSettings().workers; const seen = new Set(); for (const w of next) { // A REMOTE worker is N slots, not one — expanded HERE rather than by // turning `busy` into a counter, which would touch every grant/drain/ // summary path for one worker kind. Slot 1 keeps the BASE id (so saved // defaults, task labels and degrade marks against the configured id keep // meaning what they meant); extra slots are `${id}#2..#N`, each an // independently enable/drain/degrade-able entry sharing one remote // config. Shrinking the count retires the surplus through the exact // removed-worker path below (drain if busy, drop when idle). const slotCount = w.kind === "remote" ? this.remoteSlotCount(w) : w.kind === "llm" ? // One concurrent generation per endpoint unless the operator says // otherwise — no probe: an ollama box does not report a slot count. Math.max(1, Math.floor(w.llm?.slots ?? 1)) : 1; for (let slot = 1; slot <= slotCount; slot++) { const cfg: Worker = slot === 1 ? w : { ...w, id: `${w.id}#${slot}`, name: `${w.name} #${slot}` }; seen.add(cfg.id); const existing = this.entries.get(cfg.id); if (existing) { existing.config = cfg; existing.retireWhenIdle = false; if (applyEnabled) { // Reconcile runtime state with the (just-saved) persisted intent. // Disabling a busy worker drains it; (re-)enabling clears degraded. if (w.enabled) { existing.state = "enabled"; existing.degraded = false; existing.consecutiveFailures = 0; } else if (existing.busy) { existing.state = "draining"; } else { existing.state = "disabled"; } } // else: leave runtime state untouched (preserve operator overrides). } else { // A slot added by expansion inherits the base slot's RUNTIME state: // a probe answering mid-run must not re-enable a remote the operator // just disabled from the Workers page. const base = slot > 1 ? this.entries.get(w.id) : undefined; this.entries.set(cfg.id, { config: cfg, busy: false, state: base ? base.state === "draining" ? "enabled" : base.state : w.enabled ? "enabled" : "disabled", degraded: base?.degraded ?? false, consecutiveFailures: 0, retireWhenIdle: false, }); } } } // Workers no longer in settings: drop if idle, else retire when idle. for (const [id, entry] of this.entries) { if (seen.has(id)) continue; if (entry.busy) { entry.retireWhenIdle = true; entry.state = "draining"; } else { this.entries.delete(id); } } // Keep "pause all" sticky across a reconfigure (settings save / batch start): // re-disable everything and snapshot any newly-added worker's intended state. if (this.pausedSnapshot) { for (const [id, entry] of this.entries) { if (!this.pausedSnapshot.has(id)) { this.pausedSnapshot.set(id, entry.state); } entry.state = "disabled"; } } // A reconfigure can make new slots eligible (worker enabled/added). this.pump(); } // How many slot entries one remote config expands to. An explicit // remote.slots is the operator's number and is never second-guessed; absent, // the TTL-cached probe of the remote's /api/worker/health answers (its // enabled worker count — a payload that was already sent and discarded), and // 1 stands in until the first probe lands. The probe is fired here, // fire-and-forget, ONLY when the cache is stale — reconfigure() is // synchronous and hot (every batch start), so it never waits on the network; // a changed answer simply re-runs reconfigure with runtime state preserved. private remoteSlotCount(w: Worker): number { const explicit = w.remote?.slots; if (typeof explicit === "number" && explicit >= 1) { return Math.floor(explicit); } const cached = cachedRemoteSlots(w.id); if (remoteCapacityStale(w.id)) { void probeRemoteCapacity(w) .then((slots) => { if (slots !== null && slots !== cached) this.reconfigure(); }) .catch(() => {}); } return cached ?? 1; } // Temporary pause: disable every worker (stop granting new leases; in-flight // leases finish on their own), recording each worker's pre-pause state so // resumeAll restores it exactly. Idempotent — a second call is a no-op. pauseAll(): void { this.ensureInit(); if (this.pausedSnapshot) return; this.pausedSnapshot = new Map(); for (const [id, entry] of this.entries) { this.pausedSnapshot.set(id, entry.state); // Gracefully stop a partial-capable in-flight job so the GPU frees // between fragments (parakeet: finish the current window, cache it, exit // "paused" — resumes from cache on the next run). Mirrors // partialStopWorker: drain it so it ends disabled once the lease // releases. Other engines just stop getting new work and run their // in-flight file to completion. if (entry.busy && entry.activeStop) { entry.state = "draining"; entry.activeStop(); } else { entry.state = "disabled"; } } } // Undo pauseAll: restore each worker to the state it had when paused (a // transient "draining" restores to enabled). Wakes parked waiters. resumeAll(): void { this.ensureInit(); if (!this.pausedSnapshot) return; for (const [id, prev] of this.pausedSnapshot) { const entry = this.entries.get(id); if (entry) entry.state = prev === "draining" ? "enabled" : prev; } this.pausedSnapshot = null; this.pump(); } isPaused(): boolean { return this.pausedSnapshot !== null; } private eligible( entry: PoolEntry, requires?: readonly string[], only?: WorkerFilter, ): boolean { return ( entry.state === "enabled" && !entry.degraded && !entry.busy && workerMatches(entry.config, requires) && (!only || only(entry.config)) ); } // Highest-priority free entry matching the requirement: // (priority asc, insertion order asc). private pickFree( requires?: readonly string[], only?: WorkerFilter, ): PoolEntry | null { let best: PoolEntry | null = null; for (const entry of this.entries.values()) { if (!this.eligible(entry, requires, only)) continue; if (!best || entry.config.priority < best.config.priority) best = entry; } return best; } private grant(entry: PoolEntry): Lease { entry.busy = true; let released = false; const release = () => { if (released) return; released = true; entry.busy = false; if (entry.retireWhenIdle) { this.entries.delete(entry.config.id); } else if (entry.state === "draining") { entry.state = "disabled"; } this.pump(); }; return { worker: entry.config, release }; } // Try to satisfy parked waiters while free slots remain. The queue is // tier-ordered (see enqueueWaiter), and the grant is MATCHED, not blind: each // pass grants the earliest waiter for which a matching slot is free. Scanning // past a waiter whose requirement no free slot satisfies is what stops a GPU // waiter at the head from parking a CPU waiter forever when a CPU slot frees — // a deadlock, not a slowdown. Tier order is preserved WITHIN the set of // waiters a given slot can serve, which is the only order that means anything. private pump(): void { for (;;) { let granted = false; for (let i = 0; i < this.waiters.length; i++) { const waiter = this.waiters[i]; const entry = this.pickFree(waiter.requires, waiter.only); if (!entry) continue; this.waiters.splice(i, 1); waiter.onAbort?.(); waiter.resolve(this.grant(entry)); granted = true; break; } if (!granted) return; } } // Insert a parked waiter respecting priority: a waiter slots in BEFORE the // first waiter of strictly lower priority (so manual/foreground work jumps // ahead of queued auto/background work), preserving FIFO within a tier. Uses // the SAME compareTier comparator as the registry's scheduler — one priority // model governs both, instead of two hand-rolled copies. private enqueueWaiter(waiter: Waiter): void { let insertAt = this.waiters.length; for (let i = 0; i < this.waiters.length; i++) { if (compareTier(waiter.tier, this.waiters[i].tier) < 0) { insertAt = i; break; } } this.waiters.splice(insertAt, 0, waiter); } // Acquire a lease on the highest-priority free worker. Resolves immediately if // one is free, otherwise parks until a slot frees or a worker is enabled. A // background acquire (auto-runner transcription unit) queues behind any // foreground (manual) waiter — see enqueueWaiter. If `signal` aborts while // parked (or before), rejects with an AbortError so the caller treats it like a // cancel. // opts.tier sets the priority tier; the legacy opts.background is still // accepted and maps to the "background" tier (default "foreground"). // opts.requires is a capability requirement (see workerMatches): the lease is // granted only on a worker that matches it. A requirement NO configured worker // matches rejects immediately — parking on it would never resolve, and the // caller should treat the work as unrunnable rather than wedge. // opts.only narrows the lease further, to the workers the filter accepts (see // WorkerFilter); a filter no configured worker passes rejects the same way. acquire( signal?: AbortSignal, opts?: { tier?: SchedulerTier; background?: boolean; requires?: readonly string[]; only?: WorkerFilter; }, ): Promise { this.ensureInit(); if (signal?.aborted) { return Promise.reject(abortError()); } const requires = opts?.requires; const only = opts?.only; if ((requires && requires.length > 0) || only) { let satisfiable = false; for (const e of this.entries.values()) { // Config-level, ignoring busy/disabled/degraded on purpose: a matching // worker that is merely busy or switched off is a reason to park, not // to give up. if (workerMatches(e.config, requires) && (!only || only(e.config))) { satisfiable = true; break; } } if (!satisfiable) { return Promise.reject( new Error( requires && requires.length > 0 ? `no configured worker matches requirement [${requires.join(", ")}]` : "no configured worker can take this work", ), ); } } const entry = this.pickFree(requires, only); if (entry && this.waiters.length === 0) { return Promise.resolve(this.grant(entry)); } const tier: SchedulerTier = opts?.tier ?? (opts?.background ? "background" : "foreground"); return new Promise((resolve, reject) => { const waiter: Waiter = { resolve, reject, tier, requires, only }; if (signal) { const onAbort = () => { const idx = this.waiters.indexOf(waiter); if (idx >= 0) this.waiters.splice(idx, 1); reject(abortError()); }; waiter.onAbort = () => signal.removeEventListener("abort", onAbort); signal.addEventListener("abort", onAbort, { once: true }); } this.enqueueWaiter(waiter); // A new arrival can't create a free slot, but pump keeps the order honest if // one freed between pickFree and here. this.pump(); }); } // Non-parking claim: a matching free slot right now, or null. For dispatchers // that already have a local way to run the work (the digest/backfill fan-out): // a parked acquire inside one of their runPool slots would deadlock the batch, // so they ask, take a lease when one is free, and otherwise run locally. // // opts.kind restricts the claim to one worker kind — the fan-out paths lease // REMOTE capacity only ("llm" endpoints, "remote" unit executors) and must // never grab a local slot out from under the transcription scheduler. // // opts.taggedOnly restricts the claim to EXPLICITLY TAGGED workers. The unit // dispatcher passes it: an untagged remote worker predates the unit protocol // and matches everything under the untagged-is-universal rule, so without // this gate an existing install's transcription-only remote (possibly an // older build with no /api/worker/unit at all) would be shipped units it // cannot serve. Tagging a worker is the operator's opt-in to unit work. // // Never starves a parked waiter: if any parked waiter could take the slot this // would claim, the claim yields (returns null) and pump() serves the waiter. tryAcquire( requires: readonly string[], opts?: { kind?: WorkerKind; taggedOnly?: boolean }, ): Lease | null { this.ensureInit(); let best: PoolEntry | null = null; for (const entry of this.entries.values()) { if (opts?.kind && entry.config.kind !== opts.kind) continue; if (opts?.taggedOnly && !entry.config.tags?.length) continue; if (!this.eligible(entry, requires)) continue; if (!best || entry.config.priority < best.config.priority) best = entry; } if (!best) return null; const claimed = best; if ( this.waiters.some( (w) => workerMatches(claimed.config, w.requires) && (!w.only || w.only(claimed.config)), ) ) { return null; } return this.grant(claimed); } // How many matching slots are FREE right now — the remote term of a // dispatcher's limit() (localTerm + freeSlots(...)). Counting rather than // leasing keeps limit() pure and re-readable every poll. freeSlots( requires: readonly string[], opts?: { kind?: WorkerKind; taggedOnly?: boolean }, ): number { this.ensureInit(); let n = 0; for (const entry of this.entries.values()) { if (opts?.kind && entry.config.kind !== opts.kind) continue; if (opts?.taggedOnly && !entry.config.tags?.length) continue; if (this.eligible(entry, requires)) n++; } return n; } // --- Runtime controls (Workers page). Transient operator overrides: they // mutate the live pool's runtime `state` only, never settings.json. A restart // returns workers to their configured enabled state. They survive the // config-only reconfigure a batch does on start (see reconfigure's // applyEnabled). While a pause-all is active, they also update the snapshot so // resumeAll honors the operator's explicit choice for that worker. --- // Reflect an operator override into the pause snapshot so resumeAll restores // the operator's intent for that worker rather than its pre-pause state. private syncSnapshot(id: string, state: WorkerRuntimeState): void { if (this.pausedSnapshot) this.pausedSnapshot.set(id, state); } enableWorker(id: string): boolean { const entry = this.entries.get(id); if (!entry) return false; entry.state = "enabled"; entry.degraded = false; entry.consecutiveFailures = 0; this.syncSnapshot(id, "enabled"); this.pump(); return true; } // Hard disable: stop granting new leases immediately. In-flight leases keep // running until their controllers release (a true stop also hard-cancels the // job). State goes straight to disabled even if busy. disableWorker(id: string): boolean { const entry = this.entries.get(id); if (!entry) return false; entry.state = "disabled"; this.syncSnapshot(id, "disabled"); return true; } // Drain: stop granting NEW leases but let in-flight transcriptions finish; // flip to disabled once the last lease releases. The "free up the GPU" path. drainWorker(id: string): boolean { const entry = this.entries.get(id); if (!entry) return false; entry.state = entry.busy ? "draining" : "disabled"; this.syncSnapshot(id, "disabled"); return true; } // --- Partial-stop registry: the in-flight transcription on a worker can // register a callback to be gracefully stopped with a partial result. --- setActiveStop(id: string, fn: () => void): void { const entry = this.entries.get(id); if (entry) entry.activeStop = fn; } clearActiveStop(id: string): void { const entry = this.entries.get(id); if (entry) entry.activeStop = undefined; } // True if the worker had an active transcription that supports partial-stop and // was asked to stop. Stopping also takes the worker OUT of rotation: it's busy // finishing the current window, so drain it (→ disabled once the lease // releases, via grant().release) rather than leaving it enabled to immediately // grab the next video — "Stop & keep progress" means stop, not stop-then-go. partialStopWorker(id: string): boolean { const entry = this.entries.get(id); if (!entry?.activeStop) return false; entry.state = entry.busy ? "draining" : "disabled"; this.syncSnapshot(id, "disabled"); entry.activeStop(); return true; } // Whether a worker currently has a partial-stoppable transcription in flight. canStopPartial(id: string): boolean { return !!this.entries.get(id)?.activeStop; } // --- Failure tracking (auto-disable, Phase 5) --- markSuccess(id: string): void { const entry = this.entries.get(id); if (entry) entry.consecutiveFailures = 0; } // Returns true if this failure pushed the worker over the degrade threshold. markFailure(id: string): boolean { const entry = this.entries.get(id); if (!entry) return false; entry.consecutiveFailures++; if ( !entry.degraded && entry.consecutiveFailures >= WORKER_DEGRADE_THRESHOLD ) { entry.degraded = true; return true; } return false; } // Force a worker degraded immediately (skipped by the scheduler until // re-enabled), regardless of the consecutive-failure count. Used when a remote // is confirmed unreachable by a health ping — no point burning more attempts on // it. Returns true if this transitioned a previously-healthy worker. markDegraded(id: string): boolean { const entry = this.entries.get(id); if (!entry || entry.degraded) return false; entry.degraded = true; return true; } // True when at least one worker is enabled and not degraded (ignoring momentary // fullness — a full but enabled pool is working at capacity, not paused). Its // negation is what surfaces the "paused — waiting for a worker" state: a running // batch whose every worker is disabled/draining/degraded, so nothing can start. // `requires` narrows the question to workers matching that capability // requirement — "could a digest fan-out ever land anywhere?" is a different // question from "can transcription run?". hasEligibleWorker(requires?: readonly string[]): boolean { this.ensureInit(); for (const entry of this.entries.values()) { if ( entry.state === "enabled" && !entry.degraded && workerMatches(entry.config, requires) ) { return true; } } return false; } // Ids of every currently-enabled worker — the snapshot the Workers page "Set // as default" button persists as the launch default (see workerDefaults.ts). enabledIds(): string[] { this.ensureInit(); return Array.from(this.entries.values()) .filter((e) => e.state === "enabled") .map((e) => e.config.id); } summary(): WorkerSummary[] { this.ensureInit(); return Array.from(this.entries.values()) .map((e) => ({ id: e.config.id, name: e.config.name, kind: e.config.kind, appId: e.config.appId, priority: e.config.priority, busy: e.busy, state: e.state, degraded: e.degraded, enabled: e.config.enabled, ...(e.config.tags && e.config.tags.length > 0 ? { tags: [...e.config.tags] } : {}), device: e.config.config?.device, })) .sort((a, b) => a.priority - b.priority || a.id.localeCompare(b.id)); } } function abortError(): Error { const err = new Error("worker acquire aborted"); err.name = "AbortError"; return err; } declare global { // eslint-disable-next-line no-var var __yttWorkerPool__: WorkerPool | undefined; } export function getWorkerPool(): WorkerPool { if (!globalThis.__yttWorkerPool__) { globalThis.__yttWorkerPool__ = new WorkerPool(); } return globalThis.__yttWorkerPool__; }