commit 503d46b4352063b0b96b53881008a58243024ad8
parent b0c16f2af9326934652dd4591cff6d3b4391f0ed
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 29 Jul 2026 11:47:50 -0400
Stop the sweep loop from spinning on a failure it cannot see
The no-progress guard measured the wrong thing. It backed off only when no
channel job STARTED — but an unreachable engine makes every channel job start
normally and then throw inside its body on the probe. From the loop that is
indistinguishable from a productive pass, so it would have re-planned and marched
through all 63 channels again, as fast as they could fail, forever: filling the
job registry with its own wreckage and evicting the logs that would explain it.
Found by re-reading the loop before letting it drive an 80-day run, not by a test.
The guard now compares remaining audio-hours across passes, which is the only
measure a failing engine cannot fake, and gives up after three consecutive
barren passes. Three rather than one: a single barren pass is legitimate (every
remaining video failed a guard this time and may not next time), but three in a
row is a broken engine or a corpus that cannot progress, and both want a human.
Also extracts `classifyDigestRole` as a pure function so the role split can be
asserted without an LMDB corpus behind it, and covers it — including the case
the whole plan turns on: an unmeasured mirror counts as WORK, not as a saving,
because the sharing pass refuses to place a digest without measured alignment
and counting every mirror as free overstates the benefit by nearly half.
common 374 tests (was 356), export playwright 150/150.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
3 files changed, 142 insertions(+), 25 deletions(-)
diff --git a/common/controller/digestPlan.test.ts b/common/controller/digestPlan.test.ts
@@ -0,0 +1,83 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import type { DigestClusterRole } from "./digestSharing";
+import {
+ audioHours,
+ classifyDigestRole,
+ DIGEST_PLAN_ROLES,
+ MEASURED_SECONDS_PER_AUDIO_HOUR,
+ roleMustGenerate,
+ sweepDays,
+} from "./digestPlan";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/digestPlan.test.ts
+
+test("a video in no cluster is unclustered work", () => {
+ assert.equal(classifyDigestRole(undefined, undefined), "unclustered");
+ // Alignment is meaningless without a cluster and must not change the answer.
+ assert.equal(classifyDigestRole(undefined, true), "unclustered");
+});
+
+test("a canonical member is work regardless of alignment", () => {
+ const role: DigestClusterRole = {
+ kind: "canonical",
+ clusterId: "c1",
+ mirrors: ["a"],
+ };
+ assert.equal(classifyDigestRole(role, true), "canonical");
+ assert.equal(classifyDigestRole(role, false), "canonical");
+});
+
+// The distinction the whole plan turns on: only an ALIGNED mirror is free. The
+// sharing pass refuses to place a digest without measured alignment, so counting
+// every mirror as free overstates the saving by however many it will refuse —
+// nearly half, on the current corpus.
+test("only an explicitly-aligned mirror is free", () => {
+ const role: DigestClusterRole = {
+ kind: "mirror",
+ clusterId: "c1",
+ canonicalSlug: "x",
+ };
+ assert.equal(classifyDigestRole(role, true), "mirror-aligned");
+ assert.equal(classifyDigestRole(role, false), "mirror-unaligned");
+});
+
+// Absent means "never measured", which sharing treats as NOT aligned. The plan
+// must under-promise here, never over-promise.
+test("an unmeasured mirror counts as work, not as a saving", () => {
+ const role: DigestClusterRole = {
+ kind: "mirror",
+ clusterId: "c1",
+ canonicalSlug: "x",
+ };
+ assert.equal(classifyDigestRole(role, undefined), "mirror-unaligned");
+ assert.equal(roleMustGenerate(classifyDigestRole(role, undefined)), true);
+});
+
+test("exactly one role is free; every other role costs GPU time", () => {
+ const free = DIGEST_PLAN_ROLES.filter((r) => !roleMustGenerate(r));
+ assert.deepEqual(free, ["mirror-aligned"]);
+});
+
+test("audio-hours and the day projection agree with the measured rate", () => {
+ assert.equal(audioHours(3600), 1);
+ // One audio-hour at 90 s/audio-hour is 90 seconds of wall clock.
+ assert.equal(sweepDays(3600, 90), 90 / 86400);
+ // The corpus figure, as a regression pin on the headline: 77,298 audio-hours
+ // at the measured rate is ~80 days, which is what the whole plan is about.
+ const days = sweepDays(77_298 * 3600, MEASURED_SECONDS_PER_AUDIO_HOUR);
+ assert.ok(days > 79 && days < 82, `expected ~80 sweep days, got ${days}`);
+});
+
+// A saving counted in videos is the mistake this module exists to prevent, so
+// pin the arithmetic that makes it visible: many short mirrors are worth far
+// less than one long VOD.
+test("audio-hours, not video count, is what a saving is measured in", () => {
+ const fourThousandShorts = 4_000 * 120; // 4,000 two-minute mirrors
+ const oneVodChannel = 8_331 * 3600; // HasanAbiVODs3
+ assert.ok(
+ audioHours(oneVodChannel) > audioHours(fourThousandShorts) * 60,
+ "one VOD channel outweighs thousands of short mirrors",
+ );
+});
diff --git a/common/controller/digestPlan.ts b/common/controller/digestPlan.ts
@@ -29,6 +29,7 @@ import { loadDigest } from "../lib/digest-server";
import {
buildDigestClusterPlan,
type DigestClusterPlan,
+ type DigestClusterRole,
} from "./digestSharing";
import { resolveDigestTarget, type DigestLaneChoice } from "./digestTarget";
@@ -141,15 +142,22 @@ function mergeInto(
}
}
-// Which mirrors the alignment gate will actually accept. Read off the report's
-// per-ref `aligned`, measured at detection time. Absent means "not measured",
-// which sharing treats as NOT aligned — so absence lands in mirror-unaligned
-// and the plan under-promises rather than over-promises.
-function mirrorIsFree(
- alignedBySlug: ReadonlyMap<string, boolean>,
- slug: string,
-): boolean {
- return alignedBySlug.get(slug) === true;
+// The whole role decision, as one pure function so it can be asserted on
+// without an LMDB corpus behind it.
+//
+// `aligned` comes from the report's per-ref measurement and is deliberately
+// tri-state at the call site: absent means NOT MEASURED, which the sharing pass
+// treats as not aligned. So anything other than an explicit `true` lands in
+// mirror-unaligned and the plan under-promises rather than over-promises.
+export function classifyDigestRole(
+ clusterRole: DigestClusterRole | undefined,
+ aligned: boolean | undefined,
+): DigestPlanRole {
+ if (clusterRole?.kind === "canonical") return "canonical";
+ if (clusterRole?.kind === "mirror") {
+ return aligned === true ? "mirror-aligned" : "mirror-unaligned";
+ }
+ return "unclustered";
}
export async function buildDigestSweepPlan(
@@ -240,15 +248,10 @@ export async function buildDigestSweepPlan(
}
}
- const role = clusterPlan.bySlug.get(stat.slug);
- const planRole: DigestPlanRole =
- role?.kind === "canonical"
- ? "canonical"
- : role?.kind === "mirror"
- ? mirrorIsFree(alignedBySlug, stat.slug)
- ? "mirror-aligned"
- : "mirror-unaligned"
- : "unclustered";
+ const planRole = classifyDigestRole(
+ clusterPlan.bySlug.get(stat.slug),
+ alignedBySlug.get(stat.slug),
+ );
addTo(entry.remaining[planRole], stat.duration);
}
diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts
@@ -56,6 +56,12 @@ export const DIGEST_SWEEP_KIND = "digest-sweep";
// every channel's sidecars.
const IDLE_POLL_MS = 60_000;
+// Give up after this many consecutive passes that move no work. Three, not one:
+// a single barren pass is legitimate (every remaining video failed a guard this
+// time round and may not next time), but three in a row is a broken engine or a
+// corpus the sweep cannot make progress on, and both want a human.
+const MAX_BARREN_PASSES = 3;
+
type SweepLive = { jobId: string; startedAt: number };
type SweepSingleton = { live: SweepLive | null };
@@ -164,6 +170,14 @@ async function runSweepLoop(
drainSignal: AbortSignal,
): Promise<void> {
let pass = 0;
+ // Consecutive passes that left the remaining work UNCHANGED. This is the real
+ // termination guard, and it has to measure the plan rather than whether jobs
+ // started: an unreachable engine makes every channel job start normally and
+ // then throw inside its body, which looks exactly like progress from out here.
+ // Without this the loop would spin through all 63 channels as fast as they can
+ // fail, forever, evicting the job registry with its own wreckage.
+ let barrenPasses = 0;
+ let lastGenerateSeconds: number | null = null;
for (;;) {
if (signal.aborted || drainSignal.aborted) return;
// The operator turned the sweep off: stop cleanly rather than being
@@ -214,20 +228,37 @@ async function runSweepLoop(
background: true,
});
if (!outcome.ok) {
- // One channel failing is not the sweep failing. Log it and carry on —
- // an unreachable engine will fail every channel and the loop's own
- // no-progress guard below is what stops the spin.
+ // One channel failing to START is not the sweep failing.
onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`);
continue;
}
didWork = true;
}
- if (!didWork) {
- // Every channel refused to start (engine down, queue wedged). Back off
- // rather than spinning through 63 failures a second.
+ // Did the pass actually MOVE anything? Compared against the previous pass's
+ // remaining work, because that is the only measure a failing engine cannot
+ // fake. Some tolerance: a pass that shifts less than a minute of audio is
+ // noise, not progress.
+ const moved =
+ lastGenerateSeconds === null ||
+ lastGenerateSeconds - plan.generateSeconds > 60;
+ barrenPasses = moved ? 0 : barrenPasses + 1;
+ lastGenerateSeconds = plan.generateSeconds;
+
+ if (barrenPasses >= MAX_BARREN_PASSES) {
+ onLog(
+ `Pass ${pass} made no progress for ${MAX_BARREN_PASSES} passes in a row ` +
+ `(${audioHours(plan.generateSeconds).toFixed(0)} audio-hours still outstanding). ` +
+ `Stopping rather than spinning — check the engine and the job logs, then restart the sweep.`,
+ );
+ return;
+ }
+
+ if (!didWork || !moved) {
+ // Nothing started, or nothing moved. Back off before re-planning: a tight
+ // retry against a down engine helps nobody and buries its own diagnosis.
onLog(
- `Pass ${pass} started nothing; waiting ${IDLE_POLL_MS / 1000}s before re-planning.`,
+ `Pass ${pass} made no progress; waiting ${IDLE_POLL_MS / 1000}s before re-planning.`,
);
await sleep(IDLE_POLL_MS, signal);
}