import { getPaths, type Paths } from "../lib/paths"; import { getSettings, writeSettings, type SiteSettings, } from "../lib/settings"; import { autoPausedSlugs, autoPauseForMedia, channelPriorityFromLegacy, compileLanes, hasCompiledLaneRoots, isDefaultChannelPriority, resolveFocusSlugs, restoreAfterMedia, sanitizeChannelPriority, type AutoPauseCause, type ChannelPriority, } from "../lib/channelPriority"; import { LANES } from "../lib/autoQueueTypes"; import { siteChannelIndex } from "../lib/site"; import { inspectChannelMedia } from "../lib/channelMedia"; import { locationOfDataDir, type StorageLocation, } from "../lib/storageLocations"; import { detectLocationHealth, type HealthVerdict, type LocationHealthProbe, type VolumeBins, } from "../lib/storageVolumes"; import { NOT_ANSWERING, applyHealthTimings, healthTimings, noteLocationDetector, onPassIntervalChange, pruneLocationHealth, recordLocationHealth, registerLocationHealth, type HealthTransition, type LocationHealthState, } from "../lib/storageHealth"; import { clearRuleText, secondsText } from "../lib/storageHealthTimings"; import { listChannelConfigs } from "./channels"; import { maybeAutoRepoint, probeAllLocations } from "./storageLocations"; // THE DRIVE WENT AWAY WHILE THE CHANNEL WAS ON — now what. // // Operator ask, 2026-09-18: "Can that gracefully handle the case if the drive // disconnects while the channel is on? Maybe an automatic disable and flag?" // // Every guard the corpus has for an unreachable channel runs at the START of a // piece of work (`runManagedFunction`'s first statement, `buildChannelWork`'s // per-tick inspect, the snapshot regen, the operation batch). None of them is a // detector: they refuse work that is already being asked for, which means a // channel on a drive that vanished sits there being refused, over and over, // with nothing anywhere saying why. The flag the operator asked for is a state, // and a state needs something that looks. // // SO THIS LOOKS, ON A CADENCE, AND WRITES AT MOST ONCE PER PASS. // // Two rules, and they are what keep it from being a settings-churn machine: // // 1. IT ONLY EVER WRITES ON A TRANSITION. A pass that finds the world exactly // as the document already describes it writes nothing — no settings write, // so no pulse revision bump, so no reader re-polls. A flapping drive costs // two writes per flap, not one per tick. // 2. IT RESTORES ONLY WHAT IT PAUSED. `restoreAfterMedia` is a no-op without // an `autoPaused` record, and the one priority writer clears that record on // any manual tier change — so a drive coming back can never un-pause a // channel the operator paused on purpose in the meantime. // 3. IT TAKES TWO CONSECUTIVE DOWN PASSES TO PAUSE, AND ONE UP PASS TO // RESTORE. Availability is a bare `stat` with a blanket catch // (`storageVolumes.ts`) — an EIO on a flaky cable, or a disk that has spun // down and needs a beat to answer, reads exactly like "not mounted". One // such read would otherwise pause every channel on the drive and rewrite // the corpus's priority document. The confirmation is held IN MEMORY, not // in settings: a pending suspicion is not a fact worth persisting, and a // process restart starting the count again is the safe direction. The // asymmetry is deliberate — being slow to pause costs a few refused units // (the start-of-work guards catch those), while being slow to RESTORE costs // the operator a lane that stays off after they fixed the cable. // // IT IS A RUNNER, so `ARCHILYZER_IDLE_BOOT` must not arm it: a container // pointed at somebody else's corpus for the first time has no business // rewriting that corpus's priority document seconds after `docker compose up`. // The probe is read-only and cheap; the WRITE is the work, and idle boot // refuses work. `runStorageWatchPass({ write: false })` is the observation // without the consequence, which is what the boot pass and a test want. export type StorageWatchResult = { // Locations probed this pass. probed: number; // Channels this pass saw as unreachable for the FIRST time. They are not // paused yet; the next pass decides. Reported so a caller (and the test) can // see the confirmation working rather than infer it from silence. suspected: string[]; // Locations whose volume was found at a different mountpoint and for which an // auto re-point was queued. Empty unless a location has `autoRepoint` on. repointed: string[]; // Channels newly auto-paused, and channels restored. Both empty on a quiet // pass, which is the overwhelming majority. paused: string[]; restored: string[]; // Whether settings were written. False whenever both lists are empty. wrote: boolean; }; export type StorageWatchOpts = { paths?: Paths; bins?: VolumeBins; io?: { read: () => SiteSettings; write: (next: SiteSettings) => Promise }; // False observes and reports without touching settings — idle boot, and the // unit tests. write?: boolean; log?: (line: string) => void; }; const DEFAULT_IO = { read: getSettings, write: writeSettings }; // Channels seen down on the LAST pass and not yet paused. Module state, and // deliberately not settings: see rule 3 in the header. const suspected = new Set(); // Test seam, and the escape hatch for a process that wants a clean count. export function resetStorageWatchSuspicion(): void { suspected.clear(); } export async function runStorageWatchPass( opts: StorageWatchOpts = {}, ): Promise { const io = opts.io ?? DEFAULT_IO; const log = opts.log ?? (() => {}); const paths = opts.paths ?? getPaths(); const settings = io.read(); const locations = settings.storage.locations; const out: StorageWatchResult = { probed: 0, suspected: [], repointed: [], paused: [], restored: [], wrote: false, }; // A channel can only be auto-paused for a location's sake, so with no // locations configured there is nothing to watch — EXCEPT the records a // previous configuration left behind, which must still be restorable. Hence // the early return is on "no locations AND nothing auto-paused". const already = autoPausedSlugs(settings.channelPriority); if (locations.length === 0 && already.length === 0) return out; // REFRESH, NOT THE MEMO. The memo exists so a page render does not fork // eighteen subprocesses per click; this pass runs on a cadence measured in // minutes and its whole job is to notice a change, so an answer taken up to // ten seconds ago is not what it is asking for. const probes = await probeAllLocations(locations, opts.bins ?? paths, { refresh: true, }); out.probed = Object.keys(probes).length; // THE DRIVE CAME UP SOMEWHERE ELSE — FOLLOW IT, IF THE OPERATOR ARMED THAT. // // `mounted-elsewhere` is the one status with a remedy that moves no bytes: // the volume IS here, under a different mountpoint, and a re-point rewrites // the links. Until now only the BOOT pass took it, so a disk that came back // at a new mountpoint while the editor was up sat there while this pass // dutifully paused every channel on it — and the fix was a restart. Same // opt-in (`autoRepoint`), same preflight, same refusal-with-a-reason; what // changes is that the cadence can reach it. // // It runs BEFORE the per-channel loop and the enqueued job runs after this // pass returns, so this pass still sees (and may still suspect) the channels // on that location — which is correct: nothing has moved yet, and the // two-pass confirmation gives the re-point a whole interval to land before // anything is paused. if (opts.write !== false) { for (const loc of locations) { const probe = probes[loc.id]; if (!probe || probe.status !== "mounted-elsewhere" || !loc.autoRepoint) { continue; } const outcome = await maybeAutoRepoint({ paths, location: loc, probe, bins: opts.bins ?? paths, io, }).catch((err) => ({ started: false as const, reason: (err as Error).message, })); if (outcome.started) { out.repointed.push(loc.id); log( `[storage] "${loc.id}": the volume came up at a different mountpoint ` + `— auto re-point queued (${outcome.newRoot})`, ); } else { log(`[storage] "${loc.id}": auto re-point declined — ${outcome.reason}`); } } } const configs = await listChannelConfigs(paths); let model: ChannelPriority = settings.channelPriority; const pauseCauses = new Map(); for (const { slug, config } of configs) { const wasAutoPaused = Boolean(model.channels[slug]?.autoPaused); // ONE TIER PER CHANNEL (release 17): the drive its MEDIA is on — // `mediaDir`, or on a channel not yet migrated off the retired // whole-directory layout its `dataDir`. The pause stays the channel's (one // `autoPaused` record), not a per-lane one; its text never leaves the // corpus disk, so there is no second drive to watch. const dataDir = config.mediaDir?.trim() || config.dataDir?.trim(); if (!dataDir) { // In place. It cannot be on a drive that went away — but it CAN carry a // record from before it was moved back, and that record has to come off // or the channel stays paused for ever. suspected.delete(slug); if (wasAutoPaused) { model = restoreAfterMedia(model, slug); out.restored.push(slug); } continue; } // THE LOCATION'S PROBE FIRST, THE CHANNEL'S OWN STAT SECOND. The probe is // the cheap corpus-wide answer (one findmnt per location, not per channel); // `inspectChannelMedia` is what decides, because a location can be // available while one channel's target under it is missing — a half-done // move, a directory deleted by hand. const loc = locationOfDataDir(dataDir, locations); const probe = loc ? probes[loc.id] : undefined; const locationDown = Boolean(loc) && probe?.status !== "available"; // FRESH: this pass is the detector, and the page memo is not what it asks. // A location the health probe found not answering reads `stalled` here // without a call, and a stall is `down` like any other. const media = await inspectChannelMedia(paths, slug, config, { fresh: true, }); // `in-transition` is NEVER a reason to pause: a marker means a move is // running or was interrupted, and the relocate job is precisely the thing // that would then be refused by the state it created. // A stall is down too, on a location or not (a root typed by hand gets its // `stalled` from the watchdog alone). const down = media.status === "in-transition" ? false : locationDown || media.status === "unreachable" || media.status === "stalled"; // WHICH down it is, for the pause record's words (autoPauseReasonOf). const cause: AutoPauseCause = media.status === "stalled" || probe?.status === "stalled" ? "not-answering" : "not-there"; if (down && !wasAutoPaused) { // ONE BAD READ IS A SUSPICION, TWO IN A ROW IS A FACT. See rule 3. if (!suspected.has(slug)) { suspected.add(slug); out.suspected.push(slug); log( `[storage] ${slug}: media unreachable (${ loc ? `location "${loc.id}" is ${probe?.status ?? "unprobed"}` : media.status }) — waiting for a second pass to confirm before pausing`, ); continue; } const before = model; model = autoPauseForMedia(model, slug, new Date(), cause); pauseCauses.set(slug, cause); // autoPauseForMedia no-ops on a channel the OPERATOR already paused — // which is right, and means "nothing changed" is a normal outcome here. if (model !== before) { out.paused.push(slug); log( `[storage] ${slug}: media unreachable (${ loc ? `location "${loc.id}" is ${probe?.status ?? "unprobed"}` : media.status }) — auto-paused`, ); } continue; } if (!down) suspected.delete(slug); if (!down && wasAutoPaused) { const restoredTo = model.channels[slug]?.autoPaused?.previousTier; model = restoreAfterMedia(model, slug); out.restored.push(slug); log(`[storage] ${slug}: media reachable again, tier restored to ${restoredTo}`); } } if (out.paused.length === 0 && out.restored.length === 0) return out; if (opts.write === false) return out; // ONE WRITE PER PASS, whatever the pass found. Ten channels on one drive that // vanished is one settings write and one pulse bump, not ten. // // FROM A FRESH READ, because a probe of several locations is seconds of wall // clock during which an operator may have changed a tier in another tab — // and clobbering that with a snapshot taken before the pass started would // silently undo it. The edits are re-applied to whatever is current. const latest = io.read(); // THE LEGACY SEED, AND WHY A BACKGROUND PASS MUST PAY IT TOO. // // `laneDispatchRoot` is all-or-nothing on `isDefaultChannelPriority`: the // moment the document says ANYTHING, the stored lane trees stop being // dispatched from and compiled ones take over. So a pass that auto-paused one // channel on a corpus whose document was empty would, as a side effect, // replace the operator's hand-made lane order with an all-normal alphabetical // one — silently, at 3am, because a USB cable came loose. That is the same // trap `saveChannelPriorityAction` documents (the S2/S3 review, finding 3), // and it is worse here because nobody clicked anything. // // Same remedy, same condition: when the stored document says nothing AND the // stored trees were never compiled, derive the document the migration would // have produced FIRST and apply the pause on top of that. The corpus's // existing order survives. const stored = latest.channelPriority; const base = isDefaultChannelPriority(stored) && !hasCompiledLaneRoots(latest.autoQueue) ? channelPriorityFromLegacy( configs.map((c) => ({ slug: c.slug, config: {} })), latest.autoQueue, ) : stored; let merged: ChannelPriority = base; for (const slug of out.paused) { merged = autoPauseForMedia(merged, slug, new Date(), pauseCauses.get(slug)); } for (const slug of out.restored) merged = restoreAfterMedia(merged, slug); merged = sanitizeChannelPriority(merged); // AND THE TREES, in the same write. Two writes to one settings file race each // other, and a document that has changed tiers with trees that have not is a // corpus dispatching from an order nobody holds any more. The condition is // the writer's: a still-default document with never-compiled trees keeps its // hand-made ones. const autoQueue = { ...latest.autoQueue }; if (!isDefaultChannelPriority(merged) || hasCompiledLaneRoots(latest.autoQueue)) { const slugs = configs.map((c) => c.slug); const focusSlugs = resolveFocusSlugs(merged, siteChannelIndex(paths), slugs); const roots = compileLanes(merged, slugs, focusSlugs); for (const lane of LANES) { autoQueue[lane] = { ...latest.autoQueue[lane], root: roots[lane] }; } } await io.write({ ...latest, channelPriority: merged, autoQueue }); out.wrote = true; return out; } // Five minutes. A drive does not come and go on a timescale a person would // notice faster than that, and every pass is one findmnt per location plus two // stats per relocated channel — cheap, but not free, and this runs for the life // of the process. export const STORAGE_WATCH_INTERVAL_MS = 5 * 60_000; // --------------------------------------------------------------------------- // The health pass: is each location's drive ANSWERING // --------------------------------------------------------------------------- // // A SECOND CADENCE, AND A MUCH SHORTER ONE. The pass above asks "is the disk // here" every five minutes and pauses on two misses, which is right for a // cable pulled out. It is no help for a drive that is here and stalled — an // SMR disk in a USB enclosure resetting under a long write — because every // in-process call on that drive waits for it, and four waits stop the editor // answering at all. So every `storage.health.passIntervalMs` (15 s by default) // this reads each location's block device counters in /sys // (`detectLocationHealth`; a child `stat` of the root only where no device can // be named) and records the answer in `lib/storageHealth.ts`, which every page // and poll consults before it touches a drive. One `stalled` answer marks a // location at once; `clearAfterCleanPasses` clean answers in a row (two by // default) clear it (the rules are that module's). // // EVERY PASS APPLIES THE TIMINGS from the settings it reads (`applyHealthTimings`), // so a value changed by hand takes effect within one pass, and a changed // interval re-arms the pass's own timer (below). // // READ-ONLY AND IN MEMORY. It writes no settings and pauses nothing: the pass // above sees a stalled location as down (its probe answers `stalled` without // asking) and pauses on its own cadence. So, like the boot probe, it is armed // ABOVE the idle gate (`startStorageHealthWatch`, from instrumentation): an idle // boot has no five-minute pass, but it has the health pass — without it nothing // registers the locations, and a stall the watchdog marks is never cleared. export type StorageHealthPassOpts = { // Default: the configured locations, read from settings. locations?: readonly StorageLocation[]; io?: { read: () => SiteSettings }; // findmnt, for naming each root's block device. Default: getPaths(). bins?: Pick; // Test seam. Default: `detectLocationHealth` — the block device's counters, // or a child `stat` against `probeTimeoutMs` when no device can be named. probe?: LocationHealthProbe; now?: () => number; log?: (line: string) => void; }; export type StorageHealthPassResult = { probed: number; answers: Record; // Only the locations whose state changed. transitions: HealthTransition[]; }; // A probe's answer as a verdict. A bare state names no detector; a probe that // threw is "could not ask": `ok`, never `stalled`. function asVerdict(answer: LocationHealthState | HealthVerdict): HealthVerdict { return typeof answer === "string" ? { answer } : answer; } function statCause(): string { return `a stat of its root did not answer within ${secondsText(healthTimings().probeTimeoutMs)}`; } // Record one verdict. A verdict with no answer records nothing but the // detector that gave it. function recordVerdict( loc: StorageLocation, verdict: HealthVerdict, now: number, ): HealthTransition | null { // The counters' device, for the watchdog's slow-or-stalled check; a stat // verdict forgets it. const device = verdict.detector === "counters" ? verdict.device : verdict.detector === "stat" ? null : undefined; if (verdict.answer === null) { if (verdict.detector) noteLocationDetector(loc.id, verdict.detector, device); return null; } return recordLocationHealth(loc, verdict.answer, { now, cause: verdict.cause ?? statCause(), ...(verdict.detector ? { detector: verdict.detector } : {}), ...(device !== undefined ? { device } : {}), }); } export async function runStorageHealthPass( opts: StorageHealthPassOpts = {}, ): Promise { const log = opts.log ?? (() => {}); // The settings this pass runs on: its locations, and the drive-health // timings, applied before anything is asked. A caller that hands in the // locations reads no settings, and the timings stay as they were. const storage = opts.locations ? null : (opts.io ?? DEFAULT_IO).read().storage; if (storage) applyHealthTimings(storage.health); const locations = opts.locations ?? storage?.locations ?? []; pruneLocationHealth(locations.map((l) => l.id)); // Every configured location has an entry before anything is asked, so the // watchdog (lib/storageHealth.ts `onDrive`) can find a channel's location // even before the counters have given a first verdict. registerLocationHealth(locations); const bins = opts.bins ?? getPaths(); const probe: LocationHealthProbe = opts.probe ?? ((loc) => detectLocationHealth(loc, bins)); // Every location at once: each answer is bounded by the probe's own timer, // so the pass is too, and one stalled drive does not delay the others. const answers = await Promise.all( locations.map(async (loc) => { const verdict = await probe(loc).then(asVerdict, (): HealthVerdict => ({ answer: "ok", })); return [loc, verdict] as const; }), ); const out: StorageHealthPassResult = { probed: answers.length, answers: {}, transitions: [], }; const now = opts.now?.() ?? Date.now(); for (const [loc, verdict] of answers) { if (verdict.answer !== null) out.answers[loc.id] = verdict.answer; const t = recordVerdict(loc, verdict, now); if (!t) continue; out.transitions.push(t); if (t.to === "stalled") { log( `[storage] "${loc.id}": ${NOT_ANSWERING} — ${verdict.cause ?? statCause()}; ` + `pages and polls skip it until it answers ` + `${clearRuleText(healthTimings().clearAfterCleanPasses)}`, ); } else if (t.from === "stalled") { log(`[storage] "${loc.id}": answering again (${t.to})`); } } return out; } // ONE LOCATION, NOW: what /storage's Refresh asks before its own probe, so the // operator pressing it after doing something about the drive gets an answer // taken afterwards. It counts as one answer like any other — a stalled location // still needs two clean ones in a row — and the counters give none when their // last sample is under `minCounterIntervalMs()` old. Nothing is pruned. export async function refreshLocationHealth( loc: StorageLocation, probe?: LocationHealthProbe, ): Promise { registerLocationHealth([loc]); const ask: LocationHealthProbe = probe ?? ((l) => detectLocationHealth(l, getPaths())); const verdict = await ask(loc).then(asVerdict, (): HealthVerdict => ({ answer: "ok", })); recordVerdict(loc, verdict, Date.now()); return verdict.answer; } // --------------------------------------------------------------------------- // The cadence // --------------------------------------------------------------------------- // `storage.health.passIntervalMs`, fifteen seconds by default // (lib/storageHealth.ts says why). Each pass reads each location's device // counters in /sys (a findmnt only when the device is not known yet or its /sys // entry stopped reading; with no device, one short-lived `stat`). // A per-module-copy singleton, deliberately left so: it is not a temp-file // name (slice W folded every tmp + rename onto lib/jsonFile-server.ts, whose // state is on globalThis), and its one caller is editor/instrumentation.ts, // so only one copy ever arms it. (The health STATE the second timer writes is // on globalThis — pages in another module copy read it.) let timer: ReturnType | null = null; let healthTimer: ReturnType | null = null; let healthInFlight = false; let stopFollowingInterval: (() => void) | null = null; // THE FIVE-MINUTE PASS, ARMED ONCE PER PROCESS, below the idle gate: it writes // settings (an auto-pause). `unref()` so it never holds the event loop open — a // CLI that imports a controller must still exit. export function startStorageWatch( opts: StorageWatchOpts & { intervalMs?: number } = {}, ): boolean { if (timer) return false; const every = opts.intervalMs ?? STORAGE_WATCH_INTERVAL_MS; timer = setInterval(() => { void runStorageWatchPass(opts).catch((err) => { (opts.log ?? console.warn)( `[storage] watch pass failed: ${(err as Error).message}`, ); }); }, every); timer.unref?.(); return true; } export function stopStorageWatch(): void { if (timer) clearInterval(timer); timer = null; } // THE HEALTH PASS, ARMED ONCE PER PROCESS, above the idle gate: it is in // memory and writes nothing (see the section header). It also runs once at // once, so the locations are registered and a drive that is already stalled is // on its way to being known before the first page. // // ITS INTERVAL FOLLOWS THE SETTING. Armed at `healthTimings().passIntervalMs`, // and RE-ARMED whenever an applied change moves it — a save on /storage (at // once, from the page's module copy: the subscription is on globalThis), or a // hand edit (at the next pass, which applies what it read). An explicit // `intervalMs` (the tests') is fixed and follows nothing. export function startStorageHealthWatch( opts: { io?: { read: () => SiteSettings }; bins?: Pick; probe?: LocationHealthProbe; intervalMs?: number; log?: (line: string) => void; } = {}, ): boolean { if (healthTimer) return false; const health = () => { // One at a time: a pass is bounded by its timers, but a pass that overran // the interval must not stack a second one on top of it. if (healthInFlight) return; healthInFlight = true; void runStorageHealthPass({ io: opts.io, bins: opts.bins, probe: opts.probe, log: opts.log, }) .catch((err) => { (opts.log ?? console.warn)( `[storage] health pass failed: ${(err as Error).message}`, ); }) .finally(() => { healthInFlight = false; }); }; const arm = (every: number) => { if (healthTimer) clearInterval(healthTimer); healthTimer = setInterval(health, every); healthTimer.unref?.(); }; arm(opts.intervalMs ?? healthTimings().passIntervalMs); if (opts.intervalMs === undefined) { stopFollowingInterval = onPassIntervalChange((every) => { if (!healthTimer) return; (opts.log ?? console.log)(`[storage] health pass re-armed: every ${secondsText(every)}`); arm(every); }); } health(); return true; } export function stopStorageHealthWatch(): void { if (healthTimer) clearInterval(healthTimer); healthTimer = null; stopFollowingInterval?.(); stopFollowingInterval = null; }