commit d6f60ed9e7e5ac26078880c42170940e9549090d
parent 41a5d202d5a2b12b91ccfe86c6c9cff3b0b5db1a
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 7 Sep 2026 20:07:24 -0400
common: one executor for the operation lanes, and the yield stops keying on a field existing
Slice 1.2 of plans/one-core-phase-1.md, first commit. `digestBatch.ts` and
`backfillBatch.ts` are DELETED and `controller/operationBatch.ts` is what they
were: one classification, one unit executor, one per-lane limit, one counter.
They were the same runner written twice, and the second one's header said so —
"Structurally this is digestBatch.ts, and the three mechanics it copies are
copied for the same reasons". Copies drift, and these had. One classified a
candidate through the operation registry and the other re-implemented freshness
inline, so a digest whose cues.json was superseded was dispatched, handed to the
engine and skipped there, where the registry had already called it `deferred`.
One could count `deferred` and `blocked` and the other had no field for either,
so a channel of prerequisite-waiting videos read as finished from one button and
correctly from another. And one keyed its GPU carve-out off `laneFor` EXISTING
while the other read the declared `contendsFor`.
That last one is the fix worth naming. `k.laneFor &&` was a proxy for "this
kind's resource depends on settings, so nothing else is deciding it for us" —
which made `laneFor` load-bearing by its mere presence, so "give every operation
a laneFor defaulting to lane" was a behaviour change wearing a refactor's
clothes: it would have enrolled all of them and parked the whole lane whenever a
statically GPU-bound operation was in the run. The limit now asks
`laneYieldsToTranscription(laneForOperation(op.id) ?? op.lane)` — the
declaration itself. Identical on today's registry (no BACKFILL_QUEUE operation
declares a static GPU lane), and a fixed-GPU operation added tomorrow is caught
rather than missed. The comment on `Operation.laneFor` and the FACTS.md trap
entry now say the presence is not load-bearing, and a test pins the new rule.
Deliberately NOT collapsed: the digest unit calls `digestVideo` rather than
`op.run()`. `OperationRunOutcome` is five strings and digest's registry entry
throws away engineCalls, costUsd and warningCount — routing it through `run()`
would silently disable the SPEND CAP, since the limit reads a per-run cost that
would never leave zero. One place, named, with the fix (widen the outcome) in a
comment rather than a second executor beside it.
`countMissingDigests` and `countBackfillWork` become `countOperationWork(lane,
…)`. `runDigestChannelJob`, `runBackfillChannelJob`, the video-page digest
action and the backfill sweep's walk fallback all drive the one batch. The
`blocked is never dispatched` test moves out of `lib/operations.test.ts` and off
architecture.test.ts's allow-list: the classification is the registry's
business, what a dispatcher does with it is the dispatcher's.
Numbers: none move. The before/after diff of plans/tools/phase1-numbers.ts over
the live corpus is empty. common 896 tests (890 before, +6 laneLimit), mcp 205,
tsc clean in six packages.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
24 files changed, 2534 insertions(+), 1927 deletions(-)
diff --git a/common/architecture.test.ts b/common/architecture.test.ts
@@ -68,8 +68,6 @@ const ALLOWED: Record<string, string> = {
"registry state() freshness check; inverts with the run hook (phase 1)",
"lib/operations.test.ts -> controller/attributeOne":
"test of the above; moves with it",
- "lib/operations.test.ts -> controller/backfillBatch":
- "test of the above; moves with it",
"lib/operations.test.ts -> controller/normalizeTranscript":
"test of the above; moves with it",
diff --git a/common/bin/diarize-backfill.ts b/common/bin/diarize-backfill.ts
@@ -2,7 +2,7 @@
// Capture speaker diarization for videos whose audio is still on disk, headless.
//
// WHY THIS EXISTS. `runBackfillBatch` already accepts an `ids` list
-// (controller/backfillBatch.ts), but nothing passes one — the only caller,
+// (controller/operationBatch.ts), but nothing passes one — the only caller,
// editor/app/channels/[slug]/backfillActions.ts, sends `kindIds` and nothing
// else. So there is no way to say "diarize exactly these videos" without a
// headless entrypoint, and this is the repo's established shape for corpus work
diff --git a/common/controller/backfillBatch.test.ts b/common/controller/backfillBatch.test.ts
@@ -1,231 +0,0 @@
-import { test } from "node:test";
-import assert from "node:assert/strict";
-import { backfillLimit, candidateAction } from "./backfillBatch";
-import {
- defaultBackfill,
- defaultDiarization,
- sanitizeBackfill,
- sanitizeDiarization,
-} from "../lib/settings";
-import type { OperationClassification } from "../lib/operations";
-
-// Run with:
-// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/backfillBatch.test.ts
-//
-// The PURE function, not a live pool — same discipline as digestYield.test.ts,
-// and for the same reason: this decides whether a multi-day lane runs at all,
-// and a test that needs a GPU, a registry and a worker pool is a test nobody
-// runs before shipping.
-
-test("weight 0 is idle-only: full slots when the primary lane is quiet", () => {
- assert.equal(
- backfillLimit({ weight: 0, slots: 4, primaryBusy: false }),
- 4,
- );
- assert.equal(backfillLimit({ weight: 0, slots: 1, primaryBusy: false }), 1);
-});
-
-test("weight 0 stands aside completely while the primary lane works", () => {
- // Zero is a HOLD, not a stop: runPool idle-waits at a zero limit rather than
- // finishing, so the lane resumes the moment transcription is free without
- // re-deriving anything. That is the whole reason this is a limit and not a
- // scheduler.
- assert.equal(backfillLimit({ weight: 0, slots: 4, primaryBusy: true }), 0);
- assert.equal(backfillLimit({ weight: 0, slots: 16, primaryBusy: true }), 0);
-});
-
-test("a positive weight is a guaranteed share, busy or not", () => {
- // The point of a non-zero weight: the primary lane being busy no longer parks
- // the backfill.
- assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: true }), 2);
- assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: false }), 2);
- assert.equal(backfillLimit({ weight: 1, slots: 4, primaryBusy: true }), 4);
-});
-
-test("a small weight is a SLOW lane, never a stopped one", () => {
- // floor(4 * 0.25) is 1, and floor(1 * 0.25) is 0 — which without the floor
- // would silently turn "a quarter of the machine" into "never runs", and look
- // exactly like a wedge.
- assert.equal(backfillLimit({ weight: 0.25, slots: 4, primaryBusy: true }), 1);
- assert.equal(backfillLimit({ weight: 0.25, slots: 1, primaryBusy: true }), 1);
- assert.equal(backfillLimit({ weight: 0.01, slots: 8, primaryBusy: true }), 1);
-});
-
-test("no slots means no work, whatever the weight", () => {
- assert.equal(backfillLimit({ weight: 0, slots: 0, primaryBusy: false }), 0);
- assert.equal(backfillLimit({ weight: 1, slots: 0, primaryBusy: false }), 0);
-});
-
-test("the shipped default is idle-only", () => {
- // Stated as a test because it is a promise the feature makes: catch-up work on
- // a corpus that already exists must never slow down new arrivals.
- const d = defaultBackfill();
- assert.equal(d.weight, 0);
- assert.equal(d.enabled, false);
- assert.equal(d.allowRedownload, false);
- assert.equal(backfillLimit({ ...d, slots: d.concurrency, primaryBusy: true }), 0);
-});
-
-test("a hand-edited weight is clamped rather than rejected", () => {
- // settings.json is hand-editable. A 5 means "as much as possible", and reading
- // it as the idle-only 0 would be the opposite of the intent.
- assert.equal(sanitizeBackfill({ weight: 5 }).weight, 1);
- assert.equal(sanitizeBackfill({ weight: -3 }).weight, 0);
- assert.equal(sanitizeBackfill({ weight: "half" }).weight, 0);
- assert.equal(sanitizeBackfill({ weight: 0.25 }).weight, 0.25);
- // Concurrency shares clampPositiveInt's floor of 1, so a 0 cannot silently
- // park the lane either.
- assert.equal(sanitizeBackfill({ concurrency: 0 }).concurrency, 1);
-});
-
-test("a stale sweep scope survives sanitization as a list of slugs", () => {
- // The scope is persisted WITH the flag, so it has to round-trip: junk entries
- // are dropped, real ones kept.
- assert.deepEqual(
- sanitizeBackfill({ sweepChannels: ["a", "", 7, "b"], sweepKinds: ["diarization"] }),
- {
- ...defaultBackfill(),
- sweepChannels: ["a", "b"],
- sweepKinds: ["diarization"],
- },
- );
-});
-
-// ---------------------------------------------------------------------------
-// The candidate pull's dispatch decision.
-//
-// Untested until the duration cap needed it, and it is the single most
-// consequential branch in the file: it used to fall through to DISPATCHING any
-// state it did not recognize, so "add a state" and "run that state anyway" were
-// the same edit.
-
-const DISPATCH = { force: false, allowRedownload: false };
-
-test("only reachable work is dispatched by default", () => {
- assert.equal(candidateAction("missing", DISPATCH), "dispatch");
- assert.equal(candidateAction("stale", DISPATCH), "dispatch");
-});
-
-test("deferred is NEVER dispatched", () => {
- // The whole point of the cap. Without this the classification would be
- // computed, ignored, and the six-hour video handed to the engine — 40 minutes
- // of CPU and a kernel OOM kill, producing nothing.
- assert.equal(candidateAction("deferred", DISPATCH), "deferred");
- assert.equal(
- candidateAction("deferred", { force: true, allowRedownload: true }),
- "deferred",
- );
-});
-
-test("force redoes present work but does not overrule the cap", () => {
- assert.equal(candidateAction("present", DISPATCH), "fresh");
- assert.equal(
- candidateAction("present", { ...DISPATCH, force: true }),
- "dispatch",
- );
- // `force` means "redo work that looks done", not "ignore the cap".
- assert.equal(
- candidateAction("deferred", { ...DISPATCH, force: true }),
- "deferred",
- );
-});
-
-test("missing-input is counted unless re-download is armed", () => {
- assert.equal(candidateAction("missing-input", DISPATCH), "missing-input");
- assert.equal(
- candidateAction("missing-input", { ...DISPATCH, allowRedownload: true }),
- "dispatch",
- );
-});
-
-test("not-applicable is counted as nothing at all", () => {
- assert.equal(candidateAction("not-applicable", DISPATCH), "skip");
- assert.equal(
- candidateAction("not-applicable", { force: true, allowRedownload: true }),
- "skip",
- );
-});
-
-test("every classification has an explicit decision", () => {
- // The list is written out rather than derived so that adding a state to the
- // union without deciding what the pull does with it fails HERE as well as at
- // the compile step — a runtime backstop for the `never` check, since the
- // hazard this replaces was precisely a silent fall-through.
- const ALL: OperationClassification[] = [
- "present",
- "stale",
- // Part-done. Dispatched like stale — the runner regenerates only the
- // sections that are not fresh — and listed here so that stays a decision
- // rather than something a later edit can quietly drop.
- "partial",
- "missing",
- "missing-input",
- "deferred",
- // This list had already drifted: `blocked` was added to the union and to
- // candidateAction, and never to the backstop meant to catch exactly that.
- "blocked",
- "not-applicable",
- ];
- for (const state of ALL) {
- assert.doesNotThrow(() => candidateAction(state, DISPATCH), state);
- }
- assert.throws(
- () => candidateAction("invented" as OperationClassification, DISPATCH),
- /unhandled backfill state/,
- );
-});
-
-test("the duration cap round-trips, and 0 means off", () => {
- // 0 is a MEANINGFUL value here, not an empty one — it is how the cap is turned
- // off once windowed diarization makes it unnecessary. clampPositiveInt would
- // have floored it to 1, i.e. a one-hour cap, which is why this knob does not
- // use it.
- assert.equal(sanitizeDiarization({ maxAudioHours: 0 }).maxAudioHours, 0);
- assert.equal(sanitizeDiarization({ maxAudioHours: 6.5 }).maxAudioHours, 6.5);
- // Junk falls back to the default rather than to "no cap" — a typo in a
- // hand-edited settings.json must not silently re-arm the OOM.
- assert.equal(
- sanitizeDiarization({ maxAudioHours: "four" }).maxAudioHours,
- defaultDiarization().maxAudioHours,
- );
- assert.equal(
- sanitizeDiarization({ maxAudioHours: -2 }).maxAudioHours,
- defaultDiarization().maxAudioHours,
- );
- // The shipped default is OFF: windowed diarization removed the OOM the cap
- // was a stopgap for. This is asserted rather than assumed because turning it
- // back on by accident would silently stop diarizing long videos.
- assert.equal(defaultDiarization().maxAudioHours, 0);
-});
-
-test("an unknown engine or backend falls back to the default, never to nothing", () => {
- // The DEFAULT ENGINE IS LOAD-BEARING as a fallback, not just as a starting
- // point: it is the one every sidecar already on disk matches, so falling back
- // to it leaves the corpus fresh. Falling back to sortformer on a typo would
- // mark all of it stale and offer weeks of rework.
- assert.equal(defaultDiarization().engine, "sherpa-onnx");
- assert.equal(sanitizeDiarization({ engine: "sortformer" }).engine, "sortformer");
- assert.equal(
- sanitizeDiarization({ engine: "sortfromer" }).engine,
- defaultDiarization().engine,
- );
- assert.equal(sanitizeDiarization({ engine: 7 }).engine, defaultDiarization().engine);
-
- // The backend defaults to the GPU, which is safe only because the lane yields
- // the card rather than sharing it — see diarizationLaneFor.
- assert.equal(defaultDiarization().backend, "vulkan");
- assert.equal(sanitizeDiarization({ backend: "cpu" }).backend, "cpu");
- assert.equal(
- sanitizeDiarization({ backend: "rocm" }).backend,
- defaultDiarization().backend,
- );
-
- // Paths are plain strings and empty means "not configured", which diarizeOne
- // reports as a skip rather than a failure.
- assert.equal(sanitizeDiarization({}).sortformerBin, "");
- assert.equal(sanitizeDiarization({}).sortformerModel, "");
- assert.equal(
- sanitizeDiarization({ sortformerBin: " /opt/diarize-file " }).sortformerBin,
- "/opt/diarize-file",
- );
-});
diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts
@@ -1,937 +0,0 @@
-// Channel-scoped backfill: the runner that does the catch-up work, and the one
-// place the lane's resource SHARE is decided.
-//
-// Structurally this is digestBatch.ts, and the three mechanics it copies are
-// copied for the same reasons — re-derive eligibility from disk on every pull
-// (so a restart, a concurrent lane, or a video that became eligible mid-run are
-// all just visible), pause by returning limit() === 0 (runPool idle-WAITS at a
-// zero limit rather than finishing, so a hold is not a stop), and one job per
-// CHANNEL rather than per video (the registry keeps 100 records and the log 500;
-// 77,000 jobs would evict the history of the run that made them).
-//
-// WHAT IS NEW HERE IS THE SHARE, and it is worth being precise about what could
-// and could not be reused:
-//
-// - The queue cannot portion resources. registry.ts submits every job with
-// concurrency 1, hardcoded, and SchedulerTier only orders work WITHIN one
-// queueKey. `defaultTier` is declared, read once, and set on zero of the ~35
-// kinds — a dead field. So a distinct queueKey buys concurrency and nothing
-// else.
-// - autoQueuePolicy.ts is a real HTB/SWRR implementation with persisted
-// fairness memory, but it arbitrates BETWEEN CHANNELS INSIDE ONE RUNNER.
-// Making it arbitrate between lanes is a much larger change than this
-// feature justifies. If someone later needs true multi-lane fairness, that
-// is where it goes — not here.
-// - runPool's limit() IS a portioning knob: re-read every dispatch, with 0
-// meaning "hold", and digestYield.ts already demonstrates the exact
-// semantics we want on top of it — a yield rather than a lock, failing OPEN
-// so a bad read causes contention rather than deadlock.
-//
-// So: the queue key gives the lane concurrency, and backfillLimit() below gives
-// it a share. Nothing new in the scheduler.
-
-import path from "node:path";
-import { readFile, readdir } from "node:fs/promises";
-import type { Paths } from "../lib/paths";
-import { getSettings } from "../lib/settings";
-import { isGateHeld } from "../lib/pauseGates";
-import { runPool } from "../jobs/concurrentRunner";
-import type { TaskTracker } from "../jobs/taskHooks";
-import type { JobProgress } from "../jobs/registry";
-import { readVideoFiles } from "../lib/videoStatus";
-import {
- addOperationState,
- emptyOperationCounts,
- laneYieldsToTranscription,
- reachableOperationWork,
- resolveBackfillLaneOperations,
- type OperationClassification,
- type Operation,
- type OperationRunOutcome,
-} from "../lib/operations";
-import { transcriptionActivity } from "./digestYield";
-import { acquireLlmSlot, freeLlmSlots } from "./llmWorkers";
-import { getWorkerPool } from "../jobs/workerPool";
-import { resolveAttributionTarget } from "./attributionTarget";
-import { pingRemoteHealth } from "./remoteTranscribe";
-import { runUnitViaRemote } from "./remoteUnit";
-import { TranscribeError } from "./transcribeError";
-import { readDigestContext } from "../lib/digestContext-server";
-import type { DigestAppConfig } from "../lib/digest";
-import type { StartWorkerUnitInput } from "./workerServer";
-import { reacquireMediaFor, type ReacquireOutcome } from "./backfillReacquire";
-import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
-import { batchRecencyComparator } from "./batchRecency";
-
-// 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.
-//
-// weight 0 (the default) — IDLE-ONLY. Full slots when the primary lane is
-// quiet, zero while it works. Catch-up is by definition work on a corpus
-// that already exists, so it must never slow down new arrivals. Zero is a
-// HOLD, not a stop: runPool idle-waits, so the lane resumes the moment the
-// primary is free, having re-derived nothing.
-// weight > 0 — a guaranteed share, floored at 1. The floor is the point: a
-// small weight should mean a slow lane, not a stopped one, and
-// floor(1 * 0.25) is 0.
-export function backfillLimit(opts: {
- weight: number;
- slots: number;
- primaryBusy: boolean;
-}): number {
- const slots = Math.max(0, Math.floor(opts.slots));
- if (slots === 0) return 0;
- if (opts.weight <= 0) return opts.primaryBusy ? 0 : slots;
- const share = Math.floor(slots * Math.min(1, opts.weight));
- return Math.max(1, share);
-}
-
-export type BackfillBatchOptions = {
- channelSlug: string;
- paths: Paths;
- // Which registered kinds to run. Empty/omitted = every enabled lane kind.
- kindIds?: string[];
- // When set, only consider these video ids (intersected with what's on disk).
- ids?: string[];
- // Upload-date ordering over this channel's candidates. Absent = whatever
- // settings.backfill.order says, resolved inside the batch so every caller
- // inherits it. "listed" is today's behaviour and the default.
- order?: AutoQueueOrder;
- // Redo videos that are already present at the current identity.
- force?: boolean;
- // Stop after this many successful videos.
- limitCount?: number;
- // Re-acquire media for `missing-input` videos. Defaults to
- // settings.backfill.allowRedownload; the explicit flag is for a one-off run.
- allowRedownload?: boolean;
- setProgress?: (snap: JobProgress) => void;
- progressBaseline?: number;
- progressTarget?: number;
- onLog?: (msg: string) => void;
- signal?: AbortSignal;
- drainSignal?: AbortSignal;
- tracker?: TaskTracker;
-};
-
-export type BackfillBatchResult = {
- attempted: number;
- succeeded: number;
- // Already at the current identity when the runner got to it — a concurrent
- // lane or a previous pass got there first.
- fresh: number;
- failed: number;
- // Videos whose input was gone and that were not re-acquired (the mode is off,
- // or the disk floor refused). The number that must stay separate from the rest.
- missingInput: number;
- // Nothing to do, for a reason the cheap disk classification could not see —
- // today that is a transcript whose cues.json is stale. Kept OUT of `failed`:
- // it resolves itself once the normalize pass catches up, and counting it as a
- // failure would make a healthy lane look like a broken engine.
- skipped: number;
- // Media re-fetched for a missing-input video, and how many of those files were
- // removed again afterwards. `reacquired - reacquireCleaned` is the count of
- // files deliberately KEPT because the video is marked do-not-clean; any other
- // gap is a leak and is logged as one.
- reacquired: number;
- reacquireCleaned: number;
- // Re-acquired media deliberately LEFT on disk for the transcription lane:
- // either autoQueue.transcription would draw this video from
- // downloadedAutoSubsOnly (the hand-off), or a transcription was already
- // running on it. Both mean "left to the transcription lane"; the per-video log
- // line distinguishes them. Media kept for do-not-clean is NOT counted here —
- // that is still the unexplained-gap case the reconciliation names.
- reacquireHandedOff: number;
- reacquireFailed: number;
- // True when the run stopped taking new work because free disk fell under the
- // configured floor.
- diskFloorHit: boolean;
- // Videos the kind refused to attempt under the current configuration (the
- // diarization duration cap). Kept separate from every other number for the
- // same reason missingInput is: it is neither work done nor work failed, and
- // folding it into either would make a capped run read as a complete one.
- deferred: number;
- // Videos waiting on a prerequisite kind's output. Separate again, and for a
- // reason the others do not share: this number is expected to fall by itself
- // as the prerequisite lane runs, so an operator seeing it should wait rather
- // than change a setting.
- blocked: number;
-};
-
-// 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 kind has no opinion on this video).
- | "skip"
- // Already at the current identity.
- | "fresh"
- // Input is gone and re-acquiring is off.
- | "missing-input"
- // The kind refuses to attempt it under the current configuration.
- | "deferred"
- // Waiting on a prerequisite kind'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 kind this one dependsOn is how you unblock
- // it, and the ordering in resolveBackfillLaneOperations 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)}`);
- }
- }
-}
-
-type Candidate = { id: string; kind: Operation; target: unknown };
-
-// The lane kinds whose work bottoms out in model calls 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"]);
-
-// 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;
-}
-
-export async function runBackfillBatch(
- opts: BackfillBatchOptions,
-): Promise<BackfillBatchResult> {
- const log = opts.onLog ?? ((m: string) => console.log(m));
- const settings = getSettings();
- const kinds = resolveBackfillLaneOperations(settings, opts.kindIds);
-
- const result: BackfillBatchResult = {
- attempted: 0,
- succeeded: 0,
- fresh: 0,
- failed: 0,
- missingInput: 0,
- skipped: 0,
- reacquired: 0,
- reacquireCleaned: 0,
- reacquireHandedOff: 0,
- reacquireFailed: 0,
- diskFloorHit: false,
- deferred: 0,
- blocked: 0,
- };
-
- if (kinds.length === 0) {
- log(
- "No backfill is enabled (or none of the requested kinds is) — nothing to do.",
- );
- return result;
- }
-
- const allowRedownload =
- opts.allowRedownload ?? settings.backfill.allowRedownload;
- 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;
-
- // Upload-date ordering. Applied ONCE, here, to the frozen list the cursor
- // walks — not inside next(), which must stay O(1) per pull. Sorting `wanted`
- // is not a break of the "re-derive eligibility from disk on every pull"
- // invariant: that invariant is about kind.state(), which next() still reads
- // per item; this only decides what order the cursor reaches them in.
- //
- // Resolved inside the batch so the armed sweep, a hand-clicked channel run
- // and an ids-scoped run all inherit settings.backfill.order. "listed" gives a
- // null comparator and no sort at all, reproducing today byte-for-byte.
- const order = opts.order ?? settings.backfill.order;
- const byRecency = await batchRecencyComparator(
- opts.paths,
- opts.channelSlug,
- wantedListed,
- order,
- );
- const wanted = byRecency ? [...wantedListed].sort(byRecency) : wantedListed;
-
- // LLM fan-out, resolved ONCE per run: the primary's engine config and model
- // for the attribution kinds. 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 = kinds.map((k) => k.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 kind 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 kind (diarization) onto this
- // box's own CPU while the local term says stand aside — breaking the
- // idle-only default. A pure attribution run (the scoped sweep) 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 === kinds.length;
- let llmActive = 0;
-
- // UNIT-EXECUTOR fan-out: a TAGGED remote worker matching [kind.id,
- // contendsFor] takes whole units. Tagged-only is the opt-in gate — an
- // untagged remote predates the unit protocol (and may be an older build with
- // no /api/worker/unit), and shipping it units it cannot serve would poison
- // the sweep with refusals. The envelope's identity config is resolved ONCE
- // per run, appConfig with baseUrl STRIPPED so the executor localises its own
- // endpoint, and the channel context injected because a scratch corpus has no
- // digest-context.md to read.
- const unitConfig: StartWorkerUnitInput["config"] = {
- ...(attrResolved
- ? {
- attribution: settings.attribution,
- appConfig: stripBaseUrl(attrResolved.config),
- context: await readDigestContext(opts.paths, opts.channelSlug),
- }
- : {}),
- ...(kinds.some((k) => k.id === "diarization")
- ? { diarization: settings.diarization }
- : {}),
- };
- let unitActive = 0;
- const unitRequires = (kind: Operation): string[] => [
- kind.id,
- (kind.laneFor?.(settings) ?? kind.lane).contendsFor,
- ];
-
- // Ship one candidate to a unit executor; apply its result through the kind'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, exactly as today. 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;
- const runBackfillUnit = async (
- candidate: Candidate,
- videoDir: string,
- itemLog: (m: string) => void,
- runSignal: AbortSignal,
- ): Promise<OperationRunOutcome | null> => {
- const pool = getWorkerPool();
- for (let attempt = 0; attempt < MAX_UNIT_ATTEMPTS; attempt++) {
- const lease = pool.tryAcquire(unitRequires(candidate.kind), {
- kind: "remote",
- taggedOnly: true,
- });
- if (!lease) return null;
- unitActive++;
- try {
- const listing = await readVideoFiles(videoDir, {
- checkUntranscribable: true,
- });
- const files: Record<string, Buffer> = {};
- for (const name of candidate.kind.inputs(listing)) {
- files[name] = await readFile(path.join(videoDir, name));
- }
- const result = await runUnitViaRemote({
- worker: lease.worker,
- op: candidate.kind.id,
- channelSlug: opts.channelSlug,
- videoId: candidate.id,
- files,
- target: candidate.target,
- config: unitConfig,
- force: opts.force,
- onLog: itemLog,
- signal: runSignal,
- });
- pool.markSuccess(lease.worker.id);
- if (result.outcome === "done") {
- const applied = await candidate.kind.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";
- }
- itemLog(
- `Remote unit ${candidate.kind.id} ${candidate.id}: 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 (runSignal.aborted || opts.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);
- }
- itemLog(
- `Remote unit ${candidate.kind.id} ${candidate.id} transport failure on ${lease.worker.id}: ${String(err)} — retrying elsewhere.`,
- );
- continue;
- }
- return "failed";
- } finally {
- unitActive--;
- lease.release();
- }
- }
- return null;
- };
-
- // Resolved once per run, 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<string, unknown>();
- for (const kind of kinds) {
- targets.set(
- kind.id,
- await kind.resolveTarget({
- settings,
- paths: opts.paths,
- channelSlug: opts.channelSlug,
- }),
- );
- }
-
- log(
- `Backfill ${opts.channelSlug}: ${kinds.map((k) => k.id).join(", ")} over ` +
- `${wanted.length} video dir(s)` +
- (byRecency ? `, ${order} first` : "") +
- (allowRedownload
- ? ", re-acquiring media where it is gone."
- : ", retained media only (re-download is off)."),
- );
-
- // Ids already handed out this run, per kind. With the cursor below this keeps
- // the disk re-derivation O(n) overall rather than O(n²): the cursor only moves
- // past items that have been attempted, while eligibility for the item it stops
- // on is always re-read from disk.
- const attempted = new Set<string>();
- const key = (kindId: string, id: string) => `${kindId}\u0000${id}`;
- let cursor = 0;
- let kindIndex = 0;
-
- const next = async (): Promise<Candidate | null> => {
- for (;;) {
- if (opts.limitCount !== undefined && result.succeeded >= opts.limitCount) {
- return null;
- }
- if (kindIndex >= kinds.length) return null;
- if (cursor >= wanted.length) {
- kindIndex++;
- cursor = 0;
- continue;
- }
- const kind = kinds[kindIndex];
- const id = wanted[cursor];
- if (attempted.has(key(kind.id, id))) {
- cursor++;
- continue;
- }
- const videoDir = path.join(dataDir, id);
-
- // RE-DERIVED FROM DISK, every pull.
- const files = await readVideoFiles(videoDir, {
- checkUntranscribable: true,
- });
- const state = await kind.state({
- videoDir,
- videoId: id,
- files,
- target: targets.get(kind.id),
- settings,
- });
- attempted.add(key(kind.id, id));
- cursor++;
-
- const action = candidateAction(state, {
- force: opts.force === true,
- // 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: allowRedownload && !result.diskFloorHit,
- });
- if (action === "skip") continue;
- if (action === "fresh") {
- result.fresh++;
- continue;
- }
- if (action === "deferred") {
- result.deferred++;
- continue;
- }
- if (action === "blocked") {
- result.blocked++;
- continue;
- }
- if (action === "missing-input") {
- // Counted, not attempted. This is the population the whole
- // reachable-vs-needs-re-acquiring split exists to keep visible.
- result.missingInput++;
- continue;
- }
- return { id, kind, target: targets.get(kind.id) };
- }
- };
-
- const baseline = opts.progressBaseline ?? 0;
- const progressTarget = opts.progressTarget ?? null;
- const reportProgress = (): void => {
- if (!opts.setProgress) return;
- const current = baseline + result.succeeded + result.failed;
- opts.setProgress({
- metric: "backfills",
- initial: baseline,
- target: Math.max(progressTarget ?? current, current),
- current,
- });
- };
-
- const runOne = async (
- candidate: Candidate,
- runSignal: AbortSignal,
- ): Promise<void> => {
- const videoDir = path.join(dataDir, candidate.id);
- const task = opts.tracker?.start({
- id: candidate.id,
- label: `backfill ${candidate.kind.id} ${opts.channelSlug}/${candidate.id}`,
- kind: "backfill",
- });
- const itemLog = task ? task.onLog : log;
- // 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 pull time, but the pull and the run are separated by however long the
- // pool held us, so this is re-checked inside reacquireMediaFor.
- if (allowRedownload) {
- reacquired = await reacquireMediaFor({
- paths: opts.paths,
- channelSlug: opts.channelSlug,
- videoId: candidate.id,
- videoDir,
- onLog: itemLog,
- signal: runSignal,
- });
- if (reacquired.status === "disk-floor") {
- // Latches for the rest of the run — next() reads this and stops
- // offering videos that would need a fetch.
- if (!result.diskFloorHit) {
- log(
- `Disk floor reached — no more media will be re-acquired this run. ` +
- `Videos needing it are counted as missing-input.`,
- );
- }
- result.diskFloorHit = true;
- result.missingInput++;
- return;
- }
- if (reacquired.status === "failed" || reacquired.status === "gone") {
- result.reacquireFailed++;
- result.missingInput++;
- return;
- }
- if (reacquired.status === "fetched") result.reacquired++;
- }
-
- // Dispatch ladder: a whole unit to a tagged remote executor when one is
- // free; else the llm fan-out for the call-bound kinds; else the local
- // run exactly as today. All claims are non-parking on purpose: a parked
- // acquire inside a runPool slot would deadlock the batch.
- let outcome: OperationRunOutcome | null = await runBackfillUnit(
- candidate,
- videoDir,
- itemLog,
- runSignal,
- );
- if (outcome === null) {
- const llm =
- llmFanout && LLM_BACKFILL_OPS.has(candidate.kind.id)
- ? await acquireLlmSlot(
- candidate.kind.id,
- llmFanout.modelRequested,
- itemLog,
- )
- : null;
- if (llm) llmActive++;
- try {
- outcome = await candidate.kind.run({
- paths: opts.paths,
- videoDir,
- videoId: candidate.id,
- channelSlug: opts.channelSlug,
- target: candidate.target,
- force: opts.force,
- ...(llm
- ? { appConfig: { ...llmFanout!.config, baseUrl: llm.baseUrl } }
- : {}),
- onLog: itemLog,
- signal: runSignal,
- });
- if (llm) getWorkerPool().markSuccess(llm.workerId);
- } catch (err) {
- if (llm && !runSignal.aborted && !opts.signal?.aborted) {
- getWorkerPool().markFailure(llm.workerId);
- }
- throw err;
- } finally {
- if (llm) {
- llmActive--;
- llm.lease.release();
- }
- }
- }
- if (outcome === "done") {
- result.attempted++;
- result.succeeded++;
- } else if (outcome === "already-present") {
- result.fresh++;
- } else if (outcome === "missing-input") {
- result.missingInput++;
- } else if (outcome === "skipped") {
- result.skipped++;
- } else {
- result.attempted++;
- result.failed++;
- log(`Failed ${candidate.kind.id} ${candidate.id}: ${outcome}.`);
- }
- reportProgress();
- } catch (err) {
- if (runSignal.aborted || opts.signal?.aborted) throw err;
- result.attempted++;
- result.failed++;
- log(
- `Failed ${candidate.kind.id} ${candidate.id}: ${(err as Error)?.message ?? String(err)}`,
- );
- } finally {
- // THE FINALLY THAT KEEPS THE DISK ALIVE. A re-fetched file is removed
- // whether the backfill 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, and its cleanup was built for that and
- // never called until now.
- if (
- reacquired?.status === "fetched" ||
- reacquired?.status === "failed"
- ) {
- const out = await reacquired.cleanup();
- // Only a "fetched" contributes to the reconciliation below: the other
- // statuses were never counted in result.reacquired.
- if (reacquired.status === "fetched") {
- if (out.status === "removed") result.reacquireCleaned++;
- else if (out.status === "kept" && out.reason !== "do-not-clean") {
- result.reacquireHandedOff++;
- }
- }
- }
- task?.end();
- }
- };
-
- // Edge-triggered, so a multi-day lane logs twice per contention window rather
- // than once per poll.
- let yielding = false;
- const NEVER = new AbortController().signal;
-
- reportProgress();
-
- await runPool<Candidate>({
- next,
- run: runOne,
- limit: () => {
- // Re-read at DISPATCH time, so a settings change takes effect within one
- // poll and survives a restart with no boot hook (the downloadsPaused
- // pattern). Returning 0 makes runPool idle-wait — a hold; returning null
- // from next() would END the batch, which is not the same thing.
- const liveSettings = getSettings();
- const live = liveSettings.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(liveSettings, "backfill")) {
- if (!yielding) {
- yielding = true;
- log("Backfill lane disabled in settings — holding.");
- }
- return 0;
- }
- // 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 lanes: those are already arbitrated against transcription by the
- // same helper, and making two yielding lanes watch each other is how both
- // end up parked. The cost of that choice 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 GUARANTEED SHARE IS A CPU CONCEPT. `weight > 0` means "keep a slice of
- // the cores running even while transcription works", which is reasonable
- // when the contended resource is cores and divisible. 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.
- //
- // So a run that could dispatch GPU work is idle-only whatever the weight
- // says. Which kinds those are is DECLARED (laneFor/contendsFor), not
- // re-tested here, so the rule stays next to the queue key it belongs with.
- // The cost is bluntness — one such kind 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.
- //
- // ONLY kinds that declare laneFor are considered, and that is the fix for a
- // real regression rather than a nicety. `digest`'s declared lane is
- // `contendsFor: "gpu"`, so testing every kind's lane made this lane
- // idle-only whenever a digest was in the run — which broke the guarantee
- // that a digest runs CONCURRENTLY with a backfill rather than behind it. A
- // kind with a static GPU lane already arbitrates itself (digestBatch has
- // its own yield); laneFor marks the kinds whose resource changes under
- // them from settings, and those are the ones nothing else is deciding for.
- //
- // DIGEST NOW DECLARES A laneFor TOO, and is still not in this decision —
- // but not because of this guard. `kinds` here came through
- // resolveBackfillLaneOperations, which filters via backfillLaneOperations and so admits
- // BACKFILL_QUEUE kinds only; digest is on its own queue and can never
- // reach this array, even when asked for by name. That is the invariant
- // operations.test.ts pins.
- //
- // Which is also why laneFor MUST STAY OPTIONAL. This line keys off its
- // PRESENCE, so "give every kind a laneFor defaulting to lane" would
- // silently enrol every kind — reviving exactly the regression above for
- // any statically GPU-bound kind that IS on this queue. Add one only where
- // the lane genuinely varies with settings.
- const gpuBound = kinds.some(
- (k) => k.laneFor && laneYieldsToTranscription(k.laneFor(liveSettings)),
- );
- const limit = backfillLimit({
- weight: gpuBound ? 0 : live.weight,
- slots: live.concurrency,
- primaryBusy: activity.busy,
- });
- // The REMOTE term: llm-endpoint slots serving this run's kinds, 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-kind run gets no term.
- const remoteTerm = remoteEligible
- ? llmActive + freeLlmSlots(llmOps)
- : 0;
- // The UNIT term: tagged remote slots every kind in this run could take,
- // plus the leases it already holds. The MIN across kinds keeps it sound
- // on a mixed run — a slot justified by one kind's remote capacity must
- // not dispatch another kind onto this box while the local term holds.
- const unitTerm =
- unitActive +
- Math.min(
- ...kinds.map((k) =>
- getWorkerPool().freeSlots(unitRequires(k), {
- kind: "remote",
- taggedOnly: true,
- }),
- ),
- );
- const total = limit + remoteTerm + unitTerm;
- if (total === 0) {
- if (!yielding) {
- yielding = true;
- log(
- `Standing aside for transcription (${activity.reason}); the backfill lane resumes when it is free.`,
- );
- }
- return 0;
- }
- if (yielding) {
- yielding = false;
- log(
- limit === 0
- ? "Resuming the backfill lane on remote LLM endpoints only."
- : "Resuming the backfill lane.",
- );
- }
- return total;
- },
- signal: opts.signal ?? NEVER,
- drainSignal: opts.drainSignal ?? NEVER,
- finite: true,
- idlePollMs: 3000,
- });
-
- // 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;
-}
-
-// How much reachable backfill work a channel has right now, for sizing a
-// progress bar before starting a job. SAME derivation as the batch's own pull —
-// the two disagreeing is what makes a bar stall one short of complete forever.
-export async function countBackfillWork(
- paths: Paths,
- channelSlug: string,
- kindIds?: string[],
- // Only these videos, intersected with disk — the same scope the batch 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?: string[],
-): Promise<{
- reachable: number;
- missingInput: number;
- deferred: number;
- blocked: number;
-}> {
- const settings = getSettings();
- const kinds = resolveBackfillLaneOperations(settings, kindIds);
- if (kinds.length === 0)
- return { reachable: 0, missingInput: 0, deferred: 0, blocked: 0 };
- const dataDir = path.join(paths.channelsDir, channelSlug, "data");
- const allDirs = await readdir(dataDir).catch(() => [] as string[]);
- // Intersected with disk exactly as the batch's own pull does, so an id that
- // names nothing counts as nothing rather than as one unit of work.
- const dirs = ids
- ? ids.filter((id) => allDirs.includes(id))
- : allDirs;
- const targets = new Map<string, unknown>();
- for (const kind of kinds) {
- targets.set(
- kind.id,
- await kind.resolveTarget({ settings, paths, channelSlug }),
- );
- }
-
- // Folded through the SHARED counters rather than a private if-chain. The
- // header above promises this function derives work the same way the batch
- // does; it used to re-spell the rule instead, and had already drifted —
- // `blocked` was classified by state() and then counted by nothing here, so a
- // channel of prerequisite-waiting videos reported three zeroes and looked
- // finished. addOperationState + reachableOperationWork is the same pair the
- // channel snapshot folds with, 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 to come back.
- const counts = emptyOperationCounts();
- for (const id of dirs) {
- const videoDir = path.join(dataDir, id);
- const files = await readVideoFiles(videoDir, {
- checkUntranscribable: true,
- });
- for (const kind of kinds) {
- addOperationState(
- counts,
- await kind.state({
- videoDir,
- videoId: id,
- files,
- target: targets.get(kind.id),
- settings,
- }),
- );
- }
- }
- return {
- reachable: reachableOperationWork(counts),
- missingInput: counts.missingInput,
- deferred: counts.deferred,
- blocked: counts.blocked,
- };
-}
diff --git a/common/controller/backfillSweep.ts b/common/controller/backfillSweep.ts
@@ -38,7 +38,7 @@ import {
type OperationSnapshotEntry,
} from "../lib/operations";
import { listChannelStatsFromDisk, readChannelSnapshot } from "./channels";
-import { countBackfillWork } from "./backfillBatch";
+import { countOperationWork } from "./operationBatch";
import { runBackfillChannelJob } from "./operationJobs";
import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
import { buildRecencyKeys } from "./recencyIndex";
@@ -113,7 +113,7 @@ export type BackfillPlanEntry = SweepPlanEntry;
// corpus by how much of it is unreachable, which on the measured numbers (835
// reachable vs ~76,270 needing a re-download) means sorting by noise.
//
-// PLANNED OFF SNAPSHOTS, not off a corpus walk. countBackfillWork() is a
+// PLANNED OFF SNAPSHOTS, not off a corpus walk. countOperationWork() is a
// readdir + readVideoFiles + sidecar read for EVERY video — ~77,000 directory
// reads — and this function runs once per pass of a sweep that runs for days.
// The channel snapshots already hold these exact numbers in
@@ -165,7 +165,9 @@ export async function buildBackfillSweepPlan(opts: {
const counts = snapshot?.backfill
? countFromSnapshot(snapshot.backfill, kinds)
: ((walked++,
- await countBackfillWork(opts.paths, ch.slug, opts.kindIds)));
+ await countOperationWork("backfill", opts.paths, ch.slug, {
+ operationIds: opts.kindIds,
+ })));
if (idsBySlug && snapshot?.backfill) {
const ids = new Set<string>();
for (const id of kinds) {
diff --git a/common/controller/batchRecency.test.ts b/common/controller/batchRecency.test.ts
@@ -68,7 +68,7 @@ test("newest/oldest sort a channel's candidates by upload date", async () => {
});
test("recency COMPOSES with a prior sort rather than replacing it", async () => {
- // This is how digestBatch keeps both of its rules: sort by duration, then
+ // This is how the digest lane keeps both of its rules: sort by duration, then
// stable-sort by date. makeRecencyComparator returns 0 for two videos sharing
// a YYYYMMDD key and Array#sort is stable, so the duration order survives
// inside a day. "newest day first, shortest video within a day."
diff --git a/common/controller/batchRecency.ts b/common/controller/batchRecency.ts
@@ -2,8 +2,8 @@ import type { Paths } from "../lib/paths";
import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
import { buildRecencyKeys, makeRecencyComparator } from "./recencyIndex";
-// Recency ordering for the per-channel BATCH controllers (digestBatch,
-// backfillBatch), as distinct from the auto-queue runners.
+// Recency ordering for the per-channel BATCH controllers (the digest and
+// backfill lanes), as distinct from the auto-queue runners.
//
// The runners already sort by upload date through recencyIndex; the batches did
// not, so "newest first" meant one thing on two of the four pipelines and
diff --git a/common/controller/digestBatch.ts b/common/controller/digestBatch.ts
@@ -1,647 +0,0 @@
-// Channel-scoped digest sweep: the thing that actually runs for weeks.
-//
-// Three mechanics here are copied from existing code for specific reasons, and
-// changing any of them breaks a property a multi-week sweep depends on:
-//
-// 1. RESUME BY RE-DERIVING FROM DISK. Eligibility is re-checked against the
-// sidecar on every next() pull (the auto-runner's shape), NOT frozen into an
-// array with an index cursor (whisperBatch's shape). The cursor form does not
-// survive a restart, and it cannot see a video that became eligible mid-run
-// (a transcript that finished, a mirror that got shared to).
-// 2. PAUSE BY RETURNING limit() === 0. runPool idle-WAITS at a zero limit
-// rather than finishing (concurrentRunner.ts's documented invariant), so a
-// pause holds the job open instead of ending it. The flag is re-read from
-// settings on every pull — the downloadsPaused pattern, which needs no boot
-// hook, unlike transcriptionsPaused.
-// 3. ONE JOB PER CHANNEL, never one per video. Job logs keep the newest 500 /
-// 30 days and the in-memory registry keeps 100 records; 119,600 jobs would
-// evict everything, including the running ones' history.
-
-import path from "node:path";
-import { open, readdir } from "node:fs/promises";
-import type { Paths } from "../lib/paths";
-import { getSettings } from "../lib/settings";
-import { runPool } from "../jobs/concurrentRunner";
-import type { TaskTracker } from "../jobs/taskHooks";
-import type { JobProgress } from "../jobs/registry";
-import { isSectionFresh, type DigestSectionKind } from "../lib/digest";
-import {
- digestLaneFor,
-} from "../lib/operations";
-import { loadDigest } from "../lib/digest-server";
-import {
- resolveDigestTarget,
- type DigestLaneChoice,
-} from "./digestTarget";
-import { CUES_JSON_FILENAME } from "../lib/videoStatus";
-import { digestVideo } from "./digestVideo";
-import {
- digestGate,
- digestPreflight,
- laneSharesDuplicates,
-} from "./laneGuards";
-import { acquireLlmSlot, freeLlmSlots } from "./llmWorkers";
-import { getWorkerPool } from "../jobs/workerPool";
-import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
-import { batchRecencyComparator } from "./batchRecency";
-import {
- buildDigestClusterPlan,
- planSlugForDir,
- shareDigestToCluster,
- type DigestClusterPlan,
-} from "./digestSharing";
-
-export type DigestOrder = "shortest-first" | "longest-first";
-
-export type DigestBatchOptions = {
- channelSlug: string;
- paths: Paths;
- // Which lane to run. "local" is the default and the only one enabled unless
- // settings.digest.remoteEnabled is true; asking for "remote" while it is off is
- // an error the caller surfaces, not a silent downgrade.
- lane?: "local" | "remote";
- // When set, only consider these video ids (intersected with what's on disk).
- ids?: string[];
- sections?: DigestSectionKind[];
- // The repo only had a `reverse` flag before this; digest ordering is explicit
- // because lane routing and the pilot both depend on it. 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.
- order?: DigestOrder;
- // Upload-date ordering, applied AFTER `order` above. Absent = whatever
- // settings.digest.recencyOrder says, resolved inside the batch so the armed
- // sweep, a hand-clicked channel run and an ids-scoped run all inherit the
- // same setting instead of three call sites each remembering to pass it.
- recencyOrder?: AutoQueueOrder;
- // Duration window, for splitting the corpus between lanes (e.g. local takes
- // ≤ longTailSeconds, the metered lane takes the tail).
- minDurationSeconds?: number;
- maxDurationSeconds?: number;
- // Stop after this many successful videos. For the pilot and Stage B's sample.
- limitCount?: number;
- concurrency?: number;
- force?: boolean;
- // Skip mirrors and share the canonical member's digest to aligned ones. On by
- // default: worth ~11% of the sweep. Pass a prebuilt plan to avoid re-reading
- // the duplicates report.
- useClusters?: boolean;
- clusterPlan?: DigestClusterPlan;
- 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
- // (countMissingDigests) to size the job; the batch raises it if sharing turns
- // out to satisfy more videos than were counted missing.
- progressTarget?: number;
- onLog?: (msg: string) => void;
- signal?: AbortSignal;
- drainSignal?: AbortSignal;
- tracker?: TaskTracker;
-};
-
-export type DigestBatchResult = {
- attempted: number;
- succeeded: number;
- skipped: number;
- failed: number;
- fresh: number;
- // Cluster 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;
-};
-
-type Candidate = { id: string; duration: number };
-
-// 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;
-
-async function readDurationFast(cuesPath: string): Promise<number | null> {
- let handle: Awaited<ReturnType<typeof open>> | 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(() => {});
- }
-}
-
-export async function runDigestBatch(
- opts: DigestBatchOptions,
-): Promise<DigestBatchResult> {
- const log = opts.onLog ?? ((m: string) => console.log(m));
- const settings = getSettings();
- const digestSettings = settings.digest;
- const lane = opts.lane ?? "local";
- // ONE derivation, shared with countMissingDigests and the digest operation
- // kind the snapshot folds `backfill.digest` from — see digestTarget.ts. It
- // must match what digestVideo will
- // actually chunk with, or the batch's freshness check and the writer would
- // disagree on the identity and every video would look stale forever.
- const resolved = await resolveDigestTarget({
- paths: opts.paths,
- channelSlug: opts.channelSlug,
- lane,
- sections: opts.sections,
- });
- const { app, config, sections, modelRequested, timestampMode, context } =
- resolved;
-
- const dataDir = path.join(opts.paths.channelsDir, opts.channelSlug, "data");
-
- // BOTH fail-fasts, asked as ONE declared rule. Fail fast and loudly rather
- // than 1,100 times in a row: an unreachable engine is a configuration
- // problem, and discovering it per-item wastes the log. It still THROWS here,
- // byte-identically — the message is the rule's, and the rule is what a future
- // dispatcher asks instead of re-implementing this from memory.
- //
- // Deliberately after resolveDigestTarget, not before: `app` and `config` come
- // out of it, and probing an engine the run was never going to use would be a
- // network call for nothing.
- const preflight = await digestPreflight({
- settings,
- lane,
- app: { id: app.id, lane: app.lane, probe: () => app.probe(config) },
- });
- if (!preflight.ok) throw new Error(preflight.error);
-
- const allDirs = await readdir(dataDir).catch(() => [] as string[]);
- const onDisk = new Set(allDirs);
- const wanted = opts.ids
- ? opts.ids.filter((id) => onDisk.has(id))
- : allDirs;
-
- // WHETHER A LANE CAN SHARE A DUPLICATE'S OUTPUT IS DECLARED, not assumed from
- // the fact that this file is digestBatch. 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 clusterPlan =
- opts.useClusters === false || !laneSharesDuplicates(digestLaneFor(app.lane))
- ? null
- : (opts.clusterPlan ?? (await buildDigestClusterPlan(opts.paths)));
-
- // Durations, read once. Videos with no normalized transcript have no duration
- // and are dropped here — they are a transcription problem, not a digest one.
- const candidates: Candidate[] = [];
- let noTranscript = 0;
- let outOfWindow = 0;
- for (const id of wanted) {
- opts.signal?.throwIfAborted();
- const duration = await readDurationFast(
- path.join(dataDir, id, CUES_JSON_FILENAME),
- );
- 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({ id, duration });
- }
-
- const order = opts.order ?? "shortest-first";
- candidates.sort((a, b) =>
- order === "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 makeRecencyComparator returns 0
- // for two videos sharing a YYYYMMDD key, so this reads as "newest day first,
- // shortest video within a day" — both rules intact. With recencyOrder
- // "listed" the comparator is null and this second sort never runs, which is
- // what reproduces the historical order byte-for-byte.
- const recencyOrder = opts.recencyOrder ?? digestSettings.recencyOrder;
- const byRecency = await batchRecencyComparator(
- opts.paths,
- opts.channelSlug,
- candidates.map((c) => c.id),
- recencyOrder,
- );
- if (byRecency) candidates.sort((a, b) => byRecency(a.id, b.id));
-
- log(
- `Digest ${opts.channelSlug} (${lane} lane, ${app.id}/${modelRequested}, sections: ${sections.join(", ")}): ` +
- `${candidates.length} candidate(s) of ${wanted.length} on disk ` +
- `(${noTranscript} without a transcript, ${outOfWindow} outside the duration window), ${order}` +
- (byRecency ? `, ${recencyOrder} first.` : ".") +
- (clusterPlan
- ? ` Duplicate plan: ${clusterPlan.clusters} cluster(s), ${clusterPlan.bySlug.size} member(s) mapped.`
- : " Duplicate sharing off."),
- );
-
- const result: DigestBatchResult = {
- attempted: 0,
- succeeded: 0,
- skipped: 0,
- failed: 0,
- fresh: 0,
- shared: 0,
- misaligned: 0,
- engineCalls: 0,
- costUsd: 0,
- warnings: 0,
- spendCapped: false,
- };
-
- const freshnessTarget = resolved.target;
-
- // Ids already handed out this run. Combined with the cursor below, this is what
- // makes 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 id it stops on is always re-read from disk.
- const attempted = new Set<string>();
- let cursor = 0;
-
- // `candidate.id` is a DIRECTORY name (it comes from readdir), and a slug is
- // `${channelSlug}/${metadataId}`. Those disagree for 14.5% of the corpus —
- // every Rumble re-upload — so resolving one to the other through the plan is
- // what makes the cluster lookups below hit at all. See planSlugForDir.
- const slugOf = (videoDir: string): string =>
- planSlugForDir(clusterPlan, opts.channelSlug, videoDir);
-
- const next = async (): Promise<Candidate | null> => {
- while (cursor < candidates.length) {
- if (opts.limitCount !== undefined && result.succeeded >= opts.limitCount) {
- return null;
- }
- const candidate = candidates[cursor];
- if (attempted.has(candidate.id)) {
- cursor++;
- continue;
- }
- const videoDir = path.join(dataDir, candidate.id);
-
- // A cluster mirror is not this lane's work: its canonical member owns the
- // generation and shares the result here.
- const role = clusterPlan?.bySlug.get(slugOf(candidate.id));
- if (role?.kind === "mirror") {
- attempted.add(candidate.id);
- cursor++;
- resolvedAudioSeconds += candidate.duration;
- result.skipped++;
- log(
- `Skipping ${candidate.id}: duplicate of ${role.canonicalSlug}, which owns the digest for this cluster.`,
- );
- continue;
- }
-
- // RE-DERIVED FROM DISK, every pull. A restart, a concurrent lane, or a
- // share that landed while this job ran are all visible here.
- if (!opts.force) {
- const record = await loadDigest(videoDir);
- const allFresh = sections.every((section) =>
- isSectionFresh(record, section, freshnessTarget),
- );
- if (allFresh) {
- attempted.add(candidate.id);
- cursor++;
- resolvedAudioSeconds += candidate.duration;
- result.fresh++;
- continue;
- }
- }
- attempted.add(candidate.id);
- cursor++;
- return candidate;
- }
- return null;
- };
-
- // WIRING THE PLUMBING THAT WAS DECLARED AND NEVER USED.
- //
- // `setProgress` and `progressBaseline` were on the options type and passed by
- // digestActions.ts, but nothing in this file referenced either. 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, for two reasons the disk
- // cannot express: it knows a regenerate happened, and it knows sharing moved
- // the denominator — a canonical member's digest can satisfy a dozen mirrors
- // at once, so a sweep genuinely CHANGES its own target as it runs.
- const baseline = opts.progressBaseline ?? 0;
- const progressTarget = opts.progressTarget ?? null;
- // 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);
- let resolvedAudioSeconds = 0;
- const 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. So:
- // succeeded + shared — work this run did, and did against the target;
- // skipped — was in the target and is resolved anyway (a mirror
- // its canonical owns, or a video with no transcript);
- // without it the bar stalls short of done forever.
- // `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 — 90 fresh of 100 would report 100/10.
- const done = result.succeeded + result.shared + result.skipped;
- const remainingAudioSeconds = Math.max(
- 0,
- totalAudioSeconds - resolvedAudioSeconds,
- );
- // 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.
- const current = baseline + done;
- opts.setProgress({
- metric: "digests",
- initial: baseline,
- target: Math.max(progressTarget ?? current, current),
- current,
- remainingAudioSeconds,
- });
- };
-
- // LLM fan-out: an "llm" endpoint worker (a box running nothing but `ollama
- // serve`) can take this lane's model calls. Only for an HTTP engine with a
- // baseUrl field (ollama) — the metered CLI lane has no endpoint to swap.
- // `llmActive` counts leases THIS run holds, so limit() stays stable at
- // local + total-llm-capacity instead of sagging while slots are leased.
- const fanOutEligible = app.fields.baseUrl === true;
- let llmActive = 0;
-
- const runOne = async (
- candidate: Candidate,
- runSignal: AbortSignal,
- ): Promise<void> => {
- const task = opts.tracker?.start({
- id: candidate.id,
- label: `digest ${opts.channelSlug}/${candidate.id}`,
- kind: "digest",
- // 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 sweep ETA.
- audioSeconds: candidate.duration,
- });
- // 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 llm = fanOutEligible
- ? await acquireLlmSlot("digest", modelRequested, task ? task.onLog : log)
- : null;
- if (llm) llmActive++;
- try {
- const outcome = await digestVideo({
- paths: opts.paths,
- channelSlug: opts.channelSlug,
- videoId: candidate.id,
- sections,
- appId: app.id,
- config: llm ? { ...config, baseUrl: llm.baseUrl } : config,
- context,
- timestampMode,
- ...(digestSettings.promptVariant
- ? { promptVariant: digestSettings.promptVariant }
- : {}),
- force: opts.force,
- onLog: task ? task.onLog : opts.onLog,
- signal: runSignal,
- });
- // 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") {
- result.fresh++;
- return;
- }
- if (outcome.status === "skipped") {
- result.skipped++;
- log(`Skipped ${candidate.id}: ${outcome.reason}.`);
- return;
- }
- result.attempted++;
- result.succeeded++;
- result.engineCalls += outcome.engineCalls;
- result.costUsd += outcome.costUsd;
- result.warnings += outcome.warningCount;
-
- // 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-sweep leaves a consistent cluster.
- const role = clusterPlan?.bySlug.get(slugOf(candidate.id));
- if (role?.kind === "canonical") {
- const outcomes = await shareDigestToCluster({
- paths: opts.paths,
- clusterId: role.clusterId,
- canonicalSlug: slugOf(candidate.id),
- mirrors: role.mirrors,
- dirBySlug: clusterPlan?.dirBySlug,
- onLog: log,
- });
- for (const o of outcomes) {
- if (o.status === "shared") result.shared++;
- else if (o.status === "misaligned") result.misaligned++;
- }
- }
- reportProgress();
- } catch (err) {
- if (runSignal.aborted || opts.signal?.aborted) throw err;
- // Counted against the worker too, so a flapping endpoint degrades
- // through the pool's existing consecutive-failure machinery.
- if (llm) getWorkerPool().markFailure(llm.workerId);
- result.attempted++;
- result.failed++;
- log(
- `Failed ${candidate.id}: ${(err as Error)?.message ?? String(err)}`,
- );
- } finally {
- if (llm) {
- llmActive--;
- llm.lease.release();
- }
- resolvedAudioSeconds += candidate.duration;
- task?.end();
- }
- };
-
- // Edge-triggered, so a hold logs twice per contention window rather than once
- // per poll over a multi-week sweep. 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: "paused" | "yield" | "spend-cap" | null = null;
-
- const NEVER = new AbortController().signal;
- // 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;
- const concurrency = Math.max(1, opts.concurrency ?? defaultConcurrency);
-
- // 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<Candidate>({
- next,
- run: runOne,
- limit: () => {
- // EVERY GUARD, ASKED AS ONE DECLARED RULE. Settings are re-read at
- // DISPATCH time 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 batch. Every pause in this repo depends on
- // that distinction.
- const gate = digestGate({
- settings: getSettings(),
- appLane: app.lane,
- metered: app.metered,
- // Per-RUN, not per-settings: the spend cap is a per-job ceiling, so it
- // is read off the live accumulator rather than from disk.
- costUsd: result.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 = fanOutEligible
- ? llmActive + freeLlmSlots(["digest"])
- : 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) {
- if (heldReason !== "yield") {
- heldReason = "yield";
- log(
- `${gate.message} Remote LLM endpoint(s) keep the lane moving meanwhile.`,
- );
- }
- return remoteTerm;
- }
- // LOGGED ON THE EDGE ONLY, and the edge is tracked here rather than in
- // the rule: the rule is pure and a line per poll would bury a job log
- // over a multi-week run. An operator watching a lane sit at zero
- // throughput still has to be able to tell yielding from wedged.
- if (heldReason !== gate.reason) {
- heldReason = gate.reason;
- log(gate.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 (gate.reason === "spend-cap") result.spendCapped = true;
- return 0;
- }
- if (heldReason === "yield") {
- log("Transcription finished; resuming the digest lane.");
- }
- heldReason = null;
- return concurrency + remoteTerm;
- },
- signal: opts.signal ?? NEVER,
- drainSignal: opts.drainSignal ?? NEVER,
- finite: true,
- idlePollMs: 3000,
- });
-
- // Metered accounting is logged unconditionally when the lane is metered, even
- // at zero calls: "this run cost nothing" is information too.
- if (app.metered) {
- log(
- `Metered lane: ${result.engineCalls} model call(s), $${result.costUsd.toFixed(4)} total` +
- (digestSettings.spendCapUsd > 0
- ? ` (cap $${digestSettings.spendCapUsd.toFixed(2)})`
- : " (no cap set)"),
- );
- }
- return result;
-}
-
-// How many of a channel's videos still need a digest at the CURRENT prompt/model
-// identity. Used to scope a job's progress bar, and cheap enough to call before
-// starting one (a small sidecar read per video, no transcript parsing).
-//
-// `lane` is REQUIRED-in-spirit: it used to be absent and the local app was
-// resolved unconditionally, so a metered-lane job's progress bar was sized
-// against the LOCAL engine's identity — and, once numCtx joined the identity,
-// against the local app's context size too. Both call sites already know their
-// lane. It still defaults to "local" so the meaning of an unqualified call is
-// the same as it always was, rather than silently changing under old callers.
-export async function countMissingDigests(
- paths: Paths,
- channelSlug: string,
- ids?: ReadonlyArray<string>,
- lane: DigestLaneChoice = "local",
-): Promise<number> {
- const { target, sections } = await resolveDigestTarget({
- paths,
- channelSlug,
- lane,
- });
- const dataDir = path.join(paths.channelsDir, channelSlug, "data");
- const dirs = ids
- ? [...ids]
- : await readdir(dataDir).catch(() => [] as string[]);
- let missing = 0;
- for (const id of dirs) {
- // SAME ELIGIBILITY AS THE BATCH, and it has to be. A video with no
- // normalized transcript is a transcription problem, not a digest one:
- // runDigestBatch drops it before it is ever a candidate, so counting it here
- // sizes the progress bar against work the batch will never do and the bar
- // stalls one short of complete — for the rest of the run.
- //
- // Measured: `teamrcn` reported target 8, current 7. The eighth directory is
- // a channel-id-named folder with no transcript. Corpus-wide that is 2,987
- // videos, i.e. most channels would never show as finished, which over an
- // 81-day sweep is indistinguishable from wedged. This is exactly the
- // counter-disagreement digestTarget.ts's header was written about.
- const duration = await readDurationFast(
- path.join(dataDir, id, CUES_JSON_FILENAME),
- );
- if (duration === null || duration <= 0) continue;
- const record = await loadDigest(path.join(dataDir, id));
- const fresh = sections.every((section) =>
- isSectionFresh(record, section, target),
- );
- if (!fresh) missing++;
- }
- return missing;
-}
diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts
@@ -9,7 +9,7 @@
// Three shapes are deliberate:
//
// 1. IT LAUNCHES THE EXISTING PER-CHANNEL JOB, in sequence. It does not
-// introduce a batch over videos. digestBatch.ts is explicit that one job
+// introduce a batch over videos. operationBatch.ts is explicit that one job
// per video would be 119,600 jobs against a 100-record registry and a
// 500-record log, evicting the history of the very run it is recording.
// 63 sequential channel jobs is the same contract, driven.
diff --git a/common/controller/digestYield.ts b/common/controller/digestYield.ts
@@ -17,7 +17,7 @@
// digestApps.ts now records; until then, do not quote a contention penalty as
// fact.
//
-// The fix is deliberately NOT a scheduler. `digestBatch`'s limit() already
+// 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.
diff --git a/common/controller/laneGuards.ts b/common/controller/laneGuards.ts
@@ -18,10 +18,10 @@ import { transcriptionActivity } from "./digestYield";
// them is expressible as backfillLimit()'s single scalar."
//
// That is the whole reason there are two schedulers for one kind of work. Every
-// one of those guards lived inside digestBatch's own closure, where nothing
+// one of those guards lived inside the digest batch's own closure, where nothing
// else could ask about it — so an arbiter that wanted to dispatch digest would
// have had to re-implement all six, correctly, from memory. They live here now,
-// and digestBatch calls them, so there is exactly one definition and no second
+// 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:
diff --git a/common/controller/operationBatch.test.ts b/common/controller/operationBatch.test.ts
@@ -0,0 +1,439 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { backfillLimit, candidateAction, laneLimit } from "./operationBatch";
+import {
+ defaultBackfill,
+ defaultDigest,
+ defaultDiarization,
+ defaultSiteSettings,
+ sanitizeBackfill,
+ sanitizeDiarization,
+ type SiteSettings,
+} from "../lib/settings";
+import type { Operation, OperationClassification } from "../lib/operations";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/operationBatch.test.ts
+//
+// The PURE function, not a live pool — same discipline as digestYield.test.ts,
+// and for the same reason: this decides whether a multi-day lane runs at all,
+// and a test that needs a GPU, a registry and a worker pool is a test nobody
+// runs before shipping.
+
+test("weight 0 is idle-only: full slots when the primary lane is quiet", () => {
+ assert.equal(
+ backfillLimit({ weight: 0, slots: 4, primaryBusy: false }),
+ 4,
+ );
+ assert.equal(backfillLimit({ weight: 0, slots: 1, primaryBusy: false }), 1);
+});
+
+test("weight 0 stands aside completely while the primary lane works", () => {
+ // Zero is a HOLD, not a stop: runPool idle-waits at a zero limit rather than
+ // finishing, so the lane resumes the moment transcription is free without
+ // re-deriving anything. That is the whole reason this is a limit and not a
+ // scheduler.
+ assert.equal(backfillLimit({ weight: 0, slots: 4, primaryBusy: true }), 0);
+ assert.equal(backfillLimit({ weight: 0, slots: 16, primaryBusy: true }), 0);
+});
+
+test("a positive weight is a guaranteed share, busy or not", () => {
+ // The point of a non-zero weight: the primary lane being busy no longer parks
+ // the backfill.
+ assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: true }), 2);
+ assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: false }), 2);
+ assert.equal(backfillLimit({ weight: 1, slots: 4, primaryBusy: true }), 4);
+});
+
+test("a small weight is a SLOW lane, never a stopped one", () => {
+ // floor(4 * 0.25) is 1, and floor(1 * 0.25) is 0 — which without the floor
+ // would silently turn "a quarter of the machine" into "never runs", and look
+ // exactly like a wedge.
+ assert.equal(backfillLimit({ weight: 0.25, slots: 4, primaryBusy: true }), 1);
+ assert.equal(backfillLimit({ weight: 0.25, slots: 1, primaryBusy: true }), 1);
+ assert.equal(backfillLimit({ weight: 0.01, slots: 8, primaryBusy: true }), 1);
+});
+
+test("no slots means no work, whatever the weight", () => {
+ assert.equal(backfillLimit({ weight: 0, slots: 0, primaryBusy: false }), 0);
+ assert.equal(backfillLimit({ weight: 1, slots: 0, primaryBusy: false }), 0);
+});
+
+test("the shipped default is idle-only", () => {
+ // Stated as a test because it is a promise the feature makes: catch-up work on
+ // a corpus that already exists must never slow down new arrivals.
+ const d = defaultBackfill();
+ assert.equal(d.weight, 0);
+ assert.equal(d.enabled, false);
+ assert.equal(d.allowRedownload, false);
+ assert.equal(backfillLimit({ ...d, slots: d.concurrency, primaryBusy: true }), 0);
+});
+
+test("a hand-edited weight is clamped rather than rejected", () => {
+ // settings.json is hand-editable. A 5 means "as much as possible", and reading
+ // it as the idle-only 0 would be the opposite of the intent.
+ assert.equal(sanitizeBackfill({ weight: 5 }).weight, 1);
+ assert.equal(sanitizeBackfill({ weight: -3 }).weight, 0);
+ assert.equal(sanitizeBackfill({ weight: "half" }).weight, 0);
+ assert.equal(sanitizeBackfill({ weight: 0.25 }).weight, 0.25);
+ // Concurrency shares clampPositiveInt's floor of 1, so a 0 cannot silently
+ // park the lane either.
+ assert.equal(sanitizeBackfill({ concurrency: 0 }).concurrency, 1);
+});
+
+test("a stale sweep scope survives sanitization as a list of slugs", () => {
+ // The scope is persisted WITH the flag, so it has to round-trip: junk entries
+ // are dropped, real ones kept.
+ assert.deepEqual(
+ sanitizeBackfill({ sweepChannels: ["a", "", 7, "b"], sweepKinds: ["diarization"] }),
+ {
+ ...defaultBackfill(),
+ sweepChannels: ["a", "b"],
+ sweepKinds: ["diarization"],
+ },
+ );
+});
+
+// ---------------------------------------------------------------------------
+// The candidate pull's dispatch decision.
+//
+// Untested until the duration cap needed it, and it is the single most
+// consequential branch in the file: it used to fall through to DISPATCHING any
+// state it did not recognize, so "add a state" and "run that state anyway" were
+// the same edit.
+
+const DISPATCH = { force: false, allowRedownload: false };
+
+test("only reachable work is dispatched by default", () => {
+ assert.equal(candidateAction("missing", DISPATCH), "dispatch");
+ assert.equal(candidateAction("stale", DISPATCH), "dispatch");
+});
+
+test("deferred is NEVER dispatched", () => {
+ // The whole point of the cap. Without this the classification would be
+ // computed, ignored, and the six-hour video handed to the engine — 40 minutes
+ // of CPU and a kernel OOM kill, producing nothing.
+ assert.equal(candidateAction("deferred", DISPATCH), "deferred");
+ assert.equal(
+ candidateAction("deferred", { force: true, allowRedownload: true }),
+ "deferred",
+ );
+});
+
+test("force redoes present work but does not overrule the cap", () => {
+ assert.equal(candidateAction("present", DISPATCH), "fresh");
+ assert.equal(
+ candidateAction("present", { ...DISPATCH, force: true }),
+ "dispatch",
+ );
+ // `force` means "redo work that looks done", not "ignore the cap".
+ assert.equal(
+ candidateAction("deferred", { ...DISPATCH, force: true }),
+ "deferred",
+ );
+});
+
+test("missing-input is counted unless re-download is armed", () => {
+ assert.equal(candidateAction("missing-input", DISPATCH), "missing-input");
+ assert.equal(
+ candidateAction("missing-input", { ...DISPATCH, allowRedownload: true }),
+ "dispatch",
+ );
+});
+
+test("not-applicable is counted as nothing at all", () => {
+ assert.equal(candidateAction("not-applicable", DISPATCH), "skip");
+ assert.equal(
+ candidateAction("not-applicable", { force: true, allowRedownload: true }),
+ "skip",
+ );
+});
+
+test("every classification has an explicit decision", () => {
+ // The list is written out rather than derived so that adding a state to the
+ // union without deciding what the pull does with it fails HERE as well as at
+ // the compile step — a runtime backstop for the `never` check, since the
+ // hazard this replaces was precisely a silent fall-through.
+ const ALL: OperationClassification[] = [
+ "present",
+ "stale",
+ // Part-done. Dispatched like stale — the runner regenerates only the
+ // sections that are not fresh — and listed here so that stays a decision
+ // rather than something a later edit can quietly drop.
+ "partial",
+ "missing",
+ "missing-input",
+ "deferred",
+ // This list had already drifted: `blocked` was added to the union and to
+ // candidateAction, and never to the backstop meant to catch exactly that.
+ "blocked",
+ "not-applicable",
+ ];
+ for (const state of ALL) {
+ assert.doesNotThrow(() => candidateAction(state, DISPATCH), state);
+ }
+ assert.throws(
+ () => candidateAction("invented" as OperationClassification, DISPATCH),
+ /unhandled backfill state/,
+ );
+});
+
+test("the duration cap round-trips, and 0 means off", () => {
+ // 0 is a MEANINGFUL value here, not an empty one — it is how the cap is turned
+ // off once windowed diarization makes it unnecessary. clampPositiveInt would
+ // have floored it to 1, i.e. a one-hour cap, which is why this knob does not
+ // use it.
+ assert.equal(sanitizeDiarization({ maxAudioHours: 0 }).maxAudioHours, 0);
+ assert.equal(sanitizeDiarization({ maxAudioHours: 6.5 }).maxAudioHours, 6.5);
+ // Junk falls back to the default rather than to "no cap" — a typo in a
+ // hand-edited settings.json must not silently re-arm the OOM.
+ assert.equal(
+ sanitizeDiarization({ maxAudioHours: "four" }).maxAudioHours,
+ defaultDiarization().maxAudioHours,
+ );
+ assert.equal(
+ sanitizeDiarization({ maxAudioHours: -2 }).maxAudioHours,
+ defaultDiarization().maxAudioHours,
+ );
+ // The shipped default is OFF: windowed diarization removed the OOM the cap
+ // was a stopgap for. This is asserted rather than assumed because turning it
+ // back on by accident would silently stop diarizing long videos.
+ assert.equal(defaultDiarization().maxAudioHours, 0);
+});
+
+test("an unknown engine or backend falls back to the default, never to nothing", () => {
+ // The DEFAULT ENGINE IS LOAD-BEARING as a fallback, not just as a starting
+ // point: it is the one every sidecar already on disk matches, so falling back
+ // to it leaves the corpus fresh. Falling back to sortformer on a typo would
+ // mark all of it stale and offer weeks of rework.
+ assert.equal(defaultDiarization().engine, "sherpa-onnx");
+ assert.equal(sanitizeDiarization({ engine: "sortformer" }).engine, "sortformer");
+ assert.equal(
+ sanitizeDiarization({ engine: "sortfromer" }).engine,
+ defaultDiarization().engine,
+ );
+ assert.equal(sanitizeDiarization({ engine: 7 }).engine, defaultDiarization().engine);
+
+ // The backend defaults to the GPU, which is safe only because the lane yields
+ // the card rather than sharing it — see diarizationLaneFor.
+ assert.equal(defaultDiarization().backend, "vulkan");
+ assert.equal(sanitizeDiarization({ backend: "cpu" }).backend, "cpu");
+ assert.equal(
+ sanitizeDiarization({ backend: "rocm" }).backend,
+ defaultDiarization().backend,
+ );
+
+ // Paths are plain strings and empty means "not configured", which diarizeOne
+ // reports as a skip rather than a failure.
+ assert.equal(sanitizeDiarization({}).sortformerBin, "");
+ assert.equal(sanitizeDiarization({}).sortformerModel, "");
+ assert.equal(
+ sanitizeDiarization({ sortformerBin: " /opt/diarize-file " }).sortformerBin,
+ "/opt/diarize-file",
+ );
+});
+
+// ---------------------------------------------------------------------------
+// The layer boundary. This test used to live in lib/operations.test.ts and
+// import candidateAction across it — a lib/ test reaching into controller/,
+// carried on architecture.test.ts's allow-list as a debt. The classification is
+// the registry's business; what a DISPATCHER does with it is this file's, so
+// the assertion belongs here.
+
+test("blocked is never dispatched, and re-download cannot change that", () => {
+ // The dispatch decision is the consequential one: `blocked` must not reach a
+ // runner whatever the flags say. Note allowRedownload — the flag that DOES
+ // turn missing-input into a dispatch — is deliberately inert here.
+ for (const force of [false, true]) {
+ for (const allowRedownload of [false, true]) {
+ assert.equal(
+ candidateAction("blocked", { force, allowRedownload }),
+ "blocked",
+ `force=${force} allowRedownload=${allowRedownload}`,
+ );
+ }
+ }
+ // The contrast, so this test fails if the two ever get conflated again.
+ assert.equal(
+ candidateAction("missing-input", { force: false, allowRedownload: true }),
+ "dispatch",
+ );
+});
+
+// ---------------------------------------------------------------------------
+// laneLimit — the whole of a lane's dispatch gate, asked as one rule.
+//
+// No pool, no registry, no GPU: transcriptionActivity() fails OPEN in a bare
+// process (an unreadable pool and an unreadable registry both read as "nothing
+// there"), which is the property that makes this testable at all — and the same
+// property that means a broken read causes contention rather than deadlock.
+
+function settingsWith(over: Partial<SiteSettings>): SiteSettings {
+ return { ...defaultSiteSettings(), ...over };
+}
+
+test("the backfill lane holds when the lane switch is off, and says so", () => {
+ const verdict = laneLimit(
+ settingsWith({ backfill: { ...defaultBackfill(), enabled: false } }),
+ {
+ lane: "backfill",
+ operations: [],
+ llmOps: [],
+ remoteEligible: false,
+ llmActive: 0,
+ unitActive: 0,
+ },
+ );
+ // A HOLD, not a stop: runPool idle-waits at zero, so the lane keeps its place
+ // and resumes on the next poll after the switch comes back.
+ assert.equal(verdict.limit, 0);
+ assert.equal(verdict.hold?.reason, "paused");
+});
+
+test("an enabled backfill lane with a weight runs, primary busy or not", () => {
+ const verdict = laneLimit(
+ settingsWith({
+ backfill: { ...defaultBackfill(), enabled: true, weight: 1, concurrency: 3 },
+ }),
+ {
+ lane: "backfill",
+ operations: [],
+ llmOps: [],
+ remoteEligible: false,
+ llmActive: 0,
+ unitActive: 0,
+ },
+ );
+ assert.equal(verdict.limit, 3);
+ assert.equal(verdict.hold, null);
+});
+
+test("the GPU carve-out is keyed on contendsFor, not on laneFor existing", () => {
+ // THE RULE THIS SLICE CHANGED. It used to be `k.laneFor && ...` — the field's
+ // PRESENCE as a proxy for "this one's resource depends on settings". That
+ // made `laneFor` load-bearing by existing: giving every operation one (a
+ // refactor that looks like a no-op) would have enrolled all of them and made
+ // the lane idle-only whenever a statically GPU-bound operation was in the
+ // run. Now the declaration itself answers, so an operation with a fixed
+ // `lane: { contendsFor: "gpu" }` and NO laneFor is caught, and one with a
+ // laneFor resolving to CPU is not.
+ const settings = settingsWith({
+ backfill: { ...defaultBackfill(), enabled: true, weight: 1, concurrency: 4 },
+ });
+ const cpuOnly = laneLimit(settings, {
+ lane: "backfill",
+ operations: [fakeOperation("cpu-op", "cpu")],
+ llmOps: [],
+ remoteEligible: false,
+ llmActive: 0,
+ unitActive: 0,
+ });
+ // A weight is a CPU concept, so a CPU-bound run keeps its guaranteed share.
+ assert.equal(cpuOnly.limit, 4);
+
+ const gpuBound = laneLimit(settings, {
+ lane: "backfill",
+ operations: [fakeOperation("cpu-op", "cpu"), fakeOperation("gpu-op", "gpu")],
+ llmOps: [],
+ remoteEligible: false,
+ llmActive: 0,
+ unitActive: 0,
+ });
+ // One GPU-bound operation makes the WHOLE run idle-only, including its cheap
+ // CPU siblings — deliberately blunt, because the alternative is an OOM when
+ // two engines allocate the same 8 GB card.
+ assert.equal(gpuBound.limit, 4, "idle: nothing is transcribing in this process");
+ const weightless = laneLimit(
+ settingsWith({
+ backfill: { ...defaultBackfill(), enabled: true, weight: 1, concurrency: 4 },
+ }),
+ {
+ lane: "backfill",
+ operations: [fakeOperation("gpu-op", "gpu")],
+ llmOps: [],
+ remoteEligible: false,
+ llmActive: 0,
+ unitActive: 0,
+ },
+ );
+ // Same shape as backfillLimit({weight: 0, primaryBusy: false}) — the weight
+ // was zeroed by the carve-out, and an idle primary still gives full slots.
+ assert.equal(weightless.limit, 4);
+});
+
+test("the digest lane holds while digests are paused", () => {
+ const verdict = laneLimit(
+ settingsWith({ digest: { ...defaultDigest(), digestsPaused: true } }),
+ {
+ lane: "digest",
+ appLane: "local-gpu",
+ metered: false,
+ costUsd: 0,
+ concurrency: 1,
+ fanOutEligible: false,
+ llmActive: 0,
+ },
+ );
+ assert.equal(verdict.limit, 0);
+ assert.equal(verdict.hold?.reason, "paused");
+});
+
+test("the digest spend cap is per RUN, and parks the metered lane", () => {
+ const settings = settingsWith({
+ digest: { ...defaultDigest(), spendCapUsd: 5 },
+ });
+ const under = laneLimit(settings, {
+ lane: "digest",
+ appLane: "remote-api",
+ metered: true,
+ costUsd: 4.99,
+ concurrency: 2,
+ fanOutEligible: false,
+ llmActive: 0,
+ });
+ assert.equal(under.limit, 2);
+ const over = laneLimit(settings, {
+ lane: "digest",
+ appLane: "remote-api",
+ metered: true,
+ costUsd: 5,
+ concurrency: 2,
+ fanOutEligible: false,
+ llmActive: 0,
+ });
+ assert.equal(over.limit, 0);
+ assert.equal(over.hold?.reason, "spend-cap");
+});
+
+test("the metered digest lane does not yield the GPU it never touches", () => {
+ // laneYieldsToTranscription reads `contendsFor` off the engine's declared
+ // lane: local-gpu contends, remote-api does not. Asserted through laneLimit
+ // rather than through the declaration so the wiring is covered too.
+ const settings = settingsWith({ digest: defaultDigest() });
+ for (const appLane of ["local-gpu", "remote-api"] as const) {
+ const verdict = laneLimit(settings, {
+ lane: "digest",
+ appLane,
+ metered: appLane === "remote-api",
+ costUsd: 0,
+ concurrency: 1,
+ fanOutEligible: false,
+ llmActive: 0,
+ });
+ // Nothing is transcribing in this process, so neither holds — what is
+ // asserted is that neither holds for the WRONG reason.
+ assert.equal(verdict.hold, null, appLane);
+ assert.equal(verdict.limit, 1, appLane);
+ }
+});
+
+// A minimal registered-shaped operation. Only `id` and `lane` are read by the
+// limit; laneForOperation returns null for an id the registry does not know, so
+// the declared `lane` below is what answers — which is exactly the fallback the
+// rule relies on.
+function fakeOperation(id: string, contendsFor: "cpu" | "gpu" | "network") {
+ return {
+ id,
+ lane: { queueKey: "backfill", contendsFor },
+ } as unknown as Operation;
+}
diff --git a/common/controller/operationBatch.ts b/common/controller/operationBatch.ts
@@ -0,0 +1,1966 @@
+// 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 { 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,
+ 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.
+//
+// weight 0 (the default) — IDLE-ONLY. Full slots when the primary lane is
+// quiet, zero while it works. Catch-up is by definition work on a corpus
+// that already exists, so it must never slow down new arrivals. Zero is a
+// HOLD, not a stop: runPool idle-waits, so the lane resumes the moment the
+// primary is free, having re-derived nothing.
+// weight > 0 — a guaranteed share, floored at 1. The floor is the point: a
+// small weight should mean a slow lane, not a stopped one, and
+// floor(1 * 0.25) is 0.
+export function backfillLimit(opts: {
+ weight: number;
+ slots: number;
+ primaryBusy: boolean;
+}): number {
+ const slots = Math.max(0, Math.floor(opts.slots));
+ if (slots === 0) return 0;
+ if (opts.weight <= 0) return opts.primaryBusy ? 0 : slots;
+ const share = Math.floor(slots * Math.min(1, opts.weight));
+ return Math.max(1, share);
+}
+
+// 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 GUARANTEED SHARE IS A CPU CONCEPT. `weight > 0` means "keep a slice of the
+ // cores running even while transcription works", which is reasonable when the
+ // contended resource is cores and divisible. 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.
+ //
+ // So a run that could dispatch GPU work is idle-only whatever the weight says.
+ // 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({
+ weight: gpuBound ? 0 : cfg.weight,
+ 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<number | null> {
+ let handle: Awaited<ReturnType<typeof open>> | 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<string, number | null>();
+const DURATION_MEMO_CAP = 200_000;
+
+export async function videoDurationSeconds(
+ paths: Paths,
+ channelSlug: string,
+ videoId: string,
+): Promise<number | null> {
+ 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 because today's live value
+// comes from `settings.digest.recencyOrder` ("newest"), and slice 1.3 retires
+// that field into the lane's tree. Until then the caller passes what settings
+// say; after it, the lane's own `order` does.
+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<string, Promise<ResolvedDigestTarget>>;
+ } | 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<string, Promise<unknown>>;
+ } | null;
+};
+
+type DigestEngine = ReturnType<typeof resolveDigestChunking> & {
+ 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<void> {
+ 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<OperationRun> {
+ 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";
+ // 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. An explicit
+ // `useClusters: false` still wins — this is the default, not a lock.
+ const sharesDuplicates = laneSharesDuplicates(
+ laneForOperation(DIGEST_OPERATION_ID) ?? {
+ queueKey: "",
+ contendsFor: "cpu",
+ },
+ );
+ const clusterPlan =
+ opts.useClusters === false || !sharesDuplicates
+ ? null
+ : (opts.clusterPlan ?? (await buildDigestClusterPlan(opts.paths)));
+ const chunking = resolveDigestChunking({ lane: laneChoice });
+ 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<ResolvedDigestTarget> {
+ 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<unknown> {
+ 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<OperationUnitPlan> {
+ 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;
+ // A plan already made by the caller's pull. Omitted → classified here.
+ // Supplying it does NOT skip the re-derivation: the pull and the run are
+ // separated by however long the pool held us, so this is re-asked either way.
+ plan?: OperationUnitPlan;
+};
+
+// 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<OperationUnitReport> {
+ 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<OperationRunOutcome> {
+ 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<OperationRunOutcome> {
+ 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<OperationRunOutcome | null> {
+ 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<string, Buffer> = {};
+ 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<OperationWorkCount> {
+ 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.
+ 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) {
+ const target = run.digest
+ ? await (async () => {
+ const r = await resolveDigestFor(run, channelSlug);
+ return { target: r.target, sections: r.sections };
+ })()
+ : await resolveTargetFor(run, op, channelSlug);
+ addOperationState(
+ counts,
+ await op.state({
+ videoDir,
+ videoId: id,
+ files,
+ target,
+ 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;
+};
+
+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,
+ };
+}
+
+type Candidate = { op: Operation; id: string; duration: number | null };
+
+export async function runOperationBatch(
+ opts: OperationBatchOptions,
+): Promise<OperationBatchResult> {
+ 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;
+
+ 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<string>();
+ 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<Candidate | null> => {
+ 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;
+ }
+ // 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<void> => {
+ 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 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<Candidate>({
+ 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) {
+ heldReason = null;
+ log(
+ run.digest
+ ? "Transcription finished; resuming the digest lane."
+ : "Resuming the backfill lane.",
+ );
+ }
+ 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 an armed
+// sweep, a hand-clicked channel run and an ids-scoped run all inherit the same
+// setting instead of three call sites each remembering to pass it.
+function laneOrderFor(run: OperationRun): AutoQueueOrder {
+ return run.digest
+ ? run.settings.digest.recencyOrder
+ : run.settings.backfill.order;
+}
+
+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;
+ }
+}
diff --git a/common/controller/operationJobs.ts b/common/controller/operationJobs.ts
@@ -40,12 +40,11 @@ import {
resolveQueueKey,
} from "../lib/queueKeys";
import { DIGEST_OPERATION_ID } from "../lib/operations";
-import { countBackfillWork, runBackfillBatch } from "./backfillBatch";
import {
- countMissingDigests,
- runDigestBatch,
+ countOperationWork,
+ runOperationBatch,
type DigestOrder,
-} from "./digestBatch";
+} from "./operationBatch";
import type { DigestLaneChoice } from "./digestTarget";
import { laneForOperation } from "./operationLane";
@@ -98,11 +97,15 @@ export async function runBackfillChannelJob(
params: { kindIds, queueKey: opts.queueKey, ids },
},
fn: async (onLog, signal, setProgress, ctx) => {
- const before = await countBackfillWork(paths, channelSlug, kindIds, ids);
- const batch = await runBackfillBatch({
+ const before = await countOperationWork("backfill", paths, channelSlug, {
+ operationIds: kindIds,
+ ids,
+ });
+ const batch = await runOperationBatch({
+ lane: "backfill",
channelSlug,
paths,
- kindIds,
+ operationIds: kindIds,
ids,
setProgress,
// The bar measures THIS run, from zero. Seeding it with what is already
@@ -221,22 +224,22 @@ export async function runDigestChannelJob(
// The LANE matters: the progress target must be computed against the
// engine that is actually about to run, not whichever one the local
// setting names.
- const missing = await countMissingDigests(
- paths,
- channelSlug,
- undefined,
- lane,
- );
+ const missing = (
+ await countOperationWork("digest", paths, channelSlug, {
+ digestLane: lane,
+ })
+ ).reachable;
// The bar measures THIS RUN, from zero, and the batch reports its own
// `current`. It used to be seeded with the count of digests already on
// disk and left to the UI's disk re-count — which cannot move during a
// REGENERATION, because a regenerated digest is rewritten in place. The
// bar sat at 0% for whole jobs ({initial:39, current:39, target:72}).
- const batch = await runDigestBatch({
+ const batch = await runOperationBatch({
+ lane: "digest",
channelSlug,
paths,
- lane,
- order: isDigestOrder(opts.order) ? opts.order : undefined,
+ digestLane: lane,
+ digestOrder: isDigestOrder(opts.order) ? opts.order : undefined,
// The local lane takes everything up to the long-tail cutoff; the
// metered lane exists for the tail above it. Passing no window (the
// default) runs the whole channel on one lane, which is what a
diff --git a/common/controller/remoteUnit.ts b/common/controller/remoteUnit.ts
@@ -2,7 +2,7 @@
// diarization, …) to a unit executor over HTTP. The clone of
// remoteTranscribe.ts for the general seam: submit the kind's input files and
// the primary's identity config, poll progress, pull back the kind's declared
-// output files, clean up the remote scratch. backfillBatch calls this when a
+// output files, clean up the remote scratch. the backfill lane calls this when a
// remote worker slot matching the kind is free.
//
// Failure classes reuse TranscribeError verbatim: a network/5xx/poll problem
diff --git a/common/controller/videoOperations.test.ts b/common/controller/videoOperations.test.ts
@@ -53,7 +53,7 @@ writeSettings(settingsOn());
const { inspectVideoOperations, orderForVideoPage, shownOnVideoPage } =
await import("./videoOperations");
-const { countBackfillWork } = await import("./backfillBatch");
+const { countOperationWork } = await import("./operationBatch");
const { getPaths } = await import("../lib/paths");
const { getSettings } = await import("../lib/settings");
@@ -232,9 +232,9 @@ test("shownOnVideoPage: off with nothing on disk hides; off with a sidecar shows
assert.equal(shownOnVideoPage(withRecord!), true);
});
-// Lives here rather than in backfillBatch.test.ts because countBackfillWork
+// Lives here rather than in operationBatch.test.ts because countOperationWork
// reads settings from disk and this file already owns the settings seam.
-test("countBackfillWork sizes an ids-scoped run to those ids alone", async () => {
+test("countOperationWork sizes an ids-scoped run to those ids alone", async () => {
writeSettings(
settingsOn({
backfill: { enabled: true, weight: 1, concurrency: 1 },
@@ -247,12 +247,15 @@ test("countBackfillWork sizes an ids-scoped run to those ids alone", async () =>
const paths = getPaths();
// Both are reachable for the text lane — its input is the cue stream, which
// every transcribed video has.
- const whole = await countBackfillWork(paths, SLUG, ["attribution-text"]);
+ const whole = await countOperationWork("backfill", paths, SLUG, {
+ operationIds: ["attribution-text"],
+ });
assert.ok(whole.reachable >= 2);
// Scoped, the progress target must count what THIS run will pull.
- const scoped = await countBackfillWork(paths, SLUG, ["attribution-text"], [
- idB,
- ]);
+ const scoped = await countOperationWork("backfill", paths, SLUG, {
+ operationIds: ["attribution-text"],
+ ids: [idB],
+ });
assert.equal(scoped.reachable, 1);
});
diff --git a/common/lib/operations.test.ts b/common/lib/operations.test.ts
@@ -28,7 +28,6 @@ import {
backfillLaneOperationEntriesOf,
presentOperationWork,
} from "./operations";
-import { candidateAction } from "../controller/backfillBatch";
import {
readVideoFiles,
CUES_JSON_FILENAME,
@@ -774,25 +773,13 @@ test("attribution-diarized: no diarization.json is BLOCKED, not missing-input",
);
});
-test("blocked is never dispatched, and re-download cannot change that", () => {
- // The dispatch decision is the consequential one: `blocked` must not reach a
- // runner whatever the flags say. Note allowRedownload — the flag that DOES
- // turn missing-input into a dispatch — is deliberately inert here.
- for (const force of [false, true]) {
- for (const allowRedownload of [false, true]) {
- assert.equal(
- candidateAction("blocked", { force, allowRedownload }),
- "blocked",
- `force=${force} allowRedownload=${allowRedownload}`,
- );
- }
- }
- // The contrast, so this test fails if the two ever get conflated again.
- assert.equal(
- candidateAction("missing-input", { force: false, allowRedownload: true }),
- "dispatch",
- );
-});
+// THE DISPATCH HALF OF THIS PROPERTY LIVES IN
+// controller/operationBatch.test.ts ("blocked is never dispatched, and
+// re-download cannot change that"). It used to live here and import
+// candidateAction across the layer boundary — a lib/ test reaching into
+// controller/, carried on architecture.test.ts's allow-list as a debt. The
+// classification is this file's business; what a dispatcher does with it is
+// the dispatcher's.
test("blocked is counted, and is NOT reachable work", () => {
const counts = emptyOperationCounts();
@@ -1210,16 +1197,17 @@ test("digest declares its own lane, and it is NOT the backfill queue", () => {
});
test("the backfill lane never dispatches digest", () => {
- // backfillLaneOperations feeds backfillBatch, the channel Backfill card and the
- // dashboard instrument. Digest must be in the CATALOG and out of THAT list,
- // or it both serializes behind diarization and gets double-counted.
+ // backfillLaneOperations feeds the backfill lane's executor, the channel
+ // Backfill card and the dashboard instrument. Digest must be in the CATALOG
+ // and out of THAT list, or it both serializes behind diarization and gets
+ // double-counted.
const settings = settingsWithDiarization();
const laneIds = backfillLaneOperations(settings).map((k) => k.id);
const allIds = allOperations(settings).map((k) => k.id);
assert.ok(allIds.includes("digest"), allIds.join(","));
assert.ok(!laneIds.includes("digest"), laneIds.join(","));
- // Which means backfillBatch's kind resolution cannot reach it either, even
- // when it is asked for by name.
+ // Which means the backfill lane's operation resolution cannot reach it
+ // either, even when it is asked for by name.
assert.deepEqual(resolveBackfillLaneOperations(settings, ["digest"]), []);
});
diff --git a/common/lib/operations.ts b/common/lib/operations.ts
@@ -17,7 +17,7 @@
// and the indicator.
//
// WHAT "BACKFILL" STILL MEANS IN THIS FILE. The backfill LANE — BACKFILL_QUEUE,
-// backfillLaneOperations, backfillLaneEntriesOf, controller/backfillBatch.ts and
+// backfillLaneOperations, backfillLaneEntriesOf, controller/operationBatch.ts and
// controller/backfillSweep.ts — is ONE QUEUE that several operations share: one
// pause, one sweep, one share. It is the only thing this file still calls
// "backfill", and its persisted contracts (the `backfill` key on the snapshot,
@@ -435,14 +435,21 @@ export type Operation = {
// them. Mirrors digestLaneFor, which solved the same problem for the digest
// operation's two lanes — and which digest itself now declares.
//
- // OPTIONAL IS LOAD-BEARING, not tidiness. controller/backfillBatch.ts reads
- // this field's PRESENCE as the marker for "this kind's resource depends on
- // settings, so nothing else is deciding it for us" and makes the whole run
- // idle-only when such a kind could take the GPU. Giving every kind a laneFor
- // that defaults to `lane` would therefore not be a no-op: it would enrol
- // every kind in that rule and make the lane idle-only whenever a statically
- // GPU-bound kind was in the run — the exact regression the guard's comment
- // records. Add one only where the lane genuinely varies.
+ // ITS PRESENCE IS NO LONGER LOAD-BEARING, and that changed in phase 1 slice
+ // 1.2. It used to be: `controller/backfillBatch.ts` keyed the GPU idle-only
+ // rule off `k.laneFor &&`, i.e. off this field EXISTING, as a proxy for "this
+ // kind's resource depends on settings, so nothing else is deciding it for
+ // us". That made "give every kind a laneFor defaulting to `lane`" a
+ // behaviour change wearing a refactor's clothes — it would have enrolled
+ // every kind and made the lane idle-only whenever a statically GPU-bound kind
+ // was in the run.
+ //
+ // `controller/operationBatch.ts` now asks the DECLARATION instead:
+ // `laneYieldsToTranscription(laneForOperation(op.id) ?? op.lane)`, i.e. the
+ // resolved lane's `contendsFor`. A kind with a fixed GPU lane and no laneFor
+ // is caught (it was not before), and adding a laneFor for any other reason
+ // changes nothing. Optional, then, purely because most kinds have no choice
+ // and `lane` is the whole answer for them.
laneFor?(settings: SiteSettings): Lane;
// Ids of other kinds in this table whose output this one consumes.
//
@@ -1056,7 +1063,7 @@ function toBackfillOutcome(outcome: AttributeOneOutcome): OperationRunOutcome {
// NOTHING ABOUT DIGEST GENERATION IS REWRITTEN HERE. Every function this entry
// calls is the one the digest controller already calls — resolveDigestTarget for
// the identity, isCuesJsonFresh for the transcript gate, isSectionFresh for
-// freshness, digestVideo to do the work. If this entry and digestBatch ever
+// freshness, digestVideo to do the work. If this entry and the lane's executor ever
// disagree about whether a video is digested, that is a bug in this file, not a
// second opinion.
//
@@ -1072,7 +1079,7 @@ function toBackfillOutcome(outcome: AttributeOneOutcome): OperationRunOutcome {
// 3. ONE DEFINITION OF "digested". The snapshot, the planner and the batch all
// derive it; this is where they can converge.
//
-// IT IS DELIBERATELY NOT RUN BY backfillBatch — see backfillQueueKinds below.
+// IT IS DELIBERATELY NOT RUN ON THE BACKFILL LANE — see backfillQueueKinds below.
const digest: Operation = {
id: "digest",
label: "Digest",
@@ -1101,17 +1108,18 @@ const digest: Operation = {
// arbiter dispatches a remote digest onto DIGEST_REMOTE_QUEUE instead of the
// local key it declares.
//
- // Declaring it does NOT put digest into backfillBatch's idle-only rule, and
- // the reason is worth stating because that rule keys off laneFor's PRESENCE:
- // resolveBackfillLaneOperations is filtered through backfillLaneOperations (BACKFILL_QUEUE
- // only), so digest can never be among the `kinds` that guard inspects. See
- // controller/backfillBatch.ts, where the same fact is written from the other
- // side.
+ // Declaring it does NOT put digest into the backfill lane's idle-only rule.
+ // That used to need saying twice, because the rule keyed off laneFor's
+ // PRESENCE; it now keys off the resolved lane's `contendsFor`, and digest is
+ // still not in it for the reason that was always the real one:
+ // `operationsForLane("backfill")` admits BACKFILL_QUEUE only, so digest can
+ // never be among the operations that guard inspects. That is the invariant
+ // operations.test.ts pins.
laneFor: (settings) =>
digestLaneFor(settings.digest.remoteEnabled ? "remote-api" : "local-gpu"),
// Digests are gated by their own sweep/pause switches rather than a master
// "enabled" flag, so the feature is on whenever an app is configured. The
- // pause is honoured at DISPATCH (digestBatch's limit()), not here: a paused
+ // pause is honoured at DISPATCH (the digest lane's limit()), not here: a paused
// lane must still report how much work is outstanding.
settingsBlock: "digest",
enabled: () => true,
@@ -1248,7 +1256,7 @@ type DigestTarget = {
// The digest operation has TWO lanes, and which one a run uses follows from the
// engine it is configured with rather than from a separate setting. This is the
-// declaration; digestBatch consults it instead of re-testing the app id, so
+// declaration; the digest lane consults it instead of re-testing the app id, so
// "which lane must stand aside for transcription" is stated once.
//
// The queue keys must never be shared: registry.ts runs each key at concurrency
@@ -1309,7 +1317,7 @@ export function laneYieldsToTranscription(lane: Lane): boolean {
// worker pool provides and transcription requires.
// - A SCALAR AGAINST A POOL. backfillLimit() returns one number; the worker
// pool is per-worker with individual enable flags and priorities.
-// - PER-CHANNEL VS CROSS-CHANNEL SCOPE. backfillBatch is one job per channel;
+// - PER-CHANNEL VS CROSS-CHANNEL SCOPE. a backfill channel job is one job per channel;
// autoRunner arbitrates across every channel at once, which is the whole
// point of its policy tree.
//
@@ -1529,7 +1537,7 @@ export const OPERATIONS: readonly Operation[] = [
diarization,
attributionDiarized,
// Text-only LAST, deliberately. resolveBackfillLaneOperations preserves this order and
- // backfillBatch walks the kinds in it, so on a video that has diarization the
+ // the backfill lane walks the operations in it, so on a video that has diarization the
// cheap, better lane gets there first and the text lane then finds a record it
// must not overwrite — one wasted classification instead of ~30 model calls.
attributionText,
@@ -1557,12 +1565,12 @@ export function allOperations(settings: SiteSettings): Operation[] {
// BACKFILL_QUEUE. Three filters, and the third is new with the digest entry.
//
// The queue filter is a SAFETY RAIL, not a tidy-up. Everything downstream of
-// this function — backfillBatch's dispatch, the channel Backfill card, the
+// this function — the backfill lane's dispatch, the channel Backfill card, the
// dashboard's backfill instrument, the /channels speaker columns — treats
// these as "one lane, one job, one set of counters". Digest satisfies none of
// that:
//
-// - DISPATCH. backfillBatch runs its kinds in one job on one queue under
+// - DISPATCH. the backfill lane runs its operations in one job on one queue under
// backfillLimit(). Handing it digest would SERIALIZE the GPU digest lane
// behind CPU diarization, when the entire reason they hold separate queue
// keys is that they currently overlap.
@@ -1570,7 +1578,7 @@ export function allOperations(settings: SiteSettings): Operation[] {
// carve-out (with its CPU-worker exemption), spendCapUsd on the metered
// lane, the remoteEnabled fail-fast, shortest-first ordering,
// duplicate-cluster sharing and the engine probe() fail-fast. Those are
-// measured decisions in digestBatch's limit(), and not one of them is
+// measured decisions in the digest lane's limit(), and not one of them is
// expressible as backfillLimit()'s single scalar.
// - COUNTERS. Digest already has its own instrument, its own stage card and
// its own snapshot bucket. Folding it in here would double-count it against
@@ -1695,7 +1703,7 @@ export function resolveBackfillLaneOperations(
// Order kinds so a prerequisite is attempted before anything that declares it.
//
-// WHY THIS IS WORTH DOING AT ALL. backfillBatch walks the kinds in the order it
+// WHY THIS IS WORTH DOING AT ALL. the backfill lane walks the operations in the order it
// is given, all the way through the video list, before starting the next kind.
// So with `attribution-diarized` ahead of `diarization`, a video diarized
// during a pass becomes eligible for attribution only on whatever LATER pass
@@ -1814,7 +1822,7 @@ export function addOperationState(
// DELIBERATELY UNCHANGED by the addition of `deferred`, and unchanged again by
// `blocked`. This function is the guard: adding a state to the union raises no
// TypeScript error here (the exhaustiveness check lives on the DISPATCH
-// decision, in backfillBatch's candidateAction, which is the branch that can do
+// decision, in operationBatch's candidateAction, which is the branch that can do
// harm), so the only thing keeping capped and blocked videos out of the work
// total is that they are not added here. If a future state belongs in the
// total, it goes in on purpose.
diff --git a/common/lib/queueKeys.ts b/common/lib/queueKeys.ts
@@ -18,7 +18,7 @@ export { TRANSCRIPTION_QUEUE };
export const DIGEST_LOCAL_QUEUE = "digest:local";
export const DIGEST_REMOTE_QUEUE = "digest:remote";
-// The backfill lane (controller/backfillBatch.ts). Its own key for exactly the
+// The backfill lane (controller/operationBatch.ts). Its own key for exactly the
// reason the two digest keys have theirs: registry.ts submits every non-empty
// queueKey at concurrency 1, so a distinct key is the ONLY way to get a lane
// that runs CONCURRENTLY with transcription and the digest lanes rather than
diff --git a/common/lib/settings.ts b/common/lib/settings.ts
@@ -243,7 +243,7 @@ export type SiteSettings = {
// CHUNK, which on this corpus is ~194,000 calls, the same order as the digest
// sweep — and the digest sweep has completed 0.17% of its own. Arming both at
// once puts two multi-week lanes on the same 8 GB card with nothing to arbitrate
-// between them (backfillBatch's yield deliberately watches only the
+// between them (the backfill lane's yield deliberately watches only the
// transcription lane). Nothing here arms anything; a pilot decides whether the
// corpus-wide text-only pass is worth 25-55 GPU-days at all.
export type AttributionSettings = {
@@ -304,7 +304,7 @@ export type BackfillSettings = {
// >0 — a guaranteed share of the lane's slots, floored at 1 so a small
// weight is a slow lane rather than a stopped one.
//
- // Clamped to [0, 1]. See backfillLimit() in controller/backfillBatch.ts.
+ // Clamped to [0, 1]. See backfillLimit() in controller/operationBatch.ts.
weight: number;
// Slots the lane may use when it is not standing aside. Kept at 1 by default
// for the same reason diarization.concurrency is: this is CPU-bound work
diff --git a/editor/app/channels/[slug]/digestActions.ts b/editor/app/channels/[slug]/digestActions.ts
@@ -9,9 +9,9 @@ import {
resolveQueueKey,
} from "yt-dlp-transcript-common/lib/queueKeys";
import {
- countMissingDigests,
- runDigestBatch,
-} from "yt-dlp-transcript-common/controller/digestBatch";
+ countOperationWork,
+ runOperationBatch,
+} from "yt-dlp-transcript-common/controller/operationBatch";
import {
runManagedFunction,
type StreamActionResult,
@@ -109,11 +109,17 @@ export async function digestBucketAction(
// Same shape as the channel action: the bar measures this run from zero
// and the batch reports its own `current`, because the disk re-count it
// would otherwise use cannot see a digest rewritten in place.
- const missing = await countMissingDigests(paths, slug, cleaned, lane);
- const result = await runDigestBatch({
+ const missing = (
+ await countOperationWork("digest", paths, slug, {
+ ids: cleaned,
+ digestLane: lane,
+ })
+ ).reachable;
+ const result = await runOperationBatch({
+ lane: "digest",
channelSlug: slug,
paths,
- lane,
+ digestLane: lane,
ids: cleaned,
setProgress,
progressBaseline: 0,
diff --git a/editor/app/components/pipelines/buildBands.ts b/editor/app/components/pipelines/buildBands.ts
@@ -120,7 +120,7 @@ function addExternalBands(
]);
download.present = sumOrNull([download.present, totals.downloaded]);
// Partial downloads are reachable work like any other — the same rule
- // backfillBatch applies to `partial`.
+ // the dispatch decision applies to `partial`.
download.reachable +=
undownloaded.filter((id) => !excluded.has(id)).length +
(buckets?.partialDownloads?.length ?? 0);
diff --git a/editor/e2e/backfill.spec.ts b/editor/e2e/backfill.spec.ts
@@ -87,7 +87,7 @@ function backfillSettings(over: {
// 1, not the idle-only 0: these specs assert that work HAPPENS, and at
// weight 0 a transcription running in the same test would legitimately
// park the lane forever. The idle-only default is covered by the pure
- // unit test (common/controller/backfillBatch.test.ts), which is the right
+ // unit test (common/controller/operationBatch.test.ts), which is the right
// place for it — no pool, no GPU, no timing.
weight: 1,
concurrency: 1,
@@ -782,7 +782,7 @@ test("the diarization page shows the channel and both populations", async ({
});
// PAUSE IS A HOLD, NOT A STOP, and it is reachable from where the work is
-// watched. The hold has always existed — backfillBatch's limit() re-reads
+// watched. The hold has always existed — the backfill lane's limit() re-reads
// settings.backfill.enabled at dispatch and returns 0, so the pool idle-waits
// and the job keeps its place — but the only way to set it was the Settings
// page, which is a strange place to look for a control over a job you are
diff --git a/plans/FACTS.md b/plans/FACTS.md
@@ -2095,18 +2095,25 @@ verification, not lint.
Facts a cold agent would otherwise re-derive from the slice-2 commits (`1da6f22` → `1058b98`).
Commit bodies carry the reasoning.
-**`BackfillKind.laneFor?` is OPTIONAL AND ITS PRESENCE IS SIGNIFICANT** (`common/lib/backfillKinds.ts:387`).
-`common/controller/backfillBatch.ts` keys the GPU idle-only rule off `k.laneFor &&`, i.e. off
+**`Operation.laneFor?` is OPTIONAL AND ITS PRESENCE IS NO LONGER SIGNIFICANT** — CHANGED
+2026-09-07, phase 1 slice 1.2 (`common/lib/operations.ts`). It used to be:
+`common/controller/backfillBatch.ts` keyed the GPU idle-only rule off `k.laneFor &&`, i.e. off
the field EXISTING, as the marker for "this kind's scarce resource depends on settings, so
-nothing else is deciding it". "Give every kind a `laneFor` defaulting to `lane`" is therefore
-NOT a no-op refactor: it enrols every kind in that rule and makes the backfill lane idle-only
-whenever a statically GPU-bound kind is in the run — the exact regression the guard's comment
-records. Add one only where the lane genuinely varies. Digest and diarization declare one; no
-other kind does.
-
-**Digest declaring a `laneFor` does not put it in that guard**, and not because of the guard:
-`resolveBackfillKinds` filters through `laneBackfillKinds` (BACKFILL_QUEUE only), so digest can
-never reach `backfillBatch`'s `kinds` array even when asked for by name. `backfillKinds.test.ts`
+nothing else is deciding it" — which made "give every kind a `laneFor` defaulting to `lane`" a
+behaviour change wearing a refactor's clothes.
+
+`common/controller/operationBatch.ts` (which replaced both batch files) asks the DECLARATION
+instead: `laneYieldsToTranscription(laneForOperation(op.id) ?? op.lane)`, keyed on
+`contendsFor`. So a kind with a FIXED GPU lane and no `laneFor` is now caught (it was not
+before), and adding a `laneFor` for any other reason changes nothing. Behaviour on today's
+registry is identical — no BACKFILL_QUEUE operation declares a static GPU lane
+(`operations.ts` has diarization `cpu` with a `laneFor`, both attribution ops `network`).
+`common/controller/operationBatch.test.ts` ("the GPU carve-out is keyed on contendsFor, not on
+laneFor existing") pins the new rule.
+
+**Digest declaring a `laneFor` still does not put it in that guard**, and never did because of
+the guard: `operationsForLane("backfill")` admits BACKFILL_QUEUE only, so digest can never
+reach the backfill run's operation array even when asked for by name. `operations.test.ts`
pins it.
**`laneForOperation` lives in `common/controller/operationLane.ts`**, alone, because the arbiter
@@ -2126,7 +2133,9 @@ second, redundant gate that label-only callers would read as "no lane at all".
**`common/controller/operationJobs.ts` is the ONE per-channel runner path** — a hand-clicked run,
a swept one and an arbiter-dispatched one are the same function. `runBackfillChannelJob`,
-`runDigestChannelJob`, and `runOperationChannelJob(operation, …)` over both.
+`runDigestChannelJob`, and `runOperationChannelJob(operation, …)` over both. Since slice 1.2
+all three drive `controller/operationBatch.ts`, which is also what the auto-queue runner's
+digest and backfill lanes dispatch through — one executor, two callers.
**They do NOT drain.** They return the `StreamActionResult` and the three callers that want
sequencing (`backfillSweep`, `digestSweep`, `arbiter`) `await drainStream(result.stream)`