import { test } from "node:test"; import assert from "node:assert/strict"; import { mkdtemp, mkdir, readFile, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; import { settleRunningJobMetas } from "../jobs/bootQueuedJobs"; import { getRegistry, newJobId, type JobRecord } from "../jobs/registry"; import type { AutoQueueKind } from "../lib/autoQueueTypes"; import type { AutoRunnerInFlight } from "./autoRunner"; import { channelWriters, channelWritersRefusal, describeChannelWriter, type ChannelWritersSource, } from "./channelWriters"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/channelWriters.test.ts // // WHO IS WRITING INTO A CHANNEL, named (release 16 slice RM). The registry and // the lanes are injected: what is pinned is the reading of them — which // records count, in which order, and the sentence the move refuses with. function job(over: Partial): JobRecord { return { id: "J1", kind: "whisper-all", queueKey: "transcription", status: "running", queuedAt: 1, logPath: "/dev/null", ...over, }; } function unit( lane: AutoQueueKind, over: Partial, ): { lane: AutoQueueKind; unit: AutoRunnerInFlight } { return { lane, unit: { videoId: "v1", leafId: "leaf", channelSlug: "alpha", startedAt: 1, ...over, }, }; } function source( jobs: JobRecord[], units: { lane: AutoQueueKind; unit: AutoRunnerInFlight }[] = [], ): ChannelWritersSource { return { jobs: () => jobs, units: () => units }; } test("the 2026-09-30 case: a running Transcribe all is named by the video it is on", () => { const writers = channelWriters("realcandaceo", { source: source([ job({ id: "01M3TGYDXA7WCBGN014A90J3JW", channelSlug: "realcandaceo", tasks: [ { id: "v50t5yt", label: "v50t5yt", kind: "transcribe", startedAt: 1 }, ], }), ]), }); assert.equal(writers.length, 1); assert.equal( describeChannelWriter(writers[0]), "a transcription of v50t5yt is running (Transcribe all, job 01M3TGYDXA7WCBGN014A90J3JW)", ); assert.equal( channelWritersRefusal("realcandaceo", writers), 'Cannot move the media of "realcandaceo" now: a transcription of v50t5yt ' + "is running (Transcribe all, job 01M3TGYDXA7WCBGN014A90J3JW) — wait for " + "it or cancel it. Nothing has been touched.", ); }); test("only this channel's jobs, only running ones unless queued are asked for", () => { const jobs = [ job({ id: "A", channelSlug: "alpha", kind: "sync", status: "queued" }), job({ id: "B", channelSlug: "beta", kind: "sync" }), job({ id: "C", channelSlug: "alpha", kind: "sync", status: "done" }), job({ id: "D", channelSlug: "alpha", kind: "download-missing" }), ]; assert.deepEqual( channelWriters("alpha", { source: source(jobs) }).map((w) => w.source === "job" ? w.jobId : "", ), ["D"], ); // Running before queued, whatever the registry order. assert.deepEqual( channelWriters("alpha", { source: source(jobs), includeQueued: true }).map( (w) => (w.source === "job" ? w.jobId : ""), ), ["D", "A"], ); assert.equal( describeChannelWriter( channelWriters("alpha", { source: source(jobs), includeQueued: true })[1], ), "Sync is queued (job A)", ); }); test("the move's own kind is not a writer to refuse over", () => { const jobs = [ job({ id: "M", channelSlug: "alpha", kind: "relocate-channel-media" }), ]; assert.equal(channelWriters("alpha", { source: source(jobs) }).length, 1); assert.deepEqual( channelWriters("alpha", { source: source(jobs), ignoreKinds: ["relocate-channel-media"], }), [], ); }); test("a lane's in-process unit is a writer, and its way out is the lane", () => { const writers = channelWriters("alpha", { source: source([], [unit("transcription", { videoId: "v9" })]), }); assert.equal( channelWritersRefusal("alpha", writers), 'Cannot move the media of "alpha" now: a transcription of v9 is running ' + "(the transcription lane) — wait for it, or hold the transcription lane. " + "Nothing has been touched.", ); }); test("a download unit is named once, as the lane's, not again as its job", () => { const writers = channelWriters("alpha", { source: source( [job({ id: "U1", kind: "auto-download-unit", channelSlug: "alpha" })], [unit("download", { videoId: "v2", jobId: "U1" })], ), }); assert.equal(writers.length, 1); assert.equal( describeChannelWriter(writers[0]), "a download of v2 is running (the download lane)", ); }); test("a single-video job with no task is named by its record's video", () => { const writers = channelWriters("alpha", { source: source([ job({ id: "W", kind: "whisper-video", channelSlug: "alpha", videoId: "v3" }), ]), }); assert.equal( describeChannelWriter(writers[0]), "whisper-video of v3 is running (job W)", ); }); test("more than one writer: the first named, the rest counted; none is no refusal", () => { const writers = channelWriters("alpha", { source: source( [job({ id: "D", kind: "download-missing", channelSlug: "alpha" })], [unit("digest", { videoId: "v4" })], ), }); assert.match( channelWritersRefusal("alpha", writers) ?? "", /Download missing is running \(job D\), and 1 more writer\(s\) — wait for it or cancel it\./, ); assert.equal(channelWritersRefusal("alpha", []), null); }); // A CANCEL IS A REQUEST, NOT AN EXIT (the review's L2). A running job's cancel // marks it `cancelled` at once while its function winds down; it is a writer, // "stopping", until the registry stamps `endedAt` when the function returns. test("a job cancelled but still winding down is a writer until it has stopped", () => { const winding = job({ id: "S", channelSlug: "alpha", status: "cancelled" }); const writers = channelWriters("alpha", { source: source([winding]) }); assert.equal(writers.length, 1); assert.equal( channelWritersRefusal("alpha", writers), 'Cannot move the media of "alpha" now: Transcribe all is stopping (job S) ' + "— wait for it to stop. Nothing has been touched.", ); const stopped = { ...winding, endedAt: 2 }; assert.deepEqual(channelWriters("alpha", { source: source([stopped]) }), []); }); test("through the live registry: cancel leaves it stopping, the job's end stamps endedAt", () => { const registry = getRegistry(); const id = newJobId(); const record: JobRecord = { id, kind: "whisper-all", queueKey: `test:writers:${id}`, channelSlug: `writers-${id}`, status: "queued", queuedAt: Date.now(), logPath: "/dev/null", }; registry.register(record); registry.enqueue(record, { start: () => {}, onCancel: () => {} }); const slug = record.channelSlug as string; assert.equal(channelWriters(slug).length, 1, "running"); assert.equal(registry.cancel(id), true); assert.equal(record.status, "cancelled"); assert.equal(record.endedAt, undefined); const stopping = channelWriters(slug); assert.equal(stopping.length, 1); assert.equal(stopping[0].source === "job" && stopping[0].status, "stopping"); // The job's function returns: streamCommand finalizes it as cancelled. registry.finalize(id, "cancelled"); assert.equal(typeof record.endedAt, "number"); assert.deepEqual(channelWriters(slug), []); }); test("mediaOnly (release 17): a digest job and the digest lane are not media writers", () => { const jobs = [ job({ id: "J1", kind: "digest-channel-local", channelSlug: "alpha" }), job({ id: "J2", kind: "normalize-transcripts", channelSlug: "alpha" }), job({ id: "J3", kind: "whisper-all", channelSlug: "alpha" }), ]; const units = [ unit("digest", { videoId: "d1" }), unit("transcription", { videoId: "t1" }), ]; const all = channelWriters("alpha", { source: source(jobs, units) }); assert.equal(all.length, 5); const media = channelWriters("alpha", { source: source(jobs, units), mediaOnly: true }); assert.deepEqual( media.map((w) => (w.source === "job" ? w.jobId : `${w.lane}:${w.videoId}`)), ["J3", "transcription:t1"], ); }); // A GHOST NEVER HOLDS A MOVE (release 17 slice D0): a `running` meta a dead // process left on disk is no writer — before the boot pass closes it, and // after. test("a running meta a dead process left on disk is never a writer", async () => { const root = await mkdtemp(path.join(tmpdir(), "channel-writers-ghost-")); try { const jobsDir = path.join(root, ".jobs"); await mkdir(jobsDir, { recursive: true }); const id = newJobId(); const slug = `ghost-${id}`; await writeFile( path.join(jobsDir, `${id}.meta.json`), JSON.stringify({ id, kind: "refresh-report", queueKey: "", channelSlug: slug, status: "running", queuedAt: Date.now() - 60_000, startedAt: Date.now() - 60_000, pid: 2 ** 22 + 1, // past pid_max: no such process }), ); assert.deepEqual(channelWriters(slug, { includeQueued: true }), []); const res = await settleRunningJobMetas({ paths: { jobsDir } as Paths, bootedAt: Date.now(), isLive: (j) => getRegistry().get(j) !== undefined, }); assert.deepEqual(res.interrupted.map((j) => j.id), [id]); const meta = JSON.parse( await readFile(path.join(jobsDir, `${id}.meta.json`), "utf8"), ) as { status: string }; assert.equal(meta.status, "cancelled"); assert.deepEqual(channelWriters(slug, { includeQueued: true }), []); } finally { await rm(root, { recursive: true, force: true }); } });