import { test, beforeEach } from "node:test"; import assert from "node:assert/strict"; import { defaultPublish, type PublishSettings } from "../lib/settingsSchema"; import type { JobDoneResult } from "../jobs/streamCommand"; import { buildPublishStatus, type PublishInputs, type PublishSiteInput } from "./publishPlan"; import { publishLaneMemory, resetPublishLaneMemory } from "./publishLaneState"; import { publishLoop, runPublishPass, type PublishRunnerDeps } from "./publishRunner"; import type { EnqueueStageResult } from "./publishStages"; import type { StageRequest } from "./stages"; import type { BuiltStamp, IndexStamp } from "./stamps"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test publish/publishRunner.test.ts // // The publish lane's runner over a fake world: the "stages" are promises this // test settles, and each one moves the world the way the real stage would // (the index stamp, a built stamp). No process is spawned, nothing on disk. const MIN = 60_000; const T0 = 2_000_000_000_000; function stamp(over: Partial = {}): IndexStamp { return { v: 1, stampId: "s1", generation: 1, scannedAt: T0 - 120 * MIN, builtAt: T0 - 119 * MIN, templatesAt: T0 - 119 * MIN, commit: null, index: { shortCircuited: false, added: 0, changed: 0, removed: 0, heldChannels: [] }, stats: { shortCircuited: false, notIndexedYet: 0, notIndexable: 0 }, sites: { alpha: { siteFp: null, statsFp: null, inputSig: "a1" }, beta: { siteFp: null, statsFp: null, inputSig: "b1" } }, hubSig: "h1", ...over, }; } function built(target: string, sig: string, at: number): BuiltStamp { return { v: 1, stampId: `${target}-${at}`, target, kind: "site", indexStampId: "s1", inputSig: sig, builtAt: at, commit: null, branch: "main", runner: "local", audience: "public", corpusGeneratedAt: null, files: 1, bytes: 1, archivesStaged: 0, }; } function siteIn(id: string, policy: PublishSiteInput["policy"]): PublishSiteInput { return { siteId: id, title: id, private: false, listed: true, members: [`${id}-ch`], built: built(id, `${id[0]}1`, T0 - 118 * MIN), deployed: null, bundleProblem: null, deployProblem: null, pagesProblem: null, configChangedAt: null, policy, cloudflareProject: id, url: null, }; } type World = { inputs: PublishInputs; clock: number; settings: PublishSettings; enqueued: StageRequest[]; // The stage in flight: settle it to let the pass go on. pending: { req: StageRequest; settle: (s: JobDoneResult["status"]) => void } | null; onEnqueue?: (req: StageRequest) => void; }; function world(policyAlpha: PublishSiteInput["policy"] = "build"): World { const settings = { ...defaultPublish(), enabled: true, refreshEveryMinutes: 60 }; return { clock: T0, settings, enqueued: [], pending: null, inputs: { index: { stamp: stamp(), lastIngestDoneAt: T0 - 30 * MIN, configChangedAt: null }, ingestByChannel: { "alpha-ch": T0 - 30 * MIN }, commit: null, sites: [siteIn("alpha", policyAlpha), siteIn("beta", "off")], hub: { built: null, deployed: null, bundleProblem: null, deployProblem: null, pagesProblem: null, configChangedAt: null, policy: "off", cloudflareProject: null, url: null, }, homepage: { built: null, deployed: null, bundleProblem: null, deployProblem: null, pagesProblem: null, configChangedAt: null, policy: "off", cloudflareProject: null, url: null, mainHead: null, }, settings, jobs: [], ended: [], lane: { known: true, running: true, jobId: "lane", passRunning: false, lastCheckAt: null, nextCheckAt: null, lastPassAt: null, lastPassSummary: null, lastDecision: null, }, }, }; } // Apply what a stage that ended `done` writes. function applyDone(w: World, req: StageRequest): void { if (req.kind === "update-index") { w.inputs.index.stamp = stamp({ stampId: "s2", scannedAt: w.clock - 1000, builtAt: w.clock, sites: { alpha: { siteFp: null, statsFp: null, inputSig: "a2" }, beta: { siteFp: null, statsFp: null, inputSig: "b2" }, }, }); } else if (req.kind === "build-site") { const s = w.inputs.sites.find((x) => x.siteId === req.target)!; s.built = built(s.siteId, `${s.siteId[0]}2`, w.clock); } } function deps(w: World): PublishRunnerDeps { return { readStatus: async () => buildPublishStatus({ ...w.inputs, settings: w.settings }, w.clock), settings: () => w.settings, now: () => w.clock, sleep: async (ms) => { w.clock += ms; }, enqueue: async (req): Promise => { w.enqueued.push(req); w.onEnqueue?.(req); let settle!: (s: JobDoneResult["status"]) => void; const done = new Promise((resolve) => { settle = (status) => { w.clock += MIN; if (status === "done") applyDone(w, req); w.pending = null; resolve({ status, jobId: `J${w.enqueued.length}` }); }; }); w.pending = { req, settle }; return { ok: true, jobId: `J${w.enqueued.length}`, stream: new ReadableStream(), done }; }, }; } // Settle every stage the pass dispatches with `status` as it arrives. function autoSettle(w: World, status: (req: StageRequest) => JobDoneResult["status"] = () => "done"): void { w.onEnqueue = (req) => queueMicrotask(() => w.pending?.settle(status(req))); } const signals = () => ({ signal: new AbortController().signal, drain: new AbortController().signal }); const kinds = (w: World) => w.enqueued.map((r) => `${r.kind} ${r.target}`); beforeEach(() => resetPublishLaneMemory()); test("a pass: the index, then — re-planned against the new stamp — the policy's build; never forced", async () => { const w = world("build"); autoSettle(w); const out = await runPublishPass(deps(w), signals(), () => {}); assert.deepEqual(kinds(w), ["update-index _index", "build-site alpha"]); assert.equal(out.ended, "nothing left to do"); for (const r of w.enqueued) { assert.equal(r.force, undefined); assert.equal(r.runId, out.runId); assert.match(r.runId, /^lane-/); } // The build was planned AFTER the index ran: no indexAfter needed. assert.equal(w.enqueued[1].indexAfter, undefined); const mem = publishLaneMemory(); assert.equal(mem.passRunning, false); assert.equal(mem.lastPassAt, w.clock); assert.match(mem.lastPass?.summary ?? "", /2 stages/); }); test("a hold between stages stops dispatching; the stage in flight is never killed", async () => { const w = world("production"); w.onEnqueue = (req) => queueMicrotask(() => { // The operator holds the lane while the index update runs. if (req.kind === "update-index") w.settings = { ...w.settings, held: true }; w.pending?.settle("done"); }); const out = await runPublishPass(deps(w), signals(), () => {}); assert.deepEqual(kinds(w), ["update-index _index"]); assert.match(out.ended, /the lane is held: no next stage dispatched/); assert.equal(out.ran[0].status, "done", "the index stage ran to its end"); }); test("quiet hours beginning between stages stop dispatching too", async () => { const w = world("build"); w.onEnqueue = () => queueMicrotask(() => { const h = new Date(w.clock + MIN).getHours(); w.settings = { ...w.settings, quietHours: { start: h, end: (h + 2) % 24 } }; w.pending?.settle("done"); }); const out = await runPublishPass(deps(w), signals(), () => {}); assert.deepEqual(kinds(w), ["update-index _index"]); assert.match(out.ended, /quiet hours began/); }); test("a drain finishes the stage in flight and dispatches nothing after it", async () => { const w = world("build"); const drain = new AbortController(); w.onEnqueue = () => queueMicrotask(() => { drain.abort(); // The stage is still running when the drain lands; it ends after. setTimeout(() => w.pending?.settle("done"), 5); }); const out = await runPublishPass(deps(w), { signal: new AbortController().signal, drain: drain.signal }, () => {}); assert.deepEqual(kinds(w), ["update-index _index"]); assert.equal(out.ran[0].status, "done", "the drain waited for the stage"); assert.equal(out.ended, "the runner was drained"); }); test("a stop does not wait for the stage: it keeps running as its own job", async () => { const w = world("build"); const stop = new AbortController(); w.onEnqueue = () => queueMicrotask(() => stop.abort()); const out = await runPublishPass(deps(w), { signal: stop.signal, drain: new AbortController().signal }, () => {}); assert.equal(out.ran[0].status, "still running"); assert.equal(out.ended, "the runner was stopped"); assert.ok(w.pending, "the stage was not settled (nor killed) by the stop"); }); test("a failed index update ends the pass; a failed build drops its deploy", async () => { const w = world("production"); autoSettle(w, (r) => (r.kind === "update-index" ? "failed" : "done")); const a = await runPublishPass(deps(w), signals(), () => {}); assert.deepEqual(kinds(w), ["update-index _index"]); assert.match(a.ended, /the index update failed/); const v = world("production"); autoSettle(v, (r) => (r.kind === "build-site" ? "failed" : "done")); const b = await runPublishPass(deps(v), signals(), () => {}); assert.deepEqual(kinds(v), ["update-index _index", "build-site alpha"]); assert.equal(b.ended, "every stage left waits on a build that failed"); }); test("a pass with the production policy deploys after the build", async () => { const w = world("production"); autoSettle(w); await runPublishPass(deps(w), signals(), () => {}); assert.deepEqual(kinds(w), ["update-index _index", "build-site alpha", "deploy-site alpha"]); assert.equal(w.enqueued[2].preview, undefined); }); test("a stage someone else queued makes the pass yield", async () => { const w = world("build"); w.onEnqueue = () => queueMicrotask(() => { w.inputs.jobs = [ { id: "X", kind: "build-site", target: "beta", status: "queued", runId: "run-x", queuedAt: w.clock }, ]; w.pending?.settle("done"); }); const out = await runPublishPass(deps(w), signals(), () => {}); assert.deepEqual(kinds(w), ["update-index _index"]); assert.match(out.ended, /someone else queued/); }); test("the loop: no stamp → a pass at once; then the refresh interval gates the next index update", async () => { const w = world("build"); w.inputs.index.stamp = null; autoSettle(w); const drain = new AbortController(); let checks = 0; const d = deps(w); const lines: string[] = []; await publishLoop( { ...d, readStatus: async () => { checks++; // New data arrives after every check; stop after a few checks. w.inputs.index.lastIngestDoneAt = w.clock; w.inputs.ingestByChannel["alpha-ch"] = w.clock; if (checks > 40) drain.abort(); return d.readStatus(); }, }, { signal: new AbortController().signal, drain: drain.signal }, (l) => lines.push(l), ); const indexRuns = w.enqueued.filter((r) => r.kind === "update-index").length; assert.ok(indexRuns >= 2, `the index was updated again after the interval (${indexRuns})`); // checkEvery 10 min, refresh 60 min: never more often than every 60 min. const indexTimes = lines.filter((l) => l.includes("update-index _index →")).length; assert.equal(indexTimes, indexRuns); assert.ok(w.clock - T0 >= (indexRuns - 1) * 60 * MIN, "updates were at least refreshEveryMinutes apart"); assert.match(lines.at(-1) ?? "", /lane runner drained/); }); test("the loop ends when the lane is switched off, and a held lane never starts a pass", async () => { const w = world("build"); w.settings = { ...w.settings, held: true }; autoSettle(w); let reads = 0; const d = deps(w); await publishLoop( { ...d, readStatus: async () => { if (++reads > 3) w.settings = { ...w.settings, enabled: false }; return d.readStatus(); }, }, signals(), () => {}, ); assert.deepEqual(w.enqueued, [], "held: nothing dispatched"); assert.equal(publishLaneMemory().lastDecision, "held"); });