import { test } from "node:test"; import assert from "node:assert/strict"; import { defaultSiteSettings } from "../lib/settings"; import { createScheduler, type Scheduler } from "../jobs/scheduler"; import type { JobRecord } from "../jobs/registry"; import type { JobMeta } from "../jobs/jobMeta"; import type { Paths } from "../lib/paths"; import type { DiskGateStatus } from "../lib/diskSpace"; import type { AutoRunnerStatus } from "../controller/autoRunner"; import type { AutoQueueKind } from "../lib/autoQueueTypes"; import type { PoolReader, RegistryReader } from "./inputs"; import { buildActiveJobsPayload, liveJobRows, rowsFromJobsPage, type ActiveJobsInputs, } from "./activeJobs"; // Run with: node_modules/.bin/tsx --test common/views/activeJobs.test.ts // // This payload used to construct five singletons, stat the disk, tail a log // file and read the clock, which is why the two rules it lives by — "liveJobRows // NEVER heals" and "buildActiveJobsPayload ALWAYS heals" — were prose in a // comment and nothing else. Both are assertions below, against a real // createScheduler(); everything else is a counted stub, so "how many times did // this read the disk" is a number a test can state. const NOW = 1_700_000_000_000; // Older than jobRows' STUCK_AGE_MS (10 minutes), which is what makes a live // running record with no tasks read as possibly-stalled. const LONG_AGO = NOW - 30 * 60 * 1000; function job(over: Partial = {}): JobRecord { return { id: "j1", kind: "transcribe", queueKey: "q", status: "running", queuedAt: NOW - 1000, startedAt: NOW - 1000, logPath: "/dev/null", ...over, }; } const OK_DISK: DiskGateStatus = { ok: true, enabled: true, freeBytes: 100, thresholdBytes: 10, resumeBytes: 10, reason: "ok", message: "", }; function runnerStatus(over: Partial = {}): AutoRunnerStatus { return { kind: "transcription", running: true, jobId: null, startedAt: null, inFlight: [], activeByNode: {}, idleReason: null, ...over, }; } // A scheduler that records every complete() and otherwise IS the real thing — // the heal contract is about a side effect on live state, so a fake that only // counts would pin nothing. function watched(inner: Scheduler): { scheduler: Scheduler; completed: string[]; } { const completed: string[] = []; return { completed, scheduler: { submit: (t) => inner.submit(t), complete: (id) => { completed.push(id); inner.complete(id); }, cancel: (id) => inner.cancel(id), reorder: (id, dir) => inner.reorder(id, dir), promote: (id) => inner.promote(id), positionInQueue: (id) => inner.positionInQueue(id), queues: () => inner.queues(), }, }; } // Put an id in a running slot on `queueKey`, first come first served. function occupy(scheduler: Scheduler, id: string, queueKey = "q"): void { scheduler.submit({ id, queueKey, tier: "foreground", concurrency: 1, start: () => {}, onCancel: () => {}, }); } type Harness = { inputs: ActiveJobsInputs; completed: string[]; statCalls: string[]; metaCalls: string[]; tailCalls: string[]; }; function harness(over: { jobs?: JobRecord[]; slots?: string[]; meta?: (id: string) => JobMeta | null; disk?: DiskGateStatus; runner?: (kind: AutoQueueKind) => AutoRunnerStatus; } = {}): Harness { const jobs = over.jobs ?? []; const { scheduler, completed } = watched(createScheduler()); for (const id of over.slots ?? []) occupy(scheduler, id); const registry: RegistryReader = { list: () => jobs, get: (id: string) => jobs.find((j) => j.id === id), }; const pool: PoolReader = { summary: () => [], isPaused: () => false, canStopPartial: () => false, }; const statCalls: string[] = []; const metaCalls: string[] = []; const tailCalls: string[] = []; return { completed, statCalls, metaCalls, tailCalls, inputs: { // The view never touches paths; it is here because LiveInputs carries it. paths: {} as Paths, settings: defaultSiteSettings(), registry, scheduler, pool, now: () => NOW, channelStat: async (slug) => { statCalls.push(slug); return null; }, jobMeta: async (id) => { metaCalls.push(id); return over.meta?.(id) ?? null; }, tailLog: async (id) => { tailCalls.push(id); return "the last line"; }, disk: over.disk ?? OK_DISK, autoRunnerStatus: over.runner ?? (() => runnerStatus()), }, }; } test("builtAt is the injected clock, not the wall clock", async () => { const h = harness(); const payload = await buildActiveJobsPayload(h.inputs); assert.equal(payload.builtAt, NOW); }); test("the log tail is read for stuck rows only", async () => { // Two running records holding slots. The first has been idle since LONG_AGO // with no tasks — possibly-stalled, the soft case — and the second started a // second ago with work in flight. const h = harness({ jobs: [ job({ id: "stale", queuedAt: LONG_AGO, startedAt: LONG_AGO }), job({ id: "fresh", queueKey: "q2", tasks: [ { id: "v1", label: "v1", kind: "transcribe", startedAt: NOW - 500, }, ], }), ], slots: ["stale"], }); occupy(h.inputs.scheduler, "fresh", "q2"); const payload = await buildActiveJobsPayload(h.inputs); assert.deepEqual(h.tailCalls, ["stale"]); assert.equal( payload.jobs.find((r) => r.id === "stale")?.lastLogLine, "the last line", ); assert.equal(payload.jobs.find((r) => r.id === "fresh")?.lastLogLine, undefined); // Soft case: surfaced, never healed. assert.deepEqual(h.completed, []); assert.equal(payload.summary.stuck, 1); }); test("a channel is stat-ed once however many of its jobs are running", async () => { const h = harness({ jobs: [ job({ id: "a1", channelSlug: "alpha" }), job({ id: "a2", channelSlug: "alpha", queueKey: "q2" }), job({ id: "b1", channelSlug: "beta", queueKey: "q3" }), // No channel at all: a corpus-wide job contributes no slug. job({ id: "c1", queueKey: "q4" }), ], }); await buildActiveJobsPayload(h.inputs); assert.deepEqual(h.statCalls, ["alpha", "beta"]); }); test("a slot whose record is gone becomes a phantom row, named by its sidecar", async () => { const h = harness({ slots: ["ghost"], meta: (id) => id === "ghost" ? ({ id, kind: "digest-sweep", queueKey: "q" } as JobMeta) : null, }); const payload = await buildActiveJobsPayload(h.inputs); assert.deepEqual(h.metaCalls, ["ghost"]); const row = payload.jobs.find((r) => r.id === "ghost"); assert.ok(row, "the phantom is in the payload"); assert.equal(row.kind, "digest-sweep"); assert.equal(row.status, "evicted"); assert.equal(row.source, "slot"); assert.equal(row.stuck?.reason, "record-evicted"); // Hard case: the slot is freed, once, after the row was built. assert.deepEqual(h.completed, ["ghost"]); }); test("a terminal record holding a slot is drawn once, then healed", async () => { const h = harness({ jobs: [job({ id: "done1", status: "done", endedAt: NOW - 100 })], slots: ["done1"], }); const payload = await buildActiveJobsPayload(h.inputs); const rows = payload.jobs.filter((r) => r.id === "done1"); assert.equal(rows.length, 1, "the drift is drawn exactly once"); assert.equal(rows[0].stuck?.reason, "terminal-but-holding-slot"); // It is not ALSO in `recent`: a row in the live head is excluded there. assert.deepEqual(payload.recent.map((r) => r.id), []); assert.deepEqual(h.completed, ["done1"]); // And the slot really is free now — the scheduler is a real one. assert.deepEqual( h.inputs.scheduler.queues().flatMap((v) => v.running), [], ); }); test("liveJobRows never heals, whatever it sees", async () => { // The same drift the payload builder would have freed. A channel page reading // its own jobs must not free a scheduler slot as a side effect. const h = harness({ jobs: [ job({ id: "done1", status: "done", endedAt: NOW - 100 }), job({ id: "run1", channelSlug: "alpha" }), ], slots: ["done1"], }); const rows = await liveJobRows(h.inputs, () => true); assert.deepEqual(rows.map((r) => r.id), ["run1"]); assert.deepEqual(h.completed, []); assert.deepEqual( h.inputs.scheduler.queues().flatMap((v) => v.running), ["done1"], ); // No lanes, no disk, no meta sidecar and no log tail on this path. assert.deepEqual(h.metaCalls, []); assert.deepEqual(h.tailCalls, []); }); test("liveJobRows applies its filter to running and queued records only", async () => { const h = harness({ jobs: [ job({ id: "run1", channelSlug: "alpha" }), job({ id: "queued1", status: "queued", channelSlug: "alpha" }), job({ id: "run2", channelSlug: "beta", queueKey: "q2" }), job({ id: "done1", status: "done", channelSlug: "alpha", endedAt: NOW }), ], }); const rows = await liveJobRows(h.inputs, (j) => j.channelSlug === "alpha"); assert.deepEqual(new Set(rows.map((r) => r.id)), new Set(["run1", "queued1"])); }); test("a stopped runner reads unavailable, not idle", async () => { const h = harness({ runner: (kind) => kind === "digest" ? runnerStatus({ kind, running: false, idleReason: "stopped" }) : runnerStatus({ kind, idleReason: "no-pending" }), }); const payload = await buildActiveJobsPayload(h.inputs); const byKind = new Map(payload.lanes.map((l) => [l.label, l])); const digest = byKind.get("Digest"); assert.ok(digest); // A stopped runner will never pick anything up; "idle" would read as caught up. assert.equal(digest.state, "unavailable"); assert.equal(digest.note, "not running"); const transcribe = byKind.get("Auto-transcribe"); assert.equal(transcribe?.state, "idle"); assert.equal(transcribe?.note, "nothing pending"); assert.equal(payload.lanes.length, 4); }); test("a download runner idling on deferred videos says so on its lane", async () => { const h = harness({ runner: (kind) => kind === "download" ? runnerStatus({ kind, idleReason: "deferred" }) : runnerStatus({ kind, idleReason: "no-pending" }), }); const payload = await buildActiveJobsPayload(h.inputs); const download = payload.lanes.find((l) => l.label === "Auto-download"); assert.equal( download?.note, "every pending video was rate-limited recently and is deferred", ); }); test("the disk block is the gate the shell sampled, with `low` derived", async () => { const h = harness({ disk: { ...OK_DISK, ok: false, freeBytes: 1, reason: "below-floor", message: "full" }, }); const payload = await buildActiveJobsPayload(h.inputs); assert.equal(payload.disk.low, true); assert.equal(payload.disk.reason, "below-floor"); assert.equal(payload.disk.message, "full"); // Disabled gate: never low, whatever `ok` says. const off = harness({ disk: { ...OK_DISK, enabled: false, ok: false } }); assert.equal((await buildActiveJobsPayload(off.inputs)).disk.low, false); }); test("rowsFromJobsPage folds a listing page and nothing else", () => { const out = rowsFromJobsPage({ entries: [ { id: "j9", kind: "download", queueKey: "q", status: "done", queuedAt: NOW - 10, endedAt: NOW, inRegistry: false, replayable: false, logPath: "/dev/null", logSize: 0, }, ], hasMore: true, total: 42, }); assert.equal(out.total, 42); assert.equal(out.hasMore, true); assert.deepEqual(out.rows.map((r) => r.id), ["j9"]); });