// ONE EXECUTOR FOR THE OPERATION LANES — digest and backfill. // // This file is `digestBatch.ts` and `backfillBatch.ts` folded together. They // were structurally the same runner twice: the header of the second one said so // ("Structurally this is digestBatch.ts, and the three mechanics it copies are // copied for the same reasons"), and copies drift. They had already: one keyed // its GPU yield off the PRESENCE of `laneFor` and the other off the declared // `contendsFor`; one classified a candidate through the operation registry and // the other re-implemented freshness inline; one printed `deferred`/`blocked` // and the other could not count them at all. // // THE THREE MECHANICS, kept verbatim from both, because a multi-week lane // depends on each of them: // // 1. RE-DERIVE ELIGIBILITY FROM DISK, on every pull. Never a frozen array // with an index cursor: that form does not survive a restart and cannot // see a video that became eligible mid-run (a transcript that finished, a // mirror that got shared to). classifyOperationUnit IS that re-derivation, // and runOperationUnit re-asks it at dispatch — the pull and the run are // separated by however long the pool held us. // 2. A ZERO LIMIT IS A HOLD, NEVER A STOP. runPool idle-WAITS at limit 0, // where returning null from next() ENDS the job. Every pause in this repo // depends on that distinction, and laneLimit() returns 0 rather than any // kind of refusal for exactly that reason. // 3. ONE JOB PER CHANNEL for the manual verbs, never one per video. Job logs // keep the newest 500 / 30 days and the registry keeps 100 records; // 119,600 jobs would evict everything, including the running ones. The // auto-queue runner's per-video units are IN-PROCESS and make no job // record, which is why it may work per video and this may not. // // WHAT IS NOT HERE: a scheduler. The queue key gives a lane concurrency // (registry.ts runs every non-empty key at concurrency 1, so a distinct key is // the ONLY concurrency mechanism in this system), and laneLimit() gives it a // share. Nothing new. import path from "node:path"; import { open, readFile, readdir } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { getSettings, type SiteSettings } from "../lib/settings"; import { isGateHeld } from "../lib/pauseGates"; import { assertChannelMediaReachable, assertChannelTextReadable, readRelocationMarker, } from "../lib/channelMedia"; import { runPool } from "../jobs/concurrentRunner"; import type { TaskTracker } from "../jobs/taskHooks"; import type { JobProgress } from "../jobs/registry"; import { readVideoFiles, CUES_JSON_FILENAME } from "../lib/videoStatus"; import { DIGEST_OPERATION_ID, addOperationState, digestLaneFor, emptyOperationCounts, laneYieldsToTranscription, operationsForLane, reachableOperationWork, type Operation, type OperationClassification, type OperationRunOutcome, } from "../lib/operations"; import type { DigestAppConfig, DigestLane, DigestSectionKind } from "../lib/digest"; import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; import { getWorkerPool } from "../jobs/workerPool"; import { transcriptionActivity } from "./digestYield"; import { acquireLlmSlot, freeLlmSlots } from "./llmWorkers"; import { resolveAttributionTarget } from "./attributionTarget"; import { pingRemoteHealth } from "./remoteTranscribe"; import { runUnitViaRemote } from "./remoteUnit"; import { TranscribeError } from "./transcribeError"; import { readDigestContext } from "../lib/digestContext-server"; import type { StartWorkerUnitInput } from "./workerServer"; import { reacquireMediaFor, type ReacquireOutcome } from "./backfillReacquire"; import { batchRecencyComparator } from "./batchRecency"; import { digestGate, digestPreflight, laneSharesDuplicates, type LaneGate, } from "./laneGuards"; import { digestVideo } from "./digestVideo"; import { resolveDigestChunking, resolveDigestTarget, type DigestLaneChoice, type ResolvedDigestTarget, } from "./digestTarget"; import { laneForOperation } from "./operationLane"; import { buildDigestClusterPlan, planSlugForDir, shareDigestToCluster, type DigestClusterPlan, } from "./digestSharing"; // The two lanes with an operation work list. Not AutoQueueKind: transcription // and download are dispatched by their own executors inside the runner, and // nothing in this file can run them. export type OperationLane = "digest" | "backfill"; export type DigestOrder = "shortest-first" | "longest-first"; // --------------------------------------------------------------------------- // The dispatch decision // --------------------------------------------------------------------------- // What the candidate pull does with one classification. // // THIS USED TO BE AN IF-CHAIN INSIDE THE PULL, AND THAT WAS THE HAZARD. It // handled the states it knew about and FELL THROUGH TO DISPATCHING everything // else — so a state added to OperationState was not skipped by default, it was // RUN by default. Nothing in the repo checked OperationState exhaustively, so // widening the union raised zero TypeScript errors, and the first sign of a // missing branch would have been the six-hour video the new state existed to // avoid being handed to the engine anyway. // // Extracted as a pure function with a `never` check so that hazard is now a // COMPILE error rather than a runtime one, and so the dispatch decision — which // is the single most consequential branch in this file — can be tested without a // registry, a worker pool or a corpus. export type CandidateAction = // Hand it to the runner. | "dispatch" // Nothing to count and nothing to do (this operation has no opinion on it). | "skip" // Already at the current identity. | "fresh" // Input is gone and re-acquiring is off. | "missing-input" // The operation refuses to attempt it under the current configuration. | "deferred" // Waiting on a prerequisite operation's output. | "blocked"; export function candidateAction( state: OperationClassification, opts: { force: boolean; allowRedownload: boolean }, ): CandidateAction { switch (state) { case "not-applicable": return "skip"; case "present": // `force` is how an operator redoes a video that is already current. return opts.force ? "dispatch" : "fresh"; case "missing-input": // With re-download on, the runner re-acquires the media and then works. return opts.allowRedownload ? "dispatch" : "missing-input"; case "deferred": // NEVER dispatched, and deliberately not affected by `force`: `force` means // "redo work that looks done", not "ignore the cap". Raising the cap is how // you ask for a deferred video, and it is one edit in Settings. return "deferred"; case "blocked": // NEVER dispatched, and — like `deferred` — deliberately immune to // `force`. Forcing a video whose prerequisite has not been produced does // not make the prerequisite appear; it just hands the runner an input it // does not have. Running the operation this one dependsOn is how you // unblock it, and orderByDependencies tries to do that for you within the // same pass. // // Note what is NOT here: re-acquiring media. That is what separating this // from `missing-input` bought — allowRedownload has no bearing on a // blocked video, so the lane cannot spend a download on one. return "blocked"; case "missing": case "stale": // Part-done is reachable work like any other, and the runner already does // the right thing with it without being told: digestVideo regenerates the // sections that are not fresh and leaves the ones that are. The split // between this and `stale` is about what the work COSTS a planner, not // about whether it is dispatched — so it belongs on exactly this line, and // giving it its own CandidateAction would only invite a caller to treat it // as another do-not-run state alongside deferred and blocked. case "partial": return "dispatch"; default: { const unhandled: never = state; throw new Error(`unhandled backfill state: ${String(unhandled)}`); } } } // --------------------------------------------------------------------------- // The per-lane limit // --------------------------------------------------------------------------- // How many slots the backfill lane may use right now. // // A PURE function of the three inputs, so both directions can be asserted // without a live pool, a live registry or a GPU — the shape digestYield.test.ts // established and for the same reason: this decides whether a multi-day lane // runs at all, and a test that needs the hardware is a test nobody runs. // // idleOnly — full slots when the primary transcription lane is quiet, zero // while it works. Zero is a HOLD, not a stop: runPool idle-waits, so the // lane resumes the moment the primary is free, having re-derived nothing. // otherwise — the lane's own slots, whatever transcription is doing. // // THE INPUT USED TO BE `settings.backfill.weight`, a 0..1 scalar that meant // "idle-only" at 0 and "a guaranteed share, floored at 1" above it. Slice 1.3 // retired it: the yield is the OPERATION'S DECLARED `contendsFor` (a run that // could dispatch GPU work stands aside; one that contends for the network does // not), and the share is `concurrency` and the lane policy's `maxWorkers`. One // scalar standing for both a resource question and a size question is how a // network-bound attribution run came to park itself behind a GPU transcription // it was not competing with. export function backfillLimit(opts: { idleOnly: boolean; slots: number; primaryBusy: boolean; }): number { const slots = Math.max(0, Math.floor(opts.slots)); if (slots === 0) return 0; if (opts.idleOnly) return opts.primaryBusy ? 0 : slots; return slots; } // The live, per-RUN facts a limit cannot read from settings: leases this run // already holds, its metered spend so far, and which engine it resolved to. export type DigestLaneLive = { lane: "digest"; // 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. costUsd: number; // Local slots this run may use at once. concurrency: number; // An "llm" endpoint worker can take this lane's model calls — true only for // an HTTP engine with a baseUrl field (ollama). The metered CLI lane has no // endpoint to swap. fanOutEligible: boolean; // Leases THIS run holds, so the limit stays stable at local + total-llm // capacity instead of sagging while slots are leased. llmActive: number; }; export type BackfillLaneLive = { lane: "backfill"; // The operations in this run. Their DECLARED resource is what decides the // GPU carve-out below. operations: readonly Operation[]; // Of those, the ones whose work bottoms out in model calls an "llm" endpoint // worker can serve. llmOps: readonly string[]; remoteEligible: boolean; llmActive: number; unitActive: number; }; export type LaneLive = DigestLaneLive | BackfillLaneLive; // A lane's answer at one dispatch tick. // // `hold` is 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. It is carried beside the number rather than logged // here so the caller can log it EDGE-TRIGGERED — a line per poll would bury a // job log over a multi-week run. export type LaneLimitVerdict = { limit: number; hold: (LaneGate & { hold: true }) | null; }; // The operations whose work an "llm" endpoint worker can serve. Diarization is // deliberately absent — its work is an audio pass, which only a unit executor // (a full instance of this app) can take. const LLM_BACKFILL_OPS = new Set(["attribution-text", "attribution-diarized"]); // A tagged remote worker takes whole units matching [operation id, contendsFor]. function unitRequires(op: Operation): string[] { return [op.id, (laneForOperation(op.id) ?? op.lane).contendsFor]; } // EVERY GUARD FOR ONE LANE, ASKED AS ONE DECLARED RULE. // // Settings are re-read at DISPATCH time by the caller and handed in here, so a // pause takes effect within one poll and survives a restart with no boot hook // (the downloadsPaused pattern), and the answer is a HOLD: returning 0 makes // runPool idle-wait, where returning null from next() would END the run. export function laneLimit( settings: SiteSettings, live: LaneLive, ): LaneLimitVerdict { if (live.lane === "digest") return digestLaneLimit(settings, live); return backfillLaneLimit(settings, live); } function digestLaneLimit( settings: SiteSettings, live: DigestLaneLive, ): LaneLimitVerdict { const gate = digestGate({ settings, appLane: live.appLane, metered: live.metered, costUsd: live.costUsd, }); // The REMOTE term: leases this run already holds plus free llm slots serving // digest. Counting `llmActive` keeps the limit stable at local + // total-llm-capacity rather than sagging as slots are leased. const remoteTerm = live.fanOutEligible ? live.llmActive + freeLlmSlots([DIGEST_OPERATION_ID]) : 0; if (gate.hold) { // The yield-to-transcription hold zeroes ONLY the local term: an llm // endpoint contends for nothing on this box's GPU, so parking it too would // idle remote capacity to protect hardware it never touches. An operator // pause and the spend cap still zero BOTH terms — intent and money are // global. if (gate.reason === "yield" && remoteTerm > 0) { return { limit: remoteTerm, hold: gate }; } return { limit: 0, hold: gate }; } return { limit: Math.max(1, live.concurrency) + remoteTerm, hold: null }; } function backfillLaneLimit( settings: SiteSettings, live: BackfillLaneLive, ): LaneLimitVerdict { const cfg = settings.backfill; // The lane's gate is `enabled`, whose polarity is the opposite of every other // gate's — asked through isGateHeld so this file does not carry a second // opinion about it. if (isGateHeld(settings, "backfill")) { return { limit: 0, hold: { hold: true, reason: "paused", message: "Backfill lane disabled in settings — holding.", }, }; } // Reused verbatim, not re-derived: transcriptionActivity() already combines // "a busy local GPU-device worker" with "a job running on TRANSCRIPTION_QUEUE" // — the second signal being what covers the gaps between worker acquisitions // (audio extraction, model load, the moment between two videos), which is // exactly where a multi-minute backfill item would otherwise slip in. It fails // open in both directions. // // It sees the TRANSCRIPTION lane only. It deliberately does NOT see the digest // lane: that one is already arbitrated against transcription by the same // helper, and making two yielding lanes watch each other is how both end up // parked. The cost is that a backfill and a digest can overlap on the CPU; the // benefit is that neither can deadlock the other. const activity = transcriptionActivity(); // A RUN THAT COULD DISPATCH GPU WORK IS IDLE-ONLY. Sharing is a CPU concept: // it is reasonable when the contended resource is cores and divisible, and it // is not available for VRAM — sortformer on Vulkan holds ~4.4 GB of the same // 8 GB card parakeet is using, so "a small share" is not a slower run, it is // an out-of-memory failure of whichever lane allocates second. // // The cost is bluntness — one such operation makes the whole run idle-only, // including its cheap CPU-bound siblings — and that is the deliberate // direction to be wrong in, since the alternative is an OOM mid-sweep. // // KEYED ON `contendsFor`, off the operation's LIVE lane, and that is the fix // this slice shipped. It used to be keyed on the PRESENCE of `laneFor`, as a // proxy for "this one's resource depends on settings, so nothing else is // deciding it" — a proxy that made `laneFor` load-bearing by its existence // and would have silently enrolled every operation the day someone gave them // all one. `laneYieldsToTranscription` is the declaration itself, asked of the // lane laneForOperation resolves, so a statically GPU-bound operation on this // queue is now caught too and a `laneFor` added for any other reason changes // nothing. const gpuBound = live.operations.some((op) => laneYieldsToTranscription(laneForOperation(op.id) ?? op.lane), ); const local = backfillLimit({ idleOnly: gpuBound, slots: cfg.concurrency, primaryBusy: activity.busy, }); // The REMOTE term: llm-endpoint slots serving this run's operations, plus the // leases it already holds. transcriptionActivity zeroes only the LOCAL term — // an endpoint contends for nothing on this box — while the lane-disabled hold // above still zeroes everything (operator intent is global). See // `remoteEligible` for why a mixed-operation run gets no term. const remoteTerm = live.remoteEligible ? live.llmActive + freeLlmSlots([...live.llmOps]) : 0; // The UNIT term: tagged remote slots every operation in this run could take, // plus the leases it already holds. The MIN across operations keeps it sound // on a mixed run — a slot justified by one operation's remote capacity must // not dispatch another onto this box while the local term holds. const unitTerm = live.unitActive + (live.operations.length === 0 ? 0 : Math.min( ...live.operations.map((op) => getWorkerPool().freeSlots(unitRequires(op), { kind: "remote", taggedOnly: true, }), ), )); const total = local + remoteTerm + unitTerm; if (total === 0) { return { limit: 0, hold: { hold: true, reason: "yield", message: `Standing aside for transcription (${activity.reason}); the backfill lane resumes when it is free.`, }, }; } return { limit: total, hold: null }; } // --------------------------------------------------------------------------- // Duration, read cheaply // --------------------------------------------------------------------------- // Read a video's duration cheaply. transcript.cues.json is // {version, source, transcriptFormat, ...summary, cues} — the summary (and so // `duration`) is serialized BEFORE the multi-megabyte cues array, so the head of // the file is enough and a full parse is only the fallback. At 74k videos this is // the difference between a few seconds and reading 6.9 GB to sort a list. const DURATION_HEAD_BYTES = 8192; export async function readDurationFast( cuesPath: string, ): Promise { let handle: Awaited> | null = null; try { handle = await open(cuesPath, "r"); const buf = Buffer.alloc(DURATION_HEAD_BYTES); const { bytesRead } = await handle.read(buf, 0, DURATION_HEAD_BYTES, 0); const head = buf.subarray(0, bytesRead).toString("utf8"); const m = head.match(/"duration"\s*:\s*([0-9]+(?:\.[0-9]+)?)/); if (m) return Number(m[1]); // Short file: the whole thing is in `head`, so parse it properly. if (bytesRead < DURATION_HEAD_BYTES) { const parsed = JSON.parse(head) as { duration?: unknown }; return typeof parsed.duration === "number" ? parsed.duration : null; } return null; } catch { return null; } finally { await handle?.close().catch(() => {}); } } // A video's duration, for the lane's `cheapest` order. Memoized process-wide the // way recencyIndex memoizes its tail reads, and for the same reason: a duration // never changes, and the runner re-derives its ordering on every scheduling tick. // `null` memoizes a genuine miss (no normalized transcript) so a video with none // is not re-opened once per tick forever. const durationMemo = new Map(); const DURATION_MEMO_CAP = 200_000; export async function videoDurationSeconds( paths: Paths, channelSlug: string, videoId: string, ): Promise { const key = `${channelSlug}\u0000${videoId}`; const memo = durationMemo.get(key); if (memo !== undefined) return memo; const seconds = await readDurationFast( path.join(paths.channelsDir, channelSlug, "data", videoId, CUES_JSON_FILENAME), ); if (durationMemo.size < DURATION_MEMO_CAP) durationMemo.set(key, seconds); return seconds; } // Test seam, and the same one recencyIndex exposes. export function resetDurationMemo(): void { durationMemo.clear(); } // THE DIGEST LANE'S ORDER, COMPOSED, and stated once so a migration cannot drop // half of it. // // The live behaviour this reproduces is `digestBatch`'s two sorts: durations // ascending (or descending) FIRST, then a stable second sort by upload date. // Array#sort is stable and the recency comparator returns 0 for two videos // sharing a YYYYMMDD key, so the composition reads as "newest day first, // shortest video within a day" — DATE is primary, duration is the tiebreak, // which is the opposite of the order the two sorts are written in. // // Unknown durations sort LAST in either direction. digestBatch dropped them // from candidacy entirely (a video with no normalized transcript is a // transcription problem, not a digest one, and the registry classifies it // `blocked` or `deferred`); the runner cannot drop them, because the snapshot's // reachable ids are what it draws, so it puts them behind everything it can // price instead of letting `null` sort to the front. // // `direction` is a PARAMETER rather than a constant so this stays a pure // comparator factory. Its one caller fixes it at "newest" — the value // `settings.digest.recencyOrder` carried on the live corpus before slice 1.3 // retired that field. export function cheapestComparator(opts: { duration: (id: string) => number | null | undefined; recency?: ((a: string, b: string) => number) | null; }): (a: string, b: string) => number { return (a, b) => { if (opts.recency) { const byDate = opts.recency(a, b); if (byDate !== 0) return byDate; } const da = opts.duration(a); const db = opts.duration(b); const ka = da === null || da === undefined ? Number.POSITIVE_INFINITY : da; const kb = db === null || db === undefined ? Number.POSITIVE_INFINITY : db; if (ka !== kb) return ka - kb; return a.localeCompare(b); }; } // --------------------------------------------------------------------------- // The run context // --------------------------------------------------------------------------- // Everything a unit needs that is resolved ONCE per run, not per video: the // identity targets, the engine, the duplicate plan, the fan-out configuration. // Deriving any of it per item is how a counter and a runner end up disagreeing // about what is stale. export type OperationRun = { lane: OperationLane; paths: Paths; // The settings this run's IDENTITY was resolved against. Guards re-read live // settings at dispatch (see laneLimit); this copy is for targets only. settings: SiteSettings; // The operations this run may dispatch, in dependency order. operations: Operation[]; force: boolean; allowRedownload: boolean; log: (msg: string) => void; // Live counters the limit reads and the units mutate. live: { llmActive: number; unitActive: number; costUsd: number; // Latches for the rest of the run once free disk fell under the floor. diskFloorHit: boolean; }; // ── digest ────────────────────────────────────────────────────────────── digest: { laneChoice: DigestLaneChoice; sections?: DigestSectionKind[]; // THE ENGINE, resolved once and WITHOUT a channel — because it does not // depend on one. Only the freshness target does (it hashes the channel's // context note), which is why the two are resolved apart: a cross-channel // runner can then answer "which lane am I on, and is it metered?" — the // whole of what laneLimit needs — synchronously, on every dispatch tick. engine: DigestEngine; // Skip mirrors and share the canonical member's digest to aligned ones. clusterPlan: DigestClusterPlan | null; // Resolved per channel. byChannel: Map>; } | null; // ── backfill ──────────────────────────────────────────────────────────── backfill: { llmOps: string[]; llmFanout: { config: DigestAppConfig; modelRequested: string } | null; remoteEligible: boolean; // The unit-executor envelope MINUS the channel context, which is injected // per unit: a cross-channel runner has more than one. unitConfig: StartWorkerUnitInput["config"]; // `${op} ${channel}` -> resolved target. targets: Map>; } | null; }; type DigestEngine = ReturnType & { modelRequested: string; }; export type OpenOperationRunOptions = { lane: OperationLane; paths: Paths; // Which operations to run. Empty/omitted = every one the lane dispatches. operationIds?: readonly string[]; force?: boolean; // Re-acquire media for `missing-input` videos. Defaults to // settings.backfill.allowRedownload; the explicit flag is for a one-off run. allowRedownload?: boolean; // Digest only. digestLane?: DigestLaneChoice; sections?: DigestSectionKind[]; // Pass a prebuilt plan to avoid re-reading the duplicates report; `false` // turns sharing off for this run. On by default: worth ~11% of the sweep. useClusters?: boolean; clusterPlan?: DigestClusterPlan; onLog?: (msg: string) => void; }; // BOTH FAIL-FASTS, asked as ONE declared rule, BEFORE any pool exists. // // It THROWS, byte-identically to the batch it replaces: an unreachable engine // is a configuration problem, and discovering it per-item wastes 1,100 log // lines. `remoteEnabled` and the engine `probe()` are the two, and asking them // needs the engine only — never a channel. export async function preflightOperationRun(run: OperationRun): Promise { if (!run.digest) return; const { app, config } = run.digest.engine; const verdict = await digestPreflight({ settings: run.settings, lane: run.digest.laneChoice, app: { id: app.id, lane: app.lane, probe: () => app.probe(config) }, }); if (!verdict.ok) throw new Error(verdict.error); } export async function openOperationRun( opts: OpenOperationRunOptions, ): Promise { const settings = getSettings(); const log = opts.onLog ?? ((m: string) => console.log(m)); const wanted = opts.operationIds ?? []; // ONE RULE for what a lane may dispatch, shared with the runner's projection // (autoRunner.buildChannelWork) and the editor's operation select. Unknown // ids are dropped rather than throwing: a settings file may name an operation // from a newer build, and a stale scope must not wedge the lane. const all = operationsForLane(opts.lane, settings); const operations = wanted.length === 0 ? all : all.filter((op) => new Set(wanted).has(op.id)); const run: OperationRun = { lane: opts.lane, paths: opts.paths, settings, operations, force: opts.force === true, allowRedownload: opts.allowRedownload ?? settings.backfill.allowRedownload, log, live: { llmActive: 0, unitActive: 0, costUsd: 0, diskFloorHit: false }, digest: null, backfill: null, }; if (opts.lane === "digest") { const laneChoice = opts.digestLane ?? "local"; const chunking = resolveDigestChunking({ lane: laneChoice }); // WHETHER A LANE CAN SHARE A DUPLICATE'S OUTPUT IS DECLARED, not assumed // from the fact that this is the digest branch. A chapter list is about what // was SAID, so an aligned mirror gets the canonical member's digest for free // (~11% of the sweep); diarization and attribution are grounded in one audio // track and one set of cue timings and can share nothing. Asking the rule is // what lets a dispatcher that does not know which operation it is holding // get this right. An explicit `useClusters: false` still wins — this is the // default, not a lock. const sharesDuplicates = laneSharesDuplicates( digestLaneFor(chunking.app.lane), ); const clusterPlan = opts.useClusters === false || !sharesDuplicates ? null : (opts.clusterPlan ?? (await buildDigestClusterPlan(opts.paths))); run.digest = { laneChoice, sections: opts.sections, engine: { ...chunking, modelRequested: chunking.config.model?.trim() || chunking.app.defaultModel(), }, clusterPlan, byChannel: new Map(), }; return run; } // LLM fan-out, resolved ONCE per run: the primary's engine config and model // for the attribution operations. Only `baseUrl` will vary per leased endpoint // — baseUrl is not part of the freshness identity, so endpoint choice causes // zero churn — and only an HTTP engine with a baseUrl field (ollama) has an // endpoint to swap at all. const llmOps = operations .map((op) => op.id) .filter((id) => LLM_BACKFILL_OPS.has(id)); const attrResolved = llmOps.length > 0 ? resolveAttributionTarget("text-only", settings.attribution) : null; const llmFanout = attrResolved && attrResolved.app.fields.baseUrl ? { config: attrResolved.config, modelRequested: attrResolved.modelRequested, } : null; // The remote term is sound only when EVERY operation in this run can take an // llm lease. runPool's limit is pool-wide: on a mixed run, a slot justified by // remote capacity could dispatch a non-LLM operation (diarization) onto this // box's own CPU while the local term says stand aside — breaking the idle-only // default. A pure attribution run gets the full fan-out; a mixed run keeps // today's local limit and still fans out whatever its local slots dispatch. const remoteEligible = llmFanout !== null && llmOps.length === operations.length; // UNIT-EXECUTOR envelope, resolved ONCE per run: appConfig with baseUrl // STRIPPED so the executor localises its own endpoint. The channel CONTEXT is // deliberately not here — it is injected per unit, because a cross-channel // runner has more than one and a scratch corpus has no digest-context.md to // read at all. const unitConfig: StartWorkerUnitInput["config"] = { ...(attrResolved ? { attribution: settings.attribution, appConfig: stripBaseUrl(attrResolved.config), } : {}), ...(operations.some((op) => op.id === "diarization") ? { diarization: settings.diarization } : {}), }; run.backfill = { llmOps, llmFanout, remoteEligible, unitConfig, targets: new Map(), }; return run; } // The executor localises its own endpoint (its OLLAMA_URL); everything else in // the config IS identity and travels verbatim. function stripBaseUrl(config: DigestAppConfig): DigestAppConfig { const { baseUrl: _baseUrl, ...rest } = config; void _baseUrl; return rest; } function resolveDigestFor( run: OperationRun, channelSlug: string, ): Promise { const digest = run.digest!; let pending = digest.byChannel.get(channelSlug); if (!pending) { pending = resolveDigestTarget({ paths: run.paths, channelSlug, lane: digest.laneChoice, ...(digest.sections ? { sections: digest.sections } : {}), }); digest.byChannel.set(channelSlug, pending); } return pending; } function resolveTargetFor( run: OperationRun, op: Operation, channelSlug: string, ): Promise { const backfill = run.backfill!; const key = `${op.id}\u0000${channelSlug}`; const hit = backfill.targets.get(key); if (hit) return hit; const pending = Promise.resolve( op.resolveTarget({ settings: run.settings, paths: run.paths, channelSlug, }), ); backfill.targets.set(key, pending); return pending; } // The lane's live snapshot, for laneLimit. Built here so a caller never has to // remember which fields the limit reads — and SYNCHRONOUS, which runPool's // limit() requires: everything it needs was resolved when the run was opened. export function operationLaneLive( run: OperationRun, // A caller's own ceiling on local slots. Absent → the lane's default. concurrency?: number, ): LaneLive { if (run.digest) { const { app } = run.digest.engine; // Local lane: 1. The GPU is the bottleneck and a second concurrent // generation just thrashes the same 8 GB of VRAM. The metered lane is // network-bound, so it can overlap — but modestly, since it is paying per // call. const defaultConcurrency = app.lane === "local-gpu" ? 1 : 2; return { lane: "digest", appLane: app.lane, metered: app.metered, costUsd: run.live.costUsd, concurrency: Math.max(1, concurrency ?? defaultConcurrency), fanOutEligible: app.fields.baseUrl === true, llmActive: run.live.llmActive, }; } return { lane: "backfill", operations: run.operations, llmOps: run.backfill?.llmOps ?? [], remoteEligible: run.backfill?.remoteEligible ?? false, llmActive: run.live.llmActive, unitActive: run.live.unitActive, }; } // --------------------------------------------------------------------------- // One unit: classify, then run // --------------------------------------------------------------------------- export type OperationUnit = { op: Operation; channelSlug: string; videoId: string; }; export type OperationUnitPlan = { action: CandidateAction; // The identity this classification was made against, so the run cannot // resolve a different one. target: unknown; // Set when the video is a duplicate-cluster MIRROR whose canonical member // owns this operation's output — a `skip` with a reason worth logging. mirrorOf?: string; }; // IS THIS VIDEO THIS OPERATION'S WORK RIGHT NOW? Re-derived from disk, every // call: a restart, a concurrent lane, or a share that landed while this run was // held are all just visible here. export async function classifyOperationUnit( run: OperationRun, unit: OperationUnit, ): Promise { const videoDir = path.join( run.paths.channelsDir, unit.channelSlug, "data", unit.videoId, ); let target: unknown; if (run.digest) { const resolved = await resolveDigestFor(run, unit.channelSlug); target = { target: resolved.target, sections: resolved.sections }; // A cluster mirror is not this lane's work: its canonical member owns the // generation and shares the result here. const role = run.digest.clusterPlan?.bySlug.get( planSlugForDir(run.digest.clusterPlan, unit.channelSlug, unit.videoId), ); if (role?.kind === "mirror") { return { action: "skip", target, mirrorOf: role.canonicalSlug }; } } else { target = await resolveTargetFor(run, unit.op, unit.channelSlug); } const files = await readVideoFiles(videoDir, { checkUntranscribable: true }); const state = await unit.op.state({ videoDir, videoId: unit.videoId, files, target, settings: run.settings, }); return { action: candidateAction(state, { force: run.force, // MAKE THE FLAG TRUE. `diskFloorHit` used to be set and then ignored — // the batch kept pulling candidates and kept asking reacquireMediaFor to // refuse them, one statfs and one media re-check per video, for the rest // of a run over tens of thousands of videos. // // Once the floor is hit, re-acquisition is off for the remainder of the // run: exactly the behaviour allowRedownload:false already describes, so // those videos land in `missingInput` and stay visible as work the corpus // still owes. Note what this does NOT stop — a video whose media is // already on disk still gets diarized/attributed, because those write // kilobyte sidecars and holding them frees nothing while losing days. allowRedownload: run.allowRedownload && !run.live.diskFloorHit, }), target, }; } // What one unit did, in the vocabulary both callers fold. export type OperationUnitReport = { outcome: | "done" // Already at the current identity — classified so, or reported so by the // operation itself when the pull and the run disagreed. | "fresh" // Not this operation's work (not-applicable, or a mirror its canonical // owns), or work the operation declined for a reason the cheap disk // classification could not see — today a transcript mid-rewrite. Kept OUT // of `failed`: it resolves itself once normalize catches up, and counting // it as a failure would make a healthy lane look like a broken engine. | "skipped" | "missing-input" | "deferred" | "blocked" | "failed"; // Digest's metered accounting, which OperationRunOutcome cannot carry. engineCalls: number; costUsd: number; warnings: number; // Duplicate-cluster sharing (digest only). shared: number; misaligned: number; // Media re-acquisition (backfill only). reacquired: number; reacquireCleaned: number; reacquireHandedOff: number; reacquireFailed: number; diskFloorHit: boolean; }; function emptyReport( outcome: OperationUnitReport["outcome"], ): OperationUnitReport { return { outcome, engineCalls: 0, costUsd: 0, warnings: 0, shared: 0, misaligned: 0, reacquired: 0, reacquireCleaned: 0, reacquireHandedOff: 0, reacquireFailed: 0, diskFloorHit: false, }; } export type RunOperationUnitOptions = { onLog?: (msg: string) => void; signal?: AbortSignal; // The run's own hard-cancel signal, distinct from the per-item one: an abort // on either must propagate rather than being counted as a failure. runSignal?: AbortSignal; // NOTE: there is deliberately NO "plan already made by the caller's pull" // option. Both callers classify before dispatching — the batch to keep a // non-dispatch off a pool slot, the runner to choose WHICH operation a video // needs — and this re-asks anyway: the pull and the run are separated by // however long the pool held us, and re-deriving from disk on every pull is // the invariant, not an optimisation to be handed a cached answer. }; // RUN ONE UNIT OF ONE OPERATION ON ONE VIDEO — the whole of it, whichever lane // and whichever operation. // // Classify (from disk, now), then either count the classification or run the // work with this lane's per-video overrides: media re-acquisition and its // always-cleanup for the backfill lane, duplicate-cluster sharing for digest, // the llm fan-out and the tagged-remote unit path for both. export async function runOperationUnit( run: OperationRun, unit: OperationUnit, opts: RunOperationUnitOptions = {}, ): Promise { const log = opts.onLog ?? run.log; const signal = opts.signal ?? new AbortController().signal; const aborted = () => signal.aborted || opts.runSignal?.aborted === true; const videoDir = path.join( run.paths.channelsDir, unit.channelSlug, "data", unit.videoId, ); // RE-DERIVED FROM DISK, at dispatch. Never trusted from the pull. const plan = await classifyOperationUnit(run, unit); if (plan.action !== "dispatch") { if (plan.mirrorOf) { log( `Skipping ${unit.videoId}: duplicate of ${plan.mirrorOf}, which owns the digest for this cluster.`, ); } return emptyReport(plan.action === "skip" ? "skipped" : plan.action); } const report = emptyReport("failed"); // Re-acquisition state lives OUT here so the finally below can always see it, // whatever the run threw. let reacquired: ReacquireOutcome | null = null; try { // Only pay for a re-acquire when the input really is gone. state() said so // at classification time, but that and the run are separated by however // long the pool held us, so this is re-checked inside reacquireMediaFor. // Digest never re-acquires: its input is the transcript, which is on disk // or is somebody else's problem. if (run.backfill && run.allowRedownload && !run.live.diskFloorHit) { reacquired = await reacquireMediaFor({ paths: run.paths, channelSlug: unit.channelSlug, videoId: unit.videoId, videoDir, onLog: log, signal, }); if (reacquired.status === "disk-floor") { // Latches for the rest of the run — classifyOperationUnit reads this // and stops offering videos that would need a fetch. if (!run.live.diskFloorHit) { log( `Disk floor reached — no more media will be re-acquired this run. ` + `Videos needing it are counted as missing-input.`, ); } run.live.diskFloorHit = true; report.diskFloorHit = true; report.outcome = "missing-input"; return report; } if (reacquired.status === "failed" || reacquired.status === "gone") { report.reacquireFailed = 1; report.outcome = "missing-input"; return report; } if (reacquired.status === "fetched") report.reacquired = 1; } const outcome = await executeUnit(run, unit, plan, videoDir, { log, signal, report, }); if (outcome === "done") report.outcome = "done"; else if (outcome === "already-present") report.outcome = "fresh"; else if (outcome === "missing-input") report.outcome = "missing-input"; else if (outcome === "skipped") report.outcome = "skipped"; else { report.outcome = "failed"; log(`Failed ${unit.op.id} ${unit.videoId}: ${outcome}.`); } } catch (err) { if (aborted()) throw err; report.outcome = "failed"; log( `Failed ${unit.op.id} ${unit.videoId}: ${(err as Error)?.message ?? String(err)}`, ); } finally { // THE FINALLY THAT KEEPS THE DISK ALIVE. A re-fetched file is removed // whether the operation succeeded, failed, or threw — on a 97%-full disk a // leak here fills it. Two exceptions, both reported by cleanup(): a video // marked do-not-clean, and the hand-off to auto-transcribe (see // backfillReacquire's header). A FAILED re-acquire is cleaned too — a // download that threw can still have left a partial file, which is exactly // the leak this exists to stop. if (reacquired?.status === "fetched" || reacquired?.status === "failed") { const out = await reacquired.cleanup(); // Only a "fetched" contributes to the reconciliation: the other statuses // were never counted in `reacquired`. if (reacquired.status === "fetched") { if (out.status === "removed") report.reacquireCleaned = 1; else if (out.status === "kept" && out.reason !== "do-not-clean") { report.reacquireHandedOff = 1; } } } } return report; } // The work itself: a tagged remote executor when one is free, else the llm // fan-out for the call-bound operations, else the local run. // // All claims are non-parking on purpose: a parked acquire inside a runPool slot // would deadlock the run. async function executeUnit( run: OperationRun, unit: OperationUnit, plan: OperationUnitPlan, videoDir: string, ctx: { log: (m: string) => void; signal: AbortSignal; report: OperationUnitReport; }, ): Promise { if (run.backfill) { const viaUnit = await runViaTaggedRemote(run, unit, plan, videoDir, ctx); if (viaUnit !== null) return viaUnit; } if (run.digest) return runDigestUnit(run, unit, plan, videoDir, ctx); const fanout = run.backfill?.llmFanout ?? null; const llm = fanout && LLM_BACKFILL_OPS.has(unit.op.id) ? await acquireLlmSlot(unit.op.id, fanout.modelRequested, ctx.log) : null; if (llm) run.live.llmActive++; try { const outcome = await unit.op.run({ paths: run.paths, videoDir, videoId: unit.videoId, channelSlug: unit.channelSlug, target: plan.target, force: run.force, ...(llm ? { appConfig: { ...fanout!.config, baseUrl: llm.baseUrl } } : {}), onLog: ctx.log, signal: ctx.signal, }); if (llm) getWorkerPool().markSuccess(llm.workerId); return outcome; } catch (err) { if (llm && !ctx.signal.aborted) getWorkerPool().markFailure(llm.workerId); throw err; } finally { if (llm) { run.live.llmActive--; llm.lease.release(); } } } // THE DIGEST UNIT IS NOT `op.run()`, and that is a debt this slice names rather // than hides. // // `Operation.run()` returns an OperationRunOutcome — five strings — and digest's // registry entry throws away the three numbers the metered lane cannot work // without: engineCalls, costUsd and warningCount. Routing digest through // `op.run()` here would silently disable the SPEND CAP (laneLimit reads // `live.costUsd`, which would then never move off zero) and empty the metered // accounting line. So the digest branch calls digestVideo directly, in ONE // place, with the identity this run resolved — and the fix, when someone wants // it, is to widen OperationRunOutcome to carry a cost, not to add a second // executor beside this one. async function runDigestUnit( run: OperationRun, unit: OperationUnit, plan: OperationUnitPlan, _videoDir: string, ctx: { log: (m: string) => void; signal: AbortSignal; report: OperationUnitReport; }, ): Promise { const digest = run.digest!; const resolved = await resolveDigestFor(run, unit.channelSlug); const digestSettings = run.settings.digest; // A lease on a verified endpoint, or null → the local/default endpoint // exactly as today. ONLY baseUrl varies: model, numCtx and the rest of the // config travel verbatim, so which endpoint served a call is invisible to the // freshness identity. const fanOutEligible = resolved.app.fields.baseUrl === true; const llm = fanOutEligible ? await acquireLlmSlot( DIGEST_OPERATION_ID, resolved.modelRequested, ctx.log, ) : null; if (llm) run.live.llmActive++; try { const outcome = await digestVideo({ paths: run.paths, channelSlug: unit.channelSlug, videoId: unit.videoId, sections: resolved.sections, appId: resolved.app.id, config: llm ? { ...resolved.config, baseUrl: llm.baseUrl } : resolved.config, context: resolved.context, timestampMode: resolved.timestampMode, ...(digestSettings.promptVariant ? { promptVariant: digestSettings.promptVariant } : {}), force: run.force, onLog: ctx.log, signal: ctx.signal, }); // The call round-tripped — the endpoint is alive. Chunk-level failures are // digestVideo's per-chunk isolation and log loudly on their own. if (llm) getWorkerPool().markSuccess(llm.workerId); if (outcome.status === "fresh") return "already-present"; if (outcome.status === "skipped") { ctx.log(`Skipped ${unit.videoId}: ${outcome.reason}.`); return "skipped"; } ctx.report.engineCalls = outcome.engineCalls; ctx.report.costUsd = outcome.costUsd; ctx.report.warnings = outcome.warningCount; run.live.costUsd += outcome.costUsd; // Share to this cluster's aligned mirrors, right after the canonical // member's digest lands — so a mirror never sits un-digested waiting for a // second pass, and a crash mid-run leaves a consistent cluster. const slug = planSlugForDir( digest.clusterPlan, unit.channelSlug, unit.videoId, ); const role = digest.clusterPlan?.bySlug.get(slug); if (role?.kind === "canonical") { const outcomes = await shareDigestToCluster({ paths: run.paths, clusterId: role.clusterId, canonicalSlug: slug, mirrors: role.mirrors, dirBySlug: digest.clusterPlan?.dirBySlug, onLog: ctx.log, }); for (const o of outcomes) { if (o.status === "shared") ctx.report.shared++; else if (o.status === "misaligned") ctx.report.misaligned++; } } return "done"; } catch (err) { // Counted against the worker too, so a flapping endpoint degrades through // the pool's existing consecutive-failure machinery. if (llm && !ctx.signal.aborted) getWorkerPool().markFailure(llm.workerId); throw err; } finally { if (llm) { run.live.llmActive--; llm.lease.release(); } } } // Ship one candidate to a unit executor; apply its result through the // operation's guarded writers (trap: never a raw file copy — the unit ran // against a snapshot minutes old, and the primary's disk may have moved // meanwhile). Returns null when no tagged remote slot is free or every attempt // hit a transport failure — the caller then runs locally. The retry shape copies // transcribeOne's: transport → degrade-or-mark and try another worker; // work-class → final (the same code would fail the same way anywhere). const MAX_UNIT_ATTEMPTS = 3; async function runViaTaggedRemote( run: OperationRun, unit: OperationUnit, plan: OperationUnitPlan, videoDir: string, ctx: { log: (m: string) => void; signal: AbortSignal }, ): Promise { const pool = getWorkerPool(); for (let attempt = 0; attempt < MAX_UNIT_ATTEMPTS; attempt++) { const lease = pool.tryAcquire(unitRequires(unit.op), { kind: "remote", taggedOnly: true, }); if (!lease) return null; run.live.unitActive++; try { const listing = await readVideoFiles(videoDir, { checkUntranscribable: true, }); const files: Record = {}; for (const name of unit.op.inputs(listing)) { files[name] = await readFile(path.join(videoDir, name)); } // The channel context is injected because a scratch corpus has no // digest-context.md to read. Read per channel, memoized by the same map // the targets use. const envelope = run.backfill?.unitConfig ?? {}; const config = { ...envelope, ...(envelope.attribution ? { context: await readDigestContext(run.paths, unit.channelSlug) } : {}), }; const result = await runUnitViaRemote({ worker: lease.worker, op: unit.op.id, channelSlug: unit.channelSlug, videoId: unit.videoId, files, target: plan.target, config, force: run.force, onLog: ctx.log, signal: ctx.signal, }); pool.markSuccess(lease.worker.id); if (result.outcome === "done") { const applied = await unit.op.applyResult(videoDir, result.files); if (applied === "applied") return "done"; if (applied === "refused") { // A guard said no (the downgrade rule) — the same answer the local // runner reports as outranked/already-present. return "already-present"; } ctx.log( `Remote unit ${unit.op.id} ${unit.videoId}: unusable result payload.`, ); return "failed"; } if ( result.outcome === "already-present" || result.outcome === "missing-input" || result.outcome === "skipped" ) { return result.outcome; } return "failed"; } catch (err) { if (ctx.signal.aborted) throw err; const failureClass = err instanceof TranscribeError ? err.failureClass : "transcription"; if (failureClass === "transport") { // A transport failure must NEVER count as work failure — the item // retries on another worker or locally, and only the WORKER pays. if (!(await pingRemoteHealth(lease.worker))) { pool.markDegraded(lease.worker.id); } else { pool.markFailure(lease.worker.id); } ctx.log( `Remote unit ${unit.op.id} ${unit.videoId} transport failure on ${lease.worker.id}: ${String(err)} — retrying elsewhere.`, ); continue; } return "failed"; } finally { run.live.unitActive--; lease.release(); } } return null; } // --------------------------------------------------------------------------- // Counting a channel's work // --------------------------------------------------------------------------- export type OperationWorkCount = { // What the lane can act on WITHOUT re-acquiring media — missing + stale + // partial. The number every "how much is left?" surface leads with, and the // one a progress bar is sized against. reachable: number; // NEVER added to `reachable`. On the live corpus it is ~91x it, and one // summed "remaining" would be noise. missingInput: number; deferred: number; blocked: number; }; // How much work a channel has for a lane right now, for sizing a progress bar // before starting a job. // // SAME DERIVATION AS THE PULL — the two disagreeing is what makes a bar stall // one short of complete forever. It replaces `countMissingDigests` and // `countBackfillWork`, which were the same fold written twice and had already // drifted: the digest one counted a video with a superseded cues.json as // missing (the pull then dispatched it and the engine skipped it), where the // registry classifies it `deferred` and the pull no longer offers it at all. export async function countOperationWork( lane: OperationLane, paths: Paths, channelSlug: string, opts: { operationIds?: readonly string[]; // Only these videos, intersected with disk — the same scope the run takes. // The progress target must count what THIS run will pull, and an ids-scoped // run pulls only those; without the filter a per-video run would size its // bar to the whole channel and stall one short of complete forever. ids?: readonly string[]; digestLane?: DigestLaneChoice; sections?: DigestSectionKind[]; } = {}, ): Promise { const run = await openOperationRun({ lane, paths, operationIds: opts.operationIds, digestLane: opts.digestLane, sections: opts.sections, // Counting must not build a duplicate plan: it reads the whole duplicates // report and a mirror is `present` or `missing` on its own merits either // way, which is what the snapshot counts too. useClusters: false, onLog: () => {}, }); const empty = { reachable: 0, missingInput: 0, deferred: 0, blocked: 0 }; if (run.operations.length === 0) return empty; const dataDir = path.join(paths.channelsDir, channelSlug, "data"); const allDirs = await readdir(dataDir).catch(() => [] as string[]); // Intersected with disk exactly as the pull does, so an id that names nothing // counts as nothing rather than as one unit of work. const dirs = opts.ids ? opts.ids.filter((id) => allDirs.includes(id)) : allDirs; // Folded through the SHARED counters rather than a private if-chain, so this // path and the snapshot path can no longer disagree about what a state means, // and a state added later is counted here without anyone remembering. // Resolved ONCE per operation, not per video: the identity is a settings read // plus some string work, and deriving it per item is how a counter and a // runner end up disagreeing about what is stale. const targets = new Map(); for (const op of run.operations) { if (run.digest) { const resolved = await resolveDigestFor(run, channelSlug); targets.set(op.id, { target: resolved.target, sections: resolved.sections, }); } else { targets.set(op.id, await resolveTargetFor(run, op, channelSlug)); } } const counts = emptyOperationCounts(); for (const id of dirs) { const videoDir = path.join(dataDir, id); const files = await readVideoFiles(videoDir, { checkUntranscribable: true, }); for (const op of run.operations) { addOperationState( counts, await op.state({ videoDir, videoId: id, files, target: targets.get(op.id), settings: run.settings, }), ); } } return { reachable: reachableOperationWork(counts), missingInput: counts.missingInput, deferred: counts.deferred, blocked: counts.blocked, }; } // --------------------------------------------------------------------------- // The channel-scoped batch (the manual verbs' body) // --------------------------------------------------------------------------- export type OperationBatchOptions = { lane: OperationLane; channelSlug: string; paths: Paths; // Which operations to run. Empty/omitted = every one the lane dispatches. operationIds?: readonly string[]; // Only these video ids (intersected with what's on disk). ids?: readonly string[]; // Upload-date ordering over this channel's candidates. Absent = the lane's // configured order, resolved inside the batch so an armed sweep, a // hand-clicked channel run and an ids-scoped run all inherit the same // setting. "listed" gives a null comparator and no sort at all. order?: AutoQueueOrder; force?: boolean; allowRedownload?: boolean; limitCount?: number; concurrency?: number; setProgress?: (snap: JobProgress) => void; // Where the bar starts: how many videos already counted as done before this // run. The batch reports `initial + completed` against it rather than letting // the UI re-count files from disk, which cannot see a regeneration. progressBaseline?: number; // Where the bar ends. Supplied by the caller because it already resolves it // (countOperationWork) to size the job; the batch raises it if sharing turns // out to satisfy more videos than were counted. progressTarget?: number; onLog?: (msg: string) => void; signal?: AbortSignal; drainSignal?: AbortSignal; tracker?: TaskTracker; // ── digest only ──────────────────────────────────────────────────────── digestLane?: DigestLaneChoice; sections?: DigestSectionKind[]; // Duration ordering, applied BEFORE the date sort. Shortest-first by default: // it converts the backlog into visible coverage fastest, and the long-tail 8% // is where a prompt bug is most expensive to discover late. digestOrder?: DigestOrder; // Duration window, for splitting the corpus between lanes. minDurationSeconds?: number; maxDurationSeconds?: number; useClusters?: boolean; clusterPlan?: DigestClusterPlan; }; // ONE result for both lanes. They used to have one each and each printed a // different subset of the same outcomes — the card's had deferred and blocked, // the sweep's had skipped — so "what did that run actually do" had a different // answer depending on which button started it. export type OperationBatchResult = { attempted: number; succeeded: number; fresh: number; failed: number; // Nothing to do, for a reason the cheap disk classification could not see, or // a mirror whose canonical member owns the output. skipped: number; // Videos whose input was gone and that were not re-acquired. The number that // must stay separate from the rest. missingInput: number; // Refused under the current configuration (the diarization length cap, a // transcript that is not normalized). Neither work done nor work failed. deferred: number; // Waiting on a prerequisite operation's output. Expected to fall by itself as // the prerequisite runs, so an operator seeing it should wait rather than // change a setting. blocked: number; reacquired: number; reacquireCleaned: number; reacquireHandedOff: number; reacquireFailed: number; diskFloorHit: boolean; // Digest: mirrors that received the canonical member's digest, and mirrors // the alignment gate refused. A refusal is a normal outcome, not an error. shared: number; misaligned: number; engineCalls: number; costUsd: number; warnings: number; // True when the run stopped early because the spend cap was reached. spendCapped: boolean; // True when the run ended early because a relocation of this channel's media // started while it was running. Not a failure: the work that is left is still // there, and the next sweep picks it up once the move is done. stoppedForRelocation: boolean; }; export function emptyOperationBatchResult(): OperationBatchResult { return { attempted: 0, succeeded: 0, fresh: 0, failed: 0, skipped: 0, missingInput: 0, deferred: 0, blocked: 0, reacquired: 0, reacquireCleaned: 0, reacquireHandedOff: 0, reacquireFailed: 0, diskFloorHit: false, shared: 0, misaligned: 0, engineCalls: 0, costUsd: 0, warnings: 0, spendCapped: false, stoppedForRelocation: false, }; } type Candidate = { op: Operation; id: string; duration: number | null }; export async function runOperationBatch( opts: OperationBatchOptions, ): Promise { const log = opts.onLog ?? ((m: string) => console.log(m)); const result = emptyOperationBatchResult(); const run = await openOperationRun({ lane: opts.lane, paths: opts.paths, operationIds: opts.operationIds, force: opts.force, allowRedownload: opts.allowRedownload, digestLane: opts.digestLane, sections: opts.sections, useClusters: opts.useClusters, clusterPlan: opts.clusterPlan, onLog: log, }); if (run.operations.length === 0) { log( opts.lane === "digest" ? "The digest operation is not enabled — nothing to do." : "No backfill is enabled (or none of the requested operations is) — nothing to do.", ); return result; } // The engine fail-fast, before any pool exists. await preflightOperationRun(run); // ONE derivation of this channel's identity, shared with the classification // and with countOperationWork — see digestTarget.ts. It must match what // digestVideo will actually chunk with, or the freshness check and the writer // would disagree on the identity and every video would look stale forever. const digestResolved = run.digest ? await resolveDigestFor(run, opts.channelSlug) : null; // GUARD 4 OF FOUR (see plans/relocate-channel-media.md). Same swallow as the // snapshot's, deciding the candidate list for both operation lanes. All three // of this function's callers are already behind guard 1 today; this is here so // a fourth in-process caller that is not cannot quietly find "no candidates" // on a channel whose drive is unmounted. // // BY LANE (release 17): the digest lane reads only the text tier, so it asks // the text guard and runs while the channel's media is moving, stalled or // unmounted; the backfill lane's operations open the audio and keep the // media guard. if (opts.lane === "digest") { await assertChannelTextReadable(opts.paths, opts.channelSlug); } else { await assertChannelMediaReachable(opts.paths, opts.channelSlug); } const dataDir = path.join(opts.paths.channelsDir, opts.channelSlug, "data"); const allDirs = await readdir(dataDir).catch(() => [] as string[]); const onDisk = new Set(allDirs); const wantedListed = opts.ids ? opts.ids.filter((id) => onDisk.has(id)) : allDirs; // Ordering, applied ONCE, here, to the frozen list the cursor walks — not // inside next(), which must stay O(1) per pull. Sorting the list is not a // break of the "re-derive eligibility from disk on every pull" invariant: // that is about op.state(), which next() still reads per item; this only // decides what order the cursor reaches them in. const order = opts.order ?? laneOrderFor(run); const byRecency = await batchRecencyComparator( opts.paths, opts.channelSlug, wantedListed, order, ); let candidates: Candidate[] = []; let noTranscript = 0; let outOfWindow = 0; if (run.digest) { // Durations, read once. Videos with no normalized transcript have no // duration and are dropped here — they are a transcription problem, not a // digest one, and the registry classifies them `blocked`/`deferred` anyway. const digestOrder = opts.digestOrder ?? "shortest-first"; for (const id of wantedListed) { opts.signal?.throwIfAborted(); const duration = await videoDurationSeconds( opts.paths, opts.channelSlug, id, ); if (duration === null || duration <= 0) { noTranscript++; continue; } if ( opts.minDurationSeconds !== undefined && duration < opts.minDurationSeconds ) { outOfWindow++; continue; } if ( opts.maxDurationSeconds !== undefined && duration > opts.maxDurationSeconds ) { outOfWindow++; continue; } candidates.push({ op: run.operations[0], id, duration }); } candidates.sort((a, b) => digestOrder === "longest-first" ? b.duration! - a.duration! || a.id.localeCompare(b.id) : a.duration! - b.duration! || a.id.localeCompare(b.id), ); // Then by upload date, if asked. COMPOSED with the duration sort above // rather than replacing it: Array#sort is stable and the recency comparator // returns 0 for two videos sharing a YYYYMMDD key, so this reads as "newest // day first, shortest video within a day" — both rules intact. See // cheapestComparator, which states the same composition for the runner. if (byRecency) candidates.sort((a, b) => byRecency(a.id, b.id)); } else { const ids = byRecency ? [...wantedListed].sort(byRecency) : wantedListed; // OPERATION-MAJOR, then video: a prerequisite is attempted before anything // that declares it, all the way through the video list, so a video diarized // this pass can be attributed in the same one. candidates = run.operations.flatMap((op) => ids.map((id) => ({ op, id, duration: null })), ); } log( batchHeadline(run, opts, { resolved: digestResolved, wantedListed, candidates, order, byRecency, noTranscript, outOfWindow, }), ); // Ids already handed out this run. With the cursor below this keeps the disk // re-derivation O(n) overall rather than O(n²): the cursor only ever moves // forward past ids that have been attempted, while eligibility for the one it // stops on is always re-read from disk. const attempted = new Set(); const key = (c: Candidate) => `${c.op.id}\u0000${c.id}`; let cursor = 0; // Audio-seconds retired from the worklist, by ANY route — generated, shared // to, found fresh, or skipped as a mirror. All four remove work, and an ETA // that only counted generations would keep quoting time for videos that are // already done. const totalAudioSeconds = candidates.reduce((n, c) => n + (c.duration ?? 0), 0); let resolvedAudioSeconds = 0; const next = async (): Promise => { while (cursor < candidates.length) { if (opts.limitCount !== undefined && result.succeeded >= opts.limitCount) { return null; } const candidate = candidates[cursor]; if (attempted.has(key(candidate))) { cursor++; continue; } // GUARD 5: A RELOCATION THAT STARTED WHILE THIS BATCH WAS RUNNING. // // The reachability guard above runs ONCE, at start. The omnimirror // incident (2026-09-13) is what that misses: a relocate copy began after // this batch did, a unit wrote a sidecar into `data/` while rsync was // already past that directory, and the copy's verify then refused. A // marker is the one signal that says "these bytes are being copied right // now" — every other guard in the system already treats a channel // carrying one as in-transition, and this is the same rule applied per // pull rather than per job. // // Returning null ENDS the run cleanly (see mechanic 2 in the header: zero // is a hold, null is a stop). That is deliberate — the batch is // per-channel, so nothing further in this run could be dispatched anyway, // and the next sweep re-derives everything from disk once the move is // done. const marker = await readRelocationMarker(opts.paths, opts.channelSlug); if (marker) { log( `Stopping: a relocation (${marker.direction}) of ${opts.channelSlug} ` + `to ${marker.target} is in flight — the rest of this batch is left ` + `for the next pass.`, ); result.stoppedForRelocation = true; return null; } // RE-DERIVED FROM DISK, every pull. A restart, a concurrent lane, or a // share that landed while this job ran are all visible here. The run // re-asks the same question; this pull is what lets a non-dispatch cost a // classification rather than a pool slot. const plan = await classifyOperationUnit(run, { op: candidate.op, channelSlug: opts.channelSlug, videoId: candidate.id, }); attempted.add(key(candidate)); cursor++; if (plan.action !== "dispatch") { foldNonDispatch(result, plan, candidate, log); resolvedAudioSeconds += candidate.duration ?? 0; reportProgress(); continue; } return candidate; } return null; }; // WIRING THE PLUMBING THAT WAS DECLARED AND NEVER USED. Progress appeared to // work only because buildActiveJobs re-counts `current` from disk — and that // re-count is exactly what a REGENERATION defeats: a regenerated digest is // rewritten in place, the file count never moves, and the bar sits at 0% for // the whole job. The batch is the only thing that knows the truth: it knows a // regenerate happened, and it knows sharing moved the denominator. const baseline = opts.progressBaseline ?? 0; const progressTarget = opts.progressTarget ?? null; function reportProgress(): void { if (!opts.setProgress) return; // Counted against the TARGET, which is the count of videos that were not // fresh when the job was sized. `fresh` is deliberately EXCLUDED: a video // already fresh was never in the target, and counting it would drive the // bar past 100% on any channel that is mostly done. const done = run.digest ? result.succeeded + result.shared + result.skipped : result.succeeded + result.failed; const current = baseline + done; opts.setProgress({ metric: run.digest ? "digests" : "backfills", initial: baseline, // Never let the bar exceed its target: sharing can satisfy more videos // than the target was sized for, and a bar past 100% reads as a bug // rather than as good news. target: Math.max(progressTarget ?? current, current), current, ...(run.digest ? { remainingAudioSeconds: Math.max( 0, totalAudioSeconds - resolvedAudioSeconds, ), } : {}), }); } const runOne = async ( candidate: Candidate, runSignal: AbortSignal, ): Promise => { const task = opts.tracker?.start({ id: candidate.id, label: run.digest ? `digest ${opts.channelSlug}/${candidate.id}` : `backfill ${candidate.op.id} ${opts.channelSlug}/${candidate.id}`, kind: run.digest ? "digest" : "backfill", // What makes seconds-per-audio-hour computable. A digest's cost is // proportional to the transcript's LENGTH, not to it being one video, so // a task-count average is the wrong denominator for a run ETA. ...(candidate.duration !== null ? { audioSeconds: candidate.duration } : {}), }); const itemLog = task ? task.onLog : log; try { const report = await runOperationUnit( run, { op: candidate.op, channelSlug: opts.channelSlug, videoId: candidate.id, }, { onLog: itemLog, signal: runSignal, runSignal: opts.signal }, ); foldUnitReport(result, report); reportProgress(); } finally { resolvedAudioSeconds += candidate.duration ?? 0; task?.end(); } }; // Edge-triggered, so a hold logs twice per contention window rather than once // per poll over a multi-week run. It tracks the REASON rather than a boolean: // a lane that goes from yielding to paused has changed state and should say // so, where a boolean would stay `true` and stay silent. let heldReason: string | null = null; const laneWord = run.digest ? "digest" : "backfill"; const NEVER = new AbortController().signal; // Seed the bar before the first video finishes, so a job that spends its // first minutes on a 3-hour VOD does not look like it never started. reportProgress(); await runPool({ next, run: runOne, limit: () => { const live = operationLaneLive(run, opts.concurrency); const verdict = laneLimit(getSettings(), live); if (verdict.hold) { if (heldReason !== verdict.hold.reason) { heldReason = verdict.hold.reason; log( verdict.limit > 0 ? `${verdict.hold.message} Remote LLM endpoint(s) keep the lane moving meanwhile.` : verdict.hold.message, ); } // The one piece of state a caller keeps: the result flag surfaces a // metered run that stopped early for money rather than for work. if (verdict.hold.reason === "spend-cap") result.spendCapped = true; return verdict.limit; } if (heldReason !== null) { // NAMES WHAT ENDED, not just that something did. "Transcription // finished" after an operator lifted a pause would send them to look at // the transcription queue for something that never happened, which is // the same class of mistake the hold messages exist to prevent. const resumed = heldReason === "yield" ? `Transcription finished; resuming the ${laneWord} lane.` : `Resuming the ${laneWord} lane.`; heldReason = null; log(resumed); } return verdict.limit; }, signal: opts.signal ?? NEVER, drainSignal: opts.drainSignal ?? NEVER, finite: true, idlePollMs: 3000, }); if (run.digest) { // Metered accounting is logged unconditionally when the lane is metered, // even at zero calls: "this run cost nothing" is information too. if (run.digest.engine.app.metered) { const cap = run.settings.digest.spendCapUsd; log( `Metered lane: ${result.engineCalls} model call(s), $${result.costUsd.toFixed(4)} total` + (cap > 0 ? ` (cap $${cap.toFixed(2)})` : " (no cap set)"), ); } } // A gap between fetched and cleaned that is not explained by the hand-off or // by do-not-clean is a LEAK, and it is stated in the job log rather than left // to be discovered by a full disk. if (result.reacquired > result.reacquireCleaned) { log( `Re-acquired ${result.reacquired} file(s), removed ${result.reacquireCleaned}, ` + `handed ${result.reacquireHandedOff} to auto-transcribe. ` + `Any other difference is media kept because its video is marked "do not clean" — ` + `if that is not what you expect, check the disk.`, ); } return result; } // The lane's configured recency order, resolved INSIDE the batch so a // hand-clicked channel run, an ids-scoped run and the lane runner all inherit // the same setting instead of three call sites each remembering to pass it. // // IT COMES OFF THE LANE'S POLICY NOW — `settings.digest.recencyOrder` and // `settings.backfill.order` retired with the sweeps in slice 1.3 — so the // per-channel button and the lane runner read one field. The digest lane's own // default is `cheapest`, which is a composition (shortest inside a day, newest // day first) that this batch expresses as its own `digestOrder` plus a date // sort; the date half of it is newest-first, exactly as the retired // `recencyOrder` was on the live corpus. function laneOrderFor(run: OperationRun): AutoQueueOrder { const order = run.settings.autoQueue[run.digest ? "digest" : "backfill"].order; if (order === "cheapest") return run.digest ? "newest" : "listed"; return order ?? "listed"; } function batchHeadline( run: OperationRun, opts: OperationBatchOptions, facts: { resolved: ResolvedDigestTarget | null; wantedListed: readonly string[]; candidates: readonly Candidate[]; order: AutoQueueOrder; byRecency: unknown; noTranscript: number; outOfWindow: number; }, ): string { if (run.digest) { const resolved = facts.resolved; const digestOrder = opts.digestOrder ?? "shortest-first"; return ( `Digest ${opts.channelSlug} (${run.digest.laneChoice} lane, ` + `${resolved?.app.id}/${resolved?.modelRequested}, sections: ${resolved?.sections.join(", ")}): ` + `${facts.candidates.length} candidate(s) of ${facts.wantedListed.length} on disk ` + `(${facts.noTranscript} without a transcript, ${facts.outOfWindow} outside the duration window), ${digestOrder}` + (facts.byRecency ? `, ${facts.order} first.` : ".") + (run.digest.clusterPlan ? ` Duplicate plan: ${run.digest.clusterPlan.clusters} cluster(s), ${run.digest.clusterPlan.bySlug.size} member(s) mapped.` : " Duplicate sharing off.") ); } return ( `Backfill ${opts.channelSlug}: ${run.operations.map((k) => k.id).join(", ")} over ` + `${facts.wantedListed.length} video dir(s)` + (facts.byRecency ? `, ${facts.order} first` : "") + (run.allowRedownload ? ", re-acquiring media where it is gone." : ", retained media only (re-download is off).") ); } function foldNonDispatch( result: OperationBatchResult, plan: OperationUnitPlan, candidate: Candidate, log: (m: string) => void, ): void { switch (plan.action) { case "skip": if (plan.mirrorOf) { result.skipped++; log( `Skipping ${candidate.id}: duplicate of ${plan.mirrorOf}, which owns the digest for this cluster.`, ); } // A `not-applicable` video is counted as nothing at all — it is not this // operation's work and never was. return; case "fresh": result.fresh++; return; case "deferred": result.deferred++; return; case "blocked": result.blocked++; return; case "missing-input": // Counted, not attempted. This is the population the whole // reachable-vs-needs-re-acquiring split exists to keep visible. result.missingInput++; return; case "dispatch": return; } } function foldUnitReport( result: OperationBatchResult, report: OperationUnitReport, ): void { result.engineCalls += report.engineCalls; result.costUsd += report.costUsd; result.warnings += report.warnings; result.shared += report.shared; result.misaligned += report.misaligned; result.reacquired += report.reacquired; result.reacquireCleaned += report.reacquireCleaned; result.reacquireHandedOff += report.reacquireHandedOff; result.reacquireFailed += report.reacquireFailed; if (report.diskFloorHit) result.diskFloorHit = true; switch (report.outcome) { case "done": result.attempted++; result.succeeded++; return; case "fresh": result.fresh++; return; case "skipped": result.skipped++; return; case "missing-input": result.missingInput++; return; case "deferred": result.deferred++; return; case "blocked": result.blocked++; return; case "failed": result.attempted++; result.failed++; return; } }