import { test } from "node:test"; import assert from "node:assert/strict"; import { mkdtemp, mkdir, readFile, rm, utimes, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; import type { JobMeta } from "./jobMeta"; import type { JobSpec } from "./jobSpec"; import { INTERRUPTED_REASON, PUBLISH_RESTART_REASON, STORAGE_PASS_WAIT_MS, processIsAlive, settleAfterStoragePass, settleQueuedJobMetas, settleRunningJobMetas, waitForStoragePass, writerIsGone, type RequeueFn, } from "./bootQueuedJobs"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/bootQueuedJobs.test.ts // // THE BOOT PASS (release 9, B4b), over a temp .jobs dir with an injected // re-queue. It re-queues LITTLE: never a sync, never a meta older than 24 h, // only the newest of duplicate specs; everything else is closed `cancelled` // with the reason. An idle boot only ever cancels. Never run against a real // corpus. const BOOT = 2_000_000_000_000; async function fixture(metas: Partial[]): Promise<{ paths: Paths; root: string; read: (id: string) => Promise; log: (id: string) => Promise; }> { const root = await mkdtemp(path.join(tmpdir(), "boot-queued-")); const jobsDir = path.join(root, ".jobs"); await mkdir(jobsDir, { recursive: true }); for (const m of metas) { const meta = { kind: "whisper-all", queueKey: "platform:youtube", status: "queued", queuedAt: BOOT - 60_000, ...m, }; await writeFile( path.join(jobsDir, `${meta.id}.meta.json`), JSON.stringify(meta), ); } return { paths: { jobsDir } as Paths, root, read: async (id) => JSON.parse( await readFile(path.join(jobsDir, `${id}.meta.json`), "utf8"), ) as JobMeta, log: (id) => readFile(path.join(jobsDir, `${id}.log`), "utf8").catch(() => ""), }; } const SPEC: JobSpec = { kind: "whisper-all", slug: "teamrcn", params: {} }; const SYNC: JobSpec = { kind: "sync", slug: "teamrcn", params: { fullSweep: false }, }; const HOUR = 60 * 60 * 1000; function recordingRequeue(): { fn: RequeueFn; calls: JobSpec[] } { const calls: JobSpec[] = []; let n = 0; return { calls, fn: async (spec) => { calls.push(spec); return { ok: true, jobId: `NEW${++n}` }; }, }; } test("a fresh non-sync meta with a spec is re-queued, and the old one is closed naming the new id", async () => { const f = await fixture([{ id: "A1", spec: SPEC, channelSlug: "teamrcn" }]); try { const rq = recordingRequeue(); const lines: string[] = []; const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, log: (l) => lines.push(l), }); assert.deepEqual(rq.calls, [SPEC]); assert.deepEqual(res.requeued, [{ id: "A1", newId: "NEW1" }]); assert.deepEqual(res.cancelled, []); const meta = await f.read("A1"); assert.equal(meta.status, "cancelled"); assert.match( meta.cancelReason ?? "", /server restarted before it ran; re-queued as NEW1/, ); assert.ok(typeof meta.endedAt === "number"); assert.match(await f.log("A1"), /re-queued as NEW1/); // One line per re-queued job, then the summary. assert.equal(lines.length, 2); assert.match(lines[0], /re-queued whisper-all teamrcn: A1 -> NEW1/); assert.match(lines[1], /re-queued 1, cancelled 0/); } finally { await rm(f.root, { recursive: true, force: true }); } }); // The review's sizing of the live corpus: 346 queued syncs over 51 channels, // ~7 per channel, plus months-old whisper-all sweeps. Exactly one job may come // back from this. test("7 duplicate syncs + 1 stale whisper-all + 1 fresh non-sync: exactly one re-queue", async () => { const syncs = Array.from({ length: 7 }, (_, i) => ({ id: `S${i}`, kind: "sync", spec: SYNC, channelSlug: "teamrcn", queuedAt: BOOT - (i + 1) * HOUR, })); const f = await fixture([ ...syncs, { id: "W1", spec: { kind: "whisper-all", slug: "old", params: {} }, queuedAt: BOOT - 60 * 24 * HOUR, }, { id: "F1", spec: SPEC, queuedAt: BOOT - 2 * HOUR }, ]); try { const rq = recordingRequeue(); const lines: string[] = []; const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, log: (l) => lines.push(l), }); assert.deepEqual(rq.calls, [SPEC]); assert.deepEqual(res.requeued, [{ id: "F1", newId: "NEW1" }]); const byCat = (c: string) => res.cancelled.filter((x) => x.category === c).length; assert.equal(byCat("sync"), 7); assert.equal(byCat("stale"), 1); assert.match( (await f.read("S0")).cancelReason ?? "", /the scheduler re-derives syncs/, ); assert.match( (await f.read("W1")).cancelReason ?? "", /queued before the last restart, stale/, ); assert.match( lines.at(-1) ?? "", /re-queued 1, cancelled 8 \(sync 7, stale 1\)/, ); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("of duplicate specs only the newest is re-queued; params are compared key-order-free", async () => { const f = await fixture([ { id: "D1", spec: { kind: "whisper-all", slug: "c", params: { a: 1, b: 2 } }, queuedAt: BOOT - 3 * HOUR, }, { id: "D2", spec: { kind: "whisper-all", slug: "c", params: { b: 2, a: 1 } }, queuedAt: BOOT - 1 * HOUR, }, { id: "D3", spec: { kind: "whisper-all", slug: "c", params: { a: 9 } }, queuedAt: BOOT - 2 * HOUR, }, ]); try { const rq = recordingRequeue(); const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, }); assert.deepEqual(res.requeued.map((r) => r.id).sort(), ["D2", "D3"]); assert.deepEqual( res.cancelled.map((c) => [c.id, c.category]), [["D1", "superseded"]], ); assert.match( (await f.read("D1")).cancelReason ?? "", /superseded by a newer queued job \(D2\)/, ); // Survivors are re-queued oldest first, keeping their relative order. assert.deepEqual( res.requeued.map((r) => r.id), ["D3", "D2"], ); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("a queued meta with no spec is cancelled with the reason, not re-queued", async () => { const f = await fixture([{ id: "B1", kind: "build-export" }]); try { const rq = recordingRequeue(); const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, }); assert.equal(rq.calls.length, 0); assert.deepEqual( res.cancelled.map((c) => c.category), ["no-spec"], ); const meta = await f.read("B1"); assert.equal(meta.status, "cancelled"); assert.match(meta.cancelReason ?? "", /^server restarted before it ran/); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("an idle boot cancels every queued meta and re-queues none", async () => { const f = await fixture([{ id: "C1", spec: SPEC }, { id: "C2" }]); try { const res = await settleQueuedJobMetas({ paths: f.paths, requeue: null, bootedAt: BOOT, idleReason: "idle boot", }); assert.deepEqual(res.requeued, []); assert.deepEqual(res.cancelled.map((c) => c.id).sort(), ["C1", "C2"]); assert.match( (await f.read("C1")).cancelReason ?? "", /idle boot: not re-queued/, ); assert.equal((await f.read("C2")).status, "cancelled"); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("terminal, running, this-boot and live metas are left alone", async () => { const f = await fixture([ { id: "D1", status: "done", spec: SPEC }, { id: "D2", status: "running", spec: SPEC }, { id: "D3", spec: SPEC, queuedAt: BOOT + 5 }, { id: "D4", spec: SPEC }, ]); try { const rq = recordingRequeue(); const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, isLive: (id) => id === "D4", }); assert.equal(rq.calls.length, 0); assert.deepEqual(res, { requeued: [], cancelled: [] }); assert.equal((await f.read("D1")).status, "done"); assert.equal((await f.read("D2")).status, "running"); assert.equal((await f.read("D3")).status, "queued"); assert.equal((await f.read("D4")).status, "queued"); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("a malformed meta is skipped and the pass continues", async () => { const f = await fixture([{ id: "M2", spec: SPEC }]); try { await writeFile( path.join(f.paths.jobsDir, "M1.meta.json"), '{"id":"M1","kind":', ); await writeFile( path.join(f.paths.jobsDir, "M0.meta.json"), '{"status":"queued"}', ); const rq = recordingRequeue(); const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, }); assert.deepEqual(res.requeued, [{ id: "M2", newId: "NEW1" }]); assert.deepEqual(res.cancelled, []); assert.equal( await readFile(path.join(f.paths.jobsDir, "M1.meta.json"), "utf8"), '{"id":"M1","kind":', ); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("a refused or throwing re-queue closes the meta with the error", async () => { const f = await fixture([ { id: "E1", spec: { kind: "whisper-all", slug: "a", params: {} } }, { id: "E2", spec: { kind: "whisper-all", slug: "b", params: {} } }, ]); try { let calls = 0; const res = await settleQueuedJobMetas({ paths: f.paths, requeue: async () => { calls++; if (calls === 1) { return { ok: false, error: 'Channel "teamrcn" not found' }; } throw new Error("boom"); }, bootedAt: BOOT, }); assert.equal(res.cancelled.length, 2); assert.ok(res.cancelled.every((c) => c.category === "refused")); const reasons = [ (await f.read("E1")).cancelReason ?? "", (await f.read("E2")).cancelReason ?? "", ].join("\n"); assert.match(reasons, /re-queue refused: Channel "teamrcn" not found/); assert.match(reasons, /re-queue refused: boom/); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("a missing .jobs dir is nothing to do", async () => { const res = await settleQueuedJobMetas({ paths: { jobsDir: path.join(tmpdir(), "no-such-jobs-dir-r9") } as Paths, requeue: null, bootedAt: BOOT, }); 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; resolve: (v: string) => void; } { let resolve!: (v: string) => void; const promise = new Promise((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(() => {}); 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((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); }); // THE RUNNING PASS (release 17 slice D0). A job running when its process died // keeps `running` on disk forever; the next boot closes it `cancelled` with // INTERRUPTED_REASON — unless the process that wrote it is still alive (an // `archilyzer run` beside the editor), in which case it is left alone. const ALIVE_OTHER = 4242; // a live process that is not this one const DEAD = 4343; const SELF = 4444; const check = { selfPid: SELF, isProcessAlive: (pid: number) => pid === ALIVE_OTHER || pid === SELF, }; test("writerIsGone: no pid, this pid, a dead pid are gone; another live pid is not", () => { const m = (pid?: number) => ({ id: "X", kind: "k", queueKey: "", status: "running", queuedAt: 0, ...(pid === undefined ? {} : { pid }) }) as JobMeta; assert.equal(writerIsGone(m(), check), true, "a pre-release-17 meta"); assert.equal(writerIsGone(m(SELF), check), true, "a previous process with this pid"); assert.equal(writerIsGone(m(DEAD), check), true); assert.equal(writerIsGone(m(ALIVE_OTHER), check), false); }); test("processIsAlive: this process is; a pid past pid_max is not", () => { assert.equal(processIsAlive(process.pid), true); assert.equal(processIsAlive(2 ** 22 + 1), false); }); test("a running meta from a dead process is closed as interrupted, ended at its log's last write", async () => { const f = await fixture([ { id: "R1", kind: "refresh-report", status: "running", channelSlug: "the-quartering-rumble", pid: DEAD, startedAt: BOOT - 3 * HOUR }, { id: "R2", kind: "refresh-report", status: "running", channelSlug: "the-quartering" }, // no pid: pre-release-17 { id: "R3", kind: "whisper-all", status: "running", channelSlug: "x", pid: ALIVE_OTHER }, // an `archilyzer run` { id: "R4", kind: "whisper-all", status: "running", queuedAt: BOOT + 1, pid: SELF }, // this process's own { id: "R5", kind: "whisper-all", status: "running", pid: SELF }, // live in the registry { id: "D1", kind: "whisper-all", status: "done", pid: DEAD }, ]); try { const logMtime = BOOT - 2 * HOUR; const logFile = path.join(f.paths.jobsDir, "R1.log"); await writeFile(logFile, "Regenerating report for the-quartering-rumble…\n"); await utimes(logFile, logMtime / 1000, logMtime / 1000); const lines: string[] = []; const res = await settleRunningJobMetas({ paths: f.paths, bootedAt: BOOT, isLive: (id) => id === "R5", log: (l) => lines.push(l), ...check, }); assert.deepEqual( res.interrupted.map((j) => j.id), ["R1", "R2"], ); const r1 = await f.read("R1"); assert.equal(r1.status, "cancelled"); assert.equal(r1.cancelReason, INTERRUPTED_REASON); assert.equal(r1.endedAt, logMtime); assert.equal(r1.pid, DEAD, "the rest of the meta is kept"); assert.match(await f.log("R1"), /\[boot\] interrupted/); const r2 = await f.read("R2"); assert.equal(r2.status, "cancelled"); assert.ok(typeof r2.endedAt === "number", "no log: ended at this boot"); for (const id of ["R3", "R4", "R5"]) assert.equal((await f.read(id)).status, "running", id); assert.equal((await f.read("D1")).status, "done"); assert.equal(lines.length, 1); assert.match(lines[0], /closed as interrupted: 2 \(refresh-report 2\)/); // Idempotent: the next boot finds nothing. const again = await settleRunningJobMetas({ paths: f.paths, bootedAt: BOOT, isLive: (id) => id === "R5", ...check, }); assert.deepEqual(again.interrupted, []); } finally { await rm(f.root, { recursive: true, force: true }); } }); test("the queued pass leaves a queued meta whose writer is still alive", async () => { const f = await fixture([ { id: "Q1", spec: SPEC, pid: ALIVE_OTHER }, { id: "Q2", spec: SPEC, pid: DEAD, kind: "sync" }, ]); try { const rq = recordingRequeue(); const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, ...check, }); assert.deepEqual(rq.calls, []); assert.deepEqual(res.cancelled.map((c) => c.id), ["Q2"]); assert.equal((await f.read("Q1")).status, "queued"); } finally { await rm(f.root, { recursive: true, force: true }); } }); // RELEASE 18: a queued publish stage is never re-queued. The lane and Publish // now re-derive stages from the stamps; a stage re-queued with its run's // preconditions would wait on a run that is gone. test("a queued publish stage is cancelled as `publish`, never re-queued, even with a fresh spec", async () => { const stage: JobSpec = { kind: "publish-build-site", slug: "jeralyzer", params: { kind: "build-site", target: "jeralyzer", runId: "r1", indexAfter: BOOT - HOUR }, }; const f = await fixture([ { id: "P1", kind: "publish-build-site", queueKey: "publish", spec: stage }, { id: "P2", kind: "publish-update-index", queueKey: "publish", spec: { kind: "publish-update-index", slug: "_index" } }, { id: "F1", spec: SPEC }, ]); try { const rq = recordingRequeue(); const lines: string[] = []; const res = await settleQueuedJobMetas({ paths: f.paths, requeue: rq.fn, bootedAt: BOOT, log: (l) => lines.push(l), }); assert.deepEqual(rq.calls, [SPEC]); assert.deepEqual( res.cancelled.map((c) => [c.id, c.category]).sort(), [ ["P1", "publish"], ["P2", "publish"], ], ); assert.equal((await f.read("P1")).cancelReason, PUBLISH_RESTART_REASON); assert.match(PUBLISH_RESTART_REASON, /the publish lane re-derives stages from on-disk state/); assert.match(lines.at(-1) ?? "", /cancelled 2 \(publish 2\)/); } finally { await rm(f.root, { recursive: true, force: true }); } });