import { test } from "node:test"; import assert from "node:assert/strict"; import { mkdir, mkdtemp, readFile, rm, symlink, 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 { getRegistry, newJobId } from "./registry"; import { runManagedCommand, runManagedFunction, serialWriter, type StreamActionResult, } from "./streamCommand"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/streamCommand.test.ts // // A JOB CANCELLED WHILE STILL QUEUED (release 13 slice W1). Its start() never // runs, so the terminal sidecar write in start()'s .finally never happens; // onCancel is the only place that can write it, and it used to write nothing — // `.meta.json` kept the "queued" the enqueue wrote, and a boot pass could // re-queue a job the operator had cancelled. Each test holds its own queue key // with a job that runs until released, so the job under test only ever queues // and nothing is spawned. Temp .jobs dirs only. type Started = Extract; function ok(r: StreamActionResult): Started { if (!r.ok) throw new Error(r.error); return r; } async function jobsDir(): Promise<{ paths: Paths; root: string }> { const root = await mkdtemp(path.join(tmpdir(), "stream-command-")); return { paths: { jobsDir: path.join(root, ".jobs") } as Paths, root }; } async function readMeta(paths: Paths, id: string): Promise { try { return JSON.parse( await readFile(path.join(paths.jobsDir, `${id}.meta.json`), "utf8"), ) as JobMeta; } catch { return null; } } // The sidecar is written asynchronously and nothing hands out its promise, so // wait (briefly) for it to reach `status`; report the last one seen otherwise. async function metaReaches( paths: Paths, id: string, status: string, ): Promise { let last: JobMeta | null = null; for (let i = 0; i < 100; i++) { last = await readMeta(paths, id); if (last?.status === status) return last; await new Promise((r) => setTimeout(r, 20)); } return last; } // A job on `queueKey` that runs until released: everything submitted behind it // queues. async function hold(paths: Paths, queueKey: string) { let release!: () => void; const held = new Promise((r) => { release = r; }); const job = ok( await runManagedFunction({ kind: "test-holder", queueKey, paths, fn: () => held, }), ); return { async release() { release(); await job.done; }, }; } test("a function job cancelled while queued ends with a cancelled sidecar", async () => { const { paths, root } = await jobsDir(); const queueKey = `test:queued-cancel:${newJobId()}`; const holder = await hold(paths, queueKey); try { const job = ok( await runManagedFunction({ kind: "test-queued", queueKey, paths, fn: async () => { throw new Error("must never start"); }, }), ); assert.equal(getRegistry().get(job.jobId)?.status, "queued"); assert.equal((await metaReaches(paths, job.jobId, "queued"))?.status, "queued"); assert.equal(getRegistry().cancel(job.jobId), true); assert.deepEqual(await job.done, { status: "cancelled", jobId: job.jobId }); const meta = await metaReaches(paths, job.jobId, "cancelled"); assert.equal(meta?.status, "cancelled"); assert.equal(typeof meta?.endedAt, "number"); assert.equal(meta?.startedAt, undefined, "it never started"); } finally { await holder.release(); await rm(root, { recursive: true, force: true }); } }); test("a command job cancelled while queued ends with a cancelled sidecar", async () => { const { paths, root } = await jobsDir(); const queueKey = `test:queued-cancel-cmd:${newJobId()}`; const holder = await hold(paths, queueKey); try { const job = ok( await runManagedCommand({ kind: "test-queued-command", queueKey, paths, cwd: root, // Never run: it only queues behind the holder. command: "false", args: [], }), ); assert.equal((await metaReaches(paths, job.jobId, "queued"))?.status, "queued"); assert.equal(getRegistry().cancel(job.jobId), true); assert.equal((await job.done).status, "cancelled"); assert.equal((await metaReaches(paths, job.jobId, "cancelled"))?.status, "cancelled"); } finally { await holder.release(); await rm(root, { recursive: true, force: true }); } }); // The sidecar writer behind every job (metaWriter): each write starts only when // the previous one has settled, a write that rejects does not stall the ones // after it, and the last one rejecting is not an unhandled rejection (the test // runner would fail this test on one). test("serialWriter runs writes one at a time, in order, past a rejection", async () => { const events: string[] = []; let n = 0; const write = serialWriter(async () => { const i = ++n; events.push(`start ${i}`); // The first write is the slow one; unchained, the later ones would start // (and end) while it is still in flight. await new Promise((r) => setTimeout(r, i === 1 ? 40 : 1)); events.push(`end ${i}`); if (i === 2 || i === 4) throw new Error("a writer that rejects"); }); write(); write(); write(); write(); await new Promise((r) => setTimeout(r, 150)); assert.deepEqual(events, [ "start 1", "end 1", "start 2", "end 2", "start 3", "end 3", "start 4", "end 4", ]); }); // THE MEDIA GUARD, ASKED AGAIN AT THE START (release 16 slice RM). A media job // queued while its channel was in place, behind other work on its queue, and // started after a move wrote the channel's marker, must not write into the // tree being copied: the guard that passed at enqueue is asked again when the // queue starts the job, and the refusal — the hold's words — fails it before // `fn` runs. test("a media job queued before a move's marker and started after it refuses at its start", async () => { const { paths: base, root } = await jobsDir(); const paths = { ...base, channelsDir: path.join(root, "channels") } as Paths; await mkdir(path.join(paths.channelsDir, "alpha", "data"), { recursive: true }); const queueKey = `test:start-guard:${newJobId()}`; const holder = await hold(paths, queueKey); let ran = false; try { const job = ok( await runManagedFunction({ kind: "whisper-all", queueKey, paths, channelSlug: "alpha", fn: async () => { ran = true; }, }), ); assert.equal(getRegistry().get(job.jobId)?.status, "queued"); // The move starts while the job waits. await writeFile( path.join(paths.channelsDir, "alpha", ".relocating.json"), JSON.stringify({ target: "/mnt/platter/alpha/data", direction: "out", startedAt: new Date().toISOString(), phase: "copy", }), ); await holder.release(); assert.equal((await job.done).status, "failed"); assert.equal(ran, false, "fn never ran"); const log = await readFile(path.join(paths.jobsDir, `${job.jobId}.log`), "utf8"); assert.match( log, /Channel "alpha" is held: its media is moving \(a move is in progress or was interrupted\)/, ); } finally { await rm(root, { recursive: true, force: true }); } }); // RELEASE 17: a TEXT kind (`needsText`) asks the text guard. It runs while the // channel's media drive is unmounted, and is refused — before any record — // only where the text itself cannot be read: the retired layout. test("a text job runs over an unmounted media drive and is refused on a legacy channel", async () => { const { paths: base, root } = await jobsDir(); const paths = { ...base, channelsDir: path.join(root, "channels") } as Paths; const channelDir = path.join(paths.channelsDir, "alpha"); await mkdir(path.join(channelDir, "data"), { recursive: true }); const target = path.join(root, "platter", "alpha", "media"); await writeFile( path.join(channelDir, "config.json"), JSON.stringify({ handling: "youtube", url: "https://example.com/a", mediaDir: target }), ); await symlink(target, path.join(channelDir, "media")); // not mounted try { let ran = 0; const job = ok( await runManagedFunction({ kind: "normalize-transcripts", queueKey: `test:text-guard:${newJobId()}`, paths, channelSlug: "alpha", fn: async () => { ran += 1; }, }), ); assert.equal((await job.done).status, "done"); assert.equal(ran, 1); // A media kind on the same channel is refused, before any record. const media = await runManagedFunction({ kind: "whisper-all", queueKey: `test:text-guard:${newJobId()}`, paths, channelSlug: "alpha", fn: async () => {}, }); assert.equal(media.ok, false); assert.match((media as { error: string }).error, /does not exist \(drive not mounted\?\)/); // The retired layout: the text kind is refused too, naming the way out. await writeFile( path.join(channelDir, "config.json"), JSON.stringify({ handling: "youtube", url: "https://example.com/a", dataDir: "/mnt/platter/alpha/data" }), ); const refused = await runManagedFunction({ kind: "normalize-transcripts", queueKey: `test:text-guard:${newJobId()}`, paths, channelSlug: "alpha", fn: async () => { ran += 1; }, }); assert.equal(refused.ok, false); assert.match((refused as { error: string }).error, /archilyzer storage migrate-tier alpha/); assert.equal(ran, 1); } finally { await rm(root, { recursive: true, force: true }); } }); // THE ONE CANCEL THAT MUST NOT: the graceful-shutdown reaper cancels every // queued job only so the exit cannot promote one into a child. Nobody cancelled // it, and its `queued` sidecar is what the boot pass (bootQueuedJobs.ts) // settles or re-queues on the next start. Last in the file: it puts the // process's registry into shutdown, and replaces it afterwards. test("at shutdown a queued job's sidecar stays queued, for the boot pass", async () => { const { paths, root } = await jobsDir(); const queueKey = `test:shutdown-cancel:${newJobId()}`; const holder = await hold(paths, queueKey); try { const job = ok( await runManagedFunction({ kind: "test-queued", queueKey, paths, fn: async () => { throw new Error("must never start"); }, }), ); assert.equal((await metaReaches(paths, job.jobId, "queued"))?.status, "queued"); getRegistry().beginShutdown(); assert.equal(getRegistry().cancel(job.jobId), true); // The job itself still settles, and the registry records the cancel. assert.equal((await job.done).status, "cancelled"); assert.equal(getRegistry().get(job.jobId)?.status, "cancelled"); // Nothing rewrote the sidecar (a wrong write would land within ms). await new Promise((r) => setTimeout(r, 200)); assert.equal((await readMeta(paths, job.jobId))?.status, "queued"); } finally { await holder.release(); globalThis.__yttJobRegistry__ = undefined; await rm(root, { recursive: true, force: true }); } });