commit c395ab824ba4709c17269ed04866af48feb3dadf
parent 9331c3ec413f6ed751c771e7e38b6edd183d88ef
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 25 Sep 2026 14:05:19 -0400
common: recordDownloadBackoff writes through the shared auto-queue state (review fix)
A manual Sync / scan / video-page 429 cooldown was a disk-only
read-modify-write while every runner held the shared object, so the next
persist of any lane erased it — and nothing merged it back while the
download lane was paused, stopped or at capacity. The holder and
sharedAutoQueueState move to jobs/autoQueueState.ts on their own global
(__yttAutoQueueState__; autoRunner re-exports), liveAutoQueueState gives
writers the live object without starting a read, recordDownloadBackoff
merges the disk backoff in (later until wins) and writes through it, and
platformCooldownRemainingMs and the status builder prefer it. The e2e
reset clears the new global. +4 downloadBackoff tests (3 fail on the old
code).
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
6 files changed, 261 insertions(+), 91 deletions(-)
diff --git a/common/controller/autoRunner.test.ts b/common/controller/autoRunner.test.ts
@@ -561,6 +561,7 @@ test("an injected auto-queue state changes nothing about the answer", async () =
function freshSingletonForTest(): void {
globalThis.__yttAutoRunner__ = undefined;
+ globalThis.__yttAutoQueueState__ = undefined;
}
function stateTempPaths(): Paths {
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -63,7 +63,9 @@ import { LANES } from "../lib/autoQueueTypes";
import {
type AutoQueueKind,
type AutoQueueState,
+ liveAutoQueueState,
readAutoQueueState,
+ sharedAutoQueueState,
recordPick,
writeAutoQueueState,
} from "../jobs/autoQueueState";
@@ -243,14 +245,6 @@ type RunnerLive = {
type AutoRunnerSingleton = {
runners: Map<AutoQueueKind, RunnerLive>;
- // ONE persisted state object, shared by every lane's runner in this process.
- // See sharedAutoQueueState.
- state: AutoQueueState | null;
- stateFile: string | null;
- // The READ in flight, so every caller in the same window awaits ONE read and
- // receives ONE object. See sharedAutoQueueState.
- loading: Promise<AutoQueueState> | null;
- loadingFile: string | null;
};
declare global {
@@ -262,81 +256,17 @@ function getSingleton(): AutoRunnerSingleton {
if (!globalThis.__yttAutoRunner__) {
globalThis.__yttAutoRunner__ = {
runners: new Map(),
- state: null,
- stateFile: null,
- loading: null,
- loadingFile: null,
};
}
return globalThis.__yttAutoRunner__;
}
-// THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER, and it has to be.
-//
-// `writeAutoQueueState` serializes the WHOLE file — all four lanes — so two
-// runners each holding their own copy means every persist clobbers the other
-// lane's pick log and fairness memory with whatever that copy was read with.
-// With one slow lane that is nearly invisible (a transcription is minutes
-// apart), and it became obvious the moment two fast lanes ran together: the
-// digest lane's Recent picks panel emptied itself on every backfill dispatch,
-// and `/api/auto-queue/status` — which reads the FILE — showed zero picks for a
-// lane that was demonstrably working. `lane-runner.spec.ts`'s concurrency test
-// is what found it.
-//
-// One object mutated by both is what makes a write of it true for both. Its
-// lifetime is the auto-runner singleton's, which the e2e harness already clears
-// between specs (api/test/invalidate-cache) — so a reset corpus does not
-// inherit a previous spec's picks.
-//
-// THE CACHE IS THE PROMISE, NOT THE VALUE — and until release 8 it was the
-// value, which made the paragraph above false at exactly the moment it matters.
-// Boot (`startAutoRunnersIfEnabled`) starts all four lanes' `runLoop`s without
-// awaiting them, so all four reached the `await readAutoQueueState` here before
-// any of them had set `singleton.state`; each read its own copy and each
-// runLoop kept that private copy for its lifetime. Live, 2026-09-25: the
-// download lane persisted `fails: 13`, a later `until` and a 6 h video deferral
-// at 13:20; at 13:33:37 another lane persisted its boot-time copy of the whole
-// file (`fails: 12`, deferrals `{}`) and silently undid the pacing on disk.
-// Caching the in-flight read keyed by the state file means every caller in the
-// same window awaits one read and receives the same object. A reset (the e2e
-// harness dropping `__yttAutoRunner__`) replaces the singleton wholesale, so a
-// read still in flight resolves into the OLD singleton and cannot leak across.
-export async function sharedAutoQueueState(
- paths: Paths,
-): Promise<AutoQueueState> {
- const singleton = getSingleton();
- const file = paths.autoQueueStateFile;
- if (singleton.state && singleton.stateFile === file) {
- return singleton.state;
- }
- if (singleton.loading && singleton.loadingFile === file) {
- return singleton.loading;
- }
- const loading = readAutoQueueState(paths).then(
- (state) => {
- // A later call for a DIFFERENT file superseded this read; it owns the
- // singleton now, and this read's object goes only to its own awaiters.
- if (singleton.loading === loading) {
- singleton.state = state;
- singleton.stateFile = file;
- singleton.loading = null;
- singleton.loadingFile = null;
- }
- return state;
- },
- (err: unknown) => {
- // Never cache a failure: the next caller retries the read.
- if (singleton.loading === loading) {
- singleton.loading = null;
- singleton.loadingFile = null;
- }
- throw err;
- },
- );
- singleton.loading = loading;
- singleton.loadingFile = file;
- return loading;
-}
+// THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER — and for
+// `recordDownloadBackoff` too. The holder and `sharedAutoQueueState` live in
+// `jobs/autoQueueState.ts` (release 8) because `jobs/downloadBackoff.ts` has to
+// write through the same object and `jobs/` may not import `controller/`.
+// Re-exported here so every existing import keeps working.
+export { sharedAutoQueueState };
export type AutoRunnerStatus = {
kind: AutoQueueKind;
@@ -1044,7 +974,12 @@ export async function computeLeafPending(
// runtime.currentWeights in place. This runs on a 3-second status poll, so
// asking the live runtime would let merely HAVING the page open skew a
// round-robin group's rotation. Deep-clone first; the clone is discarded.
- const state = shared?.state ?? (await readAutoQueueState(paths));
+ // The live shared object when a runner holds one (only cloned from, below),
+ // else the file.
+ const state =
+ shared?.state ??
+ (await liveAutoQueueState(paths)) ??
+ (await readAutoQueueState(paths));
const runtime = {
currentWeights: { ...state[kind].runtime.currentWeights },
};
diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts
@@ -169,6 +169,129 @@ export async function writeAutoQueueState(
await writeJsonAtomic(paths.autoQueueStateFile, out, { mkdir: true });
}
+// ONE SHARED OBJECT PER PROCESS — the holder every lane's runner and
+// `recordDownloadBackoff` read and write through. Its own global (not the
+// auto-runner singleton, which lives in controller/ where jobs/ may not reach),
+// cleared by the e2e harness beside `__yttAutoRunner__`.
+type AutoQueueStateHolder = {
+ // The ONE persisted state object, shared by every lane's runner.
+ state: AutoQueueState | null;
+ stateFile: string | null;
+ // The READ in flight, so every caller in the same window awaits ONE read and
+ // receives ONE object.
+ loading: Promise<AutoQueueState> | null;
+ loadingFile: string | null;
+};
+
+declare global {
+ // eslint-disable-next-line no-var
+ var __yttAutoQueueState__: AutoQueueStateHolder | undefined;
+}
+
+function getHolder(): AutoQueueStateHolder {
+ if (!globalThis.__yttAutoQueueState__) {
+ globalThis.__yttAutoQueueState__ = {
+ state: null,
+ stateFile: null,
+ loading: null,
+ loadingFile: null,
+ };
+ }
+ return globalThis.__yttAutoQueueState__;
+}
+
+// THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER, and it has to be.
+//
+// `writeAutoQueueState` serializes the WHOLE file — all four lanes — so two
+// runners each holding their own copy means every persist clobbers the other
+// lane's pick log and fairness memory with whatever that copy was read with.
+// With one slow lane that is nearly invisible (a transcription is minutes
+// apart), and it became obvious the moment two fast lanes ran together: the
+// digest lane's Recent picks panel emptied itself on every backfill dispatch,
+// and `/api/auto-queue/status` — which reads the FILE — showed zero picks for a
+// lane that was demonstrably working. `lane-runner.spec.ts`'s concurrency test
+// is what found it.
+//
+// One object mutated by both is what makes a write of it true for both. Its
+// lifetime is the holder's, which the e2e harness clears between specs
+// (api/test/invalidate-cache) — so a reset corpus does not inherit a previous
+// spec's picks or backoff.
+//
+// THE CACHE IS THE PROMISE, NOT THE VALUE — and until release 8 it was the
+// value, which made the paragraph above false at exactly the moment it matters.
+// Boot (`startAutoRunnersIfEnabled`) starts all four lanes' `runLoop`s without
+// awaiting them, so all four reached the `await readAutoQueueState` here before
+// any of them had set the cached state; each read its own copy and each
+// runLoop kept that private copy for its lifetime. Live, 2026-09-25: the
+// download lane persisted `fails: 13`, a later `until` and a 6 h video deferral
+// at 13:20; at 13:33:37 another lane persisted its boot-time copy of the whole
+// file (`fails: 12`, deferrals `{}`) and silently undid the pacing on disk.
+// Caching the in-flight read keyed by the state file means every caller in the
+// same window awaits one read and receives the same object. A reset (the e2e
+// harness dropping `__yttAutoQueueState__`) replaces the holder wholesale, so a
+// read still in flight resolves into the OLD holder and cannot leak across.
+//
+// It lives in jobs/ (moved from controller/autoRunner.ts in release 8) because
+// `recordDownloadBackoff` must write THROUGH it: a manual Sync's 429 cooldown
+// written straight to disk was erased by the next persist of any lane, and never
+// merged at all while the download lane was paused, stopped or at capacity.
+export async function sharedAutoQueueState(
+ paths: Paths,
+): Promise<AutoQueueState> {
+ const holder = getHolder();
+ const file = paths.autoQueueStateFile;
+ if (holder.state && holder.stateFile === file) {
+ return holder.state;
+ }
+ if (holder.loading && holder.loadingFile === file) {
+ return holder.loading;
+ }
+ const loading = readAutoQueueState(paths).then(
+ (state) => {
+ // A later call for a DIFFERENT file superseded this read; it owns the
+ // holder now, and this read's object goes only to its own awaiters.
+ if (holder.loading === loading) {
+ holder.state = state;
+ holder.stateFile = file;
+ holder.loading = null;
+ holder.loadingFile = null;
+ }
+ return state;
+ },
+ (err: unknown) => {
+ // Never cache a failure: the next caller retries the read.
+ if (holder.loading === loading) {
+ holder.loading = null;
+ holder.loadingFile = null;
+ }
+ throw err;
+ },
+ );
+ holder.loading = loading;
+ holder.loadingFile = file;
+ return loading;
+}
+
+// The shared object if something in this process already holds (or is
+// loading) it for this file, else null — WITHOUT starting a read. For writers
+// outside the runner (`recordDownloadBackoff`): with no runner live there is no
+// in-memory copy to keep in step, and the plain disk read-modify-write is right.
+export async function liveAutoQueueState(
+ paths: Paths,
+): Promise<AutoQueueState | null> {
+ const holder = getHolder();
+ const file = paths.autoQueueStateFile;
+ if (holder.state && holder.stateFile === file) return holder.state;
+ if (holder.loading && holder.loadingFile === file) {
+ try {
+ return await holder.loading;
+ } catch {
+ return null;
+ }
+ }
+ return null;
+}
+
export function recordPick(state: AutoQueueKindState, pick: AutoQueuePick): void {
state.picks.unshift(pick);
if (state.picks.length > AUTO_QUEUE_PICK_LOG_LIMIT) {
diff --git a/common/jobs/downloadBackoff.test.ts b/common/jobs/downloadBackoff.test.ts
@@ -9,6 +9,12 @@ import {
recordDownloadBackoff,
} from "./downloadBackoff";
import { BACKOFF_BASE_MS } from "./platformBackoff";
+import {
+ emptyAutoQueueState,
+ readAutoQueueState,
+ sharedAutoQueueState,
+ writeAutoQueueState,
+} from "./autoQueueState";
// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/downloadBackoff.test.ts
//
@@ -22,9 +28,12 @@ async function withTempPaths(
const paths = {
autoQueueStateFile: path.join(dir, ".auto-queue", "state.json"),
} as Paths;
+ // No runner is live unless a test makes one (sharedAutoQueueState).
+ globalThis.__yttAutoQueueState__ = undefined;
try {
await fn(paths);
} finally {
+ globalThis.__yttAutoQueueState__ = undefined;
await rm(dir, { recursive: true, force: true });
}
}
@@ -75,3 +84,70 @@ test("remaining is 0 once the window lapses", async () => {
assert.equal(await platformCooldownRemainingMs("odysee", paths), 0);
});
});
+
+// WITH A RUNNER LIVE, THE COOLDOWN GOES THROUGH THE SHARED OBJECT (release 8).
+// Before, recordDownloadBackoff wrote the file only; the runners' in-memory
+// object never saw it, so the next persist of ANY lane wrote the object over
+// the file and the manual 429 cooldown was gone — and nothing merged it back
+// while the download lane was paused, stopped or at capacity.
+
+test("with a runner live, another lane's persist keeps a manual cooldown on disk", async () => {
+ await withTempPaths(async (paths) => {
+ // A runner holds the shared object (as every lane does from boot on).
+ const shared = await sharedAutoQueueState(paths);
+ // A clicked Sync hits a 429.
+ await recordDownloadBackoff("youtube", paths);
+ assert.ok(
+ shared.download.platformBackoff.youtube,
+ "the live object carries the cooldown",
+ );
+ // The transcription lane picks something and persists the shared object.
+ shared.transcription.runtime.currentWeights = { someLeaf: 1 };
+ await writeAutoQueueState(paths, shared);
+ const onDisk = await readAutoQueueState(paths);
+ assert.equal(onDisk.download.platformBackoff.youtube?.fails, 1);
+ assert.deepEqual(onDisk.transcription.runtime.currentWeights, {
+ someLeaf: 1,
+ });
+ assert.ok((await platformCooldownRemainingMs("youtube", paths)) > 0);
+ });
+});
+
+test("with a runner live, a cooldown already on disk is merged in, not dropped", async () => {
+ await withTempPaths(async (paths) => {
+ const shared = await sharedAutoQueueState(paths);
+ // Written to the file behind the runner's back (a later `until` than the
+ // live object knows about).
+ const behind = emptyAutoQueueState();
+ const until = Date.now() + 10 * 60_000;
+ behind.download.platformBackoff.odysee = { until, fails: 4 };
+ await writeAutoQueueState(paths, behind);
+ await recordDownloadBackoff("youtube", paths);
+ assert.deepEqual(shared.download.platformBackoff.odysee, { until, fails: 4 });
+ const onDisk = await readAutoQueueState(paths);
+ assert.deepEqual(onDisk.download.platformBackoff.odysee, { until, fails: 4 });
+ assert.ok(onDisk.download.platformBackoff.youtube);
+ });
+});
+
+test("platformCooldownRemainingMs prefers the live object over the file", async () => {
+ await withTempPaths(async (paths) => {
+ const shared = await sharedAutoQueueState(paths);
+ shared.download.platformBackoff.rumble = {
+ until: Date.now() + 60_000,
+ fails: 1,
+ };
+ // Not persisted yet — the file has no entry.
+ assert.ok((await platformCooldownRemainingMs("rumble", paths)) > 0);
+ });
+});
+
+test("with no runner live, recordDownloadBackoff stays a disk round trip and loads nothing", async () => {
+ await withTempPaths(async (paths) => {
+ await recordDownloadBackoff("youtube", paths);
+ assert.equal(globalThis.__yttAutoQueueState__?.state ?? null, null);
+ assert.equal(globalThis.__yttAutoQueueState__?.loading ?? null, null);
+ const onDisk = await readAutoQueueState(paths);
+ assert.equal(onDisk.download.platformBackoff.youtube?.fails, 1);
+ });
+});
diff --git a/common/jobs/downloadBackoff.ts b/common/jobs/downloadBackoff.ts
@@ -1,38 +1,67 @@
// Shared accessors for the per-platform download rate-limit cooldown that both
// the auto-download runner and manual sync/download paths respect.
//
-// The runner already owns this state in memory and persists it to
+// The runner owns this state in memory — ONE object shared by every lane's
+// runner (`sharedAutoQueueState`) — and persists it WHOLE to
// `.auto-queue/state.json` (the "download" kind's platformBackoff). These
-// helpers let code OUTSIDE the runner — a clicked Sync, an import — observe and
-// extend the same cooldown via a fresh read-modify-write, so a 429 hit by
-// either path pauses the other. Because auto-download units now serialize on the
-// same per-platform job queue as sync (see autoRunner.ts / registry.ts), a
-// momentarily racy write here can at most delay politeness by an iteration — it
-// can no longer cause two concurrent yt-dlp processes to hammer one source.
+// helpers let code OUTSIDE the runner — a clicked Sync, a metadata scan, the
+// video page — observe and extend the same cooldown, so a 429 hit by either path
+// pauses the other.
+//
+// WRITE THROUGH THE LIVE OBJECT WHEN THERE IS ONE (release 8). A disk-only
+// read-modify-write while a runner is live is erased by the next persist of ANY
+// lane — its in-memory object never saw the cooldown — and the download lane's
+// own merge-from-disk never runs while that lane is paused, stopped or at
+// capacity, which is exactly when an operator is riding out a rate limit. With
+// no runner live there is no in-memory copy, and the disk round trip is right.
+// Per-platform one-download-at-a-time is enforced by the job queue (see
+// autoRunner.ts / registry.ts), so a lost cooldown could only ever cost
+// politeness, never concurrency.
import { getPaths, type Paths } from "../lib/paths";
-import { readAutoQueueState, writeAutoQueueState } from "./autoQueueState";
+import {
+ type AutoQueueState,
+ liveAutoQueueState,
+ readAutoQueueState,
+ writeAutoQueueState,
+} from "./autoQueueState";
import { nextBackoff, pruneExpired } from "./platformBackoff";
// Milliseconds remaining in the platform's current cooldown window, or 0 if it
-// is not cooling down. Reads the persisted "download" backoff state fresh.
+// is not cooling down. Prefers the live shared object; else reads the file.
export async function platformCooldownRemainingMs(
platform: string,
paths: Paths = getPaths(),
): Promise<number> {
- const state = await readAutoQueueState(paths);
+ const state =
+ (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths));
const entry = state.download.platformBackoff[platform];
const now = Date.now();
return entry && entry.until > now ? entry.until - now : 0;
}
// Record a rate-limit/network failure against a platform, escalating its
-// exponential cooldown. Read-modify-write of the shared download backoff state.
+// exponential cooldown — through the live shared object when a runner holds
+// one, else as a read-modify-write of the file.
export async function recordDownloadBackoff(
platform: string,
paths: Paths = getPaths(),
): Promise<void> {
- const state = await readAutoQueueState(paths);
+ const live = await liveAutoQueueState(paths);
+ let state: AutoQueueState;
+ if (live) {
+ // Fold in whatever is on disk first (the later `until` wins, as the
+ // runner's own merge does), so a cooldown written while no runner held the
+ // object is escalated from, not forgotten.
+ const onDisk = (await readAutoQueueState(paths)).download.platformBackoff;
+ for (const [pf, e] of Object.entries(onDisk)) {
+ const cur = live.download.platformBackoff[pf];
+ if (!cur || e.until > cur.until) live.download.platformBackoff[pf] = e;
+ }
+ state = live;
+ } else {
+ state = await readAutoQueueState(paths);
+ }
const now = Date.now();
pruneExpired(state.download.platformBackoff, now);
state.download.platformBackoff[platform] = nextBackoff(
diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts
@@ -90,6 +90,12 @@ function invalidate() {
// next iteration; clearing this lets the next spec start fresh runners.
// eslint-disable-next-line @typescript-eslint/no-explicit-any
(globalThis as any).__yttAutoRunner__ = undefined;
+ // And the ONE shared auto-queue state object the runners (and
+ // recordDownloadBackoff) read and write through — it lives on its own global
+ // in jobs/autoQueueState.ts. Kept, a reset corpus would inherit the previous
+ // spec's picks, fairness weights and platform backoff.
+ // eslint-disable-next-line @typescript-eslint/no-explicit-any
+ (globalThis as any).__yttAutoQueueState__ = undefined;
// Drop the shared snapshot parse memo. It is keyed on (mtime, size), so it is
// self-invalidating in production — but resetData() rewrites the SAME fixture
// paths from the same source tree, and two specs writing an identical