import { test } from "node:test"; import assert from "node:assert/strict"; import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/snapshotScheduler.test.ts // // ONE REGENERATION AT A TIME, ONCE PER CHANNEL (release 17 slice D0). Every // refresh-report used to run on the empty queue key — all at once — and two // passes a moment apart could both enqueue the same channel. Temp corpus only. const root = await mkdtemp(path.join(tmpdir(), "snapshot-scheduler-")); process.env.TRANSCRIPTS_DIR = path.join(root, "transcripts"); process.env.SETTINGS_FILE = path.join(root, "settings.json"); await writeFile(process.env.SETTINGS_FILE, "{}"); for (const slug of ["alpha", "beta", "gamma"]) { const dir = path.join(process.env.TRANSCRIPTS_DIR, "channels", slug); await mkdir(path.join(dir, "data"), { recursive: true }); await writeFile(path.join(dir, "config.json"), JSON.stringify({ handling: "youtube" })); } // A channel relocated to a drive that is not there: its walk refuses (guard 3) // with a sentence naming the unreachable media. { const dir = path.join(process.env.TRANSCRIPTS_DIR, "channels", "unmounted"); await mkdir(dir, { recursive: true }); await writeFile( path.join(dir, "config.json"), JSON.stringify({ handling: "youtube", dataDir: path.join(root, "no-such-drive", "unmounted", "data") }), ); } const { getPaths } = await import("../lib/paths"); const { getRegistry, newJobId } = await import("./registry"); const { REFRESH_REPORT_ACTIVE, REFRESH_REPORT_QUEUE, isRefreshReportPending, refreshReportWaitNotice, requestRefreshReport, startRefreshReport, waitForRefreshReport, } = await import("./snapshotScheduler"); // A refresh-report that is RUNNING until released: a record on the queue whose // start does nothing. function holdQueue(slug: string): () => void { const registry = getRegistry(); const id = newJobId(); const holder = { id, kind: "refresh-report", queueKey: REFRESH_REPORT_QUEUE, channelSlug: slug, status: "queued" as const, queuedAt: Date.now(), logPath: "/dev/null", }; registry.register(holder); registry.enqueue(holder, { start: () => {}, onCancel: () => {} }); return () => registry.finalize(id, "done"); } test.after(async () => { await rm(root, { recursive: true, force: true }); }); test("two requests for one channel at once enqueue one regeneration", async () => { const paths = getPaths(); const [a, b] = await Promise.all([ startRefreshReport(paths, "alpha"), startRefreshReport(paths, "alpha"), ]); const started = [a, b].filter((r) => r.ok); const refused = [a, b].filter((r) => !r.ok); assert.equal(started.length, 1, "exactly one starts"); assert.equal(refused.length, 1); const r = refused[0]; assert.ok(!r.ok && r.info === true && r.error === REFRESH_REPORT_ACTIVE); const s = started[0]; assert.ok(s.ok); assert.equal((await s.done).status, "done"); assert.equal(isRefreshReportPending("alpha"), false); // Once it has finished, a new request starts a new one. const again = await startRefreshReport(paths, "alpha"); assert.ok(again.ok); await again.done; }); test("every regeneration runs on the one serial refresh-report queue", async () => { const paths = getPaths(); const a = await startRefreshReport(paths, "alpha"); const b = await startRefreshReport(paths, "beta"); assert.ok(a.ok && b.ok); for (const id of [a.jobId, b.jobId]) { assert.equal(getRegistry().get(id)?.queueKey, REFRESH_REPORT_QUEUE); } // Serial: the second is behind the first unless the first already finished. const first = getRegistry().get(a.jobId); const second = getRegistry().get(b.jobId); if (first?.status === "running") assert.equal(second?.status, "queued"); assert.equal((await a.done).status, "done"); assert.equal((await b.done).status, "done"); }); test("a running regeneration gets one queued successor, and no second", async () => { const paths = getPaths(); // A regeneration of beta that is RUNNING and stays so until released: a // record on the refresh-report queue whose start does nothing. const registry = getRegistry(); const holderId = newJobId(); const holder = { id: holderId, kind: "refresh-report", queueKey: REFRESH_REPORT_QUEUE, channelSlug: "beta", status: "queued" as const, queuedAt: Date.now(), logPath: "/dev/null", }; registry.register(holder); registry.enqueue(holder, { start: () => {}, onCancel: () => {} }); assert.equal(registry.get(holderId)?.status, "running"); // A change during the walk: it may already have passed it, so one successor // is queued behind it… const successor = await startRefreshReport(paths, "beta"); assert.ok(successor.ok, "queued behind the running one, not dropped"); assert.equal(registry.get(successor.jobId)?.status, "queued"); assert.equal(isRefreshReportPending("beta"), true); // …and a further change finds it queued: it has not started reading yet. const third = await startRefreshReport(paths, "beta"); assert.ok(!third.ok && third.info === true && third.error === REFRESH_REPORT_ACTIVE); registry.finalize(holderId, "done"); assert.equal((await successor.done).status, "done"); }); // A PERSON'S REFRESH WAITS A BOUNDED TIME (re-review R1, R2). test("a refresh behind the queue answers with where it is, once the wait runs out", async () => { const paths = getPaths(); const release = holdQueue("beta"); try { const asked = await requestRefreshReport(paths, "gamma"); assert.ok(asked.ok && asked.started); const wait = await waitForRefreshReport(paths, asked.jobId, { timeoutMs: 300, pollMs: 50 }); assert.equal(wait.state, "waiting"); assert.ok(wait.state === "waiting" && wait.ahead === 1, "one regeneration ahead of it"); assert.equal( wait.state === "waiting" && refreshReportWaitNotice(wait), `Queued behind 1 report regeneration — the report updates when it finishes (job ${asked.jobId}).`, ); // A second person's click finds the same job, queued. const again = await requestRefreshReport(paths, "gamma"); assert.deepEqual(again, { ok: true, jobId: asked.jobId, started: false }); release(); const finished = await waitForRefreshReport(paths, asked.jobId, { pollMs: 20 }); assert.deepEqual(finished, { state: "done", jobId: asked.jobId }); } finally { release(); } }); test("a failed walk is a failure with its sentence, for the click that found it queued too", async () => { const paths = getPaths(); const release = holdQueue("beta"); let queuedId = ""; try { const first = await requestRefreshReport(paths, "unmounted"); assert.ok(first.ok && first.started); queuedId = first.jobId; // Another click while it waits: it did not start this job, and still // hears how it ended. const second = await requestRefreshReport(paths, "unmounted"); assert.deepEqual(second, { ok: true, jobId: queuedId, started: false }); release(); const wait = await waitForRefreshReport(paths, queuedId, { pollMs: 20 }); assert.equal(wait.state, "failed"); assert.ok(wait.state === "failed" && /unmounted/.test(wait.error), "the walk's own sentence"); } finally { release(); } });