// The single source of truth for what an OPERATION is. // // An operation is a derived-data feature that declares what it needs — its // inputs, its per-video state probe, what one video of it costs — and the // system supplies the lane it runs on, the resource share and the indicator. // // The problem this exists for repeats: a derived-data feature lands, and the // corpus that already exists does not have what it needs. For diarization that // input is AUDIO, which cleanAudioFromTranscribed deletes once a video is // transcribed. Before this table the only catch-up was a per-channel button on // a controller written for that one feature (controller/diarizeAll.ts) — there // was no way to ask "how much of the corpus is missing this?", and no way to run // catch-up alongside new-video work without one starving the other. // // So, in the shape jobKinds.ts already uses: adding an operation should mean // adding ONE entry here, and the system supplies the lane, the resource share // and the indicator. // // WHAT "BACKFILL" STILL MEANS IN THIS FILE. The backfill LANE — BACKFILL_QUEUE, // backfillLaneOperations, backfillLaneEntriesOf and controller/operationBatch.ts // — is ONE QUEUE that several operations share: one pause, one runner, one // share. It is the only thing this file still calls "backfill", and its // persisted contracts (the `backfill` key on the snapshot, `settings.backfill`, // the `backfill-channel` job kind) keep the word because they are on disk. An operation is not a backfill; a // backfill is what the lane does to the corpus that predates an operation. // Anything named `*Kind*` outside this file's exports means "one entry of this // registry", which stays true. // // FOUR STATES, NOT TWO, and the split is the load-bearing part. Measured on this // corpus at the time of writing: 77,106 videos, 836 with media still on disk, 1 // diarized. A single "remaining" number would therefore read 77,105 — and 91x of // that is unreachable without re-downloading. The repo has already been burned by // exactly this once: common/views/widgetActionable.ts deliberately // refuses to filter on the digest work count (historically the `noDigest` // bucket) because during the backfill that is 99.87% of the corpus and counting // it would put every channel in the list forever. So // `missing` (reachable now) and `missing-input` (needs re-acquiring) are // SEPARATE numbers, everywhere, and no surface is allowed to add them together. // // STALENESS IS PROVENANCE, NOT AGE. `state` is derived from disk on every read // and never stored, and an operation that records what produced its output // compares that against what we would produce now — lib/digest.ts's // isSectionFresh, whose absent-field-equals-today's-default trick is what stops // adding a field from invalidating the whole corpus. Diarization writes that // provenance and, before this, had no comparator at all: diarizeOne // short-circuited on mere existence, so the two most likely reasons to re-run (a // threshold or model change) left everything looking done. // // SERVER-ONLY, despite living in lib/. It reads the filesystem and calls a // controller, so it is `-server.ts` in everything but name; the path is the one // the plan named. No client component imports it — the UI is handed plain // numbers off the channel snapshot, and labels as props. // // THREE ENTRIES, AND THE SECOND PAIR IS WHAT MAKES THIS AN ABSTRACTION. A // registry with one entry is a wrapper: nothing proved that "a feature declares // what it needs and the system supplies the lane, the share and the indicator" // was true. Attribution is the test of it, and it passed — registering // `attribution-diarized` and `attribution-text` lit the channel stage card, the // /channels speaker columns, the dashboard instrument and the widget strip with // ZERO UI changes, because all four iterate snapshot.backfill[operationId]. // // It also exercised the parts of the shape that one entry could not: // `missing-input` for something other than audio (diarization.json, of which // this corpus has one), and two operations writing the SAME FILE at different // quality tiers — see the ordering rule above the attribution entries. // // WHAT IS STILL NOT HERE. The other catch-up mechanisms in the repo do not fit // this per-video probe, and forcing them in would make the table lie: // controller/backfillAvailability.ts is CHANNEL-scoped (one JSON map, folded // into sync at runYtdlp.ts), and controller/normalizeAll.ts has no recorded // provenance to compare, so its "stale" is undefined. `tier` still exists // because it is what keeps them apart if they are ever added: `inline` folds // into an existing pass and `lane` gets the concurrent queue and the share. import type { Paths } from "./paths"; import type { AttributionSettings, DiarizationSettings, SiteSettings, } from "./settings"; import { BACKFILL_QUEUE, DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE, TRANSCRIPTION_QUEUE, } from "./queueKeys"; // The runner union, from the model layer. It used to come from // jobs/autoQueueState.ts as a type-only import; it now lives in lib/ outright, // which is both cycle-free and legal under the layering guard. import type { AutoQueueKind } from "./autoQueueTypes"; import { DIGEST_FILENAME, digestSectionStates, isDigestSectionKind, type DigestAppConfig, type DigestFreshnessTarget, type DigestItem, type DigestLane, type DigestProvenance, type DigestRecord, type DigestSectionKind, type DigestWarning, } from "./digest"; import { loadDigest, writeDigestSection } from "./digest-server"; import type { DigestContext } from "./digestContext-server"; import { SORTFORMER_DIARIZATION_ENGINE, diarizationTarget, isDiarizationFresh, type DiarizationBackend, type DiarizationEngineId, type DiarizationFreshnessTarget, type DiarizationRecord, } from "./diarization"; import { loadDiarization, writeDiarization } from "./diarization-server"; import { ATTRIBUTION_FILENAME, isAttributionDowngrade, isAttributionFresh, transcriptSourceOf, type AttributionFreshnessTarget, type AttributionMethod, type AttributionRecord, type AttributionTranscriptSource, } from "./attribution"; import { loadAttribution, writeAttribution } from "./attribution-server"; import { CUES_JSON_FILENAME, DIARIZATION_FILENAME, META_FILENAME, WHISPER_FILENAME, findSourceMedia, isVideoTranscribed, pickIndexTranscript, readVideoDurationSec, type VideoFiles, } from "./videoStatus"; import { pickPreferredAudio } from "./mediaFiles"; import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo"; import { resolveSavedVideo } from "./savedVideo-server"; import { diarizeOneVideo } from "../controller/diarizeOne"; // The TARGET resolver only — a settings read plus the digest app registry. // attributeOne is loaded LAZILY inside run() below: this module is imported by // controller/channelSnapshot.ts, which classifies every video of every channel, // so its eager import graph sits on the editor's hot path, and a classification // needs none of the runner's (the transcript normalizer, the markdown renderer, // the digest prompt module, the channel-context reader). See // controller/attributionTarget.ts for what this is and is not worth — it is a // structural argument, not a measured speedup. import { resolveAttributionTarget } from "../controller/attributionTarget"; import type { AttributeOneOutcome } from "../controller/attributeOne"; import { isCuesJsonFresh } from "../controller/normalizeTranscript"; // What a video's relationship to a backfill is, right now, read from disk. // // present — has it, at the identity we would produce now. // stale — has it, but from a different engine/model/threshold. // partial — has SOME of it at the current identity and not the rest. // Only meaningful for a kind whose output has parts; today // that is digest alone, whose sections are generated and // compared independently. See the digest entry for why this // is not just a nicer word for `stale`. // missing — does not have it, and the input to produce it is HERE. // missing-input — does not have it, and the input is gone. Reachable only by // re-acquiring the media, which is opt-in and bounded. // deferred — does not have it, the input is here, and the kind refuses to // attempt it under the current configuration. Today that is // only the diarization duration cap. NEVER summed into // reachable work, so a capped corpus cannot read as finished. // blocked — does not have it, and what it is waiting for is the OUTPUT OF // ANOTHER KIND IN THIS TABLE. Not the same thing as // missing-input, and conflating them was a real defect: see // below. // // `deferred` and `blocked` follow the house rule OperationRunOutcome = "skipped" // already sets on the run side: a deliberate non-action gets its own counter, // and is never folded into the work total nor reported as a failure. // // WHY `blocked` IS NOT `missing-input`. `missing-input` means one specific // thing to the rest of the system: THE MEDIA IS GONE, RE-ACQUIRE IT. It is the // population `allowRedownload` exists for, and backfillReacquire answers it by // fetching AUDIO. `attribution-diarized` waits on diarization.json — so // reporting that as missing-input told the operator ~73,000 videos needed media // re-fetched, and with re-download on, the lane would spend a download per video // fetching audio that CANNOT satisfy the wait, then cleaning it up again. The // prerequisite is not gone; it has not been produced yet, and this same table // knows how to produce it. export type OperationState = | "present" | "stale" | "partial" | "missing" | "missing-input" | "deferred" | "blocked"; // Videos this backfill has no opinion about (not transcribed, marked // untranscribable). Kept out of OperationState so it can never be counted. export type OperationClassification = OperationState | "not-applicable"; // How expensive one video is, which decides where the work runs. // // inline — microseconds to cheap I/O; folds into a pass that already walks the // corpus, and never gets a lane of its own. // lane — expensive enough to need its own queue and a resource share. // Diarization is ~500-680 s/audio-hour of CPU. export type OperationTier = "inline" | "lane"; // WHERE an operation's work runs, and what it fights with while it runs. // // This exists because "put every backfill on BACKFILL_QUEUE" is wrong, and // measurably so. registry.ts submits every non-empty queueKey at concurrency 1, // so a shared key SERIALIZES. Diarization is CPU and digest is GPU/ollama; // today they run at the same time on separate queues, and collapsing them onto // one key would idle the GPU while the CPU works and vice versa, across a // multi-week sweep. So the queue is declared per operation rather than owned by // the lane, and a scheduler dispatches one job per distinct queueKey. export type Lane = { // The registry queue key. Distinct keys are the ONLY mechanism for // concurrency between operations; the same key is the only mechanism for // serializing an operation against itself. queueKey: string; // The scarce resource one unit of this work occupies. Not decoration: it is // what decides whether the lane must stand aside for transcription. // // gpu — contends with the transcription engine for VRAM. Yields. // cpu — contends for cores. Does not yield today; its share is governed // by the backfill lane's own weight (idle-only at the default). // network — a metered or remote API. Contends with nothing local, so it // must NOT yield: doing so would park a lane that was costing // nothing to keep running. contendsFor: "gpu" | "cpu" | "network"; }; // Whether this system dispatches the operation, or merely knows about it. // // "backfill" — the backfill machinery pulls candidates and runs it. // "external" — something else owns dispatch (autoRunner, the worker pool). // Registered anyway so the catalog is complete and a dependency // on it resolves. See EXTERNAL_OPERATIONS. export type OperationDispatch = "backfill" | "external"; // WHAT AN OPERATION ACTS ON, and WHAT MAKES IT RUN. Every media-derived // operation is per-video and backlog-driven: it has a state per video, and its // work list is whichever videos lack it. Sync is the one exception the IA doc // names (plans/editor-operations-ia.md, "Where the noun breaks"): a CHANNEL has // a sync state — last synced, next due — and nothing runs it but a cadence. A // surface that draws per-video things (a band, a video-page panel, a channel // column, a worker tag) filters on `scope`, never on the id. export type OperationScope = "video" | "channel"; export type OperationTrigger = "backlog" | "cadence"; export type OperationProbe = { videoDir: string; videoId: string; // Already read by the caller. Taking it rather than re-reading is what makes // the channel snapshot's per-video classification free — see VideoFiles.entries. files: VideoFiles; target: unknown; // The live settings. Passed in rather than read here so a classification stays // a pure function of what the caller already has, and so a kind can consult a // knob that must NOT become part of its freshness identity — the diarization // duration cap is exactly that: `target` is compared by isDiarizationFresh, so // putting the cap there would mark every sidecar on disk stale the moment the // cap moved. settings: SiteSettings; }; export type OperationRunOptions = { paths: Paths; videoDir: string; videoId: string; channelSlug: string; target: unknown; // Redo a `present` video anyway (an operator forcing a regeneration). force?: boolean; // Engine-config override for the LLM-backed kinds (attribution today). The // fan-out passes the primary's resolved config with only `baseUrl` swapped // to a leased endpoint — baseUrl is not part of the freshness identity, so // which endpoint served a call can vary freely while everything that IS // identity (model, numCtx, timeouts) travels verbatim from the primary. appConfig?: DigestAppConfig; // Per-kind settings injection for a UNIT EXECUTOR (see workerServer's // startWorkerUnit). The primary resolves these and ships them in the unit // envelope; each kind's run() threads its own field into the controller's // existing `settings` override — so a bare executor's default settings // (attribution disabled, empty app config = a DIFFERENT identity) can never // leak into provenance. On the primary these stay unset and the controllers // read live settings exactly as before. attributionSettings?: AttributionSettings; diarizationSettings?: DiarizationSettings; // Channel context, injected rather than shipped as a file: the primary // already computed the note + hash, and digest pins contextHash in its // identity — injecting the object makes hash equality true by construction. digestContext?: DigestContext; onLog?: (msg: string) => void; signal?: AbortSignal; }; // Deliberately mirrors DiarizeOneOutcome's discipline: an expected condition is // an outcome, never a throw. The batch counts these and keeps going. export type OperationRunOutcome = | "done" | "already-present" | "missing-input" // Nothing to do for a reason state() could not see WITHOUT a per-video file // read — the case that matters is a transcript whose cues.json is stale, which // costs a read of the source file and cannot be paid 77,000 times per pass. // Counted separately from `failed` so a transient condition that resolves // itself does not report as a broken engine. | "skipped" | "not-configured" | "disabled" | "failed"; // WHICH PIPELINE THIS OPERATION BELONGS TO, as a declared field rather than a // list somebody maintains in a component. // // Three surfaces need this same grouping and each used to hardcode its own copy: // the channel transit line (which station does this operation live under), the // /channels table (which columns sit together), and the lane card (which // operations does this queue actually hold). Three hardcoded lists is three // places to forget when a kind is added — and the last time one was added, the // transit line kept summing three unrelated operations into one station because // nobody updated its list. // // Deriving the label from the group also means a lane holding a MIX cannot go // stale: it falls back to "Derived data" rather than naming two of its three // members. export type OperationGroup = "media" | "transcript" | "digest" | "speakers" | "sync"; // See Operation.settingsBlock. A closed union rather than `string` so the switch // that picks a settings form is exhaustive: a new block has to be named here and // then handled, instead of silently rendering nothing. export type OperationSettingsBlock = | "digest" | "diarization" | "attribution" // The WHOLE syncScheduler block — headline controls and the advanced knobs // alike. An operation's settings live on its page (slice 3's rule), and // splitting this one across two pages is what gave it two writers. | "syncScheduler"; // Group order: upstream first. The /channels columns and the transit line both // lay their pipelines out in this order, so a reader moving between the two // pages sees the same left-to-right sequence. // // `sync` is deliberately NOT here: this list is the per-video pipeline's // stations, and sync draws no channel column and no station of its own (its // channel surface is the hand-listed `playlist` bookend — see GROUP_STAGES). export const OPERATION_GROUP_ORDER: readonly OperationGroup[] = [ "media", "transcript", "digest", "speakers", ]; export function groupLabel(group: OperationGroup): string { switch (group) { case "media": return "Media"; case "transcript": return "Transcript"; case "digest": return "Digest"; case "speakers": return "Speakers"; case "sync": return "Sync"; } } export type Operation = { id: string; label: string; // One line of UI copy: what this backfill is, in the operator's terms. hint: string; // The pipeline this operation belongs to. See OperationGroup. group: OperationGroup; // The label at COLUMN width — one or two words, for a header that has to sit // above a 48px band on a table 68 rows deep. // // Declared rather than abbreviated in the component, for the same reason // `group` is: "Speaker names (from the transcript)" cannot be shortened // mechanically, and a map of abbreviations maintained next to a table is a // second place to forget when a kind is added. shortLabel: string; // WHAT ONE VIDEO OF THIS OPERATION COSTS, as a cost basis — "one audio pass // per video", "~1 model call per transcript chunk". // // "UNIT" IS NOT THE WORD FOR THIS, and the rule is one-sided. A UNIT is one // item of dispatchable work: what runOperationUnit, startWorkerUnit, // runUnitViaRemote and /api/worker/unit move around, and three of those are // persisted contracts (the route, the `worker-unit` and `auto-download-unit` // job kinds). What one video of an operation costs is its COST BASIS, and no // surface calls that a unit — a prop that once said `unit` for a population // label ("reachable", "blocked") is `population`, and prose that said "unit" // for cost says "cost basis". The dispatch side keeps its name. // // The fact this exists to surface: an operation can be armed at enormous cost // and read as a quiet row. attribution-text is reachable on 11,337 videos of // one channel and roughly 194,000 model calls corpus-wide, it has completed // ONE video, and every screen that mentioned it said only "Backfill". Printing // the cost basis beside the backlog is what makes 11,337 legible. // // Deliberately NO threshold and no editorialising. A "this is a lot" cutoff // would be a magic number the next operation gets wrong, and the operator is // the one who decides what is too expensive. costBasis: string; // What a `deferred` video of THIS kind is waiting for, and what an operator // can do about it. Belongs to the kind, not to the card: SpeakersStage used to // hardcode "too long to diarize under the current limit", which was correct // only while diarization was the sole kind that could defer. Digest defers for // an unrelated reason (no current normalized transcript), so a card summing // several kinds' `deferred` into one hardcoded sentence now states a cause // that is false for most of what it counts. // // Written as a sentence FRAGMENT completing "N videos are …", so the card // keeps ownership of the count and its pluralization. deferredHint?: string; tier: OperationTier; // Where this operation's work runs. See Lane — the queue key is what // keeps CPU and GPU operations overlapping instead of taking turns. // // This is the DECLARED lane, which for a kind with a choice means its default. // Prefer laneFor() when a live answer is needed. lane: Lane; // The lane this kind would actually use under the given settings, for the // kinds whose scarce resource is a configuration choice rather than a fact. // Diarization is one: sherpa-onnx is CPU-only, while sortformer on the Vulkan // backend holds ~4.4 GB of the same 8 GB card the transcription engine wants. // // Optional because most kinds have no choice, and `lane` is the answer for // them. Mirrors digestLaneFor, which solved the same problem for the digest // operation's two lanes — and which digest itself now declares. // // 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. // // PURELY DECLARATIVE. It does not gate anything by itself — a kind still // decides for itself, from disk, whether its input is there, and says so by // returning `blocked`. What declaring it buys is two things nothing else // could: resolveBackfillLaneOperations can order a prerequisite before its dependant // within a single pass (so a video diarized this pass can be attributed in // the same one, rather than waiting for whatever LATER pass happens to find // the sidecar on disk), and a surface can say what a blocked video is waiting // FOR rather than just that it is stuck. // // An id naming a kind that is absent or disabled is not an error: the // dependency simply imposes no ordering, and the dependant keeps reporting // `blocked` until something produces its input. dependsOn?: readonly string[]; // The feature's OWN gate. A disabled feature reports no backfill at all — // otherwise every surface would advertise catch-up work for something the // operator has switched off. enabled(settings: SiteSettings): boolean; // The settings.json block that configures this operation, when one does. Named // here so the operation page can render that block's form without a table // keyed by operation id. Digest owns `digest`; diarization owns `diarization`; // both attribution operations share `attribution` (one block, one form, drawn // on both pages). The backfill LANE's block (`settings.backfill`) is not an // operation's and is not named here — it is rendered per lane, not per operation. settingsBlock?: OperationSettingsBlock; // The identity we would produce now, resolved ONCE per run rather than per // video. `unknown` here is the one erasure point in the table: each entry // narrows it back to its own type on the line below. The alternative — making // the whole registry generic — infects every consumer with a type parameter // for no gain, since none of them look inside a target. // // ASYNC AND CHANNEL-SCOPED, since the digest entry. A digest's identity // includes the hash of the channel's context note, which is a file read — so // this can no longer be a pure function of settings. It is still resolved // ONCE PER RUN, which is the property that matters (deriving it per video is // how a counter and a runner end up disagreeing about what is stale); the // callers simply await it now. resolveTarget(ctx: OperationTargetContext): unknown | Promise; state(probe: OperationProbe): Promise; run(opts: OperationRunOptions): Promise; // --- The unit-executor contract: what one unit of this kind needs, what it // produces, and how its output lands back on the primary. --- // Filenames (within the video dir) a unit executor must be shipped, derived // from the listing the caller already has. ORDERED: the executor // materializes them in this order, and every kind lists // transcript.cues.json LAST — isCuesJsonFresh compares mtimes, and a cues // file written before its metadata/raw transcript reads as stale, making // the unit silently do nothing. inputs(files: VideoFiles): string[]; // Filenames the run writes — what the executor's result endpoint returns. outputs: readonly string[]; // Apply a unit's returned output files ON THE PRIMARY, through the guarded // writers — never a raw file copy. Both sidecar writers are // read-modify-write (writeDigestSection preserves the other section; // attribution re-checks the downgrade rule against the primary's CURRENT // disk, because the unit ran against a snapshot that is minutes old). // "invalid" = the payload is not a usable record; "refused" = a guard said // no (which is a success of the guard, not a failure of the unit). applyResult( videoDir: string, payload: Record, ): Promise; }; export type ApplyResultOutcome = "applied" | "refused" | "invalid"; export type OperationTargetContext = { settings: SiteSettings; paths: Paths; channelSlug: string; }; // --------------------------------------------------------------------------- // Unit-contract helpers, shared across the kinds. // --------------------------------------------------------------------------- // The resolved primary raw transcript — needed by every transcript-derived // kind because isCuesJsonFresh compares the cues sidecar's mtime against it. function rawTranscriptOf(files: VideoFiles): string | null { if (files.hasWhisper) return WHISPER_FILENAME; return files.ytVttFile; } // The common transcript-derived input set: metadata + raw transcript + the // existing output sidecar (for the freshness/downgrade checks) + CUES LAST — // see Operation.inputs for why the order is load-bearing. function transcriptUnitInputs( files: VideoFiles, extras: readonly string[], ): string[] { const out: string[] = []; if (files.entries.includes(META_FILENAME)) out.push(META_FILENAME); const raw = rawTranscriptOf(files); if (raw) out.push(raw); for (const name of extras) { if (files.entries.includes(name)) out.push(name); } if (files.entries.includes(CUES_JSON_FILENAME)) out.push(CUES_JSON_FILENAME); return out; } // Shared by both attribution kinds: parse, validate the same shape // loadAttribution enforces, RE-CHECK the downgrade rule against the primary's // current disk (the unit ran against a snapshot minutes old, and the diarized // lane can have landed a better record meanwhile), then the guarded writer. async function applyAttributionResult( videoDir: string, payload: Record, ): Promise { const raw = payload[ATTRIBUTION_FILENAME]; if (!raw) return "invalid"; let record: AttributionRecord; try { const parsed = JSON.parse(raw) as Partial; if ( typeof parsed?.videoId !== "string" || typeof parsed.generatedAt !== "string" || !Array.isArray(parsed.speakers) || !Array.isArray(parsed.segments) || !parsed.provenance || typeof parsed.provenance.method !== "string" ) { return "invalid"; } record = parsed as AttributionRecord; } catch { return "invalid"; } if ( isAttributionDowngrade( await loadAttribution(videoDir), record.provenance.method as AttributionMethod, ) ) { return "refused"; } await writeAttribution(videoDir, record); return "applied"; } // Digest results land SECTION-WISE through writeDigestSection, which preserves // the other section, its warnings and the history from the record already on // the primary's disk — a raw copy of the unit's file would clobber a section // the primary wrote while the unit was in flight. async function applyDigestResult( videoDir: string, payload: Record, ): Promise { const raw = payload[DIGEST_FILENAME]; if (!raw) return "invalid"; let parsed: Partial; try { parsed = JSON.parse(raw) as Partial; } catch { return "invalid"; } const sections = parsed?.sections; if (!sections || typeof sections !== "object") return "invalid"; let applied = false; for (const [section, value] of Object.entries(sections)) { if (!isDigestSectionKind(section)) continue; const v = value as { provenance?: DigestProvenance; items?: DigestItem[]; }; if (!v?.provenance || !Array.isArray(v.items)) continue; const warnings: DigestWarning[] = (parsed.warnings ?? []).filter( (w) => w.section === section, ); await writeDigestSection(videoDir, { section, items: v.items, provenance: v.provenance, warnings, }); applied = true; } return applied ? "applied" : "invalid"; } // Diarization is the ONE whole-file verbatim apply: the sidecar has a single // writer and no sections to merge, so the shape check plus the atomic writer // is the whole guard. async function applyDiarizationResult( videoDir: string, payload: Record, ): Promise { const raw = payload[DIARIZATION_FILENAME]; if (!raw) return "invalid"; try { const parsed = JSON.parse(raw) as Partial; if ( typeof parsed?.generatedAt !== "string" || !Array.isArray(parsed.turns) ) { return "invalid"; } await writeDiarization(videoDir, parsed as DiarizationRecord); return "applied"; } catch { return "invalid"; } } // Diarization: the first entry, and the reason the table exists. const diarization: Operation = { id: "diarization", label: "Speaker diarization", hint: "Speaker turns captured from the audio, written to diarization.json beside the transcript.", group: "speakers", shortLabel: "Diarize", costBasis: "one pass over the audio per video", // The wording SpeakersStage used to hardcode for every kind at once. deferredHint: "too long to diarize under the current limit — raise or clear Max audio hours on the Diarization operation's page to include them", tier: "lane", // CPU, on the shared backfill queue. Serialized against the other backfill // kinds on purpose — two channels' worth of diarization at once just thrashes // cores. How much of the machine it may take is backfillLimit()'s question, // not the queue's. // The DEFAULT engine's answer, written out rather than derived: settings.ts is // a TYPE-only import here, and pulling defaultDiarization() in as a value would // make this module's initialization depend on it at runtime. laneFor is the // live answer, and diarizationLaneFor is where the rule actually lives. lane: { queueKey: BACKFILL_QUEUE, contendsFor: "cpu" }, laneFor: (settings) => diarizationLaneFor(settings.diarization), settingsBlock: "diarization", enabled: (settings) => settings.diarization.enabled && !!settings.diarization.segModel && !!settings.diarization.embModel, // The SAME target diarizeOne's own short-circuit uses. Two derivations would // let the counter and the runner disagree about what is stale. resolveTarget: ({ settings }): DiarizationFreshnessTarget => diarizationTarget(settings.diarization), async state({ videoDir, files, target, settings }) { // Same eligibility as diarizeAll's transcribedOnly default: the capture lane // exists to pair speaker turns with a transcript, and an untranscribed // video's audio is not at risk from the cleanup sweep yet. if (!isVideoTranscribed(files) || files.isUntranscribable) { return "not-applicable"; } // Cheap negative first: no sidecar in the listing means no read at all. if (files.hasDiarization) { const record = await loadDiarization(videoDir); // A malformed file reads as ABSENT here, exactly as hasDiarization() in // diarization-server.ts treats it: a half-written sidecar must never be // what convinces anything the work is done. if (record) { return isDiarizationFresh(record, target as DiarizationFreshnessTarget) ? "present" : "stale"; } } if (!(await hasDiarizableInput(videoDir, files))) { // Read the cap here ONLY when this branch can dispatch a download: with // re-download armed, missing-input is a fetch, and diarizeOneVideo does not // enforce the cap. Otherwise stay cheap — this runs per video per job start. if ( settings.backfill.allowRedownload && (await isOverDiarizationCap(videoDir, settings.diarization)) ) { return "deferred"; } return "missing-input"; } // The duration cap, read in the would-be-`missing` branch, which is ~835 // videos corpus-wide rather than 77,000 — and in the missing-input branch // above only when re-download is armed and the branch therefore spends a // download. That gating is not optional: countBackfillWork calls state() for // every video on every job start, and readVideoDurationSec reads and parses // a file. // // Unknown duration is NOT deferred — an absent or unparseable // metadata.info.json must not silently remove a video from the work list. return (await isOverDiarizationCap(videoDir, settings.diarization)) ? "deferred" : "missing"; }, async run(opts) { // diarizeOneVideo re-reads settings when none is passed, which is what we // want: the batch may run for hours and a model change mid-run should be // picked up. `force` is how a `stale` video gets redone at all — the // existence short-circuit inside is now a freshness check, but an operator // forcing a regeneration still needs to win. const outcome = await diarizeOneVideo({ paths: opts.paths, videoDir: opts.videoDir, videoId: opts.videoId, settings: opts.diarizationSettings, force: opts.force, onLog: opts.onLog, signal: opts.signal, }); if (outcome === "diarized") return "done"; if (outcome === "already-exists") return "already-present"; if (outcome === "no-audio") return "missing-input"; if (outcome === "not-configured") return "not-configured"; if (outcome === "disabled") return "disabled"; return "failed"; }, // The one kind whose input is MEDIA: the preferred extracted audio, else the // persisted source container. Same envelope as the transcript kinds, bigger // files — and a reachable population of only ~836 videos, so modest use. inputs(files) { const out: string[] = []; if (files.entries.includes(META_FILENAME)) out.push(META_FILENAME); const audio = pickPreferredAudio(files.audioFiles) ?? findSourceMedia(files.entries); if (audio) out.push(audio); return out; }, outputs: [DIARIZATION_FILENAME], applyResult: applyDiarizationResult, }; // Is there anything on disk ffmpeg could read for this video? Mirrors // controller/diarizeOne.ts's resolveDiarizableMedia, but answered from the // listing the caller already has so the common cases cost no I/O: // extracted audio, then a persisted source container, then — only when the // pointer file is actually present — the saved-video store. async function hasDiarizableInput( videoDir: string, files: VideoFiles, ): Promise { if (files.audioFiles.length > 0) return true; if (findSourceMedia(files.entries)) return true; if (!files.entries.includes(SAVED_VIDEO_POINTER_FILENAME)) return false; // The pointer exists but the stored file may not (an unmounted backup disk), // so this last step really does have to touch the filesystem. return (await resolveSavedVideo(videoDir)) !== null; } // Is this video longer than the diarization duration cap? // // THE CAP IS OFF BY DEFAULT NOW — windowed diarization removed the OOM it // existed for. It remains because a smaller machine, or a recording longer than // anything measured here, may still want it. This function is where its two // honest limitations live. First, duration is a PROXY: the memory blowup is O(n^2) in // speech-SEGMENT count, and turn density varies 40x across this corpus, so a // sparse 7h42m video is cheaper than a dense 6h12m one. Duration is used anyway // because it is the only predictor available from metadata already on disk, for // free, before committing 45 minutes of CPU to find out the hard way. Second, // duration is the CONTAINER's, so a video whose metadata is missing or lies gets // the benefit of the doubt. // // Unknown duration therefore returns false — not deferred. Deferring on an // unreadable metadata.info.json would quietly delete work from the list on the // strength of a file that could not be parsed, which is the opposite of what a // third counter is for. async function isOverDiarizationCap( videoDir: string, diarization: SiteSettings["diarization"], ): Promise { const capHours = diarization.maxAudioHours; if (!capHours || capHours <= 0) return false; // cap off const seconds = await readVideoDurationSec(videoDir); if (seconds === null) return false; return seconds > capHours * 3600; } // --------------------------------------------------------------------------- // Attribution — the second and third entries, and the ones that make this a // registry rather than a wrapper around diarization. // // TWO KINDS, ONE FILE. Both write attribution.json, and the ordering rule in // lib/attribution.ts is the whole safety of that: the diarized lane may // overwrite a text-only record (an UPGRADE — that is what the second kind is // for), and the text lane must never overwrite a diarized one (a DOWNGRADE). // state() encodes it here and attributeOne re-checks it against disk immediately // before writing, because the pool can hold a candidate for minutes after // state() ran. It has a test; a comment would not have been enough. // // THE SPLIT IS ALSO WHAT MAKES THE UPGRADE QUEUE FREE. PLAN.md describes a // bespoke "upgrade job" for turning text-only records into diarized ones. It is // not needed: `attribution-diarized` reports a text-only record as MISSING work, // so the existing lane, sweep and indicators queue the upgrade with no new // machinery. What it needs re-acquiring media for is already // controller/backfillReacquire.ts. // // AND IT IS WHY missing-input MATTERS HERE MOST. `attribution-diarized`'s input // is diarization.json, of which this corpus has ONE. So its missing-input // population is ~73,000 videos on day one — the exact case the reachable / // needs-input split exists to stop from poisoning every surface. A single // "remaining" number would put every channel at the top of every list forever. // Shared by both attribution kinds: a video only has speakers worth naming if it // has a transcript. Not-applicable rather than missing-input, since re-acquiring // media would not help — the video needs transcribing, which is another lane's // job entirely. function attributionApplies(files: VideoFiles): boolean { return isVideoTranscribed(files) && !files.isUntranscribable; } // The identity, minus the per-video half. The diarized lane's identity also // includes the generatedAt of the diarization.json it names clusters from, and // that is a disk read — so it is added inside the one state() branch that has // already paid for the read. See AttributionProvenance.diarizationGeneratedAt. function attributionTargetFor( settings: SiteSettings, method: AttributionMethod, ): AttributionFreshnessTarget { return resolveAttributionTarget(method, settings.attribution).target; } // The transcript half of the per-video identity, for both lanes: which text the // names would be made from now. ZERO I/O — pickIndexTranscript reads the // already-loaded VideoFiles, and this runs per video per job start, where the // cost bar is the channel snapshot's own per-video classification. `{}` when // there is no transcript at all, which state() has already ruled out but which // must not be turned into a false assertion here. function transcriptSourceTarget( files: VideoFiles, ): { transcriptSource?: AttributionTranscriptSource } { const source = transcriptSourceOf(pickIndexTranscript(files)?.kind); return source ? { transcriptSource: source } : {}; } const attributionText: Operation = { id: "attribution-text", label: "Speaker names (from the transcript)", hint: "Speakers reconstructed from the transcript alone, for videos with no diarization. Cheaper to reach, worse than the diarized lane, and it never overwrites one.", group: "speakers", shortLabel: "Names·T", // The expensive one, and the reason costBasis is a field. A transcript is // many chunks; this is the only lane in the table whose unit is not the video. costBasis: "~1 model call per transcript chunk", tier: "lane", lane: { queueKey: BACKFILL_QUEUE, contendsFor: "network" }, settingsBlock: "attribution", enabled: (settings) => settings.attribution.enabled && settings.attribution.textOnlyEnabled, resolveTarget: ({ settings }) => attributionTargetFor(settings, "text-only"), async state({ videoDir, files, target }) { if (!attributionApplies(files)) return "not-applicable"; // NEVER missing-input. The input is the cue stream, and a transcribed video // has one by definition — which is exactly why this lane can reach the whole // corpus and why running it over the whole corpus costs ~194,000 model calls. if (!files.entries.includes(ATTRIBUTION_FILENAME)) return "missing"; const record = await loadAttribution(videoDir); // A malformed file reads as ABSENT, the same rule the diarization entry // uses. Here it also protects the write path: a half-written sidecar must // not be able to masquerade as a diarized record and block this lane // forever. if (!record) return "missing"; // THE DOWNGRADE RULE, in the counter as well as the runner. A diarized // record is not stale for this lane and is not work — there is simply // something better here. Reporting it as work would put this lane in a loop // of "attempt, refuse, still outstanding" across every pass of a sweep. if (isAttributionDowngrade(record, "text-only")) return "present"; return isAttributionFresh(record, { ...(target as AttributionFreshnessTarget), ...transcriptSourceTarget(files), }) ? "present" : "stale"; }, async run(opts) { // Lazy, once, at the point of actually running something. See the import // note at the top of this file. const { attributeOneVideo } = await import("../controller/attributeOne"); return toBackfillOutcome( await attributeOneVideo({ paths: opts.paths, videoDir: opts.videoDir, videoId: opts.videoId, channelSlug: opts.channelSlug, method: "text-only", force: opts.force, settings: opts.attributionSettings, appConfig: opts.appConfig, context: opts.digestContext, onLog: opts.onLog, signal: opts.signal, }), ); }, inputs: (files) => transcriptUnitInputs(files, [ATTRIBUTION_FILENAME]), outputs: [ATTRIBUTION_FILENAME], applyResult: applyAttributionResult, }; const attributionDiarized: Operation = { id: "attribution-diarized", label: "Speaker names (from the audio)", hint: "Names put to the speaker clusters in diarization.json — about one model call per video, and better than the text-only lane. Needs diarization to have run first.", group: "speakers", shortLabel: "Names·A", costBasis: "~1 model call per video", tier: "lane", // The dependency the hint has always stated in prose. Declaring it is what // turns "needs diarization to have run first" from a sentence an operator // reads into something the scheduler can order by and a counter can name. dependsOn: ["diarization"], lane: { queueKey: BACKFILL_QUEUE, contendsFor: "network" }, settingsBlock: "attribution", enabled: (settings) => settings.attribution.enabled && settings.attribution.diarizedEnabled, resolveTarget: ({ settings }) => attributionTargetFor(settings, "diarized"), async state({ videoDir, files, target }) { if (!attributionApplies(files)) return "not-applicable"; // The input is diarization.json, NOT audio. That distinction is the whole // reason this kind is cheap: the perishable input was already captured, and // what is left is a naming pass that can be redone at any time. // // Deliberately NOT gated on settings.diarization.enabled — a sidecar // captured during a past run is a perfectly good input after capture is // switched off again, and refusing to name it would strand exactly the work // the capture lane exists to protect. // // BLOCKED, NOT MISSING-INPUT. This used to say missing-input, and that was // wrong in a way that cost real work: it put ~73,000 videos into the // "re-acquire the media" population, where allowRedownload would fetch // AUDIO — which can never satisfy a wait for diarization.json — and then // delete it again. What this video is waiting for is the `diarization` kind // declared in dependsOn above, and that is a thing this table produces. if (!files.hasDiarization) return "blocked"; if (!files.entries.includes(ATTRIBUTION_FILENAME)) return "missing"; const record = await loadAttribution(videoDir); if (!record) return "missing"; // A text-only record here is THE UPGRADE QUEUE: the diarized record this // kind is responsible for genuinely does not exist yet, so it is `missing` // rather than `stale`. Both are reachable work, but the two words mean // different things to an operator reading a stage card — "stale" says // something changed under a record, "missing" says a better one was never // made. if (record.provenance.method !== "diarized") return "missing"; // Only now is the diarization read worth paying for: it is needed solely to // ask whether the clusters these names point at are still the same clusters. const diarization = await loadDiarization(videoDir); // Present in the listing but unreadable — a half-written or corrupt // sidecar. Blocked for the same reason as the branch above, and note that // the `diarization` kind reads a malformed record as ABSENT too, so it will // regenerate this file and unblock the video without anyone intervening. if (!diarization) return "blocked"; return isAttributionFresh(record, { ...(target as AttributionFreshnessTarget), diarizationGeneratedAt: diarization.generatedAt, ...transcriptSourceTarget(files), }) ? "present" : "stale"; }, async run(opts) { // Lazy, once, at the point of actually running something. See the import // note at the top of this file. const { attributeOneVideo } = await import("../controller/attributeOne"); return toBackfillOutcome( await attributeOneVideo({ paths: opts.paths, videoDir: opts.videoDir, videoId: opts.videoId, channelSlug: opts.channelSlug, method: "diarized", force: opts.force, settings: opts.attributionSettings, appConfig: opts.appConfig, context: opts.digestContext, onLog: opts.onLog, signal: opts.signal, }), ); }, inputs: (files) => transcriptUnitInputs(files, [DIARIZATION_FILENAME, ATTRIBUTION_FILENAME]), outputs: [ATTRIBUTION_FILENAME], applyResult: applyAttributionResult, }; // One mapping, shared by both kinds, so the two lanes cannot report the same // condition differently. function toBackfillOutcome(outcome: AttributeOneOutcome): OperationRunOutcome { switch (outcome) { case "attributed": return "done"; case "already-exists": // A better record already exists. Nothing to do here is the SAME answer as // "already current" for the lane's purposes, and reporting it as a failure // would make an untouched corpus look broken. case "outranked": return "already-present"; case "no-diarization": return "missing-input"; case "disabled": return "disabled"; // A transcript that is absent or about to be rewritten. Not a failure of // this lane and not something re-acquiring media fixes — it resolves itself // when the normalize pass catches up, and the next sweep pass will see it. case "no-transcript": return "skipped"; default: return "failed"; } } // --------------------------------------------------------------------------- // Digest — the fourth entry, and the one this registry was supposed to have had // from the start. // // Digests were built FIRST and never merged in, so they grew a parallel // implementation of this same idea: their own sweep, their own per-channel // batch, their own yield probe, their own cost planner, their own queue keys, // their own settings block with its own pause, and their own snapshot counter. // The backfill sweep was, in its own words, a clone of the digest one — both // retired in slice 1.3. Registering the operation is how that convergence // started. // // 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 the lane's executor ever // disagree about whether a video is digested, that is a bug in this file, not a // second opinion. // // WHAT REGISTERING ACTUALLY BUYS, today: // // 1. THE TRANSCRIPT DEPENDENCY BECOMES DECLARED. PLAN.md states in prose that // "a video cannot be digested until it has a transcript" and nothing // enforced it — an undigestable video was simply absent from every bucket. // It now reports `blocked` on `transcription`, so it is counted and named. // 2. A DECLARED LANE. The digest queue keys stop being one subsystem's private // constants and become this operation's declared lane, which is what lets a // scheduler dispatch one job per lane and keep GPU and CPU work overlapping. // 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 ON THE BACKFILL LANE — see backfillQueueKinds below. const digest: Operation = { id: "digest", label: "Digest", hint: "Chapters and tags generated from the transcript by a local or metered model. Needs a transcript first.", group: "digest", shortLabel: "Digest", costBasis: "~1 model call per transcript chunk", // See the `deferred` branch in state() below for the measurement behind this // wording. It says "run the normalize pass" and NOT "it clears itself", // because nothing automatic ever will. deferredHint: "waiting on a normalized transcript (transcript.cues.json) that nothing produces automatically — run Normalize transcripts on the channel to make them digestable", tier: "lane", // The transcript, declared. Nothing in the repo enforced this before. dependsOn: ["transcription"], // The LOCAL lane's key is the declared one because it is the default and the // only one enabled unless remoteEnabled is set. The metered lane runs on // DIGEST_REMOTE_QUEUE, and the two must never share a key: one shared key // would idle the network lane while the GPU works, across a multi-week sweep. // `contendsFor: "gpu"` is what makes this lane — and only this lane — stand // aside for transcription. lane: digestLaneFor("local-gpu"), // The live answer, for the same reason diarization has one: which lane this // runs on is a CONFIGURATION CHOICE, not a fact about the operation, and // `lane` above can only carry the default. laneForOperation asks this, so a // remote digest is dispatched onto DIGEST_REMOTE_QUEUE instead of the local // key it declares. // // 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 (the digest lane's limit()), not here: a paused // lane must still report how much work is outstanding. settingsBlock: "digest", enabled: () => true, async resolveTarget({ paths, channelSlug }) { // The existing resolver, verbatim. This is the async, channel-scoped case // OperationTargetContext exists for: a digest's identity includes the hash of // the channel's context note, which is a file read. const { resolveDigestTarget } = await import("../controller/digestTarget"); const resolved = await resolveDigestTarget({ paths, channelSlug }); return { target: resolved.target, sections: resolved.sections }; }, async state({ videoDir, files, target }) { const { target: freshness, sections } = target as DigestTarget; // Untranscribable is not-applicable, exactly as the attribution kinds treat // it: nothing will ever produce a transcript for it, so it is not blocked, // it is out of scope. if (files.isUntranscribable) return "not-applicable"; // BLOCKED, not not-applicable and not missing-input. A video with no // transcript is waiting on the transcription operation declared in // dependsOn above — a thing this system produces. Before this it was simply // invisible: absent from every digest bucket, so a channel of untranscribed // videos read as fully digested. if (!isVideoTranscribed(files)) return "blocked"; const { fresh: cuesFresh } = await isCuesJsonFresh(videoDir); if (!cuesFresh) { // NO CURRENT NORMALIZED TRANSCRIPT. Digesting now would either fail for // want of one or describe superseded text and then look fresh forever, so // the video is held back — `deferred` is the classification with the // matching meaning: not attempted, not broken, not counted as reachable // work. digestVideo reports the same condition as `skipped` rather than a // failure, for the same reason. // // THIS DOES NOT RESOLVE ITSELF, and an earlier version of this comment // said it did. Measured over the whole corpus (79,219 video dirs): // // 1,942 have NO cues.json at all ← reason "missing" // 47 have one that is superseded ← reason "stale" // // so the superseded case this branch was written for is 2.4% of what it // actually catches. The missing case is permanent: transcribeOne is the // ONLY automatic caller of normalizeTranscript, and a channel with // `handling: "youtube"` fetches subtitles with --skip-download and so // never runs it — 1,683 of the 1,942 are piratesoftware alone. It went // unnoticed because buildIndex treats cues.json as a CACHE and silently // re-parses the raw VTT when it is absent, so the published site is // correct and only this lane, which has no such fallback, can see it. // // The fix is the normalize pass, run deliberately: normalizeChannel- // Transcripts (controller/normalizeAll.ts), wired to a button on the // digest stage card next to this count. Both reasons are fixed by it, // which is why they share one classification — see CuesFreshReason. return "deferred"; } const record = await loadDigest(videoDir); // A digest SHARED from a duplicate cluster's canonical member counts as // done. The canonical member's own freshness drives regeneration and the // share is re-applied from it (isSharedFrom's contract) — so re-deriving it // here would undo ~11% of the sweep's saving. This entry IS the definition // of "digested" — the snapshot's second one (`buckets.noDigest`) is gone, // because two definitions is the exact failure this entry exists to stop. if (record?.derivedFrom != null) return "present"; // EVERY configured section must be fresh to count as done, matching // countMissingDigests and the batch. What is new is that "not all of them" // is no longer one answer. // // WHY `partial` IS NOT JUST A NICER WORD FOR `stale`. digestVideo does not // regenerate a video, it regenerates SECTIONS: its `stale` list at // digestVideo.ts:158 is `sections.filter(not fresh)`, and only those are // generated. So a video with fresh chapters and no tags is a fraction of the // cost of one with neither, and folding them together prices the work wrong // in the direction that matters — the corpus is ~77,000 videos and the // recorded surcharge for adding tags to an existing chapters pass is ~44% of // a full pass, against 100% for a genuine re-generation. // // The case is not hypothetical: `sections` is a setting, and the pre-sweep // decision is to turn tags on. The moment that happens every already-digested // video in the corpus becomes part-done at once, and without this split it // would read as `stale` — indistinguishable, on a stage card, from a // PROMPT_VERSION bump that really did invalidate everything. // // The empty-sections case still reads `present` (0 of 0 fresh), exactly as // the `every()` this replaces did, so a caller passing no sections is // unchanged. // digestSectionStates is the fold, shared with the per-video digest panel: // one definition of "this section is fresh", so a panel and a work list // cannot disagree about the same section. const fresh = digestSectionStates(record, sections, freshness).filter( (s) => s.fresh, ).length; if (fresh === sections.length) return "present"; // No record at all cannot be part-done, and it is the one branch that must // stay `missing`: it is what separates "never digested" from "digested and // superseded" everywhere downstream. if (!record) return "missing"; return fresh > 0 ? "partial" : "stale"; }, async run(opts) { // Lazy, at the point of running something: digestVideo drags in the prompt // module, the markdown renderer and the transcript normalizer, and this // module is imported by channelSnapshot on the editor's hot path. const { digestVideo } = await import("../controller/digestVideo"); const { sections } = opts.target as DigestTarget; // OperationRunOptions already carries channelSlug and videoId, which is // exactly what digestVideo takes — no reshaping of the controller. const outcome = await digestVideo({ paths: opts.paths, channelSlug: opts.channelSlug, videoId: opts.videoId, sections, // Unit-executor injection. On the primary both are unset and digestVideo // resolves exactly as before; on an executor the injected config (baseUrl // stripped — the executor localises its own endpoint) and the injected // context are what keep the identity the primary's. ...(opts.appConfig ? { config: opts.appConfig } : {}), ...(opts.digestContext ? { context: opts.digestContext } : {}), force: opts.force, onLog: opts.onLog, signal: opts.signal, }); // The outcome union already lines up almost exactly. if (outcome.status === "wrote") return "done"; if (outcome.status === "fresh") return "already-present"; return "skipped"; }, inputs: (files) => transcriptUnitInputs(files, [DIGEST_FILENAME]), outputs: [DIGEST_FILENAME], applyResult: applyDigestResult, }; type DigestTarget = { target: DigestFreshnessTarget; sections: DigestSectionKind[]; }; // 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; 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 // 1, so one key would idle the network lane while the GPU works — across a // sweep measured in weeks. export function digestLaneFor(appLane: DigestLane): Lane { return appLane === "local-gpu" ? // Competes with the transcription engine for the same VRAM. Yields. { queueKey: DIGEST_LOCAL_QUEUE, contendsFor: "gpu" } : // Metered and network-bound: it competes for nothing local, so yielding // would park a lane that costs nothing to keep running. { queueKey: DIGEST_REMOTE_QUEUE, contendsFor: "network" }; } // Which resource a diarization run competes for, which follows from the engine // it is configured with rather than from a separate setting — the same shape as // digestLaneFor, and stated here so nothing has to re-test the engine id. // // The queue key does NOT change with the engine. Diarization serializes against // itself either way, and giving the GPU variant its own key would only let two // diarizations run at once — which is precisely what must not happen when each // holds ~4.4 GB of an 8 GB card. export function diarizationLaneFor(diarization: { engine: DiarizationEngineId; backend: DiarizationBackend; }): Lane { return diarization.engine === SORTFORMER_DIARIZATION_ENGINE && diarization.backend === "vulkan" ? // Competes with the transcription engine for the same VRAM. Yields. { queueKey: BACKFILL_QUEUE, contendsFor: "gpu" } : // sherpa-onnx is ONNX/CPU, and sortformer on the CPU backend is likewise // only after cores. Contends for CPU, whose share the backfill lane's own // weight already governs. { queueKey: BACKFILL_QUEUE, contendsFor: "cpu" }; } // Whether a lane must stand aside while transcription is working. One rule, so // the digest lane and any future GPU operation cannot answer it differently. export function laneYieldsToTranscription(lane: Lane): boolean { return lane.contendsFor === "gpu"; } // --------------------------------------------------------------------------- // The operations this system knows about but does NOT dispatch. // // Registering them costs nothing and gets the dependency graph right: without a // `transcription` entry, digest.dependsOn = ["transcription"] names nothing and // no surface can say what a blocked video is waiting for. // // They stay externally dispatched DELIBERATELY, and the reasons are concrete // rather than a lack of time: // // - A SELF-REFERENTIAL YIELD. transcriptionActivity() reports busy exactly // when transcription is running, so a transcription lane that yielded to it // would yield to itself. // - NO WORKER-LEASE SURFACE. OperationRunOptions has no notion of leasing a // worker, of a tier, of draining, or of a partial stop — all of which the // 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. 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. // // The SAME shape as a catalog entry, minus the two fields it cannot have: it is // `external` by definition, and nothing in settings.json configures a download // as an operation. Declared as a projection rather than a second field list so // a new descriptor field is declared once and cannot go missing from one of the // two — which is how `scope` and `trigger` arrive here for free. export type ExternalOperation = Omit< OperationDescriptor, "dispatch" | "settingsBlock" > & { dispatch: "external" }; // THE SYNC OPERATION. Catalogued so the board, the rail and /operations/sync // come off the same table as everything else, and NOT in EXTERNAL_OPERATIONS — // that list is the media-derived pipelines the rail draws bands for (editor // buildBands.ts, EXTERNAL_BAND_IDS); sync has no per-video population. The lane // is the shape download declares: syncAction runs through runPipelineAction on // the per-platform download queue (pipelineActions.ts, downloadQueueKey). export const SYNC_OPERATION: OperationDescriptor = { id: "sync", label: "Sync", hint: "Noticing new videos: re-reading each channel's listing on its cadence and fetching what is new. Dispatched by the sync scheduler's heartbeat, and by hand per channel or all at once. A corpus that has stopped noticing new videos is not idle — it is broken.", group: "sync", shortLabel: "Sync", costBasis: "one listing fetch per channel, over the network", lane: { queueKey: "download:", contendsFor: "network" }, dispatch: "external", scope: "channel", trigger: "cadence", settingsBlock: "syncScheduler", }; // THE METADATA SCAN. A channel-scoped, backlog-fed operation that reads what // this channel's listed-but-unfetched videos are CALLED, without fetching any of // them — so the per-channel download filter can settle the ones the operator // does not want before anything is downloaded. // // On the platform download queue and contending for the network, like sync and // download: it is one metadata request per listed video, and the source counts // them the same way it counts a download's. // // THE DOWNLOAD LANE DISPATCHES IT, and that is what `runner` says. The runner // checks, before every pick, whether a filtered channel has unscanned listed // videos — exactly the backlog this entry declares — and runs the scan for it // as one channel-scoped unit on the platform download queue, holding the same // per-platform slot a download unit would. // // It goes FIRST because it decides what the downloads are. Every unscanned // non-match on a filtered channel is otherwise prefetched, rejected and // discarded one yt-dlp invocation at a time, against a source whose patience is // the scarce resource; one batch answers the whole channel. // // WHAT NAMING THE RUNNER ALSO CHANGES: `pauseLaneFor` asks `runner` first, so // this operation's console is the download lane's and the download pause now // holds the AUTO-dispatch. The operator's own Run button is unaffected — it // checks the platform cooldown and nothing else — which keeps the original // point of the exemption intact: you can still scan, while downloads are // paused, precisely to decide what the lane should fetch when it resumes. export const METADATA_SCAN_OPERATION: OperationDescriptor = { id: "metadata-scan", label: "Metadata scan", hint: "Reading the title, description and date of videos that are in the listing but not downloaded, without fetching any media. What the per-channel download filter needs in order to decide: with no scan, the only way to learn a video's title is to start downloading it.", group: "sync", shortLabel: "Scan", costBasis: "one metadata fetch per listed video, no media", lane: { queueKey: "download:", contendsFor: "network" }, dispatch: "external", scope: "channel", trigger: "backlog", runner: "download", }; export const EXTERNAL_OPERATIONS: readonly ExternalOperation[] = [ { id: "download", label: "Download", hint: "Fetching the media. Dispatched by the auto-download runner and the per-channel pipeline actions.", group: "media", shortLabel: "Download", costBasis: "one fetch per video, over the network", // Really one queue per platform (downloadQueueKey), not a single key. Named // here as the shape rather than the exact key, because the catalog's job is // the dependency graph, not dispatch. lane: { queueKey: "download:", contendsFor: "network" }, dispatch: "external", scope: "video", trigger: "backlog", runner: "download", }, { id: "transcription", label: "Transcription", hint: "Turning audio into a transcript. Dispatched by the auto-transcribe runner across the worker pool.", group: "transcript", shortLabel: "Transcribe", costBasis: "one pass over the audio per video, on a worker", lane: { queueKey: TRANSCRIPTION_QUEUE, contendsFor: "gpu" }, dependsOn: ["download"], dispatch: "external", scope: "video", trigger: "backlog", runner: "transcription", }, ]; // Every operation the console knows, dispatched here or not. The catalog — what // a dependency id resolves against, and what the board, the rail and // /operations/ enumerate. The sync scheduler is IN it now (SYNC_OPERATION), // which is why `scope` and `trigger` exist: the entries are no longer all // per-video, backlog-fed things. export type OperationDescriptor = { id: string; label: string; hint: string; group: OperationGroup; shortLabel: string; costBasis: string; lane: Lane; dependsOn?: readonly string[]; dispatch: OperationDispatch; // See OperationScope. What this operation has a state FOR — a video for // everything media-derived, a channel for sync. scope: OperationScope; // See OperationTrigger. What makes it run — a backlog for everything the // sweep and the runners feed, a cadence for sync. trigger: OperationTrigger; // The auto-queue runner that dispatches this, when one does. `dispatch` alone // cannot answer it: download and transcription are `external` WITH a runner, // and an external operation with none is a legal entry (transcode was one, // 2026-08-26 → 08-30). A console that guessed from the id would hand such an // entry the transcription runner's controls — a live Start button over the // wrong lane. Absent for every registry operation: since slice 1.2 the console // chooses a lane's console off `pauseLaneFor`, which reads the queue key. // // Typed to the auto-queue kinds on purpose; the sync heartbeat is a runner in // the IA doc's sense but not one of these, which is why /operations/sync is // chosen off `trigger`. runner?: AutoQueueKind; // See Operation.settingsBlock. Absent for every external operation but sync: // nothing in settings.json configures a download as an operation. settingsBlock?: OperationSettingsBlock; }; export function operationCatalog(): OperationDescriptor[] { return [ SYNC_OPERATION, METADATA_SCAN_OPERATION, ...EXTERNAL_OPERATIONS, ...OPERATIONS.map((k) => ({ id: k.id, label: k.label, hint: k.hint, group: k.group, shortLabel: k.shortLabel, costBasis: k.costBasis, lane: k.lane, dependsOn: k.dependsOn, dispatch: "backfill" as const, // The registry is per-video by construction — `Operation.state(probe)` is // a per-video probe — and every one of its entries is fed by a backlog of // videos that lack the output. scope: "video" as const, trigger: "backlog" as const, settingsBlock: k.settingsBlock, })), ]; } // The label for a dependency id, from anywhere in the catalog. Returns the id // itself for something unknown rather than throwing — a dangling dependency is // already tolerated everywhere else here. export function operationLabel(id: string): string { return operationCatalog().find((o) => o.id === id)?.label ?? id; } // The group an operation belongs to, or null for an id the catalog does not // know. Null rather than a fallback group: a caller grouping by this must be // able to tell "unknown" from "media", and silently filing a dangling id under // the first group would put it on the wrong station. export function operationGroup(id: string): OperationGroup | null { return operationCatalog().find((o) => o.id === id)?.group ?? null; } // What one video of an operation costs — its cost basis — in words. Empty string for an unknown // id, so a surface can print it unconditionally without a placeholder. export function operationShortLabel(id: string): string { return operationCatalog().find((o) => o.id === id)?.shortLabel ?? id; } export function operationCostBasis(id: string): string { return operationCatalog().find((o) => o.id === id)?.costBasis ?? ""; } // The label for a SET of operations — a lane card, a transit-line station, a // column group. One group means that group's name; a mix means "Derived data". // // DERIVED, never hardcoded, which is the point: a lane that gains a kind from a // different group degrades to the honest generic name instead of continuing to // advertise a label that now describes two thirds of what it holds. // The same set, named as a THING AN OPERATOR RUNS rather than as a stage on a // line — "Run speaker work", "N videos are waiting on digests". // // Two labels for one group is not duplication: "Speakers" is a station on a // transit line and has to be a noun at eyebrow width; "speaker work" is the // object of a verb and has to survive being lower-cased into a sentence. The // alternative — deriving one from the other — produces "Run speakers work", // which is why this is declared. export function groupActionLabel(group: OperationGroup): string { switch (group) { case "media": return "downloads"; case "transcript": return "transcripts"; case "digest": return "digests"; case "speakers": return "speaker work"; case "sync": return "syncs"; } } export function operationsActionLabel(ids: ReadonlyArray): string { const groups = new Set(); for (const id of ids) { const group = operationGroup(id); if (group) groups.add(group); } if (groups.size !== 1) return "derived data"; return groupActionLabel([...groups][0]); } export function operationsGroupLabel(ids: ReadonlyArray): string { const groups = new Set(); for (const id of ids) { const group = operationGroup(id); if (group) groups.add(group); } if (groups.size !== 1) return "Derived data"; return groupLabel([...groups][0]); } // The digest operation's id, named once. Surfaces that read one specific // operation off a snapshot (the digest stage card, the dashboard's coverage // instrument) need this string, and a typo in it fails the way a missing // snapshot entry does — silently, as "nothing to do". export const DIGEST_OPERATION_ID = "digest"; // The diarization operation's id, named once for the same reason. The cleanup // accounting needs it to ask a question no other surface asks: whether the lane // that would release a held video is even running (allOperations drops a kind // whose enabled() is false, so an ABSENT entry is the answer, not a zero). export const DIARIZATION_OPERATION_ID = "diarization"; // One entry per backfill known to the system. export const OPERATIONS: readonly Operation[] = [ diarization, attributionDiarized, // Text-only LAST, deliberately. resolveBackfillLaneOperations preserves this order and // 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, // Digest is in the table, but on its own lane — so it is counted and // classified by everything that reads this registry, and dispatched by none // of it. See backfillQueueKinds. digest, ]; export const OPERATION_BY_ID: Record = Object.fromEntries(OPERATIONS.map((k) => [k.id, k])); export function getOperation(id: string): Operation | undefined { return OPERATION_BY_ID[id]; } // Every registered kind whose feature is switched on, whatever lane it runs on. // The catalog view: what EXISTS and is live, for a scheduler or a dependency // lookup. Callers that mean "the backfill lane" want backfillLaneOperations below. export function allOperations(settings: SiteSettings): Operation[] { return OPERATIONS.filter((k) => k.tier === "lane" && k.enabled(settings)); } // THE BACKFILL LANE's kinds: enabled, `lane` tier, and running on the shared // 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 — 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. 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. // - GUARDS. digest carries its own pause gate, the yield-to-transcription // carve-out (with its CPU-worker exemption), spendCapUsd on the metered // lane, the remoteEnabled fail-fast, shortest-first ordering, // duplicate-cluster sharing and the engine probe() fail-fast. Those are // 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 // surfaces that are live on a 78,000-video corpus mid-sweep. // // So digest is CLASSIFIED and CATALOGUED through this registry and DISPATCHED // through its own controller. Collapsing the two sets of counters and the two // schedulers is the unified-rule-model work, not this function's job. A kind // joins this list when its lane rule can be expressed here without losing a // guard. export function backfillLaneOperations(settings: SiteSettings): Operation[] { return allOperations(settings).filter( (k) => k.lane.queueKey === BACKFILL_QUEUE, ); } // The READ-SIDE twin of backfillLaneOperations: given a snapshot's per-kind map, // return only the entries belonging to the shared backfill lane — KEYED, so a // surface can say WHICH kind a number came from. backfillLaneEntriesOf below is this // with the ids dropped, for the callers that only sum. // // The keyed form is what the corpus-wide backfill card needs. Summed, this lane // reads "77,952 reachable · 77,134 need media" — both figures correct, and // together meaningless: 99.5% of the first is attribution-text (one model call // per transcript CHUNK) and all of the second is diarization (329 runs). Adding // kinds gives a number in no unit at all, which is the mistake the header // forbids one level up for `missing` vs `missing-input`. // // THIS EXISTS BECAUSE THE SNAPSHOT MAP STOPPED BEING THE LANE. It used to be // written from backfillLaneOperations, so `Object.values(snapshot.backfill)` and "the // backfill lane" were the same set by construction, and four surfaces summed it // generically on that basis — the channel dashboard's backfill instrument, the // channel page's speakers stage and the widget's sync payload. The moment // channelSnapshot writes an entry per CATALOG operation, that identity breaks: // those four would silently absorb ~75,000 digest videos into a number that has // only ever meant diarization plus attribution. // // FILTERED BY THE DECLARATION, NOT BY ID. `key !== "digest"` would fix today and // leave the identical trap armed for the next operation registered on a lane of // its own — which is the whole direction of the unified-operations work. The // rule is the same one backfillLaneOperations applies on the write side, asked of the // registry: does this kind run on BACKFILL_QUEUE? // // An id the catalog does not know is EXCLUDED. A snapshot is a file on disk that // may have been written by an older build and may name a kind that has since // been renamed or removed; there is no lane declaration to check it against, so // it cannot be asserted to belong to this one. Deliberately not filtered on // `enabled(settings)` — these are counts already written to disk, and a feature // switched off after a snapshot was taken does not retroactively unmake the work // it recorded. export function backfillLaneOperationEntriesOf( backfill: Record | undefined | null, ): [string, T][] { if (!backfill) return []; return Object.entries(backfill).filter( ([id]) => getOperation(id)?.lane.queueKey === BACKFILL_QUEUE, ); } // The same set with the ids dropped, for the callers that only ever sum. Defined // in terms of the above rather than beside it: the filter and every word of the // rule above it must stay in ONE place, or the next surface that wants per-kind // detail copies a `key !== "digest"` in and re-arms the trap. export function backfillLaneEntriesOf( backfill: Record | undefined | null, ): T[] { return backfillLaneOperationEntriesOf(backfill).map(([, entry]) => entry); } // THE OPERATIONS A LANE DISPATCHES, in dependency order. // // The other half of the four-lane model: `bucketsForKind` says what a lane's // leaves may draw from the snapshot's buckets, and this says what they may draw // from its operation entries. A leaf naming no operation draws ALL of them, the // way a bucket-less leaf draws the default bucket union. // // OFF THE QUEUE KEY, and off `laneFor` before `lane`, for the same reason // pauseGates.ts asks `runner` first: the queue key is what actually serializes // the work, and digest's key depends on the engine it is configured with. With // `digest.remoteEnabled` on, the digest operation moves from DIGEST_LOCAL_QUEUE // to DIGEST_REMOTE_QUEUE — both of which are the digest LANE, which is why both // are named here rather than one. // // The two BUCKET lanes still get [] here, and that is not the gap it looks // like. Slice 1.5 does write `snapshot.backfill.download` and // `.transcription` — but this function returns `Operation`s, registry entries // with `state()` and `run()`, and download and transcription have neither. The // lane that dispatches them is asked for its work list by id instead, through // bucketLaneOperationId below; the runner's `run()` never asks this function // for them because it never dispatches them through operationBatch. export function operationsForLane( lane: AutoQueueKind, settings: SiteSettings, ): Operation[] { const keys = lane === "digest" ? [DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE] : lane === "backfill" ? [BACKFILL_QUEUE] : []; if (keys.length === 0) return []; return orderByDependencies( allOperations(settings).filter((op) => keys.includes((op.laneFor?.(settings) ?? op.lane).queueKey), ), ); } // THE ONE OPERATION A BUCKET LANE DISPATCHES, or null for a lane that draws // from the operation registry instead. // // The mirror of pauseGates.pauseLaneFor, which asks an operation for its lane; // this asks a lane for its operation. Both read the SAME declaration — // `ExternalOperation.runner` — so the two directions cannot drift, and neither // spells "download" or "transcription" as a literal. // // It exists because slice 1.5 gave the two bucket lanes a work list in the // snapshot (`backfill.download`, `backfill.transcription`) keyed by the // operation id, so the runner needs the id to read it and the snapshot // generator needs it to write it. There is exactly one per lane by // construction: `runner` names an AutoQueueKind, and two external operations // claiming the same runner would be two work lists for one dispatcher. export function bucketLaneOperationId(lane: AutoQueueKind): string | null { return EXTERNAL_OPERATIONS.find((op) => op.runner === lane)?.id ?? null; } // Resolve a caller-supplied list of kind ids against the registry. An empty or // absent list means "every enabled lane kind" — the sweep's scope default. // Unknown ids are dropped rather than throwing: a settings file may name a kind // from a newer build, and a stale scope must not wedge the lane. export function resolveBackfillLaneOperations( settings: SiteSettings, ids: readonly string[] | undefined, ): Operation[] { const lane = backfillLaneOperations(settings); const selected = !ids || ids.length === 0 ? lane : lane.filter((k) => new Set(ids).has(k.id)); return orderByDependencies(selected); } // Order kinds so a prerequisite is attempted before anything that declares 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 // happens to find the sidecar on disk. Ordering by the declaration collapses // that into one pass, and costs a topological sort over three entries. // // STABLE, and that is load-bearing rather than tidiness. OPERATIONS puts // attribution-text LAST on purpose (see the comment there): on a video that has // diarization, the better lane must get there first so the text lane finds a // record it must not overwrite — one wasted classification instead of ~30 model // calls. Kahn's algorithm with a queue seeded and drained in declaration order // preserves every ordering the declarations do not contradict, so that decision // survives. export function orderByDependencies(kinds: Operation[]): Operation[] { const byId = new Map(kinds.map((k) => [k.id, k])); // Only dependencies that are actually IN this selection constrain anything. A // kind that names a disabled or unselected prerequisite is not held back — // it will report `blocked` per video, which is the honest answer, rather than // being silently dropped from the run. const remaining = new Map( kinds.map((k) => [ k.id, (k.dependsOn ?? []).filter((d) => byId.has(d) && d !== k.id).length, ]), ); const dependants = new Map(); for (const k of kinds) { for (const d of k.dependsOn ?? []) { if (!byId.has(d) || d === k.id) continue; const list = dependants.get(d); if (list) list.push(k.id); else dependants.set(d, [k.id]); } } const out: Operation[] = []; const emitted = new Set(); // Repeatedly take the FIRST still-unemitted kind in declaration order whose // prerequisites are all out. Quadratic in the number of kinds, which is three. for (;;) { const next = kinds.find( (k) => !emitted.has(k.id) && (remaining.get(k.id) ?? 0) === 0, ); if (!next) break; emitted.add(next.id); out.push(next); for (const id of dependants.get(next.id) ?? []) { remaining.set(id, (remaining.get(id) ?? 1) - 1); } } // A CYCLE leaves entries unemitted. Append them in declaration order rather // than throwing or dropping them: a mis-declared dependency should degrade to // the old behaviour (run in table order), never wedge the lane or silently // stop a backfill from running at all. for (const k of kinds) if (!emitted.has(k.id)) out.push(k); return out; } // Per-kind counts, the shape every indicator reads. `missing` and `missingInput` // are never summed — see the header. export type OperationCounts = { missing: number; stale: number; missingInput: number; // Work the kind is refusing to attempt under the current configuration (the // diarization duration cap). A THIRD number, alongside the other two that are // never summed. Snapshots written before this field existed do not carry it, // so every read site needs `?? 0` — `.toLocaleString()` on undefined throws. deferred: number; // Waiting on a prerequisite kind's output. A FOURTH number, and the same rule // applies: never summed with the others, and `?? 0` at every read site, // because every snapshot currently on disk predates it. // // This number should FALL on its own as the prerequisite lane works, which is // the whole difference from missingInput — that one only falls if an operator // turns re-download on. blocked: number; // Reachable work that is PART DONE. Unlike deferred and blocked, this one IS // summed into reachableOperationWork — it is work the lane can do today. It is // split out of `stale` because the two cost different amounts and want // different decisions: see the digest entry's state(). // // `?? 0` at every read site, like deferred and blocked before it. Every // snapshot currently on disk predates this field. partial: number; }; export function emptyOperationCounts(): OperationCounts { return { missing: 0, stale: 0, partial: 0, missingInput: 0, deferred: 0, blocked: 0, }; } // Fold one classification into a counts record. Central so no surface invents // its own accounting: `present` and `not-applicable` add to nothing, which is // what makes these counts a WORK LIST rather than a coverage measure. export function addOperationState( counts: OperationCounts, state: OperationClassification, ): void { if (state === "missing") counts.missing++; else if (state === "stale") counts.stale++; else if (state === "partial") counts.partial++; else if (state === "missing-input") counts.missingInput++; else if (state === "deferred") counts.deferred++; else if (state === "blocked") counts.blocked++; } // What the lane can act on WITHOUT re-acquiring media. The number every "how // much is left?" surface should lead with. // // 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 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. // // A blocked video is emphatically not reachable work: there is nothing this // lane can do about it this pass. Counting it would make a corpus with one // diarization and 73,000 waiting attributions report 73,000 jobs ready to run. // // `partial` IS in the total, and that is the on-purpose case the paragraph // above reserves. A part-done video is work the lane can pick up right now and // the run will write to it; leaving it out would make a corpus mid-tags-backfill // report less work than it has. Splitting it from `stale` is about what the // work COSTS, not about whether it is reachable — so the sum is unchanged from // what it would have been before the split, which is the property that keeps // this a refinement rather than a behaviour change. // // `?? 0` because every snapshot on disk predates the field, and undefined would // poison the sum to NaN rather than merely under-report. export function reachableOperationWork(counts: OperationCounts): number { return counts.missing + counts.stale + (counts.partial ?? 0); } // What a channel snapshot stores per kind: the three counts, plus the ids of the // REACHABLE work only. // // The asymmetry is deliberate. A stage card has to list what it would act on, so // those ids have to be somewhere the render path can read without walking the // corpus (there is a guard test forbidding exactly that). But `missingInput` is // ~76,000 videos corpus-wide, and writing that list into all 66 snapshots would // put tens of megabytes of ids on disk to say a number we already have. So: ids // for the actionable half, a count for the other. export type OperationSnapshotEntry = OperationCounts & { // Exactly the reachable set — missing + stale + partial — sorted. Never // includes missing-input, deferred or blocked. Kept equal to // reachableOperationWork(entry) by a test, because a policy leaf hands this // list out as work while the cards render the count. ids: string[]; // How many videos this operation has an OPINION about: everything it did not // classify not-applicable. The denominator, and deliberately not part of // OperationCounts — those are a work list, and mixing a coverage measure into // them is what would let a surface add "done" to "to do". // // It is stored rather than derived because `present` is the one classification // addOperationState throws away, so nothing downstream can reconstruct the // total from the counts alone. With it, present = eligible - (every work // count), which is what presentOperationWork below computes. // // Optional: every snapshot written before this field lacks it, and a reader // that cannot tell how many videos were considered must say so rather than // divide by a zero it invented. eligible?: number; }; // How many videos this operation is DONE with, derived from the stored // denominator minus every work state. Returns null when the snapshot predates // `eligible`, because the honest answer there is "unknown" — a 0 would render as // "nothing digested" on a fully digested channel. export function presentOperationWork( entry: OperationSnapshotEntry, ): number | null { if (entry.eligible == null) return null; return Math.max( 0, entry.eligible - (entry.missing + entry.stale + (entry.partial ?? 0) + entry.missingInput + (entry.deferred ?? 0) + (entry.blocked ?? 0)), ); }