commit 2fd9f84421cdc3e18bf453d3f65ceda42eacc7d9
parent 3b2491cc00dc617930926b67f9fef0d5e36e04b3
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 26 Sep 2026 02:30:48 -0400
jobs: the boot pass stops waiting for the storage pass after 60 s
The queued-meta settle (release 9, B4b) waited for the storage boot pass
with no bound. Every findmnt in that pass has a 3 s timeout, but the stat of
a location's root and its statfs do not, and on a hung network mount they
never return: the settle never ran and every stale meta stayed `queued`.
`settleAfterStoragePass` waits at most STORAGE_PASS_WAIT_MS (60 s: a healthy
pass takes milliseconds, its bounded worst case is ~12 s a location) and then
settles anyway, with a `[boot] storage pass still running after 60 s …` line,
and a second line when the pass finally ends. The race is over a derived
promise, so the storage pass itself is never cancelled; a re-point it
enqueued runs to the end. A pass that throws counts as finished.
bootQueuedJobs.test.ts +5: a pass that never resolves (the settle runs after
the bound and logs it), a pass finished by hand after the timeout (its own
value intact, the late line), a pass that finishes in time (waited for,
nothing logged), a pass that throws, and the 60 s constant.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
3 files changed, 224 insertions(+), 30 deletions(-)
diff --git a/common/jobs/bootQueuedJobs.test.ts b/common/jobs/bootQueuedJobs.test.ts
@@ -6,7 +6,13 @@ import path from "node:path";
import type { Paths } from "../lib/paths";
import type { JobMeta } from "./jobMeta";
import type { JobSpec } from "./jobSpec";
-import { settleQueuedJobMetas, type RequeueFn } from "./bootQueuedJobs";
+import {
+ STORAGE_PASS_WAIT_MS,
+ settleAfterStoragePass,
+ settleQueuedJobMetas,
+ waitForStoragePass,
+ type RequeueFn,
+} from "./bootQueuedJobs";
// Run with:
// pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/bootQueuedJobs.test.ts
@@ -336,3 +342,112 @@ test("a missing .jobs dir is nothing to do", async () => {
});
assert.deepEqual(res, { requeued: [], cancelled: [] });
});
+
+// THE WAIT ON THE STORAGE PASS IS BOUNDED (release 10, L2). A probe stuck on a
+// hung mount (a `stat` that never returns) used to hold the settle forever, so
+// every stale meta stayed `queued`. Timeouts are shortened here; production
+// waits STORAGE_PASS_WAIT_MS.
+
+// A pass we finish by hand, so a test can prove the timeout did not.
+function deferred(): {
+ promise: Promise<string>;
+ resolve: (v: string) => void;
+} {
+ let resolve!: (v: string) => void;
+ const promise = new Promise<string>((r) => {
+ resolve = r;
+ });
+ return { promise, resolve };
+}
+
+test("a storage pass that never finishes: the settle runs after the timeout, and says so", async () => {
+ const f = await fixture([{ id: "A1", spec: SPEC, channelSlug: "teamrcn" }]);
+ try {
+ const rq = recordingRequeue();
+ const lines: string[] = [];
+ const never = new Promise<void>(() => {});
+ const started = Date.now();
+ const res = await settleAfterStoragePass(never, {
+ paths: f.paths,
+ requeue: rq.fn,
+ bootedAt: BOOT,
+ log: (l) => lines.push(l),
+ waitMs: 50,
+ });
+ assert.ok(Date.now() - started >= 45, "it did wait, for the bound");
+ assert.deepEqual(res.requeued, [{ id: "A1", newId: "NEW1" }]);
+ assert.equal((await f.read("A1")).status, "cancelled");
+ // The timeout line comes first, then the settle's own lines.
+ assert.match(lines[0], /^\[boot\] storage pass still running after 0 s; settling queued jobs without it/);
+ assert.match(lines.at(-1) ?? "", /re-queued 1, cancelled 0/);
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("the timeout does not cancel the storage pass: it finishes later, with its own result, and is logged", async () => {
+ const pass = deferred();
+ const lines: string[] = [];
+ let t = 1_000;
+ const outcome = await waitForStoragePass(pass.promise, {
+ timeoutMs: 20,
+ log: (l) => lines.push(l),
+ now: () => t,
+ });
+ assert.equal(outcome, "timed-out");
+ assert.equal(lines.length, 1);
+ // Another reader of the pass (the storage pass's own caller) still gets its
+ // value when it ends — the race was over a derived promise.
+ t = 1_000 + 95_000;
+ pass.resolve("probed 2");
+ assert.equal(await pass.promise, "probed 2");
+ await new Promise((r) => setImmediate(r));
+ assert.equal(lines.length, 2);
+ assert.match(
+ lines[1],
+ /^\[boot\] storage pass finished 95 s after the queued-job pass began waiting \(it stopped waiting at 0 s\)$/,
+ );
+});
+
+test("a storage pass that finishes in time is waited for, and nothing is logged about it", async () => {
+ const f = await fixture([{ id: "A1", spec: SPEC, channelSlug: "teamrcn" }]);
+ try {
+ let passDone = false;
+ const pass = new Promise<void>((r) =>
+ setTimeout(() => {
+ passDone = true;
+ r();
+ }, 30),
+ );
+ const lines: string[] = [];
+ let requeuedAfterPass: boolean | undefined;
+ await settleAfterStoragePass(pass, {
+ paths: f.paths,
+ requeue: async () => {
+ requeuedAfterPass = passDone;
+ return { ok: true, jobId: "NEW1" };
+ },
+ bootedAt: BOOT,
+ log: (l) => lines.push(l),
+ waitMs: 5_000,
+ });
+ assert.equal(requeuedAfterPass, true, "the settle waited for the pass");
+ assert.equal(lines.filter((l) => /storage pass/.test(l)).length, 0);
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("a storage pass that throws counts as finished: no timeout, the settle runs", async () => {
+ const lines: string[] = [];
+ const outcome = await waitForStoragePass(
+ Promise.reject(new Error("findmnt exploded")),
+ { timeoutMs: 5_000, log: (l) => lines.push(l) },
+ );
+ assert.equal(outcome, "done");
+ assert.deepEqual(lines, []);
+});
+
+test("the production bound is 60 s", () => {
+ assert.equal(STORAGE_PASS_WAIT_MS, 60_000);
+});
diff --git a/common/jobs/bootQueuedJobs.ts b/common/jobs/bootQueuedJobs.ts
@@ -80,7 +80,7 @@ export function specKey(spec: JobSpec): string {
});
}
-export async function settleQueuedJobMetas(opts: {
+export type SettleQueuedOpts = {
paths: Paths;
// null: cancel, never re-queue (an idle boot, or the e2e test server).
requeue: RequeueFn | null;
@@ -91,7 +91,91 @@ export async function settleQueuedJobMetas(opts: {
log?: (line: string) => void;
idleReason?: string;
maxAgeMs?: number;
-}): Promise<BootQueuedResult> {
+};
+
+// THE WAIT ON THE STORAGE BOOT PASS, BOUNDED (release 10, L2).
+//
+// The settle waits for the storage pass (release 9 review, LOW-1) so a channel
+// whose location is mid-autoRepoint is reachable when its job is re-queued. It
+// waited with no bound, and one probe in that pass is not bounded either: every
+// findmnt has FINDMNT_TIMEOUT_MS (3 s), but the `stat` of a location's root and
+// the statfs for its free space are plain syscalls, and on a hung network mount
+// they never return. The boot pass then never ran, and every `queued` meta
+// stayed `queued` on /jobs for the life of the process.
+//
+// 60 s. A healthy pass takes milliseconds (findmnt answers in under 10 ms here).
+// Its bounded worst case is two or three 3 s findmnts per location — identity,
+// fstab, or where the uuid is mounted — plus, for a location being re-pointed,
+// the preflight's second probe and a stat per channel on it: ~12 s a location.
+// 60 s covers several locations at that worst case; past it the pass is stuck on
+// a syscall that will not answer, and waiting longer buys nothing.
+//
+// On a timeout the settle runs anyway, and the cost is bounded: a job re-queued
+// for a channel whose drive the pass has not re-pointed yet meets the media
+// guard (jobs/jobKinds.ts `needsMedia`) — refused at submission, which closes
+// the old meta `cancelled` with the error, or failed when it starts, with the
+// reason in its log. Either way /jobs says so, and Retry is one click. The
+// storage pass is NOT cancelled (it has no signal, and a re-point it already
+// enqueued must finish): it runs on, and a line says when it finally ends.
+export const STORAGE_PASS_WAIT_MS = 60_000;
+
+export async function waitForStoragePass(
+ pass: Promise<unknown>,
+ opts: {
+ timeoutMs?: number;
+ log?: (line: string) => void;
+ now?: () => number;
+ } = {},
+): Promise<"done" | "timed-out"> {
+ const timeoutMs = opts.timeoutMs ?? STORAGE_PASS_WAIT_MS;
+ const log = opts.log ?? (() => {});
+ const now = opts.now ?? Date.now;
+ const startedAt = now();
+ // A pass that throws has still finished; the settle does not care how. This
+ // is a DERIVED promise — racing it never touches the pass itself.
+ const finished = pass.then(
+ () => "done" as const,
+ () => "done" as const,
+ );
+ let timer: ReturnType<typeof setTimeout> | undefined;
+ const timedOut = new Promise<"timed-out">((resolve) => {
+ timer = setTimeout(() => resolve("timed-out"), timeoutMs);
+ });
+ const outcome = await Promise.race([finished, timedOut]);
+ clearTimeout(timer);
+ if (outcome === "timed-out") {
+ const secs = (ms: number) => `${Math.round(ms / 1000)} s`;
+ log(
+ `[boot] storage pass still running after ${secs(timeoutMs)}; settling queued jobs without it ` +
+ `(a hung mount? a job re-queued for a channel it has not re-pointed yet ` +
+ `will be refused as unreachable, and /jobs will say so)`,
+ );
+ void finished.then(() =>
+ log(
+ `[boot] storage pass finished ${secs(now() - startedAt)} after the queued-job pass began waiting ` +
+ `(it stopped waiting at ${secs(timeoutMs)})`,
+ ),
+ );
+ }
+ return outcome;
+}
+
+// What instrumentation.ts runs: the settle, after the storage pass or after
+// STORAGE_PASS_WAIT_MS, whichever comes first.
+export async function settleAfterStoragePass(
+ storagePass: Promise<unknown>,
+ opts: SettleQueuedOpts & { waitMs?: number },
+): Promise<BootQueuedResult> {
+ await waitForStoragePass(storagePass, {
+ timeoutMs: opts.waitMs,
+ log: opts.log,
+ });
+ return settleQueuedJobMetas(opts);
+}
+
+export async function settleQueuedJobMetas(
+ opts: SettleQueuedOpts,
+): Promise<BootQueuedResult> {
const log = opts.log ?? (() => {});
const maxAgeMs = opts.maxAgeMs ?? REQUEUE_MAX_AGE_MS;
const result: BootQueuedResult = { requeued: [], cancelled: [] };
diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts
@@ -119,11 +119,14 @@ export async function register() {
// too, but there it ONLY cancels — an idle boot must not resume work — and so
// does the e2e test server, whose leftover metas belong to a previous run's
// fixture. It waits for the storage pass above, so a channel being
- // re-pointed is reachable when its job is re-queued. See
+ // re-pointed is reachable when its job is re-queued — for at most
+ // STORAGE_PASS_WAIT_MS (60 s, release 10): a probe stuck on a hung mount
+ // must not keep every stale meta `queued` for the life of the process. The
+ // timeout is logged, and the storage pass itself runs on untouched. See
// common/jobs/bootQueuedJobs.ts. Lazy, voided, best-effort: never blocks
// readiness.
try {
- const { settleQueuedJobMetas } = await import(
+ const { settleAfterStoragePass } = await import(
"yt-dlp-transcript-common/jobs/bootQueuedJobs"
);
const { getPaths } = await import("yt-dlp-transcript-common/lib/paths");
@@ -132,31 +135,23 @@ export async function register() {
);
const testServer = process.env.EDITOR_TEST_ROUTES === "1";
const cancelOnly = idle || testServer;
- void storagePass
- .then(() =>
- settleQueuedJobMetas({
- paths: getPaths(),
- bootedAt,
- isLive: (id) => getRegistry().get(id) !== undefined,
- log: (line) => console.log(line),
- idleReason: idle
- ? "idle boot"
- : testServer
- ? "test server"
- : undefined,
- requeue: cancelOnly
- ? null
- : async (spec) => {
- const { runJobSpec } = await import("./app/jobs/runJobSpec");
- const res = await runJobSpec(spec);
- if (!res.ok) return { ok: false, error: res.error };
- // Nobody reads this stream; release it as the ops routes do.
- void res.stream.cancel();
- return { ok: true, jobId: res.jobId };
- },
- }),
- )
- .catch(() => {});
+ void settleAfterStoragePass(storagePass, {
+ paths: getPaths(),
+ bootedAt,
+ isLive: (id) => getRegistry().get(id) !== undefined,
+ log: (line) => console.log(line),
+ idleReason: idle ? "idle boot" : testServer ? "test server" : undefined,
+ requeue: cancelOnly
+ ? null
+ : async (spec) => {
+ const { runJobSpec } = await import("./app/jobs/runJobSpec");
+ const res = await runJobSpec(spec);
+ if (!res.ok) return { ok: false, error: res.error };
+ // Nobody reads this stream; release it as the ops routes do.
+ void res.stream.cancel();
+ return { ok: true, jobId: res.jobId };
+ },
+ }).catch(() => {});
} catch {
/* a boot pass that fails to start must not block server readiness */
}