import { beforeEach, test } from "node:test"; import assert from "node:assert/strict"; import { mkdir, mkdtemp, 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 { SiteSettings } from "../lib/settings"; import { autoPauseReasonOf, compileLanes, sanitizeChannelPriority, } from "../lib/channelPriority"; import { LANES } from "../lib/autoQueueTypes"; import { refreshLocationHealth, resetStorageWatchSuspicion, runStorageHealthPass, runStorageWatchPass, startStorageHealthWatch, startStorageWatch, stopStorageHealthWatch, stopStorageWatch, } from "./storageWatch"; import { inspectChannelMedia } from "../lib/channelMedia"; import { applyHealthTimings, healthTimings, locationHealth, resetStorageHealth, type LocationHealthState, } from "../lib/storageHealth"; // THE CONFIRMATION COUNT IS MODULE STATE (see storageWatch.ts rule 3), so each // case starts from a clean one — otherwise the second test inherits the first // test's suspicions and pauses on what should be its first pass. The health // state (lib/storageHealth.ts) is process state for the same reason. beforeEach(() => { resetStorageWatchSuspicion(); resetStorageHealth(); applyHealthTimings(); }); // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/storageWatch.test.ts // // NO SUBPROCESS. `bins` is passed as a Paths-shaped object whose findmnt points // at a binary that does not exist, so `probeLocation` fails open to "identity // unknown" — which is the container case and exactly what this pass must keep // working under. What decides here is the ROOT's existence and the channel's // own link, which is what the pass actually reads. // Lane policies whose roots the compiler already wrote. See the seed comment in // storageWatch.ts: a corpus whose trees were NEVER compiled keeps its hand-made // order, and the watch must not be the thing that switches it over. const COMPILED_LANES = (() => { const roots = compileLanes(sanitizeChannelPriority({ channels: {} }), [], []); return Object.fromEntries( LANES.map((lane) => [lane, { enabled: false, root: roots[lane] }]), ); })(); type H = { paths: Paths; root: string; io: { read: () => SiteSettings; write: (n: SiteSettings) => Promise }; writes: number; }; async function withTmp(fn: (h: H) => Promise): Promise { const dir = await mkdtemp(path.join(tmpdir(), "ttb-storagewatch-")); const transcriptsDir = path.join(dir, "corpus"); const paths = { transcriptsDir, channelsDir: path.join(transcriptsDir, "channels"), findmntBin: path.join(dir, "no-such-findmnt"), udisksctlBin: path.join(dir, "no-such-udisksctl"), } as Paths; const root = path.join(dir, "platter"); await mkdir(paths.channelsDir, { recursive: true }); await mkdir(root, { recursive: true }); const h: H = { paths, root, writes: 0, io: { read: () => settings, write: async (n) => { settings = n; h.writes += 1; }, }, }; let settings = { storage: { locations: [{ id: "cold", label: "Cold", root, autoRepoint: false }], defaultLocationId: "cold", }, channelPriority: sanitizeChannelPriority({ channels: {} }), // The compiled-lane flag the writer's seed condition reads. Compiled here, // so the seed path is not taken and the test is about the watch and not // about the migration — `hasCompiledLaneRoots` is exercised by // channelPriorityCompile.test.ts. autoQueue: COMPILED_LANES, } as unknown as SiteSettings; try { await fn(h); } finally { await rm(dir, { recursive: true, force: true }); } } // A channel whose media tier is relocated (release 17): `channels//media` // a link to `//media`, recorded as `mediaDir`, with the target // present or not. Its `data/` (the text) is on the corpus disk throughout. async function seedRelocated( h: H, slug: string, opts: { targetExists: boolean }, ): Promise { const channelDir = path.join(h.paths.channelsDir, slug); await mkdir(channelDir, { recursive: true }); const target = path.join(h.root, slug, "media"); if (opts.targetExists) await mkdir(target, { recursive: true }); await writeFile( path.join(channelDir, "config.json"), JSON.stringify({ handling: "youtube", mediaDir: target }), ); await symlink(target, path.join(channelDir, "media")); } function tierOf(h: H, slug: string): string | undefined { return h.io.read().channelPriority.channels[slug]?.tier; } // ONE BAD READ IS A SUSPICION, TWO IN A ROW IS A FACT — availability is a bare // stat with a blanket catch, so an EIO or a spun-down disk reads exactly like // "not mounted". Most cases here are about what happens once a drive really is // gone, so they run the confirming pair and assert on the second. async function twoPasses(h: H) { const first = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths, }); const second = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths, }); return { first, second }; } test("a channel whose target is gone is auto-paused, once, in one write", async () => { await withTmp(async (h) => { await seedRelocated(h, "gone-a", { targetExists: false }); await seedRelocated(h, "gone-b", { targetExists: false }); await seedRelocated(h, "fine", { targetExists: true }); // RUN THE PASSES ONE AT A TIME HERE, not through twoPasses(): the write // count between them is exactly what this case is about. const pass = () => runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths }); const first = await pass(); // The first pass only suspects — and writes NOTHING, which is the point: // one flaky stat must not rewrite the corpus's priority document. assert.deepEqual(first.suspected.sort(), ["gone-a", "gone-b"]); assert.deepEqual(first.paused, []); assert.equal(first.wrote, false); assert.equal(h.writes, 0); // The second confirms. const second = await pass(); assert.deepEqual(second.paused.sort(), ["gone-a", "gone-b"]); assert.deepEqual(second.restored, []); assert.equal(second.wrote, true); // TWO CHANNELS, ONE WRITE. Ten on a drive that vanished must be one pulse // bump, not ten. assert.equal(h.writes, 1); assert.equal(tierOf(h, "gone-a"), "paused"); assert.equal(tierOf(h, "fine"), undefined); assert.equal( h.io.read().channelPriority.channels["gone-a"].autoPaused?.reason, "storage", ); // A QUIET PASS WRITES NOTHING. The document already describes the world. const third = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths }); assert.deepEqual(third.paused, []); assert.equal(third.wrote, false); assert.equal(h.writes, 1); }); }); test("the drive coming back restores the tier it overwrote", async () => { await withTmp(async (h) => { await seedRelocated(h, "away", { targetExists: false }); const settings = h.io.read(); await h.io.write({ ...settings, channelPriority: sanitizeChannelPriority({ channels: { away: { tier: "low", rank: 4 } }, }), }); h.writes = 0; await twoPasses(h); assert.equal(tierOf(h, "away"), "paused"); await mkdir(path.join(h.root, "away", "media"), { recursive: true }); const back = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths }); assert.deepEqual(back.restored, ["away"]); assert.equal(tierOf(h, "away"), "low"); assert.equal(h.io.read().channelPriority.channels.away.rank, 4); assert.equal(h.io.read().channelPriority.channels.away.autoPaused, undefined); }); }); // THE OPERATOR'S WORD WINS. A channel they paused is not the machine's to // claim, and one they resumed while the drive was still away must not be // re-paused... by a restore. (Re-pausing it on a LATER pass is correct: the // drive is still gone. What must not happen is the restore un-pausing a // deliberate pause, which is what the no-record no-op guarantees.) test("a manually paused channel is never claimed by the watch", async () => { await withTmp(async (h) => { await seedRelocated(h, "off", { targetExists: false }); const settings = h.io.read(); await h.io.write({ ...settings, channelPriority: sanitizeChannelPriority({ channels: { off: { tier: "paused" } }, }), }); h.writes = 0; const { second } = await twoPasses(h); assert.deepEqual(second.paused, []); assert.equal(second.wrote, false); assert.equal(h.writes, 0); assert.equal( h.io.read().channelPriority.channels.off.autoPaused, undefined, ); }); }); // A MARKER MEANS A MOVE IS RUNNING OR WAS INTERRUPTED, and the relocate job is // precisely the thing an auto-pause would then be refusing. test("a channel mid-relocation is not auto-paused", async () => { await withTmp(async (h) => { await seedRelocated(h, "moving", { targetExists: false }); await writeFile( path.join(h.paths.channelsDir, "moving", ".relocating.json"), JSON.stringify({ target: path.join(h.root, "moving", "media"), direction: "out", startedAt: "", phase: "copy", }), ); const { second } = await twoPasses(h); assert.deepEqual(second.paused, []); assert.equal(second.wrote, false); }); }); // A CHANNEL MOVED BACK IN PLACE CARRIES NO DRIVE TO BE AWAY, but it can carry a // record from before — which has to come off or it stays paused for ever. test("a record on an in-place channel is restored", async () => { await withTmp(async (h) => { const channelDir = path.join(h.paths.channelsDir, "home"); await mkdir(path.join(channelDir, "data"), { recursive: true }); await writeFile( path.join(channelDir, "config.json"), JSON.stringify({ handling: "youtube" }), ); const settings = h.io.read(); await h.io.write({ ...settings, channelPriority: sanitizeChannelPriority({ channels: { home: { tier: "paused", autoPaused: { reason: "storage", since: "2026-09-20T00:00:00.000Z", previousTier: "normal", }, }, }, }), }); h.writes = 0; const r = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths }); assert.deepEqual(r.restored, ["home"]); assert.equal(h.io.read().channelPriority.channels.home, undefined); }); }); // `write: false` IS IDLE BOOT. It observes and reports; the write is the work, // and idle boot refuses work. test("write: false reports the transition and changes nothing", async () => { await withTmp(async (h) => { await seedRelocated(h, "gone", { targetExists: false }); // The first pass only suspects, whatever `write` says. const opts = { paths: h.paths, io: h.io, bins: h.paths, write: false }; assert.deepEqual((await runStorageWatchPass(opts)).suspected, ["gone"]); const r = await runStorageWatchPass(opts); assert.deepEqual(r.paused, ["gone"]); assert.equal(r.wrote, false); assert.equal(h.writes, 0); assert.equal(tierOf(h, "gone"), undefined); }); }); test("no locations and nothing auto-paused is a free pass", async () => { await withTmp(async (h) => { const settings = h.io.read(); await h.io.write({ ...settings, storage: { locations: [], defaultLocationId: "" }, }); h.writes = 0; const r = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths }); assert.deepEqual(r, { probed: 0, suspected: [], repointed: [], paused: [], restored: [], wrote: false, }); }); }); // ONE BAD READ MUST NOT PAUSE A TIER. Availability is a bare `stat` with a // blanket catch (storageVolumes.ts), so an EIO on a flaky cable or a disk that // has spun down and needs a beat to answer is indistinguishable from "not // mounted" — and pausing on it rewrites the corpus's priority document for a // drive that is fine. test("a drive that blips for one pass is never paused", async () => { await withTmp(async (h) => { await seedRelocated(h, "blip", { targetExists: false }); const first = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths, }); assert.deepEqual(first.suspected, ["blip"]); assert.deepEqual(first.paused, []); assert.equal(h.writes, 0); // It answers on the next pass. Nothing was ever paused, and the suspicion // is dropped — so a LATER real outage starts its own two-pass count rather // than pausing immediately on the strength of a blip an hour ago. await mkdir(path.join(h.root, "blip", "media"), { recursive: true }); const second = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths, }); assert.deepEqual(second.paused, []); assert.deepEqual(second.restored, []); assert.equal(second.wrote, false); assert.equal(h.writes, 0); // Prove the suspicion really was dropped: the drive going away again takes // two fresh passes. await rm(path.join(h.root, "blip"), { recursive: true, force: true }); assert.deepEqual( (await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths })) .paused, [], ); assert.deepEqual( (await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths })) .paused, ["blip"], ); }); }); // RESTORE STAYS SINGLE-PASS, and the asymmetry is the point: being slow to // pause costs a few refused units (the start-of-work guards catch those), while // being slow to restore leaves a lane off after the operator fixed the cable. test("the restore needs only one good pass", async () => { await withTmp(async (h) => { await seedRelocated(h, "back", { targetExists: false }); await twoPasses(h); assert.equal(tierOf(h, "back"), "paused"); h.writes = 0; await mkdir(path.join(h.root, "back", "media"), { recursive: true }); const r = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths, }); assert.deepEqual(r.restored, ["back"]); assert.equal(h.writes, 1); }); }); // --------------------------------------------------------------------------- // The health pass (15 s): is the drive ANSWERING // --------------------------------------------------------------------------- // // The probe is injected: `probeLocationHealth` (a child `stat` raced against a // 3 s timer) has its own tests in lib/storageHealthProbe.test.ts. Here the // answers are scripted, one per pass. function scripted(answers: LocationHealthState[]) { let i = 0; return async () => answers[Math.min(i++, answers.length - 1)]; } test("one missed probe stalls the location; pages then answer 'stalled' without asking", async () => { await withTmp(async (h) => { await seedRelocated(h, "slow", { targetExists: true }); const lines: string[] = []; const r = await runStorageHealthPass({ io: h.io, probe: scripted(["stalled"]), log: (l) => lines.push(l), }); assert.deepEqual(r.answers, { cold: "stalled" }); // Registered first (answering), so the miss is a transition from ok. assert.deepEqual(r.transitions, [{ id: "cold", from: "ok", to: "stalled" }]); assert.equal(locationHealth("cold")?.state, "stalled"); assert.match(lines.join("\n"), /"cold": drive not answering/); // The target is there and would answer, but nothing asks it. const media = await inspectChannelMedia(h.paths, "slow"); assert.equal(media.status, "stalled"); // The five-minute pass sees the location as down (its probe answers // "stalled" without a stat) and suspects the channel, as for any outage. const w = await runStorageWatchPass({ paths: h.paths, io: h.io, bins: h.paths }); assert.deepEqual(w.suspected, ["slow"]); assert.equal(h.writes, 0); }); }); test("the stall clears only after two clean probes in a row", async () => { await withTmp(async (h) => { await seedRelocated(h, "slow", { targetExists: true }); const probe = scripted(["stalled", "ok", "stalled", "ok", "ok"]); const lines: string[] = []; const pass = () => runStorageHealthPass({ io: h.io, probe, log: (l) => lines.push(l) }); await pass(); await pass(); // one clean answer assert.equal(locationHealth("cold")?.state, "stalled"); await pass(); // missed again: the count starts over await pass(); // one clean assert.equal(locationHealth("cold")?.state, "stalled"); assert.equal((await inspectChannelMedia(h.paths, "slow")).status, "stalled"); const last = await pass(); // two clean in a row assert.deepEqual(last.transitions, [{ id: "cold", from: "stalled", to: "ok" }]); assert.match(lines.at(-1) ?? "", /"cold": answering again/); assert.equal((await inspectChannelMedia(h.paths, "slow")).status, "ok"); }); }); test("a location no longer configured is forgotten", async () => { await withTmp(async (h) => { await runStorageHealthPass({ io: h.io, probe: scripted(["stalled"]) }); assert.equal(locationHealth("cold")?.state, "stalled"); await runStorageHealthPass({ locations: [], probe: scripted(["ok"]) }); assert.equal(locationHealth("cold"), undefined); }); }); test("a probe that throws is 'could not ask': ok, never stalled", async () => { await withTmp(async (h) => { const r = await runStorageHealthPass({ io: h.io, probe: async () => { throw new Error("spawn failed"); }, }); assert.deepEqual(r.answers, { cold: "ok" }); }); }); test("the health pass is armed on its own, runs once at once, and stops; the five-minute watch arms no health pass", async () => { await withTmp(async (h) => { let asked = 0; const armed = startStorageHealthWatch({ io: h.io, probe: async () => { asked += 1; return "stalled"; }, log: () => {}, }); try { assert.equal(armed, true); // Armed once per process. assert.equal(startStorageHealthWatch({ io: h.io }), false); for (let i = 0; i < 50 && locationHealth("cold")?.state !== "stalled"; i++) { await new Promise((r) => setTimeout(r, 10)); } assert.equal(asked, 1); assert.equal(locationHealth("cold")?.state, "stalled"); } finally { stopStorageHealthWatch(); } assert.equal(startStorageHealthWatch({ io: h.io, probe: async () => "ok" }), true); stopStorageHealthWatch(); // The watch (below the idle gate) runs nothing at arm time and asks no drive. resetStorageHealth(); assert.equal(startStorageWatch({ paths: h.paths, io: h.io, bins: h.paths, write: false }), true); stopStorageWatch(); assert.equal(locationHealth("cold"), undefined); }); }); test("a Refresh asks one location now: it counts as one answer, and prunes nothing", async () => { await withTmp(async (h) => { const other = { id: "other", label: "Other", root: "/elsewhere", autoRepoint: false }; await runStorageHealthPass({ locations: [h.io.read().storage.locations[0], other], probe: scripted(["stalled"]), }); const cold = h.io.read().storage.locations[0]; assert.equal(await refreshLocationHealth(cold, async () => "ok"), "ok"); // One clean answer is not two. assert.equal(locationHealth("cold")?.state, "stalled"); assert.equal(locationHealth("other")?.state, "stalled"); await refreshLocationHealth(cold, async () => "ok"); assert.equal(locationHealth("cold")?.state, "ok"); assert.equal(locationHealth("other")?.state, "stalled"); }); }); test("the pass registers every location, records a verdict's detector, and a verdict with no answer changes nothing", async () => { await withTmp(async (h) => { const verdicts = [ { answer: null, detector: "counters" as const, device: "sdz1" }, { answer: "stalled" as const, detector: "counters" as const, device: "sdz1", cause: "its disk (sdz1) had 1 request(s) in flight and completed none in 15 s", }, ]; let i = 0; const probe = async () => verdicts[i++]; const lines: string[] = []; const first = await runStorageHealthPass({ io: h.io, probe, log: (l) => lines.push(l) }); // Registered, answering, and the detector named — with no verdict yet. assert.deepEqual(first.answers, {}); assert.deepEqual(first.transitions, []); assert.equal(locationHealth("cold")?.state, "ok"); assert.equal(locationHealth("cold")?.detector, "counters"); const second = await runStorageHealthPass({ io: h.io, probe, log: (l) => lines.push(l) }); assert.deepEqual(second.transitions, [{ id: "cold", from: "ok", to: "stalled" }]); assert.match(String(locationHealth("cold")?.cause), /its disk \(sdz1\)/); assert.match(lines.join("\n"), /"cold": drive not answering — its disk \(sdz1\)/); }); }); test("a counters verdict records its device (for the watchdog); a stat verdict forgets it", async () => { await withTmp(async (h) => { const verdicts = [ { answer: null, detector: "counters" as const, device: "sdz1" }, { answer: "ok" as const, detector: "stat" as const }, ]; let i = 0; const probe = async () => verdicts[i++]; await runStorageHealthPass({ io: h.io, probe }); assert.equal(locationHealth("cold")?.device, "sdz1"); await runStorageHealthPass({ io: h.io, probe }); assert.equal(locationHealth("cold")?.device, undefined); assert.equal(locationHealth("cold")?.detector, "stat"); }); }); test("a stall auto-pauses after two passes, and says the drive is not answering (not that it is not there)", async () => { await withTmp(async (h) => { await seedRelocated(h, "slow", { targetExists: true }); await runStorageHealthPass({ io: h.io, probe: async () => "stalled" }); await twoPasses(h); const entry = h.io.read().channelPriority.channels.slow; assert.equal(entry?.tier, "paused"); assert.equal(entry?.autoPaused?.cause, "not-answering"); const reason = autoPauseReasonOf(h.io.read().channelPriority, "slow"); assert.match(String(reason), /drive that is not answering/); assert.match(String(reason), /when the drive answers again/); // The sanitizer keeps the cause; a record without one reads as not there. const kept = sanitizeChannelPriority(h.io.read().channelPriority); assert.equal(kept.channels.slow?.autoPaused?.cause, "not-answering"); const old = sanitizeChannelPriority({ channels: { a: { tier: "paused", autoPaused: { reason: "storage", since: "", previousTier: "low" } }, }, }); assert.match(String(autoPauseReasonOf(old, "a")), /drive that is not there/); }); }); // ── the timings are settings (release 15 slice DT) ───────────────────────── test("DT: every pass applies the timings it reads — the clear count and the log line follow storage.health", async () => { await withTmp(async (h) => { await seedRelocated(h, "slow", { targetExists: true }); h.io.read().storage.health = { clearAfterCleanPasses: 3, budgetMs: 5_000 }; const probe = scripted(["stalled", "ok", "ok", "ok"]); const lines: string[] = []; const pass = () => runStorageHealthPass({ io: h.io, probe, log: (l) => lines.push(l) }); await pass(); assert.equal(healthTimings().budgetMs, 5_000, "applied before anything was asked"); assert.match(lines.join("\n"), /until it answers 3 times in a row/); await pass(); await pass(); assert.equal(locationHealth("cold")?.state, "stalled", "two clean passes are not three"); await pass(); assert.equal(locationHealth("cold")?.state, "ok"); // A pass handed its locations reads no settings: the timings stay. await runStorageHealthPass({ locations: h.io.read().storage.locations, probe: scripted(["ok"]) }); assert.equal(healthTimings().clearAfterCleanPasses, 3); }); }); test("DT: a changed pass interval re-arms the armed pass; an explicit interval follows nothing", async () => { await withTmp(async (h) => { let asked = 0; const lines: string[] = []; startStorageHealthWatch({ io: h.io, probe: async () => { asked += 1; return "ok"; }, log: (l) => lines.push(l), }); try { assert.equal(asked, 1, "one pass at once"); // The save on /storage: written to settings, then applied at once. h.io.read().storage.health = { passIntervalMs: 5_000 }; applyHealthTimings(h.io.read().storage.health); assert.match(lines.join("\n"), /health pass re-armed: every 5 s/); // At the default 15 s nothing would run for another 15 s; re-armed at // 5 s, the next pass comes within about 5 s. const started = Date.now(); while (asked < 2 && Date.now() - started < 7_000) { await new Promise((r) => setTimeout(r, 50)); } assert.equal(asked, 2, `a second pass after ${Date.now() - started} ms`); assert.ok(Date.now() - started >= 4_500); } finally { stopStorageHealthWatch(); } // Stopped: a later change re-arms nothing. const before = lines.length; applyHealthTimings({ passIntervalMs: 20_000 }); assert.equal(lines.length, before); // Armed with an explicit interval, a change is not followed. startStorageHealthWatch({ io: h.io, probe: async () => "ok", intervalMs: 60_000, log: (l) => lines.push(l) }); try { applyHealthTimings({ passIntervalMs: 6_000 }); assert.equal(lines.some((l) => /re-armed/.test(l) && /6 s/.test(l)), false); } finally { stopStorageHealthWatch(); } }); });