Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 05a6e47378b245380e86a7afed8af8f2e9552289
parent 1815abeaae2cfb01b4c69b81a2caa3841b970c78
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun,  9 Aug 2026 17:11:22 -0400

Put digest in the catalog, on its own declared lane

Digests were built first and never merged in, so they are a second
implementation of this registry's idea: own sweep, own per-channel
batch, own yield probe, own cost planner, own queue keys, own pause, own
counters. backfillSweep.ts says in its header that it is a clone of
digestSweep.ts. This registers the operation.

Nothing about digest generation is rewritten. The entry calls the
functions the controller already calls -- resolveDigestTarget,
isCuesJsonFresh, isSectionFresh, digestVideo -- including the
shared-from-a-duplicate-cluster rule, because two definitions of
"digested" is the exact failure this entry exists to prevent.

Two things are true that were not:

  - THE TRANSCRIPT DEPENDENCY IS DECLARED. PLAN.md:146 states "a video
    cannot be digested until it has a transcript" in prose and nothing
    enforced it. It is now dependsOn ["transcription"] and reports
    `blocked`, which is why transcription and download get
    descriptor-only catalog entries: without them the id resolves to
    nothing and no surface can name what a video is waiting for. They
    stay externally dispatched for four concrete reasons, written down
    at EXTERNAL_OPERATIONS -- chiefly that transcriptionActivity()
    reports busy exactly when transcription runs, so a transcription
    lane would yield to itself.

  - THE LANE IS DECLARED PER OPERATION. The obvious "unification" --
    everything on BACKFILL_QUEUE -- would be a regression: registry.ts
    runs each key at concurrency 1, so it would make the GPU digest lane
    wait on CPU diarization across a multi-week sweep. laneBackfillKinds
    now filters on the queue key so backfillBatch can never dispatch
    digest and can never drop its guards; digestBatch's yield rule reads
    contendsFor off the declared lane instead of re-testing the app id.

resolveTarget became async and channel-scoped, because a digest's
identity includes the hash of the channel's context note. Still resolved
once per run, which is the property that matters.

Deliberately NOT done here: collapsing noDigest into
snapshot.backfill.digest, and merging digest into the backfill
indicators. Both change live coverage numbers on a 78,000-video corpus
for no gain this branch can verify, and both are the unified-rule-model
work.

  common 646/646; backfill+digest+attribution+diarization+disk-space
  e2e 41/41 including a new spec that fails if the two lanes ever
  serialize; tsc clean both packages; editor build clean.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

Diffstat:
Mcommon/controller/backfillBatch.ts | 18++++++++++++++++--
Mcommon/controller/channelSnapshot.ts | 6+++++-
Mcommon/controller/digestBatch.ts | 11++++++++++-
Mcommon/lib/backfillKinds.test.ts | 215+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/lib/backfillKinds.ts | 357+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Meditor/CHANGELOG.md | 1+
Meditor/e2e/backfill.spec.ts | 96+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
7 files changed, 688 insertions(+), 16 deletions(-)

diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts @@ -244,7 +244,16 @@ export async function runBackfillBatch( // 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, kind.resolveTarget(settings)); + 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 ` + @@ -530,7 +539,12 @@ export async function countBackfillWork( const dataDir = path.join(paths.channelsDir, channelSlug, "data"); const dirs = await readdir(dataDir).catch(() => [] as string[]); const targets = new Map<string, unknown>(); - for (const kind of kinds) targets.set(kind.id, kind.resolveTarget(settings)); + for (const kind of kinds) { + targets.set( + kind.id, + await kind.resolveTarget({ settings, paths, channelSlug }), + ); + } let reachable = 0; let missingInput = 0; diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -385,7 +385,11 @@ export async function generateChannelSnapshot( const backfillKinds = laneBackfillKinds(backfillSettings); const backfillTargets: Record<string, unknown> = {}; for (const kind of backfillKinds) { - backfillTargets[kind.id] = kind.resolveTarget(backfillSettings); + backfillTargets[kind.id] = await kind.resolveTarget({ + settings: backfillSettings, + paths, + channelSlug: slug, + }); } const limit = pLimit(SNAPSHOT_VIDEO_CONCURRENCY); diff --git a/common/controller/digestBatch.ts b/common/controller/digestBatch.ts @@ -25,6 +25,10 @@ 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, + laneYieldsToTranscription, +} from "../lib/backfillKinds"; import { loadDigest } from "../lib/digest-server"; import { resolveDigestTarget, @@ -463,8 +467,13 @@ export async function runDigestBatch( // batch, so the sweep resumes the moment the card is free without // re-deriving anything. Only the local (GPU) lane yields — the metered // lane is network-bound and competes for nothing here. + // WHICH LANE YIELDS IS NOW DECLARED, not re-tested here. digestLaneFor + // maps the engine to the digest operation's registered lane, and + // laneYieldsToTranscription reads `contendsFor: "gpu"` off it. Same + // decision, same behaviour — but stated once, next to the queue key it + // belongs with, so a future GPU operation cannot answer it differently. if ( - app.lane === "local-gpu" && + laneYieldsToTranscription(digestLaneFor(app.lane)) && getSettings().digest.yieldToTranscription ) { const activity = transcriptionActivity(); diff --git a/common/lib/backfillKinds.test.ts b/common/lib/backfillKinds.test.ts @@ -5,9 +5,14 @@ import os from "node:os"; import { mkdtemp, mkdir, writeFile, rm, utimes } from "node:fs/promises"; import { getBackfillKind, + allBackfillKinds, laneBackfillKinds, resolveBackfillKinds, orderByDependencies, + operationCatalog, + operationLabel, + digestLaneFor, + laneYieldsToTranscription, addBackfillState, emptyBackfillCounts, reachableBackfillWork, @@ -23,7 +28,8 @@ import { } from "./videoStatus"; import { CUES_FILE_VERSION } from "../controller/normalizeTranscript"; import { attributeOneVideo } from "../controller/attributeOne"; -import { OLLAMA_DIGEST_APP_ID } from "./digest"; +import { DIGEST_FILENAME, OLLAMA_DIGEST_APP_ID } from "./digest"; +import { BACKFILL_QUEUE, DIGEST_LOCAL_QUEUE } from "./queueKeys"; import { ATTRIBUTION_FILENAME, ATTRIBUTION_PROMPT_VERSION, @@ -47,6 +53,19 @@ import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo"; const diarization = getBackfillKind("diarization")!; +// resolveTarget is channel-scoped since the digest entry (a digest's identity +// includes the hash of the channel's context note). None of the kinds tested +// here reads paths/channelSlug, so a stub is honest as well as convenient. +function targetCtx(settings: SiteSettings) { + return { + settings, + paths: { channelsDir: "/nonexistent" } as unknown as Parameters< + typeof diarization.resolveTarget + >[0]["paths"], + channelSlug: "test-channel", + }; +} + function settingsWithDiarization( over: Partial<SiteSettings["diarization"]> = {}, ): SiteSettings { @@ -148,7 +167,7 @@ async function classify( videoDir: dir, videoId: "vid1", files, - target: diarization.resolveTarget(settings), + target: await diarization.resolveTarget(targetCtx(settings)), settings, }); } finally { @@ -490,7 +509,7 @@ async function classifyAttr( videoDir: dir, videoId: "vid1", files, - target: kind.resolveTarget(settings), + target: await kind.resolveTarget(targetCtx(settings)), settings, }); } finally { @@ -824,6 +843,196 @@ async function writeTranscriptFixture(dir: string): Promise<void> { await utimes(path.join(dir, CUES_JSON_FILENAME), base + 10, base + 10); } +// --------------------------------------------------------------------------- +// The digest entry. Registered but NOT dispatched by the backfill lane — see +// laneBackfillKinds — so these tests pin the classification, which is the half +// that is now shared, and the lane declaration, which is what keeps the GPU and +// CPU lanes from serializing. +// --------------------------------------------------------------------------- + +const digestKind = getBackfillKind("digest")!; + +// The digest freshness target, hand-built rather than resolved: resolveTarget +// reads the channel's context note through paths, and these cases are about +// what state() does with a record, not about how the identity is derived. +const DIGEST_TARGET = { + target: { + appId: OLLAMA_DIGEST_APP_ID, + model: "llama3:8b", + promptVersion: 2, + contextHash: "ctx-1", + }, + sections: ["chapters"] as const, +}; + +function digestSidecar(over: Record<string, unknown> = {}): string { + return JSON.stringify({ + videoId: "vid1", + digestSchemaVersion: 1, + sections: { + chapters: { + items: [{ start: 0, title: "Intro" }], + provenance: { + appId: OLLAMA_DIGEST_APP_ID, + model: "llama3:8b", + modelRequested: "llama3:8b", + promptVersion: 2, + contextHash: "ctx-1", + }, + }, + }, + ...over, + }); +} + +async function classifyDigest(opts: { + transcript?: boolean; + cues?: boolean; + staleCues?: boolean; + sidecar?: string; +}): Promise<BackfillClassification> { + const { dir, cleanup } = await fixture({ transcript: opts.transcript }); + try { + if (opts.cues !== false && opts.transcript !== false) { + await writeTranscriptFixture(dir); + if (opts.staleCues) { + // cues.json OLDER than the raw transcript: the transcript changed + // underneath and the normalize pass owes this video a rewrite. + const base = Date.now() / 1000; + await utimes(path.join(dir, CUES_JSON_FILENAME), base - 100, base - 100); + } + } + if (opts.sidecar) { + await writeFile(path.join(dir, DIGEST_FILENAME), opts.sidecar); + } + const files = await readVideoFiles(dir, { checkUntranscribable: true }); + return await digestKind.state({ + videoDir: dir, + videoId: "vid1", + files, + target: DIGEST_TARGET, + settings: defaultSiteSettings(), + }); + } finally { + await cleanup(); + } +} + +test("digest: a video with no transcript is BLOCKED on transcription", async () => { + // PLAN.md states "a video cannot be digested until it has a transcript" in + // prose and nothing enforced it. Now it is declared (dependsOn) and reported. + assert.equal(await classifyDigest({ transcript: false }), "blocked"); + assert.deepEqual(digestKind.dependsOn, ["transcription"]); + // And the dependency resolves to something nameable, which is the whole + // reason the externally-dispatched operations are in the catalog at all. + assert.equal(operationLabel("transcription"), "Transcription"); +}); + +test("digest: a fresh sidecar is present, a mismatched one is stale", async () => { + assert.equal( + await classifyDigest({ sidecar: digestSidecar() }), + "present", + ); + // A model change is work. This is the case the snapshot's old + // "does an ai-digest.json exist?" test got wrong, reading "All digested" + // while the batch reported the whole channel as stale. + assert.equal( + await classifyDigest({ + sidecar: digestSidecar({ + sections: { + chapters: { + items: [{ start: 0, title: "Intro" }], + provenance: { + appId: OLLAMA_DIGEST_APP_ID, + model: "llama3:70b", + modelRequested: "llama3:70b", + promptVersion: 2, + contextHash: "ctx-1", + }, + }, + }, + }), + }), + "stale", + ); +}); + +test("digest: no sidecar at all is missing, not stale", async () => { + assert.equal(await classifyDigest({}), "missing"); +}); + +test("digest: a stale cues.json defers rather than digesting superseded text", async () => { + // Digesting now would describe text that is about to be rewritten, and would + // then look fresh forever. Not a failure — the normalize pass fixes it — so + // it must not be counted as reachable work or as a broken engine. + assert.equal(await classifyDigest({ staleCues: true }), "deferred"); +}); + +test("digest: a digest shared from a duplicate cluster counts as done", async () => { + // Worth ~11% of the sweep. If this entry disagreed with the snapshot's + // noDigest bucket here, every mirror would be regenerated. + assert.equal( + await classifyDigest({ + sidecar: digestSidecar({ derivedFrom: { videoId: "canonical" } }), + }), + "present", + ); +}); + +test("digest declares its own lane, and it is NOT the backfill queue", () => { + // THE PROPERTY THAT KEEPS THEM CONCURRENT. registry.ts runs every non-empty + // queueKey at concurrency 1, so sharing a key with diarization would make the + // GPU lane wait on the CPU lane and vice versa. + assert.notEqual(digestKind.lane.queueKey, BACKFILL_QUEUE); + assert.equal(digestKind.lane.queueKey, DIGEST_LOCAL_QUEUE); + assert.equal(getBackfillKind("diarization")!.lane.queueKey, BACKFILL_QUEUE); + // The two digest lanes are distinct from each other for the same reason. + assert.notEqual( + digestLaneFor("local-gpu").queueKey, + digestLaneFor("remote-api").queueKey, + ); + // Only the GPU lane stands aside for transcription. Yielding the metered lane + // would park something that costs nothing to keep running. + assert.equal(laneYieldsToTranscription(digestLaneFor("local-gpu")), true); + assert.equal(laneYieldsToTranscription(digestLaneFor("remote-api")), false); +}); + +test("the backfill lane never dispatches digest", () => { + // laneBackfillKinds 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. + const settings = settingsWithDiarization(); + const laneIds = laneBackfillKinds(settings).map((k) => k.id); + const allIds = allBackfillKinds(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. + assert.deepEqual(resolveBackfillKinds(settings, ["digest"]), []); +}); + +test("the catalog covers every operation, dispatched here or not", () => { + const ids = operationCatalog().map((o) => o.id); + for (const id of [ + "download", + "transcription", + "diarization", + "attribution-diarized", + "attribution-text", + "digest", + ]) { + assert.ok(ids.includes(id), `${id} missing from ${ids.join(",")}`); + } + // Every declared dependency resolves to a catalogued operation. This is the + // check that would have caught digest.dependsOn naming something that did not + // exist. + for (const op of operationCatalog()) { + for (const dep of op.dependsOn ?? []) { + assert.ok(ids.includes(dep), `${op.id} depends on unknown ${dep}`); + } + } +}); + // The accounting rule the whole feature turns on: reachable work and // needs-re-acquiring are never added together. test("counts keep reachable work and needs-re-acquiring apart", () => { diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts @@ -60,6 +60,19 @@ import type { Paths } from "./paths"; import type { SiteSettings } from "./settings"; import { + BACKFILL_QUEUE, + DIGEST_LOCAL_QUEUE, + DIGEST_REMOTE_QUEUE, + TRANSCRIPTION_QUEUE, +} from "./queueKeys"; +import { + isSectionFresh, + type DigestFreshnessTarget, + type DigestLane, + type DigestSectionKind, +} from "./digest"; +import { loadDigest } from "./digest-server"; +import { diarizationTarget, isDiarizationFresh, type DiarizationFreshnessTarget, @@ -92,6 +105,7 @@ import { diarizeOneVideo } from "../controller/diarizeOne"; // 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. // @@ -142,6 +156,40 @@ export type BackfillClassification = BackfillState | "not-applicable"; // Diarization is ~500-680 s/audio-hour of CPU. export type BackfillCostTier = "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 BackfillLane = { + // 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 BackfillDispatch = "backfill" | "external"; + export type BackfillProbe = { videoDir: string; videoId: string; @@ -192,6 +240,9 @@ export type BackfillKind = { // One line of UI copy: what this backfill is, in the operator's terms. hint: string; tier: BackfillCostTier; + // Where this operation's work runs. See BackfillLane — the queue key is what + // keeps CPU and GPU operations overlapping instead of taking turns. + lane: BackfillLane; // Ids of other kinds in this table whose output this one consumes. // // PURELY DECLARATIVE. It does not gate anything by itself — a kind still @@ -216,24 +267,42 @@ export type BackfillKind = { // 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. - resolveTarget(settings: SiteSettings): unknown; + // + // 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: BackfillTargetContext): unknown | Promise<unknown>; state(probe: BackfillProbe): Promise<BackfillClassification>; run(opts: BackfillRunOptions): Promise<BackfillRunOutcome>; }; +export type BackfillTargetContext = { + settings: SiteSettings; + paths: Paths; + channelSlug: string; +}; + // Diarization: the first entry, and the reason the table exists. const diarization: BackfillKind = { id: "diarization", label: "Speaker diarization", hint: "Speaker turns captured from the audio, written to diarization.json beside the transcript.", 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. + lane: { queueKey: BACKFILL_QUEUE, contendsFor: "cpu" }, 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 => + resolveTarget: ({ settings }): DiarizationFreshnessTarget => diarizationTarget(settings.diarization), async state({ videoDir, files, target, settings }) { // Same eligibility as diarizeAll's transcribedOnly default: the capture lane @@ -383,9 +452,10 @@ const attributionText: BackfillKind = { 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.", tier: "lane", + lane: { queueKey: BACKFILL_QUEUE, contendsFor: "network" }, enabled: (settings) => settings.attribution.enabled && settings.attribution.textOnlyEnabled, - resolveTarget: (settings) => attributionTargetFor(settings, "text-only"), + 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 @@ -435,9 +505,10 @@ const attributionDiarized: BackfillKind = { // 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" }, enabled: (settings) => settings.attribution.enabled && settings.attribution.diarizedEnabled, - resolveTarget: (settings) => attributionTargetFor(settings, "diarized"), + 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 @@ -526,6 +597,239 @@ function toBackfillOutcome(outcome: AttributeOneOutcome): BackfillRunOutcome { } } +// --------------------------------------------------------------------------- +// 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. +// backfillSweep.ts is, in its own words, a clone of digestSweep.ts. Registering +// the operation is how that convergence starts. +// +// 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 +// 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 BY backfillBatch — see backfillQueueKinds below. +const digest: BackfillKind = { + id: "digest", + label: "Digest", + hint: "Chapters and tags generated from the transcript by a local or metered model. Needs a transcript first.", + 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"), + // 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 + // lane must still report how much work is outstanding. + enabled: () => true, + async resolveTarget({ paths, channelSlug }) { + // The existing resolver, verbatim. This is the async, channel-scoped case + // BackfillTargetContext 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) { + // The raw transcript changed under cues.json, so the normalize pass owes + // this video a rewrite. Digesting now would describe superseded text and + // then look fresh forever. It resolves itself, which is why digestVideo + // reports it as `skipped` rather than a failure — and `deferred` is the + // classification with the matching meaning: not attempted, not broken, + // not counted as reachable work. + 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. Same rule as the snapshot's + // noDigest bucket, deliberately, because two definitions of "digested" is + // the exact failure this entry exists to stop. + if (record?.derivedFrom != null) return "present"; + // EVERY configured section must be fresh, matching countMissingDigests and + // the batch. A record with fresh chapters and stale tags is stale. + return sections.every((section) => + isSectionFresh(record, section, freshness), + ) + ? "present" + : record + ? "stale" + : "missing"; + }, + 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; + // BackfillRunOptions 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, + 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"; + }, +}; + +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; digestBatch 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): BackfillLane { + 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" }; +} + +// 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: BackfillLane): 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. BackfillRunOptions 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. backfillBatch is one job per channel; +// autoRunner arbitrates across every channel at once, which is the whole +// point of its policy tree. +export type ExternalOperation = { + id: string; + label: string; + hint: string; + lane: BackfillLane; + dependsOn?: readonly string[]; + dispatch: "external"; +}; + +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.", + // 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:<platform>", contendsFor: "network" }, + dispatch: "external", + }, + { + id: "transcription", + label: "Transcription", + hint: "Turning audio into a transcript. Dispatched by the auto-transcribe runner across the worker pool.", + lane: { queueKey: TRANSCRIPTION_QUEUE, contendsFor: "gpu" }, + dependsOn: ["download"], + dispatch: "external", + }, +]; + +// Every media-derived operation, dispatched here or not. The catalog — what a +// dependency id resolves against, and what a future scheduler enumerates. +export type OperationDescriptor = { + id: string; + label: string; + hint: string; + lane: BackfillLane; + dependsOn?: readonly string[]; + dispatch: BackfillDispatch; +}; + +export function operationCatalog(): OperationDescriptor[] { + return [ + ...EXTERNAL_OPERATIONS, + ...BACKFILL_KINDS.map((k) => ({ + id: k.id, + label: k.label, + hint: k.hint, + lane: k.lane, + dependsOn: k.dependsOn, + dispatch: "backfill" as const, + })), + ]; +} + +// 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; +} + // One entry per backfill known to the system. export const BACKFILL_KINDS: readonly BackfillKind[] = [ diarization, @@ -535,6 +839,10 @@ export const BACKFILL_KINDS: readonly BackfillKind[] = [ // 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 BACKFILL_KIND_BY_ID: Record<string, BackfillKind> = @@ -544,12 +852,43 @@ export function getBackfillKind(id: string): BackfillKind | undefined { return BACKFILL_KIND_BY_ID[id]; } -// The kinds that get the concurrent lane, filtered by their own feature gate. -// Both filters matter: `inline` kinds must never take a lane slot, and a -// disabled feature must report and run nothing. +// 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 laneBackfillKinds below. +export function allBackfillKinds(settings: SiteSettings): BackfillKind[] { + return BACKFILL_KINDS.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 — backfillBatch's dispatch, the channel Backfill card, the +// dashboard's backfill instrument, /actionable's backfill rows — 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 +// 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 digestsPaused, 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 digestBatch'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 laneBackfillKinds(settings: SiteSettings): BackfillKind[] { - return BACKFILL_KINDS.filter( - (k) => k.tier === "lane" && k.enabled(settings), + return allBackfillKinds(settings).filter( + (k) => k.lane.queueKey === BACKFILL_QUEUE, ); } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **Digests are now part of the same catalogue as every other derived-data job.** Digests were built first and never folded in, so they grew a complete parallel copy of the same idea — their own sweep, their own per-channel runner, their own pause switch, their own counters. They are now a registered operation alongside speaker capture and speaker naming, which makes two things true that were not. **The transcript requirement is declared rather than merely written down**: a video with no transcript reports as *waiting on transcription* and names it, where the plan document had stated the rule in prose and nothing enforced it. And **each operation now declares which lane it runs in**, which is what keeps the graphics-card work and the processor work overlapping instead of taking turns — the obvious way to "unify" them would have been to put everything in one queue, and that would have made a weeks-long sweep take roughly twice as long. There is a test that fails if they ever start queueing behind each other. Digest generation itself is untouched: every guard it had — the pause, standing aside for transcription, the spend cap, shortest-first ordering, sharing between duplicate videos — still runs exactly as before. - **"Waiting for an earlier step" is no longer reported as "the media is gone".** Putting names to speakers from the audio needs the speaker-turn capture to have run first. When it hadn't, those videos were reported as *needing their media re-acquired* — which on this corpus is around **73,000 videos**, filed under the one heading that means "fetch the audio again". That was wrong twice over: it was the wrong number in front of the operator, and with re-downloading switched on the system would have spent one download per video fetching audio that **cannot** satisfy the wait, then deleted it again. There is now a distinct **blocked** state for work waiting on an earlier step, counted on its own, never added to the work-to-do total, never handed to a runner, and — crucially — completely unaffected by the re-download setting. The channel card names what each blocked video is waiting for and says plainly that there is nothing to do, because unlike the other numbers **this one falls by itself** as the earlier step runs. - **Backfills now declare what they depend on, and run in that order.** A video whose speaker turns were captured during a pass used to become eligible for speaker naming only on some *later* pass that happened to notice the new file on disk. Prerequisites are now declared and ordered, so both happen in the same pass. A mis-declared or circular dependency degrades to the old ordering rather than stalling anything. - **The automatic downloader now stops when the disk is nearly full — it never did before.** Every download you start by clicking something has checked free space for a long time. The one thing that runs unattended, for days, choosing downloads by itself, did not check at all: there was no mention of disk anywhere in it. That is the process most likely to fill a disk and the least likely to have anyone watching while it does. It now consults the same floor as everything else and simply **goes idle** rather than stopping, so it picks up again on its own once space is free — no restart, nothing to remember. Four more paths that write large files were checked and gated the same way: **re-downloading a single video** (its near-identical sibling, "re-download to archive", already checked — this one had been missed), **persisting kept videos**, which now checks *before each video* instead of once at the start, since it writes full video containers in a loop and the twentieth should not be relying on the first one's headroom, **re-downloading truncated audio**, checked before the old file is deleted rather than after, and the **backfill batch**, whose "disk floor reached" flag was being set and then ignored while it kept asking for more work. Derived files — digests, speaker diarization, attribution — are deliberately **not** gated: they are kilobytes, holding them back frees nothing, and it would throw away days of a multi-week run for no gain. diff --git a/editor/e2e/backfill.spec.ts b/editor/e2e/backfill.spec.ts @@ -8,6 +8,7 @@ import { resetData, resolvePath, writeSettings, + writeDigestVideo, generateReport, } from "./helpers"; @@ -557,6 +558,101 @@ test("a backfill runs concurrently with a transcription", async ({ return seen; }); +// THE PROPERTY REGISTERING DIGEST MUST NOT BREAK. +// +// Digest is now a registered operation in the same table as diarization, and +// the obvious way to "unify" them would have been to put both on BACKFILL_QUEUE. +// That would be a regression, not a simplification: registry.ts runs every +// non-empty queueKey at concurrency 1, so one shared key makes the GPU digest +// lane wait on CPU diarization and vice versa — across a sweep measured in +// weeks. The lane is declared per operation precisely so this keeps working. +test("a digest runs concurrently with a backfill, not behind it", async ({ + page, + request, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings({ + ...backfillSettings(), + digest: { + localAppId: "ollama-direct", + remoteAppId: "claude-code", + sections: ["chapters"], + }, + }); + // ENOUGH backfill work that the lane is provably still running when the + // digest is submitted. The fake diarizer takes ~0.2s per video, so the three + // fixture videos finish in well under a second — faster than a click — and + // the test would measure nothing. 40 dirs buys ~8s of overlap window. + for (const id of ["vidA", "vidB"]) await seedTranscript(id); + for (let i = 0; i < 40; i++) { + const id = `slowdiar${String(i).padStart(2, "0")}`; + const dir = resolvePath(`test-transcripts/channels/${SLUG}/data/${id}`); + await mkdir(dir, { recursive: true }); + await writeFile(`${dir}/audio.m4a`, `fake audio ${id}\n`); + await writeFile( + `${dir}/transcript.json`, + '{"transcription":[{"text":"hi"}]}\n', + ); + } + // Digest work: a video with a real normalized transcript. + await writeDigestVideo({ channelSlug: SLUG, videoId: "digestme" }); + + await generateReport(page, SLUG); + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + await expect + .poll( + async () => + activeJobs(request).then((j) => + j.some((x) => x.kind === "backfill-channel" && x.status === "running"), + ), + { timeout: 30_000 }, + ) + .toBe(true); + + await page.getByRole("button", { name: "Digest stage summary" }).click(); + await page.getByRole("button", { name: "Digest channel" }).click(); + + const seen = await expect + .poll( + async () => { + const jobs = await activeJobs(request); + const backfill = jobs.find((j) => j.kind === "backfill-channel"); + const digest = jobs.find((j) => j.kind.startsWith("digest-channel")); + // A digest sitting QUEUED behind a running backfill is exactly the + // failure a shared queue key produces, so record that state distinctly + // rather than letting it time out as a generic "never ran". + if (digest?.status === "running" && backfill?.status === "running") { + expect(digest.queueKey).not.toBe(backfill.queueKey); + return "concurrent"; + } + return digest?.status ?? "absent"; + }, + { timeout: 30_000, intervals: [100] }, + ) + .toBe("concurrent"); + + // Settle, so neither lane's children leak into the next spec. + await expect + .poll( + async () => { + const jobs = await activeJobs(request); + return jobs.filter( + (j) => + j.kind.startsWith("digest-channel") || + j.kind === "backfill-channel", + ).length; + }, + { timeout: 90_000, intervals: [250] }, + ) + .toBe(0); + return seen; +}); + // (10) /actionable lists the channel, keeps the two numbers in separate columns, // and — the part that is hand-exhaustive and so easy to miss — does not hide the // whole page behind "nothing pending".