commit 14bc961a1f4b12b27f18684296c0d7bb7f944075
parent ca4be1962ddb125caf249a0d4ec6b6d3bd4aa100
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 8 Sep 2026 00:58:43 -0400
common: every lane gets a work list in the snapshot, and the two bucket lanes get theirs
generateChannelSnapshot now writes backfill.download and backfill.transcription
— the fold of each lane's default buckets, in priority order, unsorted —
through one foldBucketLaneEntry beside foldBackfillEntry, so the invariant
ids.length === reachableOperationWork holds for them too. bucketLaneWorkIds is
that fold and is also the runner's fallback, which is what makes the entry
invisible: a regenerated snapshot and an old one hand a leaf the same array.
laneForOperation resolves an EXTERNAL operation's declared lane instead of
answering null; sync keeps its null. ChannelSnapshot.backfill is aliased
ChannelSnapshotOperations, and the disk key is unchanged.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
8 files changed, 597 insertions(+), 62 deletions(-)
diff --git a/common/controller/channelSnapshot.test.ts b/common/controller/channelSnapshot.test.ts
@@ -6,9 +6,11 @@ import {
digestWorkOf,
emptyHeldAudio,
foldBackfillEntry,
+ foldBucketLaneEntry,
} from "./channelSnapshot";
import {
emptyOperationCounts,
+ presentOperationWork,
reachableOperationWork,
type OperationClassification,
} from "../lib/operations";
@@ -83,6 +85,64 @@ test("ids are sorted, because a snapshot is compared byte-for-byte", () => {
});
// ---------------------------------------------------------------------------
+// foldBucketLaneEntry: the same entry for a lane that has no classifier.
+//
+// Slice 1.5. Extracted for exactly the reason foldBackfillEntry was: inside
+// generateChannelSnapshot it could only be reached with lmdb, an archive reader
+// and a corpus on disk, and its failure mode — a work list that disagrees with
+// its own count — looks like success.
+
+const SOURCE = {
+ buckets: {
+ partialDownloads: ["p1"],
+ downloadedNoTranscript: ["t2", "t1"],
+ failedListed: ["t1", "t3"],
+ noTranscript: ["n1", "n2"],
+ downloadedAutoSubsOnly: ["a1"],
+ },
+ undownloadedIds: ["zz9", "aa1"],
+};
+
+test("a bucket lane's entry keeps the reachable invariant too", () => {
+ for (const lane of ["download", "transcription"] as const) {
+ const entry = foldBucketLaneEntry(lane, { ...SOURCE, present: 100 });
+ assert.equal(entry.ids.length, reachableOperationWork(entry));
+ assert.equal(entry.stale, 0);
+ assert.equal(entry.partial, 0);
+ assert.equal(entry.blocked, 0);
+ assert.equal(entry.deferred, 0);
+ }
+});
+
+test("a bucket lane's ids are the default union, unsorted", () => {
+ // partialDownloads before undownloadedIds, and undownloadedIds in PLAYLIST
+ // order — sorting it would silently reorder the auto-download queue.
+ assert.deepEqual(
+ foldBucketLaneEntry("download", { ...SOURCE, present: 100 }).ids,
+ ["p1", "zz9", "aa1"],
+ );
+ // downloadedNoTranscript then failedListed, `t1` claimed once, and the opt-in
+ // auto-captions bucket nowhere in it.
+ assert.deepEqual(
+ foldBucketLaneEntry("transcription", { ...SOURCE, present: 100 }).ids,
+ ["t2", "t1", "t3"],
+ );
+});
+
+test("missingInput is transcription's alone, and present round-trips", () => {
+ // Download's input is the channel listing, which is never missing; a video
+ // with no audio is the TRANSCRIPTION lane's missing input, and the download
+ // lane's ordinary work.
+ const dl = foldBucketLaneEntry("download", { ...SOURCE, present: 40 });
+ assert.equal(dl.missingInput, 0);
+ assert.equal(presentOperationWork(dl), 40);
+
+ const tr = foldBucketLaneEntry("transcription", { ...SOURCE, present: 7 });
+ assert.equal(tr.missingInput, 2);
+ assert.equal(presentOperationWork(tr), 7);
+});
+
+// ---------------------------------------------------------------------------
// digestWorkOf: one reader, one definition.
test("digestWorkOf reads the registry entry, the split included", () => {
diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts
@@ -28,6 +28,7 @@ import { loadDigest } from "../lib/digest-server";
import {
addOperationState,
allOperations,
+ bucketLaneOperationId,
emptyOperationCounts,
presentOperationWork,
reachableOperationWork,
@@ -36,6 +37,12 @@ import {
type OperationClassification,
type OperationSnapshotEntry,
} from "../lib/operations";
+import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes";
+import {
+ bucketIdsFrom,
+ bucketLaneWorkIds,
+ type BucketSource,
+} from "../jobs/autoQueuePolicy";
import { isExcludedFromTruncatedCheck } from "../lib/excludeTruncatedCheck-server";
import { loadDownloadOutcome } from "../lib/downloadOutcome-server";
import type { Paths } from "../lib/paths";
@@ -56,6 +63,15 @@ import { readTranscriptCoverage } from "./normalizeTranscript";
import { resolveVttProvenance } from "../lib/subtitleProvenance";
import { isIncompleteTranscript } from "../lib/transcriptCoverage";
+// THE PER-OPERATION WORK LISTS, keyed by operation id.
+//
+// The name the map has wanted since it stopped being the backfill lane: it is
+// `operations`, and since slice 1.5 it holds one entry for EVERY lane —
+// `digest`, the three speaker kinds, and now `download` and `transcription`.
+// The DISK KEY is still `backfill`, and stays: renaming it would make every
+// snapshot on disk unreadable to say something a type alias says for free.
+export type ChannelSnapshotOperations = Record<string, OperationSnapshotEntry>;
+
export type AvailabilitySnapshot = {
byStatus: Record<Availability, string[]>;
unchecked: string[];
@@ -101,9 +117,18 @@ export type ChannelSnapshot = {
// 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.
//
+ // SINCE SLICE 1.5 IT CARRIES EVERY LANE'S WORK LIST, the two bucket lanes
+ // included: `backfill.download` and `backfill.transcription` are the fold of
+ // those lanes' default buckets (see bucketLaneWorkIds), so "which videos does
+ // lane X have to do" has ONE answer shape for all four lanes. Those two are
+ // written by hand rather than by state() — they have no registry entry — and
+ // the generator says how.
+ //
// Optional: snapshots written before this existed lack it, and readers default
- // to {}.
- backfill?: Record<string, OperationSnapshotEntry>;
+ // to {}. The two new entries are optional in the same way and for a stronger
+ // reason: no live snapshot is regenerated by the slice that added them, so
+ // every reader falls back to the bucket fold until a channel's next regen.
+ backfill?: ChannelSnapshotOperations;
buckets: {
noTranscript: string[];
downloadedNoTranscript: string[];
@@ -437,6 +462,54 @@ export function foldBackfillEntry(
return { ...counts, ids: ids.sort(), eligible };
}
+// THE SAME ENTRY FOR A BUCKET LANE, which has no per-video classifier to fold.
+//
+// download and transcription are EXTERNAL_OPERATIONS: registered for the
+// dependency graph, dispatched by their own runners, and carrying no
+// `Operation.state()`. So their four numbers are STATED rather than counted,
+// and only the ones that mean something are non-zero:
+//
+// ids = bucketLaneWorkIds — the lane's default buckets, in priority
+// order, deduped, UNSORTED. The one function the runner also
+// falls back to, which is what makes writing this entry
+// invisible: a regenerated snapshot's list and an old
+// snapshot's fold are the same array.
+// missing = ids.length. `stale` and `partial` stay 0 — a download is
+// fetched or it is not, and a transcript is never part done —
+// which keeps the invariant the test above pins
+// (ids.length === reachableOperationWork) true here too.
+// missingInput = transcription's videos with no audio yet, work the DOWNLOAD
+// lane has to do first. Download's input is the listing, which
+// is never missing, so it is 0.
+// eligible = present + every work count, so presentOperationWork() gives
+// back totals.downloaded / totals.transcribed exactly. Videos
+// excluded from download and untranscribable ones are in
+// NEITHER half: they are this operation's not-applicable, the
+// same exclusion foldBackfillEntry applies.
+//
+// THESE ENTRIES ARE NOT THE /channels BANDS, and must not become them.
+// buildBands keeps its own fold for these two lanes because its numbers are
+// different ones: its transcription `reachable` is downloadedNoTranscript
+// ALONE — the retry bucket is not in it — and its `blocked` is derived from
+// noTranscript, which is this entry's missingInput. Pointing the bands here
+// would move two rendered numbers. See editor/app/components/pipelines/
+// buildBands.ts.
+export function foldBucketLaneEntry(
+ lane: AutoQueueKind,
+ source: BucketSource & { present: number },
+): OperationSnapshotEntry {
+ const ids = bucketLaneWorkIds(lane, source);
+ const missingInput =
+ lane === "transcription" ? bucketIdsFrom(source, "noTranscript").length : 0;
+ return {
+ ...emptyOperationCounts(),
+ missing: ids.length,
+ missingInput,
+ ids,
+ eligible: source.present + ids.length + missingInput,
+ };
+}
+
// The digest work list. ONE DEFINITION: the operation registry's entry.
//
// The digest layer once grew its own counter (`buckets.noDigest`) beside the
@@ -1140,6 +1213,50 @@ export async function generateChannelSnapshot(
const corruptSourceSet = new Set(corruptSource);
const corruptFullSourceSet = new Set(corruptFullSource);
+ const snapshotBuckets: ChannelSnapshot["buckets"] = {
+ noTranscript: noTranscript.sort(),
+ downloadedNoTranscript: downloadedNoTranscript.sort(),
+ wrongFormatAudio: wrongFormatAudio.sort(),
+ multipleAudioFormats: multipleAudioFormats.sort(),
+ transcribedWithAudio: transcribedWithAudio.sort(),
+ untranscribable: untranscribable.sort(),
+ noMetadata: noMetadata.sort(),
+ failedListed: failedListed.filter(
+ (id) =>
+ !excludedById.has(id) &&
+ !corruptSourceSet.has(id) &&
+ !corruptFullSourceSet.has(id),
+ ),
+ missingFromArchive: missingFromArchive.sort(),
+ duplicateDirs: [],
+ partialDownloads: partialDownloads.sort(),
+ corruptSource: corruptSource.sort(),
+ corruptFullSource: corruptFullSource.sort(),
+ nonStandardVtt: nonStandardVtt.sort(),
+ skippedByFilter: skippedByFilter.sort(),
+ incompleteTranscript: incompleteTranscript.sort(),
+ shortAudio: shortAudio.sort(),
+ autoSubsOnly: autoSubsOnly.sort(),
+ downloadedAutoSubsOnly: downloadedAutoSubsOnly.sort(),
+ supersededAutoSubs: supersededAutoSubs.sort(),
+ needsCookies: needsCookies.sort(),
+ digestWarnings: digestWarnings.sort(),
+ };
+
+ // THE TWO BUCKET LANES GET A WORK LIST TOO — one entry per lane, so
+ // `snapshot.backfill[op].ids` is where EVERY lane's dispatch list lives and
+ // the runner has one way to ask. See foldBucketLaneEntry.
+ const bucketLaneEntries: Record<string, OperationSnapshotEntry> = {};
+ for (const lane of LANES) {
+ const opId = bucketLaneOperationId(lane);
+ if (!opId) continue;
+ bucketLaneEntries[opId] = foldBucketLaneEntry(lane, {
+ buckets: snapshotBuckets as Record<string, string[] | undefined>,
+ undownloadedIds,
+ present: lane === "transcription" ? transcribed : downloaded,
+ });
+ }
+
const snapshot: ChannelSnapshot = {
generatedAt: new Date().toISOString(),
totals: {
@@ -1147,37 +1264,9 @@ export async function generateChannelSnapshot(
transcribed,
downloaded,
},
- buckets: {
- noTranscript: noTranscript.sort(),
- downloadedNoTranscript: downloadedNoTranscript.sort(),
- wrongFormatAudio: wrongFormatAudio.sort(),
- multipleAudioFormats: multipleAudioFormats.sort(),
- transcribedWithAudio: transcribedWithAudio.sort(),
- untranscribable: untranscribable.sort(),
- noMetadata: noMetadata.sort(),
- failedListed: failedListed.filter(
- (id) =>
- !excludedById.has(id) &&
- !corruptSourceSet.has(id) &&
- !corruptFullSourceSet.has(id),
- ),
- missingFromArchive: missingFromArchive.sort(),
- duplicateDirs: [],
- partialDownloads: partialDownloads.sort(),
- corruptSource: corruptSource.sort(),
- corruptFullSource: corruptFullSource.sort(),
- nonStandardVtt: nonStandardVtt.sort(),
- skippedByFilter: skippedByFilter.sort(),
- incompleteTranscript: incompleteTranscript.sort(),
- shortAudio: shortAudio.sort(),
- autoSubsOnly: autoSubsOnly.sort(),
- downloadedAutoSubsOnly: downloadedAutoSubsOnly.sort(),
- supersededAutoSubs: supersededAutoSubs.sort(),
- needsCookies: needsCookies.sort(),
- digestWarnings: digestWarnings.sort(),
- },
+ buckets: snapshotBuckets,
digestEngines,
- backfill: backfillCounts,
+ backfill: { ...backfillCounts, ...bucketLaneEntries },
undownloadedIds,
excludedFromDownload,
keptCount: keptIds.size,
diff --git a/common/controller/laneForOperation.test.ts b/common/controller/laneForOperation.test.ts
@@ -29,8 +29,12 @@ const SETTINGS_FILE = path.join(ROOT, "settings.json");
process.env.SETTINGS_FILE = SETTINGS_FILE;
const { laneForOperation } = await import("./operationLane");
-const { DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE, BACKFILL_QUEUE } =
- await import("../lib/queueKeys");
+const {
+ DIGEST_LOCAL_QUEUE,
+ DIGEST_REMOTE_QUEUE,
+ BACKFILL_QUEUE,
+ TRANSCRIPTION_QUEUE,
+} = await import("../lib/queueKeys");
after(() => rm(ROOT, { recursive: true, force: true }));
@@ -77,19 +81,31 @@ test("diarization resolves through its laneFor too", () => {
assert.equal(laneForOperation("diarization")?.contendsFor, "cpu");
});
-test("a KNOWN external operation still has no lane", () => {
- // download and transcription are in operationCatalog() — they have to be, or
- // digest.dependsOn = ["transcription"] names nothing — but they are dispatched
- // by their OWN runners off their own queue keys. laneForOperation asks
- // getOperation, not the catalog, precisely so they come back null.
+test("a KNOWN external operation resolves its DECLARED lane", () => {
+ // Slice 1.5 flipped this case. download and transcription are in
+ // operationCatalog() — they have to be, or digest.dependsOn =
+ // ["transcription"] names nothing — and they are dispatched by their OWN
+ // runners off their own queue keys. That is a fact about the DISPATCHER, not
+ // about whether they have a lane: they declare one, and every four-lane
+ // surface asking "which lane does this run on" was getting null for the two
+ // lanes that carry all of the live dispatch.
//
- // A catalog lookup would hand them a lane, and a caller would reserve the
- // backfill queue for work no backfill operation can do. That is a different
- // failure from an unknown id, and only this case would catch it.
+ // The old fear was that a caller would "reserve the backfill queue for work no
+ // backfill operation can do" — which these two assertions are what rules out:
+ // neither declared lane is BACKFILL_QUEUE.
writeSettings({});
- assert.equal(laneForOperation("download"), null);
- assert.equal(laneForOperation("transcription"), null);
- // The unknown-id case, for contrast: same answer, different reason.
+ const download = laneForOperation("download");
+ assert.equal(download?.queueKey, "download:<platform>");
+ assert.equal(download?.contendsFor, "network");
+ assert.notEqual(download?.queueKey, BACKFILL_QUEUE);
+ const transcription = laneForOperation("transcription");
+ assert.equal(transcription?.queueKey, TRANSCRIPTION_QUEUE);
+ assert.equal(transcription?.contendsFor, "gpu");
+ assert.notEqual(transcription?.queueKey, BACKFILL_QUEUE);
+ // SYNC KEEPS ITS NULL. It is catalogued beside them and is a per-channel
+ // cadence, not a per-video pipeline — the same answer pauseLaneFor gives it.
+ assert.equal(laneForOperation("sync"), null);
+ // The unknown-id case: same answer, different reason.
assert.equal(laneForOperation("sortformer"), null);
});
diff --git a/common/controller/operationLane.ts b/common/controller/operationLane.ts
@@ -8,7 +8,11 @@
// dispatch's private business.
import { getSettings } from "../lib/settings";
-import { getOperation, type Lane } from "../lib/operations";
+import {
+ EXTERNAL_OPERATIONS,
+ getOperation,
+ type Lane,
+} from "../lib/operations";
// The lane an operation runs on, or null when the registry does not know it.
//
@@ -20,13 +24,28 @@ import { getOperation, type Lane } from "../lib/operations";
// places and only digest's copy was live: diarization's laneFor was invisible
// to every caller.
//
-// getOperation, NOT operationCatalog(), and that is deliberate. Today
-// getOperation("download") is undefined, so download and transcription get no
-// lane — which is correct, because they are dispatched by their own runners off
-// their own queue keys. A catalog lookup would hand them a lane and a caller
-// would reserve the backfill queue for work no backfill operation can do.
+// THE REGISTRY FIRST, THEN THE EXTERNAL DECLARATIONS — slice 1.5.
+//
+// This used to answer null for download and transcription, on the argument that
+// a lane for them would let a caller "reserve the backfill queue for work no
+// backfill operation can do". That was never what a catalog lookup would have
+// done: their DECLARED lanes are `download:<platform>` and the transcription
+// queue, neither of which is BACKFILL_QUEUE. What the null actually cost was the
+// question every four-lane surface now asks — "which lane does this operation
+// run on" — coming back unanswered for the two lanes carrying 100% of live
+// dispatch.
+//
+// EXTERNAL_OPERATIONS only, not operationCatalog(): sync is catalogued beside
+// them and is a per-CHANNEL cadence, not a per-video pipeline, so it keeps the
+// null it has always had here and pauseGates.pauseLaneFor keeps saying so too.
+//
+// The one caller that reserves a queue key from this answer is
+// controller/operationJobs.runOperationChannelJob, and it is unaffected: it is
+// only ever handed a registry operation id (the /channels group buttons, the
+// channel page, the video page), and an external id would already have found no
+// kind to run.
export function laneForOperation(operation: string): Lane | null {
const kind = getOperation(operation);
- if (!kind) return null;
- return kind.laneFor?.(getSettings()) ?? kind.lane;
+ if (kind) return kind.laneFor?.(getSettings()) ?? kind.lane;
+ return EXTERNAL_OPERATIONS.find((op) => op.id === operation)?.lane ?? null;
}
diff --git a/common/jobs/autoQueuePolicy.test.ts b/common/jobs/autoQueuePolicy.test.ts
@@ -6,9 +6,12 @@ import {
type AutoQueueRuntime,
type ChannelWork,
buildPendingByLeaf,
+ bucketIdsFrom,
+ bucketLaneWorkIds,
bucketsForKind,
defaultAutoQueue,
defaultBucketsForPolicy,
+ defaultDrawsForPolicy,
isGroup,
optInBucketsForKind,
selectableBucketsForKind,
@@ -19,6 +22,7 @@ import {
sanitizeAutoQueueOrder,
selectNextWork,
} from "./autoQueuePolicy";
+import { bucketLaneOperationId } from "../lib/operations";
// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/autoQueuePolicy.test.ts
// (or `node_modules/.bin/tsx --test common/jobs/autoQueuePolicy.test.ts` from the repo root)
@@ -1142,3 +1146,224 @@ test("policyDrawsBucket: an empty tree draws nothing", () => {
false,
);
});
+
+// ---------------------------------------------------------------------------
+// Slice 1.5: one work list per lane.
+//
+// The bucket lanes' default draw stopped being "walk these two bucket names"
+// and became "draw the lane's work list", projected from
+// snapshot.backfill[op].ids with bucketLaneWorkIds as its fallback. The whole
+// claim of the change is that NOTHING A LEAF DRAWS MOVES, so these tests are
+// equivalence tests: the same tree, over the same corpus, under both
+// projections.
+
+// A snapshot-shaped corpus with every trap the two projections could differ on:
+// a video in two default buckets (dedup), an opt-in bucket that must stay at
+// the tail, and a channel whose undownloadedIds are NOT sorted (playlist order).
+const CORPUS = [
+ {
+ slug: "alpha",
+ platform: "youtube" as const,
+ buckets: {
+ partialDownloads: ["p1"],
+ downloadedNoTranscript: ["t2", "t1"],
+ failedListed: ["t1", "t3"],
+ downloadedAutoSubsOnly: ["a1"],
+ autoSubsOnly: ["a2"],
+ },
+ undownloadedIds: ["zz9", "aa1", "mm5"],
+ },
+ {
+ slug: "beta",
+ platform: "rumble" as const,
+ buckets: {
+ partialDownloads: [],
+ downloadedNoTranscript: [],
+ failedListed: ["t9"],
+ downloadedAutoSubsOnly: [],
+ autoSubsOnly: [],
+ },
+ undownloadedIds: ["b1"],
+ },
+];
+
+// Every pick, as `<leaf>/<video>` — the units the runner would dispatch, in the
+// order it would dispatch them. `drain` above returns leaf ids alone, which
+// would let two projections that pick different VIDEOS look equal.
+function drainUnits(
+ root: AutoQueueGroup,
+ pending: Record<string, string[]>,
+): string[] {
+ const runtime = emptyAutoQueueRuntime();
+ const out: string[] = [];
+ for (let i = 0; i < 1000; i++) {
+ const pick = selectNextWork(root, pending, runtime);
+ if (!pick) break;
+ out.push(`${pick.leafId}/${pick.videoId}`);
+ pending[pick.leafId] = pending[pick.leafId].filter(
+ (id) => id !== pick.videoId,
+ );
+ }
+ return out;
+}
+
+// The OLD projection: every selectable bucket, drawn by name.
+function bucketProjection(kind: "download" | "transcription"): ChannelWork[] {
+ return CORPUS.map((c) => ({
+ slug: c.slug,
+ platform: c.platform,
+ buckets: Object.fromEntries(
+ selectableBucketsForKind(kind).map((name) => [
+ name,
+ [...bucketIdsFrom(c, name)],
+ ]),
+ ),
+ }));
+}
+
+// The NEW projection, exactly as buildChannelWork builds it: the same
+// selectable buckets (a `match.bucket` leaf still needs them) plus the lane's
+// work list under the operation's id, read from the snapshot entry and folded
+// from the buckets when there is none.
+//
+// `entryPresent: false` is the corpus this slice actually ships against — no
+// live snapshot is regenerated, so every channel takes the fallback until its
+// next regen.
+function operationProjection(
+ kind: "download" | "transcription",
+ opts: { entryPresent: boolean },
+): ChannelWork[] {
+ const opId = bucketLaneOperationId(kind)!;
+ return bucketProjection(kind).map((ch, i) => {
+ const stored = opts.entryPresent
+ ? // What generateChannelSnapshot wrote for this channel.
+ { ids: bucketLaneWorkIds(kind, CORPUS[i]) }
+ : undefined;
+ return {
+ ...ch,
+ buckets: {
+ ...ch.buckets,
+ [opId]: stored?.ids ?? bucketLaneWorkIds(kind, CORPUS[i]),
+ },
+ };
+ });
+}
+
+const EQUIV_ROOT: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [
+ // A bucket leaf FIRST: retry buckets are a real UI, and the video it claims
+ // must not be claimed a second time by the catch-all below it.
+ {
+ id: "retry",
+ match: { type: "channel", value: "alpha", bucket: "failedListed" },
+ weight: 1,
+ maxWorkers: null,
+ },
+ { id: "alpha", match: { type: "channel", value: "alpha" }, weight: 1, maxWorkers: null },
+ { id: "all", match: { type: "all" }, weight: 1, maxWorkers: null },
+ ],
+};
+
+for (const kind of ["download", "transcription"] as const) {
+ for (const replaceAutoSubs of [false, true]) {
+ test(`the ${kind} lane draws the same list under both projections (replaceAutoSubs=${replaceAutoSubs})`, () => {
+ const policy = { replaceAutoSubs };
+ const opId = bucketLaneOperationId(kind)!;
+ const before = buildPendingByLeaf(
+ EQUIV_ROOT,
+ bucketProjection(kind),
+ defaultBucketsForPolicy(kind, policy),
+ );
+ const after = buildPendingByLeaf(
+ EQUIV_ROOT,
+ operationProjection(kind, { entryPresent: true }),
+ defaultDrawsForPolicy(kind, policy, opId),
+ );
+ assert.deepEqual(after, before);
+
+ // And the picks, in order, which is the property the runner actually has:
+ // the same first N units — same leaves AND same videos.
+ const beforeUnits = drainUnits(EQUIV_ROOT, structuredClone(before));
+ assert.ok(beforeUnits.length > 0, "the fixture must produce work");
+ assert.deepEqual(
+ drainUnits(EQUIV_ROOT, structuredClone(after)),
+ beforeUnits,
+ );
+ });
+ }
+}
+
+test("a snapshot with no lane entry falls back to the identical list", () => {
+ // THE MIGRATION PATH. No live snapshot is regenerated by slice 1.5, so every
+ // channel reads its work list through this fallback until its next regen —
+ // and the fallback is the same fold the generator writes.
+ for (const kind of ["download", "transcription"] as const) {
+ const opId = bucketLaneOperationId(kind)!;
+ const policy = { replaceAutoSubs: false };
+ const withEntry = buildPendingByLeaf(
+ EQUIV_ROOT,
+ operationProjection(kind, { entryPresent: true }),
+ defaultDrawsForPolicy(kind, policy, opId),
+ );
+ const withoutEntry = buildPendingByLeaf(
+ EQUIV_ROOT,
+ operationProjection(kind, { entryPresent: false }),
+ defaultDrawsForPolicy(kind, policy, opId),
+ );
+ assert.deepEqual(withoutEntry, withEntry);
+ }
+});
+
+test("bucketLaneWorkIds is the default bucket union, in order, deduped", () => {
+ // Download's default order is partialDownloads THEN undownloadedIds — the
+ // half-fetched file before the one not started — and undownloadedIds keeps
+ // playlist order, unsorted.
+ assert.deepEqual(bucketLaneWorkIds("download", CORPUS[0]), [
+ "p1",
+ "zz9",
+ "aa1",
+ "mm5",
+ ]);
+ // Transcription's is downloadedNoTranscript then failedListed, and `t1` is in
+ // both: claimed once, under the higher-priority bucket.
+ assert.deepEqual(bucketLaneWorkIds("transcription", CORPUS[0]), [
+ "t2",
+ "t1",
+ "t3",
+ ]);
+ // The opt-in buckets are never folded in — they are a policy switch, not a
+ // corpus fact.
+ assert.ok(!bucketLaneWorkIds("transcription", CORPUS[0]).includes("a1"));
+ // An operation lane has no buckets, so it has no bucket work list.
+ assert.deepEqual(bucketLaneWorkIds("digest", CORPUS[0]), []);
+});
+
+test("defaultDrawsForPolicy keeps the opt-in buckets at the tail", () => {
+ assert.deepEqual(
+ [...defaultDrawsForPolicy("transcription", { replaceAutoSubs: true }, "transcription")],
+ ["transcription", ...optInBucketsForKind("transcription")],
+ );
+ assert.deepEqual(
+ [...defaultDrawsForPolicy("transcription", { replaceAutoSubs: false }, "transcription")],
+ ["transcription"],
+ );
+ // An operation lane has no work-list name and falls straight through to the
+ // (empty) bucket defaults — its draw is `defaultOperations`, not this.
+ assert.deepEqual(
+ [...defaultDrawsForPolicy("digest", { replaceAutoSubs: false }, null)],
+ [...defaultBucketsForPolicy("digest", { replaceAutoSubs: false })],
+ );
+});
+
+test("bucketIdsFrom knows undownloadedIds is a TOP-LEVEL field", () => {
+ assert.deepEqual([...bucketIdsFrom(CORPUS[0], "undownloadedIds")], [
+ "zz9",
+ "aa1",
+ "mm5",
+ ]);
+ assert.deepEqual([...bucketIdsFrom(CORPUS[0], "partialDownloads")], ["p1"]);
+ assert.deepEqual([...bucketIdsFrom({}, "partialDownloads")], []);
+ assert.deepEqual([...bucketIdsFrom({}, "undownloadedIds")], []);
+});
diff --git a/common/jobs/autoQueuePolicy.ts b/common/jobs/autoQueuePolicy.ts
@@ -136,6 +136,93 @@ export function defaultBucketsForPolicy(
: bucketsForKind(kind);
}
+// WHERE A BUCKET'S IDS LIVE IN A SNAPSHOT. One place, because the answer is not
+// uniform: `undownloadedIds` is a TOP-LEVEL field and every other bucket is
+// under `snapshot.buckets`, and that asymmetry was spelled out at each of the
+// two places that walked the list.
+export type BucketSource = {
+ buckets?: Record<string, string[] | undefined> | null;
+ undownloadedIds?: string[] | null;
+};
+
+export function bucketIdsFrom(
+ source: BucketSource,
+ name: string,
+): readonly string[] {
+ if (name === "undownloadedIds") return source.undownloadedIds ?? [];
+ return source.buckets?.[name] ?? [];
+}
+
+// A BUCKET LANE'S WHOLE WORK LIST, folded from its default buckets in priority
+// order with duplicates dropped — the exact list a bucket-less leaf on this lane
+// would draw before the opt-in buckets are appended.
+//
+// ONE FOLD, TWO CALLERS, and that is the point of slice 1.5. The snapshot
+// GENERATOR writes it as `backfill.download.ids` / `backfill.transcription.ids`
+// so every lane's work list is one array in the snapshot, and the runner reads
+// it back — falling back to this same fold for a snapshot written before those
+// entries existed. Because both sides are this function, the migration cannot
+// change what a leaf draws: the entry a fresh snapshot carries is byte-identical
+// to the fold the old one gets.
+//
+// UNSORTED, DELIBERATELY. `undownloadedIds` is playlist order — newest first,
+// which is the auto-download runner's queue — and sorting it here would silently
+// reorder the download lane. The dedup preserves first-seen order for the same
+// reason.
+//
+// The OPT-IN buckets are NOT here. They enter only through
+// defaultDrawsForPolicy, at the tail, exactly as defaultBucketsForPolicy has
+// always appended them — a snapshot entry that folded them in would make
+// `replaceAutoSubs` a property of the corpus instead of of the policy.
+export function bucketLaneWorkIds(
+ kind: AutoQueueKind,
+ source: BucketSource,
+): string[] {
+ const out: string[] = [];
+ const seen = new Set<string>();
+ for (const name of bucketsForKind(kind)) {
+ for (const id of bucketIdsFrom(source, name)) {
+ if (seen.has(id)) continue;
+ seen.add(id);
+ out.push(id);
+ }
+ }
+ return out;
+}
+
+// THE NAMES A BUCKET-LESS LEAF DRAWS, in priority order — what the runner
+// actually projects, as opposed to what `defaultBucketsForPolicy` describes.
+//
+// The two differ on a bucket lane and only there. Since slice 1.5 the runner
+// draws that lane's whole work list as ONE list, named for the operation the
+// lane dispatches (`laneWorkList`, from operations.bucketLaneOperationId) and
+// projected from `snapshot.backfill[op].ids` with the fold above as its
+// fallback. The result is identical to walking the default buckets one by one —
+// bucketLaneWorkIds IS that walk — so nothing a leaf draws moves; what changes
+// is that the list has a name, a home in the snapshot, and one definition.
+//
+// It stays in the BUCKET claim space (buildPendingByLeaf's `\0id`, not
+// `op\0id`), because on a bucket lane the operation and the bucket union are the
+// same work: a video claimed by an explicit `failedListed` leaf must not be
+// claimed AGAIN by a catch-all leaf drawing the operation. That is why the id is
+// projected into `ChannelWork.buckets` rather than into `.operations` — the
+// operation claim space exists to keep diarization and attribution apart on ONE
+// video, and there is no second operation here to keep apart from.
+//
+// `defaultBucketsForPolicy` is unchanged and still answers the question
+// policyDrawsBucket asks ("would the runner draw this bucket"), which is about
+// the POPULATION and not about the projection.
+export function defaultDrawsForPolicy(
+ kind: AutoQueueKind,
+ policy: Pick<AutoQueuePolicy, "replaceAutoSubs">,
+ laneWorkList: string | null,
+): readonly string[] {
+ if (!laneWorkList) return defaultBucketsForPolicy(kind, policy);
+ return policy.replaceAutoSubs
+ ? [laneWorkList, ...optInBucketsForKind(kind)]
+ : [laneWorkList];
+}
+
// Would this runner kind, under this policy, draw `bucket` for this channel?
//
// The question a lane outside the runner has to ask before it hands work over:
diff --git a/common/lib/operations.test.ts b/common/lib/operations.test.ts
@@ -26,8 +26,10 @@ import {
type Operation,
backfillLaneEntriesOf,
backfillLaneOperationEntriesOf,
+ bucketLaneOperationId,
presentOperationWork,
} from "./operations";
+import { pauseLaneFor } from "./pauseGates";
import {
readVideoFiles,
CUES_JSON_FILENAME,
@@ -1591,10 +1593,27 @@ function settingsWithBackfillLane(): SiteSettings {
return { ...d, attribution: a.attribution };
}
-test("operationsForLane: the runner lanes draw no operations yet", () => {
- // Not an opinion about the future: download and transcription are external
- // entries with no state(), so there is no snapshot.backfill[op].ids for a
- // leaf to draw. Slice 1.5 writes those entries.
+test("bucketLaneOperationId is pauseLaneFor read backwards", () => {
+ // The lane -> operation direction, off the SAME declaration
+ // (ExternalOperation.runner) that pauseLaneFor reads the other way. Neither
+ // spells "download" or "transcription" as a literal, so the two cannot drift.
+ assert.equal(bucketLaneOperationId("download"), "download");
+ assert.equal(bucketLaneOperationId("transcription"), "transcription");
+ // The operation lanes draw from the registry, not from a single external
+ // entry, so there is nothing for this to answer.
+ assert.equal(bucketLaneOperationId("digest"), null);
+ assert.equal(bucketLaneOperationId("backfill"), null);
+ // Round-trips through the other direction for every lane that has one.
+ for (const lane of ["download", "transcription"] as const) {
+ assert.equal(pauseLaneFor(bucketLaneOperationId(lane)!), lane);
+ }
+});
+
+test("operationsForLane: the bucket lanes draw no registry operations", () => {
+ // Slice 1.5 gave download and transcription a snapshot entry, but not a
+ // registry one: they still have no state() and no run(), so this function —
+ // which returns Operations — still has nothing to hand back. The lane asks
+ // bucketLaneOperationId for its work list instead.
const s = settingsWithBackfillLane();
assert.deepEqual(operationsForLane("transcription", s), []);
assert.deepEqual(operationsForLane("download", s), []);
diff --git a/common/lib/operations.ts b/common/lib/operations.ts
@@ -1662,11 +1662,13 @@ export function backfillLaneEntriesOf<T>(
// to DIGEST_REMOTE_QUEUE — both of which are the digest LANE, which is why both
// are named here rather than one.
//
-// The two runner lanes get [] today. They are not opinions about the future:
-// download and transcription are EXTERNAL_OPERATIONS with no `state`, so there
-// is no `snapshot.backfill[op].ids` for a leaf to draw. Slice 1.5 writes those
-// entries, and this function starts answering for them then with no caller
-// changing.
+// The two BUCKET lanes still get [] here, and that is not the gap it looks
+// like. Slice 1.5 does write `snapshot.backfill.download` and
+// `.transcription` — but this function returns `Operation`s, registry entries
+// with `state()` and `run()`, and download and transcription have neither. The
+// lane that dispatches them is asked for its work list by id instead, through
+// bucketLaneOperationId below; the runner's `run()` never asks this function
+// for them because it never dispatches them through operationBatch.
export function operationsForLane(
lane: AutoQueueKind,
settings: SiteSettings,
@@ -1685,6 +1687,24 @@ export function operationsForLane(
);
}
+// THE ONE OPERATION A BUCKET LANE DISPATCHES, or null for a lane that draws
+// from the operation registry instead.
+//
+// The mirror of pauseGates.pauseLaneFor, which asks an operation for its lane;
+// this asks a lane for its operation. Both read the SAME declaration —
+// `ExternalOperation.runner` — so the two directions cannot drift, and neither
+// spells "download" or "transcription" as a literal.
+//
+// It exists because slice 1.5 gave the two bucket lanes a work list in the
+// snapshot (`backfill.download`, `backfill.transcription`) keyed by the
+// operation id, so the runner needs the id to read it and the snapshot
+// generator needs it to write it. There is exactly one per lane by
+// construction: `runner` names an AutoQueueKind, and two external operations
+// claiming the same runner would be two work lists for one dispatcher.
+export function bucketLaneOperationId(lane: AutoQueueKind): string | null {
+ return EXTERNAL_OPERATIONS.find((op) => op.runner === lane)?.id ?? null;
+}
+
// Resolve a caller-supplied list of kind ids against the registry. An empty or
// absent list means "every enabled lane kind" — the sweep's scope default.
// Unknown ids are dropped rather than throwing: a settings file may name a kind