import type { SiteSettings } from "../lib/settings"; import { digestLaneFor, laneYieldsToTranscription, type Lane, } from "../lib/operations"; import type { DigestLane } from "../lib/digest"; import { isGateHeld } from "../lib/pauseGates"; import { transcriptionActivity } from "./digestYield"; // THE LANE'S RULES, declared once, in the two shapes a dispatcher can act on. // // This exists because of a sentence in backfillLaneOperations' own header: // // "digest carries its own pause gate, the yield-to-transcription carve-out (with // its CPU-worker exemption), spendCapUsd on the metered lane, the // remoteEnabled fail-fast, the engine probe() fail-fast … and not one of // them is expressible as backfillLimit()'s single scalar." // // That is the reason there used to be two schedulers for one kind of work. Every // one of those guards lived inside the digest batch's own closure, where nothing // else could ask about it — so a second dispatcher wanting to run digest would // have had to re-implement all six, correctly, from memory. They live here now, // and operationBatch calls them, so there is exactly one definition and no second // opinion to drift. // // TWO SHAPES, AND THE SPLIT IS FORCED BY WHEN THEY CAN BE ASKED: // // * a PREFLIGHT is answered before any pool exists. `remoteEnabled` and the // engine `probe()` both THROW today, before runPool is constructed, and // they have to: an unreachable engine is a configuration problem, and // discovering it per-item wastes 1,100 log lines. A dispatcher needs to ask // that question before it commits a job, not after. // * a GATE is answered at DISPATCH time, on every pull. It returns a HOLD, // never a stop — runPool idle-waits at a zero limit, whereas next() // returning null ends the job. Every pause in this repo depends on that // distinction, and stating it here is what keeps the next dispatcher from // re-deciding it. // // Nothing here reads disk or mutates settings. The spend cap is the one rule // that needs per-RUN state, and it takes it as an argument rather than reaching // for a module-level accumulator — see DigestGateInput.costUsd. // A gate's answer. `hold` carries the reason in words, because an operator // watching a lane sit at zero throughput has to be able to tell YIELDING from // WEDGED, and those look identical from outside. export type LaneGate = | { hold: false } | { hold: true; // Machine-readable, for a surface that wants to branch. reason: "paused" | "yield" | "spend-cap"; // Human-readable, logged on the EDGE only — a line per poll would bury a // job log over a multi-week run. message: string; }; export const GO: LaneGate = { hold: false }; export type PreflightVerdict = { ok: true } | { ok: false; error: string }; // --- Digest ------------------------------------------------------------------ export type DigestPreflightInput = { settings: SiteSettings; // Which lane was ASKED for, not which one is configured. lane: "local" | "remote"; app: { id: string; lane: DigestLane; // The engine's own reachability check. Awaited here so the caller has one // place to ask "can this lane run at all". probe: () => Promise; }; }; // Can this digest lane run at all? // // Both answers are refusals rather than holds, and deliberately: a disabled // metered lane and an unreachable engine are configuration problems that will // not resolve on their own, so parking the pool on them would idle forever // while looking busy. export async function digestPreflight( input: DigestPreflightInput, ): Promise { if (input.lane === "remote" && !input.settings.digest.remoteEnabled) { return { ok: false, error: "The metered digest lane is disabled (settings.digest.remoteEnabled). Enable it in Settings before running it.", }; } if (!(await input.app.probe())) { return { ok: false, error: `Digest engine ${input.app.id} is not reachable. ` + (input.app.lane === "local-gpu" ? "Is the ollama service running (systemctl status ollama)?" : "Is the claude CLI installed and on PATH (set CLAUDE_BIN)?"), }; } return { ok: true }; } export type DigestGateInput = { settings: SiteSettings; // The engine's lane, which decides both the queue key and whether this run // contends for the GPU. Passed rather than re-derived so a caller that has // already resolved its app cannot disagree with one that has not. appLane: DigestLane; metered: boolean; // Cumulative metered spend SO FAR IN THIS RUN. The cap is per-job, so it // cannot be read from settings alone — and passing it in is what keeps this // function pure enough to test. costUsd: number; }; // Should the digest lane hold right now? // // Order matters and is the order it was in: an explicit pause beats a yield // beats a spend cap. A paused lane that reported "yielding" would send an // operator to look at the transcription queue for a switch they turned off // themselves. export function digestGate(input: DigestGateInput): LaneGate { const digest = input.settings.digest; // Re-read at DISPATCH time so a pause takes effect within one poll and // survives a restart with no boot hook — the downloadsPaused pattern. The // gate itself is defined once, in lib/pauseGates.ts. if (isGateHeld(input.settings, "digest")) { return { hold: true, reason: "paused", message: "Digests are paused; the lane is holding.", }; } // Step aside for whisper. WHICH LANE YIELDS IS DECLARED, not re-tested here: // digestLaneFor maps the engine to the digest operation's registered lane and // laneYieldsToTranscription reads `contendsFor: "gpu"` off it, so a future GPU // operation cannot answer this differently. // // The CPU carve-out lives one level down, in transcriptionActivity(): a // worker pinned to `device: "cpu"` is not GPU contention, and treating it as // such once stalled the digest lane for transcription competing for zero // shaders. A worker with NO device set still triggers the yield — the unknown // case fails safe. if ( laneYieldsToTranscription(digestLaneFor(input.appLane)) && digest.yieldToTranscription ) { const activity = transcriptionActivity(); if (activity.busy) { return { hold: true, reason: "yield", message: `Yielding the GPU to transcription (${activity.reason}); the digest lane will resume when it is free.`, }; } } if ( input.metered && digest.spendCapUsd > 0 && input.costUsd >= digest.spendCapUsd ) { return { hold: true, reason: "spend-cap", message: `Spend cap reached ($${input.costUsd.toFixed(2)} of ` + `$${digest.spendCapUsd.toFixed(2)}) — parking the metered lane.`, }; } return GO; } // --- What a lane does with a duplicate cluster ------------------------------- // Whether an operation's output can be SHARED between two recordings of the // same event rather than generated twice. // // Digest can: a chapter list is about what was said, so an aligned mirror gets // the canonical member's digest for free (worth ~11% of the sweep). Diarization // and attribution cannot — they are grounded in a specific audio track and a // specific set of cue timings, and a mirror's clock is not the same one. // // Declared here rather than inferred from an id, so a dispatcher can ask the // question without knowing which operation it is holding. `useClusters: false` // on a run still overrides it; this is the DEFAULT, not a lock. export function laneSharesDuplicates(lane: Lane): boolean { return ( lane.queueKey === digestLaneFor("local-gpu").queueKey || lane.queueKey === digestLaneFor("remote-api").queueKey ); }