// GPU arbitration between the digest lane and the transcription lane. // // `queueKeys.ts` deliberately puts `digest:local` on a DIFFERENT registry queue // from TRANSCRIPTION_QUEUE, so the two never serialize against each other. That // reasoning is right for digest-local vs digest-remote (GPU-bound vs // network-bound), but its consequence for digest-local vs transcription is that // **ollama and whisper run concurrently on the same 8 GB card**. // // That is not hypothetical, but the size of it is a HYPOTHESIS, not a // measurement. The bake-off's 11.2 s per chunk was ENGINE time on an idle box; // the 102-video validation run's 24.9 s per chunk was WALL time on a box also // running auto-transcribe. Those are not the same quantity, so their 2.22× ratio // is an upper bound on contention that also contains model-load, prefill and — // crucially — this file's own 3 s poll idling. Re-priced per chunk the corpus is // ~25 days idle to ~55 contended (see controller/digestPlan.ts). Splitting that // gap needs a production run with the per-call load/prefill/decode breakdown that // digestApps.ts now records; until then, do not quote a contention penalty as // fact. // // The fix is deliberately NOT a scheduler. the digest lane's limit() already // returns 0 to idle-wait on a pause, and runPool treats a zero limit as "hold, // don't finish" — so the digest lane can step aside using machinery that is // already proven, with no priority system, no new queue and no boot hook. // // This is a YIELD, not a lock. A digest already in flight is not interrupted, so // the two can still overlap for the length of one generation if transcription // starts just after a dispatch. Interrupting would waste the partial work; the // next dispatch sees the busy lane and holds. import { getWorkerPool, type WorkerSummary } from "../jobs/workerPool"; import { getRegistry, type JobRecord } from "../jobs/registry"; import { TRANSCRIPTION_QUEUE } from "../lib/queueKeys"; import { getSettings } from "../lib/settings"; export type TranscriptionActivity = { busy: boolean; // Which signal fired, for the log line. An operator watching a sweep sit at // zero throughput needs to be told it is yielding rather than wedged. reason: "local-worker" | "queued-job" | null; }; // Does this busy worker actually compete for GPU shaders? // // The original check was `kind === "local"` alone, and that was a BUG. `local` // distinguishes "on this box" from "delegated to another machine" — it says // nothing about which processor the engine uses. This box runs three enabled local // parakeet workers: one on the engine's default device and two pinned to // `device: "cpu"`. At parallelTranscriptions 2 the CPU pair alone was enough to // hold the digest lane at zero throughput for work using no shaders at all. // // The unknown case fails SAFE, in the direction that costs throughput rather than // correctness: only an explicit "cpu" is treated as non-contending. No device set // means parakeet-cli's own default, which may be the GPU. export function workerContendsForGpu( w: Pick, yieldToCpuWorkers: boolean, ): boolean { if (w.kind !== "local") return false; if (yieldToCpuWorkers) return true; return w.device?.trim().toLowerCase() !== "cpu"; } // Is the local transcription lane using the GPU right now? The whole decision, as // a PURE function of the two signals — so both directions of the CPU-worker fix can // be asserted without a live pool, a live registry or a GPU. transcriptionActivity() // below is the thin I/O wrapper that feeds it. // // Two signals, because one alone leaves a hole: // - a busy LOCAL worker on a GPU device is the transcription engine actually // running on this box's card (`remote` workers delegate to another machine // and compete for nothing here, so they are deliberately excluded, as are // CPU-pinned workers unless `yieldToCpuWorkers` says otherwise); // - a RUNNING job on TRANSCRIPTION_QUEUE covers the GAPS BETWEEN worker // acquisitions — audio extraction, model load, the moment between two // videos in a batch. Those gaps are exactly where a multi-minute digest // generation would otherwise slip in and hold the card. // // The second signal is deliberately consulted ONLY when no local worker is busy. // It is a proxy for "a phase the pool cannot see", and if a worker IS busy then // the pool has already told us the truth — including which DEVICE. Without that // condition the device check above would be dead code on this box: a CPU-pinned // transcription is a running job on the queue, so the digest lane would go on // yielding to it through the gap signal and the fix would change nothing. // // Known, accepted narrow window: if one worker is busy on CPU while a second job // is loading a model onto the GPU, this reports not-busy and one digest chunk may // overlap it. A yield is not a lock — an in-flight generation was never // interrupted either — and the alternative reinstates the bug. export function evaluateTranscriptionActivity(opts: { workers: readonly Pick[]; transcriptionJobRunning: boolean; yieldToCpuWorkers: boolean; }): TranscriptionActivity { let anyLocalBusy = false; for (const w of opts.workers) { if (!w.busy || w.kind !== "local") continue; anyLocalBusy = true; if (workerContendsForGpu(w, opts.yieldToCpuWorkers)) { return { busy: true, reason: "local-worker" }; } } if (anyLocalBusy) return { busy: false, reason: null }; if (opts.transcriptionJobRunning) return { busy: true, reason: "queued-job" }; return { busy: false, reason: null }; } // The I/O wrapper the digest batch actually calls: read the live pool, the live // registry and the setting, then hand them to the pure decision above. export function transcriptionActivity(): TranscriptionActivity { // Read once: a setting flipping mid-scan would produce an answer that matches // neither configuration. let yieldToCpuWorkers = false; try { yieldToCpuWorkers = getSettings().digest.yieldToCpuWorkers; } catch { /* unreadable settings: the sanitizer's default is false, so keep it */ } // Each read is guarded separately and falls back to "nothing there", i.e. fail // OPEN: a pool or registry that cannot be read must not wedge the digest lane. // The cost of being wrong here is contention, not deadlock. let workers: readonly WorkerSummary[] = []; try { workers = getWorkerPool().summary(); } catch { /* fail open */ } let transcriptionJobRunning = false; try { transcriptionJobRunning = getRegistry() .list() .some( (job) => job.queueKey === TRANSCRIPTION_QUEUE && job.status === "running", ); } catch { /* fail open */ } return evaluateTranscriptionActivity({ workers, transcriptionJobRunning, yieldToCpuWorkers, }); } // --- Is a transcription running on THIS video right now? -------------------- // // A second consumer of the same registry, asking a narrower question than the // lane-level one above: not "is the transcription lane busy" but "is some job // transcribing this exact video". The backfill's cleanup uses it as defense in // depth before unlinking a re-acquired audio file — deleting audio out from // under a running engine costs that video's whole run and appends it to // failed-transcriptions. // // Registry TASKS are the right signal because they cover BOTH producers: the // auto runner and the manual whisper-all batch both go through transcribeOne's // tracker (transcribeOne.ts -> jobs/taskHooks.ts -> registry.addTask), whereas // getAutoRunnerStatus() sees only the runner and would miss the manual batch. // // Known blind spot, accepted: an inline transcribe inside a download // (downloadOneManaged) passes no tracker and registers no task, so it is // invisible here. The primary mechanism is the policy hand-off below it, not // this check. export function evaluateTranscribingVideo( jobs: ReadonlyArray>, videoId: string, ): boolean { return jobs.some( (job) => job.status === "running" && (job.tasks ?? []).some( (t) => t.kind === "transcribe" && t.id === videoId, ), ); } // Thin I/O wrapper, same shape as transcriptionActivity() above. Fails OPEN to // false: "not transcribing" simply falls through to the caller's own policy // decision, which is the mechanism that actually protects the file. An // unreadable registry must not change what the backfill does. export function isTranscribingVideo(videoId: string): boolean { try { return evaluateTranscribingVideo(getRegistry().list(), videoId); } catch { return false; } }