Archilyzer · Source

archilyzer

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

commit bbcce85dd22484fdf1f7cdc262a4851c55e238ba
parent e8b29c4bb75afe561ab4a3cfd23fc22418ffe933
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun,  9 Aug 2026 19:33:00 -0400

Give the digest work list one definition, and let the registry carry it

The digest kind's state() was called by nothing. channelSnapshot resolved its
per-video probes from laneBackfillKinds(), which filters on the backfill queue
key and so excludes digest by construction — Phase C registered the operation
and moved no number anywhere. It now iterates allBackfillKinds(): a work list
the snapshot does not carry is one nothing can count, select or schedule.

The hazard that made this worth designing rather than just doing: four surfaces
summed Object.values(snapshot.backfill) on the assumption that the map WAS the
backfill lane, so an entry per catalog operation would have folded ~75,000
digest videos into the dashboard instrument, /actionable's backfill rows and the
widget. They now use laneEntriesOf(), which asks the REGISTRY whether a kind runs
on BACKFILL_QUEUE. A hardcoded `key !== "digest"` would have fixed today and left
the same trap armed for the next operation on a lane of its own.

Measured on the corpus, read-only, before changing any surface — and it
contradicts what the plan predicted. The count does not rise by ~1,530, it FALLS
by 414: +1,516 from the 11 channels whose snapshots predate the noDigest bucket,
-1,933 now correctly deferred because their cues.json is stale. That second
number is concentrated: piratesoftware alone is 1,683, so 94% of that channel's
apparent digest backlog was work digestVideo would have taken and immediately put
back. blocked lands at exactly 1,631 (the untranscribed count) and eligible at
78,885 of 78,963, i.e. exactly the 78 untranscribable.

partial is a fifth state because digestVideo regenerates only the sections that
are not fresh, so "has chapters, needs tags" costs a fraction of "has neither"
and both used to report as stale. It IS reachable work — the split is about cost,
never about whether the lane can act — and it made candidateAction's `never`
check fire, which is what that check is for.

Two of the plan's premises were wrong and are corrected here. backfillSweep was
never comparing scope on " "; it already joined on a literal NUL and was right —
the raw byte is what made the file binary to grep and got it misread. Both that
line and backfillBatch's key now use the \u0000 escape. And the registry does not
expose `present`: addBackfillState discards it by design, so `eligible` is stored
on the snapshot entry and presentBackfillWork() returns null — never 0 — for a
snapshot that predates it.

Two latent defects found on the way. countBackfillWork re-spelled the reachable
rule instead of sharing it and had already lost `blocked`, so a channel of
prerequisite-waiting videos reported three zeroes and read as finished. And the
runtime backstop list in backfillBatch.test.ts was itself missing "blocked" — the
one test whose job is catching an unwired state had missed one.

digestWorkOf() keeps reading noDigest as a fallback. A third of the snapshots on
disk are months old, and without it every un-regenerated channel would read as
fully digested the moment this landed.

Also, asked for mid-session: a diarization progress bar. A backfill task was
handed the DOWNLOAD parser, which matches nothing in diarization output, so a
40-minute windowed run showed an indeterminate bar while the job bar sat at 0/9
videos. diarize-sherpa.py was already printing "window 3/9 ... in 210.3s" and
nothing read it; the new parser turns that into a fraction plus an ETA measured
from the windows already done, chained ahead of the download parser so a
re-acquiring backfill keeps its download progress.

Verified: 664/664 common unit tests (was 646), tsc clean in common and editor,
editor production build exit 0, and 39/39 across backfill/digest/attribution/
disk-space — including a new spec that asserts BOTH that snapshot.backfill.digest
exists with the right ids and that the backfill instrument does not move because
of it. Asserting only the second half would have passed vacuously.

Not done, and it needs the live editor restarted: regenerating the 66 snapshots.
They only gain a digest entry through an editor running this code, and the one on
:3001 is a production build of the old one, currently supervising a diarization.

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

Diffstat:
Mcommon/controller/backfillBatch.test.ts | 7+++++++
Mcommon/controller/backfillBatch.ts | 63++++++++++++++++++++++++++++++++++++++++++++++-----------------
Mcommon/controller/backfillSweep.ts | 2+-
Acommon/controller/channelSnapshot.test.ts | 153+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/channelSnapshot.ts | 210++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mcommon/jobs/progressParsers.test.ts | 81++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/progressParsers.ts | 66++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/taskHooks.ts | 14++++++++++++--
Mcommon/lib/backfillKinds.test.ts | 131++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/lib/backfillKinds.ts | 154+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Meditor/CHANGELOG.md | 5+++++
Meditor/app/actionable/lib/loadActionable.ts | 40+++++++++++++++++++++++++++++++---------
Meditor/app/api/widget/actionable/route.ts | 13++++++++-----
Meditor/app/api/widget/sync/route.ts | 43+++++++++++++++++++++++++++++++++++++++----
Meditor/app/channels/[slug]/components/stages/DigestStage.tsx | 59++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Meditor/app/channels/[slug]/lib/stageStatus.ts | 48++++++++++++++++++++++++++++++++++++++----------
Meditor/app/channels/[slug]/page.tsx | 16+++++++++++++++-
Meditor/app/components/dashboard/PipelineBand.tsx | 17+++++++++++++++--
Meditor/app/components/dashboard/types.ts | 10+++++++---
Meditor/e2e/backfill.spec.ts | 88+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mplans/FACTS.md | 65+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mplans/STATE.md | 150+++++++++++++++++++++++++++++++++++++++++++++++++------------------------------
22 files changed, 1282 insertions(+), 153 deletions(-)

diff --git a/common/controller/backfillBatch.test.ts b/common/controller/backfillBatch.test.ts @@ -154,9 +154,16 @@ test("every classification has an explicit decision", () => { const ALL: BackfillClassification[] = [ "present", "stale", + // Part-done. Dispatched like stale — the runner regenerates only the + // sections that are not fresh — and listed here so that stays a decision + // rather than something a later edit can quietly drop. + "partial", "missing", "missing-input", "deferred", + // This list had already drifted: `blocked` was added to the union and to + // candidateAction, and never to the backstop meant to catch exactly that. + "blocked", "not-applicable", ]; for (const state of ALL) { diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts @@ -39,6 +39,9 @@ import type { TaskTracker } from "../jobs/taskHooks"; import type { JobProgress } from "../jobs/registry"; import { readVideoFiles } from "../lib/videoStatus"; import { + addBackfillState, + emptyBackfillCounts, + reachableBackfillWork, resolveBackfillKinds, type BackfillClassification, type BackfillKind, @@ -194,6 +197,14 @@ export function candidateAction( return "blocked"; case "missing": case "stale": + // Part-done is reachable work like any other, and the runner already does + // the right thing with it without being told: digestVideo regenerates the + // sections that are not fresh and leaves the ones that are. The split + // between this and `stale` is about what the work COSTS a planner, not + // about whether it is dispatched — so it belongs on exactly this line, and + // giving it its own CandidateAction would only invite a caller to treat it + // as another do-not-run state alongside deferred and blocked. + case "partial": return "dispatch"; default: { const unhandled: never = state; @@ -268,7 +279,7 @@ export async function runBackfillBatch( // past items that have been attempted, while eligibility for the item it stops // on is always re-read from disk. const attempted = new Set<string>(); - const key = (kindId: string, id: string) => `${kindId} + const key = (kindId: string, id: string) => `${kindId}\u0000${id}`; let cursor = 0; let kindIndex = 0; @@ -532,10 +543,16 @@ export async function countBackfillWork( paths: Paths, channelSlug: string, kindIds?: string[], -): Promise<{ reachable: number; missingInput: number; deferred: number }> { +): Promise<{ + reachable: number; + missingInput: number; + deferred: number; + blocked: number; +}> { const settings = getSettings(); const kinds = resolveBackfillKinds(settings, kindIds); - if (kinds.length === 0) return { reachable: 0, missingInput: 0, deferred: 0 }; + if (kinds.length === 0) + return { reachable: 0, missingInput: 0, deferred: 0, blocked: 0 }; const dataDir = path.join(paths.channelsDir, channelSlug, "data"); const dirs = await readdir(dataDir).catch(() => [] as string[]); const targets = new Map<string, unknown>(); @@ -546,26 +563,38 @@ export async function countBackfillWork( ); } - let reachable = 0; - let missingInput = 0; - let deferred = 0; + // Folded through the SHARED counters rather than a private if-chain. The + // header above promises this function derives work the same way the batch + // does; it used to re-spell the rule instead, and had already drifted — + // `blocked` was classified by state() and then counted by nothing here, so a + // channel of prerequisite-waiting videos reported three zeroes and looked + // finished. addBackfillState + reachableBackfillWork is the same pair the + // channel snapshot folds with, so this path and the snapshot path can no + // longer disagree about what a state means, and a state added later is + // counted here without anyone remembering to come back. + const counts = emptyBackfillCounts(); for (const id of dirs) { const videoDir = path.join(dataDir, id); const files = await readVideoFiles(videoDir, { checkUntranscribable: true, }); for (const kind of kinds) { - const state = await kind.state({ - videoDir, - videoId: id, - files, - target: targets.get(kind.id), - settings, - }); - if (state === "missing" || state === "stale") reachable++; - else if (state === "missing-input") missingInput++; - else if (state === "deferred") deferred++; + addBackfillState( + counts, + await kind.state({ + videoDir, + videoId: id, + files, + target: targets.get(kind.id), + settings, + }), + ); } } - return { reachable, missingInput, deferred }; + return { + reachable: reachableBackfillWork(counts), + missingInput: counts.missingInput, + deferred: counts.deferred, + blocked: counts.blocked, + }; } diff --git a/common/controller/backfillSweep.ts b/common/controller/backfillSweep.ts @@ -379,7 +379,7 @@ export async function startBackfillSweep( const channelScope = opts.channelSlugs ?? settings.backfill.sweepChannels; const kindScope = opts.kindIds ?? settings.backfill.sweepKinds; const same = (a: string[], b: string[]) => - a.join(" + a.join("\u0000") === b.join("\u0000"); if ( !settings.backfill.sweepEnabled || !same(settings.backfill.sweepChannels, channelScope) || diff --git a/common/controller/channelSnapshot.test.ts b/common/controller/channelSnapshot.test.ts @@ -0,0 +1,153 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { digestWorkOf, foldBackfillEntry } from "./channelSnapshot"; +import { + emptyBackfillCounts, + reachableBackfillWork, + type BackfillClassification, +} from "../lib/backfillKinds"; + +// The snapshot's accounting, tested where it can actually be reached. +// +// generateChannelSnapshot needs lmdb, an archive reader and a corpus on disk, so +// the fold inside it was never exercised by anything. These two functions are +// the parts every downstream surface trusts, and both have a failure mode that +// looks exactly like success: a work list that disagrees with its own count, and +// a coverage number that reads "all done" over a channel nothing has looked at. + +function fold(states: BackfillClassification[]) { + return foldBackfillEntry(states.map((state, i) => ({ id: `v${i}`, state }))); +} + +test("ids are EXACTLY the reachable set — the invariant nothing else checks", () => { + // A policy leaf hands `ids` out as the work to do while the cards render the + // count. The two drifting apart is a progress bar that stalls one short of + // complete forever, and this pair has drifted once already (countBackfillWork + // silently dropped `blocked`). + const entry = fold([ + "missing", + "stale", + "partial", + "present", + "blocked", + "deferred", + "missing-input", + "not-applicable", + ]); + assert.equal(entry.ids.length, reachableBackfillWork(entry)); + assert.deepEqual(entry.ids, ["v0", "v1", "v2"]); + // And the three that must never be reachable are still counted, just not there. + assert.equal(entry.blocked, 1); + assert.equal(entry.deferred, 1); + assert.equal(entry.missingInput, 1); +}); + +test("eligible counts every video the operation had an opinion about", () => { + // not-applicable is the ONLY exclusion: an untranscribable video is out of + // scope, everything else is either done or outstanding. Without this the + // `present` half is unrecoverable, because addBackfillState discards it. + const entry = fold([ + "present", + "present", + "missing", + "blocked", + "not-applicable", + "not-applicable", + ]); + assert.equal(entry.eligible, 4); + // 4 eligible - 1 missing - 1 blocked = 2 present, which is the number no + // counter stores directly. + assert.equal(digestWorkOf({ backfill: { digest: entry } }).present, 2); +}); + +test("a video no kind classified is skipped, not counted as done", () => { + const entry = foldBackfillEntry([ + { id: "a", state: "missing" }, + { id: "b", state: undefined }, + ]); + assert.equal(entry.eligible, 1); + assert.deepEqual(entry.ids, ["a"]); +}); + +test("ids are sorted, because a snapshot is compared byte-for-byte", () => { + const entry = foldBackfillEntry( + ["zz", "aa", "mm"].map((id) => ({ id, state: "missing" as const })), + ); + assert.deepEqual(entry.ids, ["aa", "mm", "zz"]); +}); + +// --------------------------------------------------------------------------- +// digestWorkOf: one reader, three generations of snapshot. + +test("digestWorkOf prefers the registry entry over the legacy bucket", () => { + // When both exist the registry wins — it is the classification the RUNNER + // uses, and the bucket has no cues-staleness or transcript gate, so the two + // genuinely disagree. Two definitions of "digested" is the failure this whole + // change exists to end. + const work = digestWorkOf({ + backfill: { + digest: { + ...emptyBackfillCounts(), + missing: 2, + partial: 1, + blocked: 4, + deferred: 3, + ids: ["a", "b", "c"], + eligible: 10, + }, + }, + buckets: { noDigest: ["x", "y", "z", "w", "v", "u", "t"] }, + }); + assert.equal(work.source, "registry"); + assert.equal(work.reachable, 3); + assert.deepEqual(work.ids, ["a", "b", "c"]); + assert.equal(work.blocked, 4); + assert.equal(work.deferred, 3); + assert.equal(work.partial, 1); +}); + +test("digestWorkOf falls back to noDigest so a stale snapshot is not 'all done'", () => { + // A THIRD of the channels on disk carry snapshots months old. Without this + // every one of them would report zero digest work until the regen reached it, + // which is the single most dangerous way for a coverage number to be wrong. + const work = digestWorkOf({ buckets: { noDigest: ["a", "b"] } }); + assert.equal(work.source, "bucket"); + assert.equal(work.reachable, 2); + assert.deepEqual(work.ids, ["a", "b"]); + // The bucket cannot know these, and says so via `source` rather than by + // implying a split it never computed. + assert.equal(work.blocked, 0); + assert.equal(work.deferred, 0); + assert.equal(work.present, null); + assert.equal(work.eligible, null); +}); + +test("digestWorkOf reports UNKNOWN, not zero coverage, when it cannot tell", () => { + for (const snapshot of [null, undefined, {}, { buckets: {} }]) { + const work = digestWorkOf(snapshot); + assert.equal(work.source, "none"); + assert.equal(work.reachable, 0); + // The distinction that matters: no work outstanding AND no idea how much is + // done. A 0 here would render as "0% digested" on every surface. + assert.equal(work.present, null); + assert.equal(work.eligible, null); + } +}); + +test("digestWorkOf survives an entry written before `partial` and `eligible`", () => { + const work = digestWorkOf({ + backfill: { + digest: { + missing: 4, + stale: 1, + missingInput: 0, + deferred: 0, + blocked: 0, + ids: ["a", "b", "c", "d", "e"], + } as never, + }, + }); + assert.equal(work.reachable, 5); + assert.equal(work.partial, 0); + assert.equal(work.present, null); +}); diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -27,8 +27,11 @@ import { loadDigest } from "../lib/digest-server"; import { isSectionFresh } from "../lib/digest"; import { addBackfillState, + allBackfillKinds, emptyBackfillCounts, - laneBackfillKinds, + presentBackfillWork, + reachableBackfillWork, + DIGEST_KIND_ID, type BackfillClassification, type BackfillSnapshotEntry, } from "../lib/backfillKinds"; @@ -78,17 +81,26 @@ export type ChannelSnapshot = { // local lane is actually carrying the corpus. Optional: older snapshots lack // it; readers default to {}. digestEngines?: Record<string, number>; - // Per-backfill-kind work counts, keyed by kind id (see lib/backfillKinds.ts). + // Per-OPERATION work counts, keyed by kind id (see lib/backfillKinds.ts). // Beside `totals` and NOT in `buckets`, following the digestEngines precedent // above for the same reason: buckets is a closed literal of `string[]` id // lists, and this is per-kind counts. // - // THE THREE NUMBERS ARE NEVER SUMMED. `missing` + `stale` is work the lane can - // do today; `missingInput` needs the media re-acquired and, measured on this - // corpus, is 91x larger. A single "remaining" figure here would put every - // channel permanently at the top of every list — which is the documented + // THE NUMBERS ARE NEVER SUMMED. `missing` + `stale` + `partial` is work the + // lane can do today; `missingInput` needs the media re-acquired and, measured + // on this corpus, is 91x larger. A single "remaining" figure here would put + // every channel permanently at the top of every list — which is the documented // reason /api/widget/actionable refuses to filter on `noDigest`. // + // THIS MAP IS NO LONGER THE BACKFILL LANE, and anything reading it generically + // must say which lane it means. It is written from allBackfillKinds — every + // enabled operation in the catalog — because a work list the snapshot does not + // carry is a work list nothing can select, count or schedule. The digest + // operation runs on its own queue key and is the first entry here that the + // backfill lane must not touch. Read it with laneEntriesOf() to get the lane, + // or by id to get one operation; a bare Object.values() over this map now + // means "every operation", which on this corpus is a ~75,000-video difference. + // // Optional: snapshots written before this existed lack it, and readers default // to {}. backfill?: Record<string, BackfillSnapshotEntry>; @@ -300,6 +312,125 @@ export function normalizeMaybeMissing( return { ids: raw?.ids ?? [], checkedAt: raw?.checkedAt ?? "" }; } +// Fold one operation's per-video classifications into the record a snapshot +// stores for it. +// +// EXTRACTED SO THE INVARIANT CAN BE TESTED. It is four lines of accounting, but +// they are the four that every downstream surface trusts, and inside +// generateChannelSnapshot they could only be exercised by standing up lmdb, an +// archive reader and a corpus on disk — which is to say, never. The rule that +// matters and that nothing else checks: +// +// ids.length === reachableBackfillWork(entry) +// +// A policy leaf hands `ids` out as the work to do while the stage cards render +// the count, so the two drifting apart is a progress bar that stalls one short +// of complete forever. That is not hypothetical — countBackfillWork re-spelled +// this same rule and had already lost `blocked` from it. +export function foldBackfillEntry( + videos: Iterable<{ id: string; state: BackfillClassification | undefined }>, +): BackfillSnapshotEntry { + const counts = emptyBackfillCounts(); + const ids: string[] = []; + // The denominator: every video this kind had an opinion about. `present` is + // the classification addBackfillState deliberately discards, so without this + // no reader can tell "0 outstanding because it is all done" from "0 + // outstanding because there was nothing here". + let eligible = 0; + for (const { id, state } of videos) { + if (!state) continue; + if (state !== "not-applicable") eligible++; + addBackfillState(counts, state); + if (state === "missing" || state === "stale" || state === "partial") { + ids.push(id); + } + } + return { ...counts, ids: ids.sort(), eligible }; +} + +// The digest work list, from whichever of the two places this snapshot has it. +// +// ONE READER, because the alternative is what this whole change exists to undo. +// The digest layer grew its own counter (`buckets.noDigest`) before the +// operation registry existed, and for a while the two disagreed by an entire +// channel — the stage card read "All digested" while the batch reported +// everything stale. There is now one classification, in the registry, and this +// is the only function allowed to decide where a surface reads it from. +// +// THE FALLBACK IS NOT DEFENSIVENESS, IT IS THE MIGRATION. A third of the +// channels on disk carry snapshots months old, and 11 of them predate the +// `noDigest` bucket itself — so there are three generations in play at once: +// snapshots with `backfill.digest` (written after this change), snapshots with +// only `noDigest`, and snapshots with neither. Without the middle case every +// channel that has not been regenerated yet would read as fully digested, which +// is the single most dangerous way for a coverage number to be wrong. `source` +// is returned so a surface can say which generation it is looking at rather +// than quietly presenting a guess as a measurement. +export type DigestWork = { + // Videos needing digest work that the lane can do right now. + ids: string[]; + reachable: number; + // Zero for a `noDigest`-derived answer: that bucket has no cues-staleness + // gate and no transcript gate, so it cannot distinguish these. + blocked: number; + deferred: number; + partial: number; + // Videos with a current digest, and how many were considered. Null when the + // snapshot cannot say — never 0, which would render as "none digested". + present: number | null; + eligible: number | null; + source: "registry" | "bucket" | "none"; +}; + +export function digestWorkOf( + snapshot: + | Pick<ChannelSnapshot, "backfill"> & { + buckets?: Partial<ChannelSnapshot["buckets"]>; + } + | null + | undefined, +): DigestWork { + const entry = snapshot?.backfill?.[DIGEST_KIND_ID]; + if (entry) { + return { + ids: entry.ids ?? [], + reachable: reachableBackfillWork(entry), + blocked: entry.blocked ?? 0, + deferred: entry.deferred ?? 0, + partial: entry.partial ?? 0, + present: presentBackfillWork(entry), + eligible: entry.eligible ?? null, + source: "registry", + }; + } + const bucket = snapshot?.buckets?.noDigest; + if (bucket) { + return { + ids: bucket, + reachable: bucket.length, + // The bucket genuinely does not know these. Reporting 0 is not a guess: + // it is the honest statement that this snapshot carries no such split, + // and `source` says why. + blocked: 0, + deferred: 0, + partial: 0, + present: null, + eligible: null, + source: "bucket", + }; + } + return { + ids: [], + reachable: 0, + blocked: 0, + deferred: 0, + partial: 0, + present: null, + eligible: null, + source: "none", + }; +} + // The snapshot READER lives in ./channels — reading a report shouldn't require // loading the machinery that generates one (this module pulls in lmdb, the // archive reader and the digest layer). Re-exported here so the name stays @@ -376,13 +507,28 @@ export async function generateChannelSnapshot( (e) => !e.toLowerCase().startsWith("youtube"), ); - // Which backfills are live, and what identity each would produce right now. + // Which operations are live, and what identity each would produce right now. // Resolved ONCE per channel — a settings read and some string work — so the // per-video probe below is a comparison rather than a derivation, exactly as - // digestTarget is above. An empty list is the default (every backfill feature - // ships off), and it makes the whole indicator free. + // digestTarget is above. + // + // allBackfillKinds, NOT laneBackfillKinds. The snapshot's job is to carry a + // work list for every operation in the catalog, not for one lane: nothing can + // count, select or schedule work the snapshot does not record, which is why + // registering the digest kind in Phase C changed no number anywhere — its + // state() was called by nothing. The lane filter belongs on the READ side, + // where a surface says which lane it means (laneEntriesOf). + // + // COST, measured rather than assumed, because this runs per video over ~79,000 + // of them: the digest classification is 0.34 ms/video on a 125-video channel + // and 0.38 ms/video on a 773-video one — 0.7x to 1.2x the readVideoFiles call + // directly above, which this path already pays. About 30 s across the whole + // corpus, spread over 67 per-channel regenerations. The re-listing inside + // isCuesJsonFresh is the obvious thing to fold into `files.entries` and it is + // deliberately NOT done: at this cost it would be an optimization with no + // measurement behind it. const backfillSettings = getSettings(); - const backfillKinds = laneBackfillKinds(backfillSettings); + const backfillKinds = allBackfillKinds(backfillSettings); const backfillTargets: Record<string, unknown> = {}; for (const kind of backfillKinds) { backfillTargets[kind.id] = await kind.resolveTarget({ @@ -451,13 +597,17 @@ export async function generateChannelSnapshot( isVideoTranscribed(files) && !files.isUntranscribable ? await loadDigest(dir) : null; - // Backfill state, per registered kind. Costs NOTHING when no backfill - // feature is enabled — `backfillKinds` is empty then, and this loop does - // not run — which is the normal case and why this can live in the - // snapshot's hot path at all. When one IS enabled it is one small - // sidecar read per transcribed video, the same shape and cost as the - // digest read directly above; the classification itself reuses `files` - // rather than re-listing the directory. + // Work state, per registered operation. + // + // This USED to cost nothing in the default configuration, because every + // backfill feature ships off and `backfillKinds` was therefore empty. + // That is no longer true and the change is deliberate: the digest + // operation is enabled whenever an app is configured (its pause lives at + // dispatch, not here, so a paused lane still reports its outstanding + // work), so this loop now always runs at least once per video. See the + // measured per-video cost where backfillKinds is resolved above — it is + // roughly the readVideoFiles call this classification reuses rather than + // repeats. const backfill: Record<string, BackfillClassification> = {}; for (const kind of backfillKinds) { backfill[kind.id] = await kind.state({ @@ -850,19 +1000,23 @@ export async function generateChannelSnapshot( } : undefined; - // Fold the per-video classifications into per-kind counts. Only the reachable - // half carries ids — see BackfillSnapshotEntry for why missing-input does not. + // Fold the per-video classifications into per-operation counts. Only the + // reachable half carries ids — see BackfillSnapshotEntry for why missing-input + // does not. + // + // `ids` MUST stay exactly the set reachableBackfillWork counts, because a + // policy leaf consumes this list as the work to hand out while the stage cards + // render the number: a list and a total that disagree is a progress bar that + // stalls one short of complete forever. `partial` is reachable, so it is in + // both. channelSnapshot.test.ts pins `ids.length === reachableBackfillWork()` + // for exactly this reason: the two derivations have drifted apart once already + // (countBackfillWork silently dropped `blocked`), and nothing about the shape + // here would have caught it. const backfillCounts: Record<string, BackfillSnapshotEntry> = {}; for (const kind of backfillKinds) { - const counts = emptyBackfillCounts(); - const ids: string[] = []; - for (const v of perVideo) { - const state = v.backfill[kind.id]; - if (!state) continue; - addBackfillState(counts, state); - if (state === "missing" || state === "stale") ids.push(v.id); - } - backfillCounts[kind.id] = { ...counts, ids: ids.sort() }; + backfillCounts[kind.id] = foldBackfillEntry( + perVideo.map((v) => ({ id: v.id, state: v.backfill[kind.id] })), + ); } const corruptSourceSet = new Set(corruptSource); diff --git a/common/jobs/progressParsers.test.ts b/common/jobs/progressParsers.test.ts @@ -1,6 +1,9 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { createDownloadProgressParser } from "./progressParsers"; +import { + createDiarizeProgressParser, + createDownloadProgressParser, +} from "./progressParsers"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/progressParsers.test.ts @@ -59,3 +62,79 @@ test("media download fraction is byte-based from the start", () => { ); assert.equal(u?.fraction, 0.25); }); + +// --------------------------------------------------------------------------- +// Diarization. The lines are real output from scripts/diarize-sherpa.py. + +test("diarization reports window progress with a measured ETA", () => { + const p = createDiarizeProgressParser(); + // The plan line makes the bar determinate before the first window finishes — + // which on a 6-hour file is several minutes of otherwise-blank spinner. + const planned = p.feed( + "diarize-sherpa: 6.20 h -> 9 window(s) of 45 min (15s overlap), windowed mode", + ); + assert.equal(planned?.fraction, 0); + assert.equal(planned?.detail, "0/9 windows"); + + const first = p.feed( + "diarize-sherpa: window 1/9 (0-45 min): 412 turns, 6 local speaker(s) in 300.0s", + ); + // Printed AFTER the window is processed, so 1/9 is one window DONE. + assert.equal(first?.fraction, 1 / 9); + // 8 windows left at the one measured time: 2400s = 40:00. + assert.equal(first?.detail, "window 1/9 · ETA 40:00"); + + // The mean moves with the second sample: (300+180)/2 = 240 × 7 = 1680 = 28:00. + const second = p.feed( + "diarize-sherpa: window 2/9 (45-90 min): 300 turns, 4 local speaker(s) in 180.0s", + ); + assert.equal(second?.fraction, 2 / 9); + assert.equal(second?.detail, "window 2/9 · ETA 28:00"); +}); + +test("diarization drops the ETA on the last window rather than showing 0:00", () => { + const p = createDiarizeProgressParser(); + p.feed("diarize-sherpa: window 1/2 (0-45 min): 10 turns, 2 local speaker(s) in 60.0s"); + const last = p.feed( + "diarize-sherpa: window 2/2 (45-90 min): 10 turns, 2 local speaker(s) in 60.0s", + ); + assert.equal(last?.fraction, 1); + assert.equal(last?.detail, "window 2/2"); +}); + +test("a skipped window leaves the bar alone rather than inventing a step", () => { + const p = createDiarizeProgressParser(); + const first = p.feed( + "diarize-sherpa: window 1/4 (0-45 min): 10 turns, 2 local speaker(s) in 100.0s", + ); + assert.equal(first?.fraction, 0.25); + // The REAL zero-samples line (diarize-sherpa.py) prints a bare index with no + // "/N" and no elapsed time: "window 2 decoded zero samples; skipping". So it + // matches nothing here and the bar holds until window 3 reports — which is + // the correct outcome, and the reason this asserts the actual string rather + // than a plausible-looking one. + assert.equal( + p.feed("diarize-sherpa: window 2 decoded zero samples; skipping"), + null, + ); + // The mean is untouched by the skip: still 100s, 2 windows left = 3:20. + const third = p.feed( + "diarize-sherpa: window 3/4 (90-135 min): 10 turns, 2 local speaker(s) in 100.0s", + ); + assert.equal(third?.detail, "window 3/4 · ETA 1:40"); +}); + +test("the diarize parser ignores everything that is not its own output", () => { + const p = createDiarizeProgressParser(); + // Non-windowed short files print no per-window line: indeterminate is honest. + assert.equal(p.feed("diarize-sherpa: 3.2 min decoded, diarizing…"), null); + // And a backfill task that is re-acquiring media must fall through to the + // download parser untouched — this parser must not claim yt-dlp's lines. + assert.equal( + p.feed( + "DLOM_PROGRESS status=downloading downloaded=500 total=1000 frag=NA/NA speed=1000 eta=1", + ), + null, + ); + assert.equal(p.feed("[download] 50.0% of 1.00GiB at 2.31MiB/s ETA 00:30"), null); +}); diff --git a/common/jobs/progressParsers.ts b/common/jobs/progressParsers.ts @@ -372,6 +372,72 @@ export function createChoughProgressParser(): { // Progress = completed segment index / total segments. Once a segment has // finished the line carries a per-video ETA (avg measured time per segment × // remaining segments), which we surface verbatim in the detail. +// scripts/diarize-sherpa.py, windowed mode. It already printed everything a bar +// needs and nothing was reading it: +// +// diarize-sherpa: 6.20 h -> 9 window(s) of 45 min (15s overlap), windowed mode +// diarize-sherpa: window 3/9 (90-135 min): 412 turns, 6 local speaker(s) in 210.3s +// +// WHY THIS IS WORTH A PARSER AT ALL. A backfill task got the DOWNLOAD parser, +// which matches nothing in this output, so a diarization showed an indeterminate +// bar — for up to 46 minutes on the 6-hour files, while the job-level bar sat at +// "0/9 videos" because no video had finished yet. Both readouts were technically +// honest and together they said nothing: the operator could not tell a working +// diarization from a wedged one, which is the exact case windowing was built for. +// +// The per-window line is printed AFTER the window is processed and carries its +// elapsed time, so `i` is a COMPLETED count and the times give a measured ETA +// rather than a guess. Windows are near-identical in cost (same duration, same +// model), which is what makes a flat mean good enough here — the parakeet parser +// takes its ETA from the producer for the same reason the producer knows better +// there and does not here. +// +// Short files take the non-windowed path and print no per-window line at all. +// They stay indeterminate on purpose: there is no progress to report, and the +// whole run is over in the time the windowed path spends on one window. +export function createDiarizeProgressParser(): { + feed: (line: string) => ProgressUpdate | null; +} { + let elapsedSum = 0; + let elapsedCount = 0; + return { + feed(line: string): ProgressUpdate | null { + if (!line.includes("diarize-sherpa:")) return null; + // The plan line, before any window has run. Worth consuming so the bar is + // determinate from the start rather than jumping from a spinner to 1/9. + const planned = line.match(/->\s*(\d+)\s+window\(s\)/); + if (planned) { + const n = Number.parseInt(planned[1], 10); + if (Number.isFinite(n) && n > 0) { + return { fraction: 0, detail: `0/${n} windows` }; + } + return null; + } + const m = line.match(/window\s+(\d+)\s*\/\s*(\d+)/i); + if (!m) return null; + const i = Number.parseInt(m[1], 10); + const n = Number.parseInt(m[2], 10); + if (!Number.isFinite(n) || n <= 0 || !Number.isFinite(i)) return null; + // "skipping" is the zero-samples window: it reports an index but no time, + // and counting its (absent) duration would drag the mean down. + const secs = line.match(/in\s+([\d.]+)s\b/); + if (secs) { + const v = Number.parseFloat(secs[1]); + if (Number.isFinite(v) && v > 0) { + elapsedSum += v; + elapsedCount++; + } + } + let detail = `window ${i}/${n}`; + if (elapsedCount > 0 && i < n) { + const eta = Math.round((elapsedSum / elapsedCount) * (n - i)); + if (eta > 0) detail += ` · ETA ${formatClock(eta)}`; + } + return { fraction: clamp01(i / n), detail }; + }, + }; +} + export function createParakeetProgressParser(): { feed: (line: string) => ProgressUpdate | null; } { diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts @@ -1,6 +1,9 @@ import type { JobTaskKind } from "./registry"; import type { JobRunContext } from "./streamCommand"; -import { createDownloadProgressParser } from "./progressParsers"; +import { + createDiarizeProgressParser, + createDownloadProgressParser, +} from "./progressParsers"; import { getTranscriptionApp } from "../lib/transcriptionApps"; // A handle for one in-flight sub-operation. `onLog` is a drop-in replacement @@ -78,6 +81,13 @@ export function makeTaskTracker( transcribeParser || kind === "digest" ? null : createDownloadProgressParser(); + // Backfill tasks are the diarization lane (and the attribution lanes, + // which produce no progress output and simply never match). Chained + // AHEAD of the download parser rather than replacing it: a backfill task + // really can download, when allowRedownload re-acquires media for a + // missing-input video, and that phase's progress is worth keeping. + const diarizeParser = + kind === "backfill" ? createDiarizeProgressParser() : null; let ended = false; const onLog = (line: string) => { forwardLog(line); @@ -92,7 +102,7 @@ export function makeTaskTracker( // generated, N failed" while looking like an engine problem. const update = transcribeParser ? transcribeParser.feed(part) - : downloadParser?.feed(part); + : (diarizeParser?.feed(part) ?? downloadParser?.feed(part)); if (update) ctx.updateTask(id, update); } }; diff --git a/common/lib/backfillKinds.test.ts b/common/lib/backfillKinds.test.ts @@ -18,6 +18,8 @@ import { reachableBackfillWork, type BackfillClassification, type BackfillKind, + laneEntriesOf, + presentBackfillWork, } from "./backfillKinds"; import { candidateAction } from "../controller/backfillBatch"; import { @@ -890,6 +892,10 @@ async function classifyDigest(opts: { cues?: boolean; staleCues?: boolean; sidecar?: string; + // Which sections must ALL be fresh. Defaults to the single-section shape the + // live corpus is configured with; the two-section form is what makes a + // part-done digest possible at all. + sections?: readonly string[]; }): Promise<BackfillClassification> { const { dir, cleanup } = await fixture({ transcript: opts.transcript }); try { @@ -910,7 +916,9 @@ async function classifyDigest(opts: { videoDir: dir, videoId: "vid1", files, - target: DIGEST_TARGET, + target: opts.sections + ? { ...DIGEST_TARGET, sections: opts.sections } + : DIGEST_TARGET, settings: defaultSiteSettings(), }); } finally { @@ -1052,6 +1060,7 @@ test("counts keep reachable work and needs-re-acquiring apart", () => { assert.deepEqual(counts, { missing: 2, stale: 1, + partial: 0, missingInput: 3, deferred: 0, blocked: 0, @@ -1074,6 +1083,7 @@ test("deferred is counted, and is NOT reachable work", async () => { assert.deepEqual(counts, { missing: 1, stale: 1, + partial: 0, missingInput: 0, deferred: 2, blocked: 0, @@ -1083,3 +1093,122 @@ test("deferred is counted, and is NOT reachable work", async () => { // deferred into the total. assert.equal(reachableBackfillWork(counts), 2); }); + +// --------------------------------------------------------------------------- +// PARTIAL: some sections at the current identity, some not. + +test("digest: some sections fresh and some not is PARTIAL, not stale", async () => { + // THE CASE THIS EXISTS FOR. `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 is part-done at once. Without the split all ~77,000 would read as + // `stale` — indistinguishable on a stage card from a PROMPT_VERSION bump that + // really did invalidate everything, when in fact digestVideo would regenerate + // only the tags. + assert.equal( + await classifyDigest({ + sidecar: digestSidecar(), + sections: ["chapters", "tags"], + }), + "partial", + ); + // The same sidecar against the section list it was made for is DONE, which is + // what makes the line above about the sections and not about the sidecar. + assert.equal( + await classifyDigest({ sidecar: digestSidecar(), sections: ["chapters"] }), + "present", + ); +}); + +test("digest: partial is REACHABLE work, unlike deferred and blocked", () => { + // The distinction is about what the work COSTS, never about whether the lane + // can do it. A part-done video is dispatched exactly like a stale one, so the + // sum must be unchanged by splitting them — otherwise a corpus mid-tags- + // backfill reports less work than it has. + const split = emptyBackfillCounts(); + for (const st of ["missing", "partial", "stale"] as BackfillClassification[]) { + addBackfillState(split, st); + } + assert.equal(reachableBackfillWork(split), 3); + const blockedAndDeferred = emptyBackfillCounts(); + for (const st of ["blocked", "deferred"] as BackfillClassification[]) { + addBackfillState(blockedAndDeferred, st); + } + assert.equal(reachableBackfillWork(blockedAndDeferred), 0); +}); + +test("reachableBackfillWork survives a snapshot written before `partial`", () => { + // Every snapshot on disk predates the field. Reading it as undefined and + // adding it would produce NaN, which renders as "NaN" and sorts unpredictably + // — strictly worse than under-reporting. + const old = { + missing: 2, + stale: 1, + missingInput: 5, + deferred: 0, + blocked: 0, + } as unknown as ReturnType<typeof emptyBackfillCounts>; + assert.equal(reachableBackfillWork(old), 3); +}); + +// --------------------------------------------------------------------------- +// laneEntriesOf: the read-side twin of laneBackfillKinds. + +test("laneEntriesOf keeps a digest entry OUT of the lane's sums", () => { + // THE REGRESSION THIS WHOLE DESIGN EXISTS TO PREVENT. Four surfaces summed + // Object.values(snapshot.backfill) on the assumption that the map WAS the + // backfill lane. Now that the snapshot carries an entry per catalog operation, + // that assumption would fold ~75,000 digest videos into the dashboard's + // backfill instrument, /actionable's backfill rows and the widget. + const backfill = { + diarization: { ...emptyBackfillCounts(), missing: 3, ids: [], eligible: 3 }, + digest: { + ...emptyBackfillCounts(), + missing: 75_000, + ids: [], + eligible: 75_000, + }, + }; + const lane = laneEntriesOf(backfill); + assert.equal(lane.length, 1); + assert.equal( + lane.reduce((n, e) => n + reachableBackfillWork(e), 0), + 3, + ); +}); + +test("laneEntriesOf filters by the DECLARATION, not by a hardcoded id", () => { + // `key !== "digest"` would pass the test above and leave the identical trap + // armed for the next operation registered on a lane of its own. The rule is + // the queue key, asked of the registry — so every lane kind is in, and an id + // the catalog does not know is out. + const laneIds = laneBackfillKinds(settingsWithDiarization()).map((k) => k.id); + const backfill: Record<string, ReturnType<typeof emptyBackfillCounts> & { ids: string[] }> = {}; + for (const id of [...laneIds, "digest", "some-kind-from-a-newer-build"]) { + backfill[id] = { ...emptyBackfillCounts(), missing: 1, ids: [] }; + } + assert.equal(laneEntriesOf(backfill).length, laneIds.length); + assert.ok(laneIds.length > 0, "expected at least one lane kind enabled"); + // Empty and absent are both simply nothing, never a throw: a snapshot may + // predate the field entirely. + assert.deepEqual(laneEntriesOf(undefined), []); + assert.deepEqual(laneEntriesOf({}), []); +}); + +test("presentBackfillWork says UNKNOWN rather than zero on an old snapshot", () => { + // A 0 here would render as "nothing digested" on a fully digested channel, + // which is the most dangerous direction for a coverage number to be wrong. + assert.equal( + presentBackfillWork({ ...emptyBackfillCounts(), ids: [] }), + null, + ); + assert.equal( + presentBackfillWork({ + ...emptyBackfillCounts(), + missing: 2, + blocked: 1, + ids: [], + eligible: 10, + }), + 7, + ); +}); diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts @@ -111,6 +111,11 @@ import { isCuesJsonFresh } from "../controller/normalizeTranscript"; // // 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. @@ -139,6 +144,7 @@ import { isCuesJsonFresh } from "../controller/normalizeTranscript"; export type BackfillState = | "present" | "stale" + | "partial" | "missing" | "missing-input" | "deferred" @@ -685,15 +691,37 @@ const digest: BackfillKind = { // 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) => + // 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. + const fresh = sections.filter((section) => isSectionFresh(record, section, freshness), - ) - ? "present" - : record - ? "stale" - : "missing"; + ).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 @@ -830,6 +858,12 @@ export function operationLabel(id: string): string { return operationCatalog().find((o) => o.id === id)?.label ?? id; } +// 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_KIND_ID = "digest"; + // One entry per backfill known to the system. export const BACKFILL_KINDS: readonly BackfillKind[] = [ diarization, @@ -892,6 +926,40 @@ export function laneBackfillKinds(settings: SiteSettings): BackfillKind[] { ); } +// The READ-SIDE twin of laneBackfillKinds: given a snapshot's per-kind map, +// return only the entries belonging to the shared backfill lane. +// +// THIS EXISTS BECAUSE THE SNAPSHOT MAP STOPPED BEING THE LANE. It used to be +// written from laneBackfillKinds, 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, +// /actionable's two backfill functions 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 laneBackfillKinds 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 laneEntriesOf<T>( + backfill: Record<string, T> | undefined | null, +): T[] { + if (!backfill) return []; + return Object.entries(backfill) + .filter(([id]) => getBackfillKind(id)?.lane.queueKey === BACKFILL_QUEUE) + .map(([, entry]) => entry); +} + // 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 @@ -987,10 +1055,25 @@ export type BackfillCounts = { // 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 reachableBackfillWork — 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 emptyBackfillCounts(): BackfillCounts { - return { missing: 0, stale: 0, missingInput: 0, deferred: 0, blocked: 0 }; + 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 @@ -1002,6 +1085,7 @@ export function addBackfillState( ): 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++; @@ -1021,8 +1105,19 @@ export function addBackfillState( // 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 reachableBackfillWork(counts: BackfillCounts): number { - return counts.missing + counts.stale; + return counts.missing + counts.stale + (counts.partial ?? 0); } // What a channel snapshot stores per kind: the three counts, plus the ids of the @@ -1035,6 +1130,43 @@ export function reachableBackfillWork(counts: BackfillCounts): number { // 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 BackfillSnapshotEntry = BackfillCounts & { - // missing + stale, sorted. Never includes missing-input. + // Exactly the reachable set — missing + stale + partial — sorted. Never + // includes missing-input, deferred or blocked. Kept equal to + // reachableBackfillWork(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 + // BackfillCounts — 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 + // addBackfillState throws away, so nothing downstream can reconstruct the + // total from the counts alone. With it, present = eligible - (every work + // count), which is what presentBackfillWork 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 presentBackfillWork( + entry: BackfillSnapshotEntry, +): 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)), + ); +} diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,11 @@ # Changelog ## [Unreleased] +- **There is now one definition of "digested", and every screen reads it.** The digest layer kept its own private tally of what still needed doing, separate from the one the digest runner actually uses — and the two could disagree by an entire channel. They are now the same number. Two consequences you will see immediately. **Channels whose reports were months old start reporting digest work at all**: eleven of them predated the old tally entirely and had been quietly reading as "nothing to do" (they account for about 1,500 videos, most of them in one channel). And **videos whose transcript is being rewritten are no longer offered as work**: they were being handed to the digest runner, which looked at them and immediately put them back. On the current corpus that is 1,933 videos, and 1,683 of them are in a single channel — so nearly all of that channel's apparent digest backlog was work that could never have started. The overall count goes *down* slightly as a result, from 75,613 to 75,199, which is the counter becoming honest rather than anything being skipped. +- **A video with no transcript is now visible instead of invisible.** It used to be absent from every digest count, which meant a channel of untranscribed videos read as fully digested. The Digest card now says how many are waiting on transcription and, when that is all that is left, says so instead of "All digested". These are counted separately and never added to the work total — there is nothing you can do about them from that screen, and they clear themselves as transcription catches up. +- **Part-finished digests are counted apart from unstarted ones.** A digest is generated in sections, and a re-run only regenerates the sections that are actually out of date. So a video that already has its chapters and only needs its topic tags is a fraction of the work of one with neither — but both used to be reported with the same word. They are now separate, which matters most the moment topic tags are switched on: every already-digested video in the archive becomes part-finished at once, and without the distinction that would have looked identical to a change that invalidated the whole archive. +- **The digest coverage percentage on the dashboard was dividing by the wrong number.** It counted every video folder, including the ones with no transcript and the ones marked as impossible to transcribe — so the bar could never have reached 100% however long a sweep ran. It now divides by the videos that can actually carry a digest. +- **Long speaker-capture jobs now show real progress instead of a spinner.** Diarizing a six-hour recording takes around forty minutes and is processed in windows; the tool already reported which window it was on and nothing was reading it. The task now shows a moving bar with a time remaining estimate measured from the windows already done. Short recordings are not processed in windows and still show an indeterminate bar, which is honest — they finish in about the time a single window takes. - **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. diff --git a/editor/app/actionable/lib/loadActionable.ts b/editor/app/actionable/lib/loadActionable.ts @@ -1,11 +1,16 @@ import type { Paths } from "yt-dlp-transcript-common/lib/paths"; import { cache } from "react"; +import { + laneEntriesOf, + reachableBackfillWork, +} from "yt-dlp-transcript-common/lib/backfillKinds"; import type { ChannelBrief } from "yt-dlp-transcript-common/controller/channels"; import { getChannelBriefs, getDuplicateReport, } from "../../lib/requestCache"; import { + digestWorkOf, excludedDownloadIdSet, type ChannelSnapshot, } from "yt-dlp-transcript-common/controller/channelSnapshot"; @@ -125,12 +130,19 @@ export function actionableCleanExtraFormatsCount(row: ActionableRow): number { return row.snapshot?.buckets.multipleAudioFormats?.length ?? 0; } -// The digest layer's two work lists. `noDigest` is "has no digest at the -// current identity" — the backfill's denominator. `digestWarnings` is "the -// model produced something a human should look at", which includes the total -// failures that write no section and so are invisible to any count of files. +// The digest layer's two work lists. The first is "has no digest at the current +// identity" — missing, stale or part-done. `digestWarnings` is "the model +// produced something a human should look at", which includes the total failures +// that write no section and so are invisible to any count of files. +// +// Read through digestWorkOf rather than off `buckets.noDigest` directly: the +// registry's classification is the one the runner uses, and it additionally +// excludes videos with no transcript (blocked) and videos whose cues.json is +// stale (deferred) — work that was being offered here and then declined by the +// runner. The bucket remains the fallback for a snapshot written before the +// registry entry existed. export function actionableNoDigestCount(row: ActionableRow): number { - return row.snapshot?.buckets.noDigest?.length ?? 0; + return digestWorkOf(row.snapshot).reachable; } export function actionableDigestWarningsCount(row: ActionableRow): number { @@ -146,16 +158,26 @@ export function actionableDigestWarningsCount(row: ActionableRow): number { // the reachable count corpus-wide, so filtering on it would put every channel in // the list forever. That is not a hypothetical: it is the documented reason // /api/widget/actionable refuses to filter on `noDigest`. +// laneEntriesOf, not Object.values: the snapshot map is every catalog operation +// now, and digest is one of them. These two functions decide whether a channel +// appears in the BACKFILL section at all, so folding a ~75,000-video operation +// that runs on another queue into them would put every channel in the list +// forever — the same trap /api/widget/actionable documents for `noDigest`, hit +// from the other direction. export function actionableBackfillCount(row: ActionableRow): number { - const entries = Object.values(row.snapshot?.backfill ?? {}); - return entries.reduce((n, e) => n + e.missing + e.stale, 0); + return laneEntriesOf(row.snapshot?.backfill).reduce( + (n, e) => n + reachableBackfillWork(e), + 0, + ); } export function actionableBackfillMissingInputCount( row: ActionableRow, ): number { - const entries = Object.values(row.snapshot?.backfill ?? {}); - return entries.reduce((n, e) => n + e.missingInput, 0); + return laneEntriesOf(row.snapshot?.backfill).reduce( + (n, e) => n + e.missingInput, + 0, + ); } // Estimated bytes each cleanup would reclaim (default 0 for snapshots written diff --git a/editor/app/api/widget/actionable/route.ts b/editor/app/api/widget/actionable/route.ts @@ -13,11 +13,14 @@ export type WidgetActionableChannel = { slug: string; undownloaded: number; untranscribed: number; - // Videos with no digest at the CURRENT identity — the backfill's per-channel - // work list. Reported but NOT used to decide whether a channel "needs work": - // during the backfill that is 99.87% of the corpus, so counting it would put - // every channel in the list forever and drown the two buckets a human can - // actually act on today. + // Videos with no current digest — missing, stale or part-done — now read from + // the operation registry rather than the `noDigest` bucket. Reported but STILL + // NOT used to decide whether a channel "needs work", and the reasoning is + // unchanged by the swap: during the backfill this is ~99.87% of the corpus, so + // counting it would put every channel in the list forever and drown the two + // buckets a human can actually act on today. The registry's classification + // removes untranscribed and cues-stale videos from the number, which makes it + // smaller and more honest — nowhere near small enough to filter on. noDigest: number; }; diff --git a/editor/app/api/widget/sync/route.ts b/editor/app/api/widget/sync/route.ts @@ -1,9 +1,14 @@ import { NextResponse } from "next/server"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { digestCountOf } from "yt-dlp-transcript-common/controller/channels"; +import { digestWorkOf } from "yt-dlp-transcript-common/controller/channelSnapshot"; import { getChannelBriefs } from "../../../lib/requestCache"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; -import { laneBackfillKinds } from "yt-dlp-transcript-common/lib/backfillKinds"; +import { + laneBackfillKinds, + laneEntriesOf, + reachableBackfillWork, +} from "yt-dlp-transcript-common/lib/backfillKinds"; import { buildScheduleView } from "yt-dlp-transcript-common/jobs/syncScheduler"; import { readSchedulerState } from "yt-dlp-transcript-common/jobs/syncSchedulerState"; @@ -34,7 +39,20 @@ export type WidgetSyncPayload = { // channel page can show. digest: { digested: number; // videos carrying a non-empty ai-digest.json - videos: number; // videos in the corpus (the denominator) + videos: number; // every video dir in the corpus + // Videos that CAN be digested: transcribed and not untranscribable, summed + // from each channel's digest work list. The honest denominator — `videos` + // counts video dirs, including the ~1,700 that have no transcript and never + // will, so a coverage percentage against it can never reach 100%. + // + // NULL UNTIL EVERY CHANNEL CAN REPORT IT, deliberately. This is summed + // across channels, and a channel whose snapshot predates the digest registry + // entry contributes videos to `digested` but nothing to this — so a partial + // sum would be a denominator smaller than its own numerator, which is a + // worse lie than the one it replaces. A null says "not yet knowable" and the + // band falls back to `videos`; it becomes a number when the snapshots have + // all been regenerated. + eligible: number | null; channelsWithAny: number; // channels the layer has reached at all paused: boolean; // settings.digest.digestsPaused sweeping: boolean; // a corpus-wide sweep is armed @@ -104,6 +122,11 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> { let digested = 0; let videos = 0; let channelsWithAny = 0; + // The eligible denominator, and whether it is complete. One channel that + // cannot report voids the whole sum — see the payload type for why a partial + // one would be worse than no answer at all. + let digestEligible = 0; + let digestEligibleKnown = true; // Summed off the same snapshots, in the same pass. Kept apart all the way // through — see the payload type. let backfillReachable = 0; @@ -113,8 +136,19 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> { const n = digestCountOf(c.snapshot); digested += n; if (n > 0) channelsWithAny++; - for (const entry of Object.values(c.snapshot?.backfill ?? {})) { - backfillReachable += entry.missing + entry.stale; + // Eligibility, minus what is waiting on a transcript: a blocked video is + // eligible in principle and cannot be digested today, and including it would + // make the coverage bar sag every time the downloader finds new videos — + // which reads as digest progress going backwards. + const work = digestWorkOf(c.snapshot); + if (work.eligible == null) digestEligibleKnown = false; + else digestEligible += Math.max(0, work.eligible - work.blocked); + // The backfill LANE only. The snapshot's per-kind map now carries every + // catalog operation, digest included, and the widget's backfill strip has + // only ever meant diarization plus attribution — the digest coverage figure + // it shows beside this one is computed separately, from digestCountOf above. + for (const entry of laneEntriesOf(c.snapshot?.backfill)) { + backfillReachable += reachableBackfillWork(entry); backfillNeedsMedia += entry.missingInput; } } @@ -130,6 +164,7 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> { }, digest: { digested, + eligible: digestEligibleKnown ? digestEligible : null, videos, channelsWithAny, paused: settings.digest.digestsPaused, diff --git a/editor/app/channels/[slug]/components/stages/DigestStage.tsx b/editor/app/channels/[slug]/components/stages/DigestStage.tsx @@ -23,8 +23,26 @@ import { VideoIdList } from "../VideoIdList"; type Props = { slug: string; existingQueues: string[]; - // Videos with a transcript but no digest at the current identity. + // Videos with a transcript but no current digest — missing, stale or + // part-done. Exactly what "Digest channel" would act on. noDigestIds: string[]; + // Videos with NO TRANSCRIPT YET. Never added to the count above: there is + // nothing this button can do about them, and they come back on their own as + // the transcription lane runs. They were invisible before the digest + // operation declared its dependency — absent from every bucket, so a channel + // of untranscribed videos read as fully digested. + blocked: number; + // Videos whose cues.json is stale, so the normalize pass owes them a rewrite. + // Digesting one now would describe superseded text and then look fresh + // forever, so the runner skips them — this is the count of work that would + // have been offered and immediately declined. + deferred: number; + // Part of the reachable count above, not an addition to it: videos that have + // SOME configured section at the current identity and not the rest. Broken + // out because a re-run regenerates only the missing sections, so these cost a + // fraction of an undigested video and an operator sizing a sweep needs to + // know how much of the work is a top-up. + partial: number; // The two lanes' default queue keys. They are DIFFERENT on purpose: the local // lane is GPU-bound and the metered lane is network-bound, and the registry // hardcodes concurrency 1 per key, so a shared key would serialize them. @@ -45,6 +63,9 @@ export function DigestStage({ slug, existingQueues, noDigestIds, + blocked, + deferred, + partial, localQueueKey, remoteQueueKey, remoteEnabled, @@ -91,6 +112,42 @@ export function DigestStage({ time. Hand corrections live in a separate{" "} <code>ai-digest.overrides.json</code> and are never overwritten. </p> + {partial > 0 && ( + <p + aria-label="digest partial" + className="mt-1 text-sm text-muted-foreground" + > + {partial.toLocaleString()} of them already{" "} + {partial === 1 ? "has" : "have"} some sections at the current + settings — a re-run generates only the ones still outstanding, so{" "} + {partial === 1 ? "it costs" : "they cost"} a fraction of an + undigested video. + </p> + )} + {deferred > 0 && ( + <p + aria-label="digest deferred" + className="mt-1 text-sm text-muted-foreground" + > + {deferred.toLocaleString()} more{" "} + {deferred === 1 ? "video is" : "videos are"} waiting on a transcript + rewrite — {deferred === 1 ? "its" : "their"} <code>cues.json</code>{" "} + is older than the transcript it came from, so digesting now would + describe superseded text. Nothing to do here; the normalize pass + clears {deferred === 1 ? "it" : "them"}. + </p> + )} + {blocked > 0 && ( + <p + aria-label="digest blocked" + className="mt-1 text-sm text-muted-foreground" + > + {blocked.toLocaleString()}{" "} + {blocked === 1 ? "video has" : "videos have"} no transcript yet and{" "} + {blocked === 1 ? "is" : "are"} waiting on transcription — nothing to + do here. + </p> + )} </div> <VideoIdList diff --git a/editor/app/channels/[slug]/lib/stageStatus.ts b/editor/app/channels/[slug]/lib/stageStatus.ts @@ -1,9 +1,14 @@ import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; import { + digestWorkOf, excludedDownloadIdSet, type ChannelSnapshot, } from "yt-dlp-transcript-common/controller/channelSnapshot"; import type { JobRecord } from "yt-dlp-transcript-common/jobs/registry"; +import { + laneEntriesOf, + reachableBackfillWork, +} from "yt-dlp-transcript-common/lib/backfillKinds"; export type SnapshotBuckets = ChannelSnapshot["buckets"]; @@ -341,14 +346,20 @@ export function computeStageStatuses( }), }; - // Videos with a transcript whose digest is missing OR stale against the local - // lane's current identity — the snapshot's noDigest bucket now computes that - // full freshness target (channelSnapshot.ts), so this label and that bucket - // finally mean the same thing. Counted as pending work rather than merely - // informational: unlike the auto-captions lane, every transcribed video is - // eventually meant to have one. + // Videos whose digest is missing, stale or part-done against the local lane's + // current identity. Read from the operation registry via digestWorkOf, which + // falls back to the legacy `noDigest` bucket for a snapshot written before + // this — so a channel that has not been regenerated yet still reports its work + // rather than reading as fully digested. + // + // Counted as pending work rather than merely informational: unlike the + // auto-captions lane, every transcribed video is eventually meant to have one. + // Videos BLOCKED on transcription are deliberately not in this number — there + // is nothing the digest lane can do about them — and the stage card names them + // separately. const digestRunning = runningByStage.has("digest"); - const digestPending = buckets.noDigest.length; + const digestWork = digestWorkOf(snapshot); + const digestPending = digestWork.reachable; const digest: StageStatus = { id: "digest", title: "Digest", @@ -364,7 +375,18 @@ export function computeStageStatuses( "transcript needs a digest", "transcripts need a digest", ) - : "All digested at the current settings.", + : // "All digested" MUST NOT be said over a channel that simply has + // nothing to digest yet. Before the registry classified them, videos + // with no transcript were absent from every digest bucket, so a + // channel of untranscribed videos read as finished — the exact failure + // declaring the transcription dependency exists to end. + digestWork.blocked > 0 + ? pluralize( + digestWork.blocked, + "video is waiting on a transcript", + "videos are waiting on transcripts", + ) + : "All digested at the current settings.", tone: pickTone({ running: digestRunning, pending: digestPending, @@ -381,9 +403,15 @@ export function computeStageStatuses( // for work that cannot be done without an opt-in re-download — precisely the // trap /api/widget/actionable documents for `noDigest`. const backfillRunning = runningByStage.has("backfill"); - const backfillEntries = Object.values(snapshot.backfill ?? {}); + // laneEntriesOf, not Object.values. The snapshot's per-kind map carries every + // operation in the catalog now, including digest — which runs on its own queue + // key, has its own stage card directly above, and would otherwise add ~75,000 + // videos to this instrument on the measured corpus. The filter is by the kind's + // declared lane rather than by its id, so the next operation registered on a + // lane of its own does not re-arm the same trap. + const backfillEntries = laneEntriesOf(snapshot.backfill); const backfillPending = backfillEntries.reduce( - (n, e) => n + e.missing + e.stale, + (n, e) => n + reachableBackfillWork(e), 0, ); const backfillMissingInput = backfillEntries.reduce( diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx @@ -23,6 +23,7 @@ import { listPostFetchersFor } from "./socialActions"; import { NoReportYet } from "./components/NoReportYet"; import type { ChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot"; import { + digestWorkOf, excludedDownloadIdSet, normalizeAvailability, normalizeExcludedFromDownload, @@ -318,7 +319,13 @@ export default async function ChannelDetailPage({ const transcodeApplies = config.handling === "transcribe" && !!config.audioFormat; // Enabled lane backfills, for the stage card. A settings read, no I/O. + // + // laneBackfillKinds, still: this card is the BACKFILL lane's, and the snapshot + // now carries an entry for every catalog operation including digest, which has + // its own card and its own queue key. const backfillKinds = laneBackfillKinds(settings); + // The digest operation's work list, from the registry or the legacy bucket. + const digestWork = digestWorkOf(snapshot); const stageOrder: StageId[] = [ "configure", @@ -406,7 +413,14 @@ export default async function ChannelDetailPage({ <DigestStage slug={slug} existingQueues={existingQueues} - noDigestIds={buckets.noDigest} + // From the operation registry, falling back to the legacy `noDigest` + // bucket for a snapshot written before it (digestWorkOf). A third of the + // channels on disk are months old, so the fallback is what stops them + // reading as fully digested until the regen reaches them. + noDigestIds={digestWork.ids} + blocked={digestWork.blocked} + deferred={digestWork.deferred} + partial={digestWork.partial} localQueueKey={DIGEST_LOCAL_QUEUE} remoteQueueKey={DIGEST_REMOTE_QUEUE} remoteEnabled={settings.digest.remoteEnabled} diff --git a/editor/app/components/dashboard/PipelineBand.tsx b/editor/app/components/dashboard/PipelineBand.tsx @@ -49,11 +49,24 @@ export function PipelineBand({ const disk = jobs?.disk ?? null; const diskLow = disk?.low ?? false; const digest = sync?.digest ?? null; + // THE DENOMINATOR IS ELIGIBLE VIDEOS, NOT EVERY VIDEO DIRECTORY. `videos` + // counts every directory in the corpus, ~1,700 of which have no transcript or + // are marked untranscribable and so can never carry a digest — against that, + // this bar could not reach 100% however long the sweep ran. `eligible` comes + // from the digest operation's own work list (transcribed, not untranscribable, + // not waiting on transcription), so it is the same population the sweep + // actually walks. + // + // It falls back to `videos` while null, which is the state until every + // channel's snapshot has been regenerated — see the sync payload. That keeps + // the bar reading exactly as it did before rather than blank. + const digestDenominator = + digest === null ? 0 : (digest.eligible ?? digest.videos); // Deliberately not rounded up. At 0.13% a "1%" would be a lie of the kind // that makes an 80-day backfill look nearly begun. const digestPct = - digest && digest.videos > 0 - ? (digest.digested / digest.videos) * 100 + digest && digestDenominator > 0 + ? (digest.digested / digestDenominator) * 100 : null; const backfill = sync?.backfill ?? null; diff --git a/editor/app/components/dashboard/types.ts b/editor/app/components/dashboard/types.ts @@ -11,8 +11,12 @@ export type DashboardChannel = { hasUrl: boolean; undownloaded: number; untranscribed: number; - // Videos with no digest at the CURRENT identity. The backfill's per-channel - // denominator, and the only per-channel number that makes corpus coverage - // legible while a multi-week sweep is running. + // Videos needing digest work at the CURRENT identity — missing, stale or + // part-done — from the operation registry. The only per-channel number that + // makes corpus coverage legible while a multi-week sweep is running. + // + // Videos with no transcript are NOT in it: they are classified as waiting on + // transcription and counted separately, so this number is what the digest + // lane could actually act on today. noDigest: number; }; diff --git a/editor/e2e/backfill.spec.ts b/editor/e2e/backfill.spec.ts @@ -722,3 +722,91 @@ test("the dashboard pauses and resumes the backfill lane", async ({ page }) => { await resume.click(); await expect.poll(laneEnabled, { timeout: 30_000 }).toBe(true); }); + +// (N) THE REGRESSION THIS STEP'S DESIGN EXISTS TO PREVENT. +// +// The channel snapshot now carries a work-list entry for EVERY catalog +// operation, not just the ones on the backfill lane — that is what makes the +// digest kind's state() reachable by anything at all. Four surfaces used to sum +// `Object.values(snapshot.backfill)` on the assumption that the map WAS the +// lane: the channel Backfill card, /actionable's two backfill functions and the +// widget sync payload. On the measured corpus the digest entry is ~75,000 +// videos against the lane's 9, so getting this wrong does not shift a number +// slightly — it replaces it. +// +// The spec asserts BOTH halves, and the first half is what stops it passing +// vacuously: if the digest entry were simply never written, every "unchanged" +// assertion below would still pass while the feature was missing entirely. +test("a digest entry in the snapshot does not move the backfill instrument", async ({ + page, + request, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + + // Both videos get a real, digestable transcript (fresh cues.json against a + // backdated vtt — a stale one would classify `deferred` and prove nothing). + // vidA keeps its audio, so diarization is reachable; vidB's is removed, so + // diarization needs it re-acquired. Digest applies to both. + await writeDigestVideo({ channelSlug: SLUG, videoId: "vidA" }); + await writeDigestVideo({ channelSlug: SLUG, videoId: "vidB" }); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + + await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { + force: true, + }); + await generateReport(page, SLUG); + + // HALF ONE: the digest entry is really there, with the videos in it. + const snapshot = await readJson<{ + backfill?: Record< + string, + { missing?: number; ids?: string[]; eligible?: number; blocked?: number } + >; + }>(`test-transcripts/channels/${SLUG}/snapshot.json`); + const digestEntry = snapshot.backfill?.digest; + expect(digestEntry, "snapshot.backfill.digest must exist").toBeTruthy(); + expect(digestEntry?.missing).toBe(2); + expect(digestEntry?.ids?.slice().sort()).toEqual(["vidA", "vidB"]); + // vidC is in the fixture with audio and NO transcript, so it is waiting on + // transcription: counted, named, and never in the work list. Before the digest + // operation declared that dependency it was in no bucket at all, which is how + // a channel of untranscribed videos read as fully digested. + expect(digestEntry?.blocked).toBe(1); + // The denominator the dashboard's coverage percentage divides by: every video + // the operation has an opinion about, which includes the blocked one (it + // becomes digestable the moment it is transcribed) and excludes nothing else + // here. 3 = 2 missing + 1 blocked. + expect(digestEntry?.eligible).toBe(3); + // And the lane kind is still counted separately, on its own entry — vidC is + // NOT-APPLICABLE to diarization (untranscribed), so the two operations + // legitimately disagree about the same video. That is the whole reason these + // are per-operation entries rather than one shared number. + expect(snapshot.backfill?.diarization?.missing).toBe(1); + + // HALF TWO: every surface that sums the lane ignores it. + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + const backfillSection = page.getByLabel("backfill section"); + // 1 — vidA's diarization. NOT 3, which is what folding the digest entry in + // would produce here. + await expect(backfillSection.getByRole("heading")).toContainText( + "Backfill derived data (1)", + ); + + // The digest card reports its own work, off the same snapshot entry. + await page.getByRole("button", { name: "Digest stage summary" }).click(); + await expect(page.getByLabel("digest section").getByRole("heading")).toContainText( + "Generate digests (2)", + ); + + // The widget payload the dashboard instrument renders: reachable is the lane + // only, and stays 1. + const sync = await request.get(`${baseUrl}/api/widget/sync`); + const body = await sync.json(); + expect(body.backfill.reachable).toBe(1); + expect(body.backfill.needsMedia).toBe(1); +}); diff --git a/plans/FACTS.md b/plans/FACTS.md @@ -1929,3 +1929,68 @@ Tier 1 only (ffprobe container read), **953 media files across 68 channels in 83 **The scan deliberately looks at every MEDIA-EXTENSION file, not only recognized ones.** The file it was written for is precisely one the fixed predicates exclude — scanning only recognized outputs would have missed it. + +--- + +## Phase D step 1 — the digest counter collapse (verified 2026-08-09) + +**The digest kind's `state()` was dead code as shipped.** `channelSnapshot.ts` resolved its +per-video probes from `laneBackfillKinds()`, which filters on +`lane.queueKey === BACKFILL_QUEUE` and therefore excludes digest by construction. So Phase C +registered the operation and moved no number anywhere. It now iterates `allBackfillKinds()`, +and the lane filter moved to the READ side as `laneEntriesOf()` +(`common/lib/backfillKinds.ts`). + +**Four generic sum sites assumed `snapshot.backfill` WAS the backfill lane**, and all four +would have absorbed ~75,000 digest videos: +`editor/app/channels/[slug]/lib/stageStatus.ts:384`, +`editor/app/actionable/lib/loadActionable.ts:150` and `:158`, +`editor/app/api/widget/sync/route.ts:116`. + +**`addBackfillState` discards `present` by design** (`common/lib/backfillKinds.ts`), so the +counts cannot yield a coverage denominator. That is why `BackfillSnapshotEntry.eligible` is +stored rather than derived, and why `presentBackfillWork()` returns `null` — not 0 — for a +snapshot written before the field. + +**`digestVideo.ts:158` regenerates only the sections that are not fresh** +(`stale = sections.filter((s) => force || !isSectionFresh(...))`). This is what makes +`partial` a cost distinction rather than a label: a video needing only tags is a fraction of +one needing everything. + +**`candidateAction` in `common/controller/backfillBatch.ts` is still the ONLY `never` check in +the backfill system.** Adding `partial` to `BackfillState` produced a compile error there and +nowhere else — `addBackfillState` and (formerly) `countBackfillWork` are unguarded if-chains. + +**`countBackfillWork` had already drifted before this change.** It re-spelled the reachable +rule as `if (state === "missing" || state === "stale")` and counted `blocked` nowhere, despite +its own header promising the same derivation as the batch. It now folds through +`addBackfillState` + `reachableBackfillWork` and returns `blocked` as a fourth field. + +**`backfillSweep.ts` and `backfillBatch.ts` contained LITERAL NUL BYTES** (2 and 1 +respectively), used correctly as join/key separators. This is why `grep` treated both as +binary, and why the scope comparison at `backfillSweep.ts:382` was repeatedly misread as +`join(" ")` — it was already joining on a NUL. Both now use the `\u0000` escape; +`digestSweep.ts` always did, and neither file needs `grep -a` any more. + +**Corpus census, 2026-08-09** (read-only walk of all 68 channel dirs, 20 s): + +| measure | value | +| --- | --- | +| video dirs / `eligible` for digest | 78,963 / 78,885 (so exactly **78** untranscribable) | +| `blocked` (no transcript) | **1,631** | +| `deferred` (stale `cues.json`) | **1,933**, of which `piratesoftware` alone is **1,683** | +| new reachable | **75,199** | +| old `noDigest` (54 channels that have the bucket) | 75,613 | +| `present` (digested at the current identity) | 122 | +| snapshots with a `backfill` block | 44 of 66 | +| channels predating the `noDigest` bucket | 11 (~1,516 videos) | + +**The plan predicted the count would RISE by ~1,530; it FALLS by 414.** `deferred` was +underestimated. `digest.state()` costs **0.34 ms/video** (125-video channel) to **0.38 +ms/video** (773-video channel), i.e. 0.7x–1.2x the `readVideoFiles` the snapshot already pays. + +**`scripts/diarize-sherpa.py` already emitted per-window progress and nothing read it** +(`window {i}/{n} ... in {elapsed}s`, plus a plan line `-> N window(s)`). A `backfill` task was +given `createDownloadProgressParser()` by `common/jobs/taskHooks.ts`, which matches none of +it. The zero-samples line prints a BARE index (`window 2 decoded zero samples`) with no `/N`, +so it correctly matches nothing. diff --git a/plans/STATE.md b/plans/STATE.md @@ -3,66 +3,102 @@ The working memory for the local-AI derived-corpus work. Rewritten at the end of every session, before context is cleared. See [`README.md`](README.md) for the protocol. -**Last updated:** 2026-08-09 — **the media-derived jobs got one catalog, declared -dependencies and one disk rule.** Same branch `feat/diarization-oom-wall`, three commits, -all verified. Phase D is designed only, in -[`unified-operations-model.md`](unified-operations-model.md). - -**A: THE UNATTENDED DOWNLOADER HAD NO DISK CHECK AT ALL** (`1f82296`). Every download a -person starts by clicking has preflighted for a long time; `autoRunner` — the one path that -dispatches for days unobserved — contained zero references to disk. The disk was at 98% -(29 GB) when this started. There is now one shared `diskGate` with **hysteresis** (resume -needs `minFreeDiskGB + resumeMarginGB`, default margin 2 GB, or the first resumed download -drops back under the floor and the pipeline flaps) and a **reason string**. Three modes, and -the split is the point: `enforce` reads+writes the latch (unattended loops), `observe` reads -it (UI polls — so a dashboard cannot show green while the pipeline is held), `manual` -ignores it (the operator is standing right there). Four more byte-writing holes closed: -`downloadVideoPipelineAction` (whose sibling already called the `lowDiskError` defined in -the same file), `persistKept` per-item rather than once, the truncated-audio re-fetch checked -BEFORE it deletes the stub, and `backfillBatch` where `diskFloorHit` was set and then -ignored. Derived sidecars stay ungated on purpose — kilobytes, and holding them costs days. -The dashboard's red "downloads paused" could only ever mean the manual toggle; it is two -instruments now. **The new auto-runner spec was checked against a disabled gate before being -believed — it goes red with 2 picks dispatched.** - -**B: `blocked` IS NOT `missing-input`** (`67d2ff2`). `attribution-diarized` waits on -diarization.json and reported that as `missing-input` — which means one thing to the rest of -the system: *the media is gone, re-acquire it*. So ~73,000 videos sat in the re-acquire -population, and with `allowRedownload` on the lane would have spent a download each fetching -AUDIO, which can never satisfy that wait, then deleted it again. `blocked` follows the -`deferred` precedent exactly: counted, never summed into `reachableBackfillWork`, never -dispatched, immune to `force` — and, the actual fix, unaffected by `allowRedownload`. -`dependsOn` is declared and `resolveBackfillKinds` topologically sorts by it, **stably**, -because attribution-text must stay last (the better lane has to reach a diarized video first -or the text lane spends ~30 model calls on a record it must not write). Two existing -assertions pinned the old behaviour and now pin the fix. - -**C: DIGEST IS IN THE CATALOG, ON ITS OWN LANE** (`f9c15d5`). Nothing about digest generation -is rewritten — the entry calls `resolveDigestTarget` / `isCuesJsonFresh` / `isSectionFresh` / -`digestVideo`, including the shared-from-a-duplicate-cluster rule, because two definitions of -"digested" is the exact failure it exists to prevent. Two things are now true that were not: -the transcript dependency is **declared** (`dependsOn: ["transcription"]`, so a video with no -transcript reports `blocked` and names what it waits for — PLAN.md:146 stated this in prose -and nothing enforced it), and each operation **declares its lane**. 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` filters on the queue key so `backfillBatch` -can never dispatch digest and can never drop its guards. **`backfill.spec.ts` now fails if -the two lanes ever serialize.** `resolveTarget` became async and channel-scoped (a digest's -identity includes the hash of the channel context note), still resolved once per run. - -**DELIBERATELY NOT DONE, and it is step 1 of Phase D:** collapsing `noDigest` into -`snapshot.backfill.digest`, and merging digest into the backfill indicators. Both change live -coverage numbers on a 78,000-video corpus, and `blocked` additionally removes untranscribed -videos from the digest count — real improvements, but not ones this branch could verify. - -**Coordination note (again):** `settings.backfill.enabled` was set to false while -`backfillKinds.ts` / `backfillBatch.ts` were edited, and **restored to true afterwards**. -A backfill sweep was running at the time (`sweeping: true`). Free disk moved from 29 GB to -61 GB during the session — something outside this work freed ~32 GB. +**Last updated:** 2026-08-09 (later) — **Phase D step 1: the digest counters collapsed into the +operation registry.** Branch `feat/digest-in-registry`, off a fast-forward merge of +`feat/diarization-oom-wall` into `main` (8 commits, no conflicts). + +**THE DIGEST KIND'S `state()` WAS CALLED BY NOTHING, AND NOW IT IS THE ONLY DEFINITION.** +`channelSnapshot` iterated `laneBackfillKinds()`, which excludes digest by queue key, so +Phase C registered the operation and changed no number anywhere. It now iterates +`allBackfillKinds()` — a work list the snapshot does not carry is a work list nothing can +count, select or schedule, and that is the prerequisite for every later step. + +**THE HAZARD WAS REAL AND IS CLOSED BY DECLARATION, NOT BY AN ID CHECK.** Four surfaces summed +`Object.values(snapshot.backfill)` on the assumption that the map WAS the backfill lane +(`stageStatus.ts`, `loadActionable.ts` ×2, `api/widget/sync`). All four now use +`laneEntriesOf()`, which filters on `lane.queueKey === BACKFILL_QUEUE` asked of the registry — +so the next operation on its own lane does not re-arm the same trap. `backfill.spec.ts` has a +spec that asserts BOTH halves: that `snapshot.backfill.digest` exists with the right ids, and +that the Backfill card still reads 1 while the Digest card reads 2. Asserting only the second +half would have passed vacuously if the entry were never written. + +**MEASURED, AND IT CONTRADICTS THE PLAN'S PREDICTION.** The plan expected the corpus digest +count to RISE by ~1,530. It FALLS by 414, and the decomposition closes exactly: + +| | videos | +| --- | --- | +| old `noDigest` total (54 channels that have the bucket) | 75,613 | +| + the 11 pre-bucket channels, now counted | +1,516 | +| − newly `deferred` (stale `cues.json`) | −1,933 | +| = new reachable | **75,199** | + +`deferred` is far larger than the plan assumed and it is CONCENTRATED: **`piratesoftware` +alone has 1,683**, so 94% of that channel's 1,797-video apparent digest backlog was work +`digestVideo` would have taken and immediately put back. Everything else the plan stated held: +`blocked` = **1,631** exactly (the untranscribed count), `eligible` = 78,885 of 78,963 videos, +i.e. exactly **78** untranscribable. Coverage is **122** digested at the current identity. + +**COST GATE PASSED, SO THE OPTIMIZATION WAS NOT DONE.** `digest.state()` adds +`isCuesJsonFresh()` per video: **0.34 ms/video** on a 125-video channel, **0.38 ms/video** on +a 773-video one — 0.7x to 1.2x the `readVideoFiles` call the snapshot already pays, ~30 s +across the whole corpus spread over 67 regenerations. `isCuesJsonFresh` re-lists a directory +the caller already listed; folding that into `files.entries` is deliberately NOT done, because +at this cost it would be an optimization with no measurement behind it. + +**A FIFTH STATE: `partial`.** Digests are generated per SECTION and `digestVideo` regenerates +only the sections that are not fresh (`digestVideo.ts:158`), so "has chapters, needs tags" is +a fraction of the cost of "has neither" and the two were both reported as `stale`. It is +summed into `reachableBackfillWork` — the split is about what work COSTS, never about whether +it is reachable — and it made the `never` check in `candidateAction` fire, which is exactly +what that check is for. Not reachable on the live corpus today (`sections: ["chapters"]`), and +it becomes the dominant case the moment the pre-sweep decision to enable tags lands: every +already-digested video goes part-done at once, and without the split that reads identically to +a `PROMPT_VERSION` bump that invalidated everything. + +**TWO THINGS THE PLAN GOT WRONG, CORRECTED IN PLACE:** + +- **`backfillSweep.ts:382` was never comparing on `" "`.** It already joins on a LITERAL NUL + BYTE and was correct. The plan misread it because a raw NUL makes the file binary to `grep` + — which is the actual defect. Both that line and `backfillBatch.ts:282` now use the + `\u0000` escape (`digestSweep`'s form), so neither file needs `grep -a` any more. +- **`present` is NOT available from `BackfillCounts`** — `addBackfillState` discards it by + design, so the plan's "the registry now gives it directly" was false. Fixed by storing + `eligible` (videos the operation had an opinion about) on `BackfillSnapshotEntry`, keeping + the counts a pure work list; `presentBackfillWork()` derives present from it and returns + **null**, never 0, on a snapshot that predates the field. + +**TWO LATENT DEFECTS FOUND ON THE WAY, BOTH FIXED.** `countBackfillWork` re-spelled the +reachable rule instead of sharing it and had ALREADY drifted — it classified `blocked` and +then counted it nowhere, so a channel of prerequisite-waiting videos reported three zeroes and +read as finished. It now folds through `addBackfillState` + `reachableBackfillWork`. And the +runtime backstop list in `backfillBatch.test.ts` was missing `"blocked"` — the one test whose +job is to catch a state nobody wired up had itself missed one. + +**STEP 4 IS NOT DONE AND CANNOT BE DONE FROM HERE.** Regenerating the 66 snapshots only takes +effect through an editor running this code, and the live editor on :3001 is a PRODUCTION build +(`pnpm start`) serving the old one, currently supervising a diarization. It needs a restart on +the new build, then a regen. The `noDigest` fallback in `digestWorkOf()` is what makes the +interim safe: an un-regenerated channel reports its old bucket rather than reading as fully +digested. + +**ALSO SHIPPED, ASKED FOR MID-SESSION: a diarization progress bar.** A `backfill` task got the +DOWNLOAD parser, which matches nothing in diarization output, so a 40-minute windowed +diarization showed an indeterminate bar while the job bar sat at "0/9 videos". `diarize-sherpa.py` +was already printing `window 3/9 ... in 210.3s` and nothing read it. +`createDiarizeProgressParser` turns that into a fraction plus an ETA measured from the windows +already done, chained AHEAD of the download parser so a re-acquiring backfill keeps its +download progress. Short files take the non-windowed path, print no per-window line, and stay +indeterminate on purpose. + +**Coordination note:** NO pause was needed this time and none was taken. The plan's instruction +to pause the backfill assumed a `next dev` server that hot-reloads; the live editor is +`next start`, which serves a compiled build and cannot pick up source edits at all. The +backfill sweep ran throughout (`sweepEnabled: true`), untouched, and `settings.json` was read +but never written. --- + **Previous session:** 2026-08-08 — **the OOM wall is gone, and it was removed rather than worked around.** Branch `feat/diarization-oom-wall`, four parts, all verified.