Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 0f375893eda5ae97fd457cd616b0bd2a220f3621
parent bac70970c5ddf868d9b8c7c02524a07d7d31174d
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Tue, 29 Sep 2026 21:52:34 -0400

common: a stalled drive is asked out of process and then not asked in-process

A new lib/storageHealth.ts holds each storage location's health (ok, stalled,
absent, with since) in one map per process on globalThis. probeLocationHealth
(lib/storageVolumes.ts) runs `stat` on the root as a child process raced
against a 3 s timer; the timer's answer is `stalled`, and the child is killed
without being waited for. The storage watch arms a second, 15 s health pass
beside its 5-minute pass (and runs one at arm time): one miss marks a location
stalled at once, two clean answers in a row clear it, a re-pointed root starts
over. Nothing is persisted and nothing is paused by it.

The gate, before any filesystem call on a stalled location:
- inspectChannelMedia answers a new status `stalled` (detail "drive not
  answering (location "<label>", since HH:MM)"); with a config passed in it
  makes no call at all. HELD_REASON gains its words, so the index and stats
  builds hold such a channel; assertChannelMediaReachable refuses it.
- probeLocation and its memo answer a new probe status `stalled` ("Not
  answering") with no stat, statfs or findmnt.
- volumeFreeBytes reads the location as unknown; readChannelStat (the 1 s job
  list poll's full data/ walk) returns no counts; the recency layer skips the
  metadata tail reads for such channels and does not remember them as misses;
  a move onto such a root is refused before its stat; the saved-video store
  reads unreachable.

inspectChannelMedia results are memoised for 5 s per channel (slug + the
configured dataDir), asked after the gate. `fresh: true` bypasses the memo and
is passed by every caller that decides from the answer: the start-of-work
guard, both movers, the index and stats builds, the watch, eviction, the
re-point preflight and doctor. The channel mover's marker writes and clears,
clearRelocationMarker and a re-point clear it.

Tests: storageHealth (the rules), storageHealthProbe (a child that never
answers is stalled on the timer; the 3 s default; fail-open), storageStall
(every gated caller makes no call on the drive, through a spy on node:fs; the
memo), storageWatch (the health pass, recovery after two clean probes, arming).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

Diffstat:
Mcommon/bin/doctor.ts | 4+++-
Mcommon/controller/buildIndex.ts | 8++++++--
Mcommon/controller/buildStats.ts | 4+++-
Mcommon/controller/channels.ts | 8+++++++-
Mcommon/controller/evictClipWindows.ts | 4+++-
Mcommon/controller/recencyIndex.ts | 24+++++++++++++++++++++---
Mcommon/controller/relocateChannelMedia.ts | 31+++++++++++++++++++++++++++++--
Mcommon/controller/relocateSavedVideos.ts | 18++++++++++++++++++
Mcommon/controller/storageLocations.ts | 21+++++++++++++++++++--
Acommon/controller/storageStall.test.ts | 360+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/storageWatch.test.ts | 127+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/controller/storageWatch.ts | 168+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcommon/lib/channelMedia.ts | 157+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/lib/channelMediaHold.ts | 1+
Acommon/lib/storageHealth.test.ts | 129+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/storageHealth.ts | 231+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/storageHealthProbe.test.ts | 112+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/storageVolumes.ts | 100++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/views/pipeline/stageStatus.ts | 3++-
Mcommon/views/storage.ts | 1+
Meditor/app/components/MediaLocationBadge.tsx | 11+++++++----
21 files changed, 1486 insertions(+), 36 deletions(-)

diff --git a/common/bin/doctor.ts b/common/bin/doctor.ts @@ -121,7 +121,9 @@ export async function collectDoctorReport(deps: DoctorDeps): Promise<DoctorRepor const unreachable: string[] = []; const { inspectChannelMedia } = await import("../lib/channelMedia"); for (const slug of channelSlugs) { - const loc = await inspectChannelMedia(paths, slug); + const loc = await inspectChannelMedia(paths, slug, undefined, { + fresh: true, + }); if (loc.status !== "ok" && loc.status !== "in-place") { unreachable.push(`${slug}: ${loc.status} — ${loc.detail ?? ""}`.trim()); } diff --git a/common/controller/buildIndex.ts b/common/controller/buildIndex.ts @@ -345,7 +345,9 @@ async function scanSource( // An unmounted drive is not an empty channel (lib/channelMedia.ts): the // readdir below would fail, and every record the channel has would be // removed as gone. - const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg); + const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { + fresh: true, + }); if (isMediaHeld(media.status)) { held.set(ch.name, heldReason(media, cfg.dataDir, locations)); continue; @@ -462,7 +464,9 @@ async function scanSource( // Asked again after the walk: a drive that went away DURING it leaves the // videos after that point missing from this scan, which would remove them. // Three syscalls a channel. - const after = await inspectChannelMedia({ channelsDir }, ch.name, cfg); + const after = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { + fresh: true, + }); if (isMediaHeld(after.status)) { held.set(ch.name, heldReason(after, cfg.dataDir, locations)); continue; diff --git a/common/controller/buildStats.ts b/common/controller/buildStats.ts @@ -283,7 +283,9 @@ async function scanSource( channels.set(ch.name, cfg); // An unmounted drive is not an empty channel (lib/channelMedia.ts): the // readdir below would fail and every one of its stats would be removed. - const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg); + const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { + fresh: true, + }); if (isMediaHeld(media.status)) { held.set(ch.name, heldReason(media, cfg.dataDir, locations)); continue; diff --git a/common/controller/channels.ts b/common/controller/channels.ts @@ -16,7 +16,7 @@ import { readVideoFiles, } from "../lib/videoStatus"; import { loadDigest } from "../lib/digest-server"; -import { readRelocationMarker } from "../lib/channelMedia"; +import { channelMediaStall, readRelocationMarker } from "../lib/channelMedia"; // TYPE-ONLY, and it must stay that way: ./channelSnapshot imports // readChannelConfig from this module, and it drags in the snapshot generator's // whole dependency graph (lmdb, the archive reader, the digest layer). A value @@ -189,6 +189,12 @@ export async function readChannelStat( ): Promise<ChannelStat | null> { const config = await readChannelConfig(paths, slug); if (!config) return null; + // A WALK OF `data/` ON A DRIVE THAT IS NOT ANSWERING IS NOT STARTED. The + // one-second job-list poll asks this for every channel with a job listed, and + // on a stalled drive each readdir and stat in the walk would hold an I/O + // thread until the drive came back. No counts is what a caller already + // handles (the row draws no progress bar). + if (channelMediaStall(config)) return null; const channelDir = path.join(paths.channelsDir, slug); const counts = await countDataFiles(path.join(channelDir, "data")); return { diff --git a/common/controller/evictClipWindows.ts b/common/controller/evictClipWindows.ts @@ -96,7 +96,9 @@ async function evictChannel( // fine. (The JOB also declares `needsMedia: true`, which covers a // single-channel run before it starts; this covers the corpus-wide one, // where there is no slug for that guard to check.) - const media = await inspectChannelMedia(opts.paths, slug); + const media = await inspectChannelMedia(opts.paths, slug, undefined, { + fresh: true, + }); if (media.status !== "ok" && media.status !== "in-place") { out.skipped.push( `${slug}: media ${media.status}${media.detail ? ` (${media.detail})` : ""} — nothing was touched`, diff --git a/common/controller/recencyIndex.ts b/common/controller/recencyIndex.ts @@ -7,6 +7,8 @@ import { mapConcurrent } from "../lib/concurrency"; import { extractVideoId } from "../lib/videoId"; import { uploadKeyFor } from "./keptVideos"; import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; +import type { ChannelConfig } from "../lib/channelConfig"; +import { channelMediaStall } from "../lib/channelMedia"; // Upload-date lookup for the auto-queue's "newest first" ordering. // @@ -189,11 +191,18 @@ async function readTailUploadDate(file: string): Promise<string | null> { // Date the ids in `wanted` from their on-disk metadata. Removes each id it // keys from `wanted`, like interpolateFromPlaylist. +// +// `stalled` names the channels whose media is on a drive that is not answering +// (lib/storageHealth.ts). Their ids are not read — each read would hold an I/O +// thread until the drive came back, 32 at a time — and NOT memoized as misses +// either: they fall through to layers 3 and 4 for now, and a later refresh with +// the drive answering reads them. async function datesFromMetadata( paths: Paths, owner: ReadonlyMap<string, string>, wanted: Set<string>, out: Map<string, RecencyKey>, + stalled: ReadonlySet<string> = new Set(), ): Promise<void> { const todo: string[] = []; for (const id of wanted) { @@ -205,7 +214,8 @@ async function datesFromMetadata( } continue; } - if (!owner.has(id)) continue; + const slug = owner.get(id); + if (slug === undefined || stalled.has(slug)) continue; todo.push(id); if (todo.length >= TAIL_READS_PER_BUILD) { if (!tailCapLogged) { @@ -312,7 +322,12 @@ export function interpolateFromPlaylist( export type BuildRecencyKeysArgs = { paths: Paths; // The channels whose work is in play — the runner's own channel-meta list. - meta: ReadonlyArray<{ slug: string }>; + // The config, when the caller holds it, is what tells layer 2 which channels + // are on a drive that is not answering (it skips their reads). + meta: ReadonlyArray<{ + slug: string; + config?: Pick<ChannelConfig, "dataDir"> | null; + }>; // Every video id the caller might sort. Ids outside this set are not keyed. candidateIds: ReadonlySet<string>; // videoId -> owning channel slug. Required for the tail-read layer, which has @@ -468,7 +483,10 @@ export async function buildRecencyKeys({ // auto-transcribe, whose entire candidate set is by definition absent from the // transcript index. if (owner && missing.size > 0) { - await datesFromMetadata(paths, owner, missing, out); + const stalled = new Set( + meta.filter((m) => channelMediaStall(m.config)).map((m) => m.slug), + ); + await datesFromMetadata(paths, owner, missing, out, stalled); } // Layer 3: playlist interpolation, for videos with nothing on disk at all. diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts @@ -38,9 +38,15 @@ import { type VolumeBins, } from "../lib/storageVolumes"; import { getFreeBytes } from "../lib/diskSpace"; +import { + NOT_ANSWERING, + sinceText, + stalledLocation, +} from "../lib/storageHealth"; import { getSettings, type SiteSettings } from "../lib/settings"; import { formatBytes } from "../lib/format"; import { + forgetChannelMedia, inspectChannelMedia, relocatedDataDir, relocationMarkerPath, @@ -282,6 +288,16 @@ export async function relocationRootPresenceProblem( const named = locationForRoot(r, storage.locations); const where = named ? ` (location "${named.id}")` : ""; + // A destination whose drive is not answering is refused WITHOUT the stat + // below: that stat would wait on the drive, and a move onto it would too. + const stall = named ? stalledLocation(named) : null; + if (stall) { + return ( + `The destination root ${r}${where}: ${NOT_ANSWERING} ` + + `(${sinceText(stall.since)}). Wait for it to answer, or check the drive.` + ); + } + let isDir = false; try { isDir = (await stat(r)).isDirectory(); @@ -373,16 +389,23 @@ async function sweepParked( // The channel's marker, at `channels/<slug>/.relocating.json`. The FILE is the // contract — see relocateDir.ts — and these three are the channel's name for it. +// +// Each write and the clear also drop the channel from `inspectChannelMedia`'s +// five-second page memo. Every phase change (copy → swap → reclaim) writes the +// marker AFTER the link and the config it changes, so a page asks the disk +// again the moment the move has done something it would see. async function writeMarker( paths: Paths, slug: string, marker: RelocationMarker, ): Promise<void> { await writeDirMarker(relocationMarkerPath(paths, slug), marker); + forgetChannelMedia(slug); } async function clearMarker(paths: Paths, slug: string): Promise<void> { await clearDirMarker(relocationMarkerPath(paths, slug)); + forgetChannelMedia(slug); } async function readMarkerRaw( @@ -601,7 +624,9 @@ async function moveOut(args: { // every guard and exactly wrong here: the rerun that finishes an interrupted // move is the one caller allowed to see it. A marker for a DIFFERENT target // was already refused above. - const location = await inspectChannelMedia(paths, slug, args.config); + const location = await inspectChannelMedia(paths, slug, args.config, { + fresh: true, + }); if (!args.resumed && location.status !== "in-place") { throw new Error( `Channel "${slug}" is not in a movable state: ${ @@ -838,7 +863,9 @@ async function moveBack(args: { // // `in-transition` is allowed because a marker is what a resume carries, and a // rerun is the caller this precondition must not refuse. - const location = await inspectChannelMedia(paths, slug, args.config); + const location = await inspectChannelMedia(paths, slug, args.config, { + fresh: true, + }); if ( !args.resumed && location.status !== "ok" && diff --git a/common/controller/relocateSavedVideos.ts b/common/controller/relocateSavedVideos.ts @@ -20,6 +20,11 @@ import { } from "../lib/settings"; import type { RelocationMarker, RelocationPhase } from "../lib/channelMedia"; import { + NOT_ANSWERING, + sinceText, + stalledLocation, +} from "../lib/storageHealth"; +import { relocatedSavedVideosDir, savedVideosMarkerPath, } from "../lib/savedVideoStore"; @@ -156,6 +161,19 @@ export async function inspectSavedVideosStore( detail: `the store links to ${target}, but no storage location is recorded for it`, }; } + // The store's drive is not answering: reported without the stat below, + // which would wait on it (the /storage page asks on every render). The + // page skips the store's size walk for an unreachable store, too. + const stall = loc ? stalledLocation(loc) : null; + if (stall) { + return { + dir, + locationId, + target, + status: "unreachable", + detail: `${NOT_ANSWERING} (location "${stall.label}", ${sinceText(stall.since)})`, + }; + } if (!(await isDirectory(target))) { return { dir, diff --git a/common/controller/storageLocations.ts b/common/controller/storageLocations.ts @@ -24,7 +24,12 @@ import { type StorageLocationProbe, type VolumeBins, } from "../lib/storageVolumes"; -import { inspectChannelMedia, relocatedDataDir } from "../lib/channelMedia"; +import { + forgetChannelMedia, + inspectChannelMedia, + relocatedDataDir, +} from "../lib/channelMedia"; +import { stalledLocation } from "../lib/storageHealth"; import { relocationQueueKey } from "../lib/queueKeys"; import { runManagedFunction, @@ -247,6 +252,12 @@ export async function volumeFreeBytes(opts: { out[INTERNAL_LOCATION_ID] = Number.isFinite(corpus) ? corpus : undefined; await Promise.all( opts.locations.map(async (loc) => { + // NOT ASKED WHILE ITS DRIVE IS NOT ANSWERING: each stat and the statfs + // below would hold an I/O thread until it did. "—", like unmounted. + if (stalledLocation(loc)) { + out[loc.id] = undefined; + return; + } const st = await stat(loc.root).catch(() => null); if (!st?.isDirectory()) { out[loc.id] = undefined; @@ -528,7 +539,9 @@ export async function preflightRepoint(opts: { const missing: string[] = []; for (const slug of slugs) { const config = await readChannelConfig(opts.paths, slug); - const media = await inspectChannelMedia(opts.paths, slug, config); + const media = await inspectChannelMedia(opts.paths, slug, config, { + fresh: true, + }); if (media.status === "in-transition") { base.problems.push( `${slug}: a media relocation is in flight or was interrupted ` + @@ -795,6 +808,10 @@ export async function repointStorageLocation(opts: { } resetStorageProbeMemo(); + // Every channel on it has a new link and a new dataDir: the page memo's keys + // already differ, and this drops the old answers rather than letting them + // age out. + forgetChannelMedia(); log( `Done. "${loc.label}" is at ${pre.newRoot}; ${pre.channels.length} ` + `channel(s) re-pointed.`, diff --git a/common/controller/storageStall.test.ts b/common/controller/storageStall.test.ts @@ -0,0 +1,360 @@ +// A STALLED DRIVE IS NOT ASKED: every page-and-poll path that would touch a +// storage location's drive in-process asks the health state first +// (lib/storageHealth.ts) and, on a stalled location, answers WITHOUT the call. +// +// No test stalls a real drive. The drive here is an ordinary temp directory +// that answers every call at once; the health state is TOLD it is stalled. +// Every node:fs and node:fs/promises call this file's code makes is recorded +// (the buildIndex.test.ts spy, reads included), so a gate that let one call +// through shows up as a recorded path under the drive — it cannot hide behind +// a call that happened to answer quickly. +// +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/storageStall.test.ts + +import { after, beforeEach, test } from "node:test"; +import assert from "node:assert/strict"; +import { createRequire, syncBuiltinESMExports } from "node:module"; +import { + existsSync, + mkdirSync, + mkdtempSync, + rmSync, + symlinkSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import type { Paths } from "../lib/paths"; +import type { StorageLocation } from "../lib/storageLocations"; +import { + assertChannelMediaReachable, + ChannelMediaUnreachableError, + CHANNEL_MEDIA_MEMO_MS, + clearRelocationMarker, + forgetChannelMedia, + inspectChannelMedia, +} from "../lib/channelMedia"; +import { HELD_REASON, isMediaHeld } from "../lib/channelMediaHold"; +import { + recordLocationHealth, + resetStorageHealth, +} from "../lib/storageHealth"; +import { + probeLocation, + probeLocationMemo, + resetStorageProbeMemo, +} from "../lib/storageVolumes"; +import { volumeFreeBytes } from "./storageLocations"; +import { readChannelStat } from "./channels"; +import { buildRecencyKeys, clearRecencyCache } from "./recencyIndex"; +import { relocationRootPresenceProblem } from "./relocateChannelMedia"; +import { generateChannelSnapshot } from "./channelSnapshot"; +import { inspectSavedVideosStore } from "./relocateSavedVideos"; +import type { SiteSettings } from "../lib/settings"; + +// ── the spy ──────────────────────────────────────────────────────────────── +type Call = { fn: string; path: string }; +let calls: Call[] = []; +{ + const req = createRequire(import.meta.url); + const fsCjs = req("node:fs") as Record<string, unknown>; + const fspCjs = req("node:fs/promises") as Record<string, unknown>; + const NAMES = [ + "access", "appendFile", "chmod", "copyFile", "cp", "lstat", "mkdir", + "open", "opendir", "readdir", "readFile", "readlink", "realpath", "rename", + "rm", "rmdir", "stat", "statfs", "symlink", "unlink", "utimes", "writeFile", + ]; + const asPath = (v: unknown) => + typeof v === "string" + ? v + : v instanceof URL + ? fileURLToPath(v) + : Buffer.isBuffer(v) + ? v.toString() + : null; + const wrap = (mod: Record<string, unknown>, name: string) => { + const fn = mod[name]; + if (typeof fn !== "function") return; + mod[name] = function (this: unknown, ...args: unknown[]) { + const p = asPath(args[0]); + if (p !== null) calls.push({ fn: name, path: path.resolve(p) }); + return (fn as (...a: unknown[]) => unknown).apply(this, args); + }; + }; + for (const n of NAMES) { + wrap(fspCjs, n); + wrap(fsCjs, n); + wrap(fsCjs, `${n}Sync`); + } + wrap(fsCjs, "existsSync"); + wrap(fsCjs, "createReadStream"); + syncBuiltinESMExports(); +} + +// ── the corpus and the drive ─────────────────────────────────────────────── +const ROOT = mkdtempSync(path.join(tmpdir(), "ttb-stall-")); +after(() => rmSync(ROOT, { recursive: true, force: true })); + +const CORPUS = path.join(ROOT, "corpus"); +const DRIVE = path.join(ROOT, "drive"); +const SLUG = "on-drive"; +const paths = { + transcriptsDir: CORPUS, + channelsDir: path.join(CORPUS, "channels"), + lmdbPath: path.join(CORPUS, "index.mdb"), + savedVideosDir: path.join(CORPUS, "saved-videos"), + jobsDir: path.join(CORPUS, "jobs"), + // A findmnt that leaves a mark if anything runs it. + findmntBin: path.join(ROOT, "findmnt-ran.sh"), + udisksctlBin: path.join(ROOT, "no-udisksctl"), +} as unknown as Paths; +const LOC: StorageLocation = { + id: "usb", + label: "USB drive", + root: DRIVE, + autoRepoint: false, +}; +const TARGET = path.join(DRIVE, SLUG, "data"); +const LINK = path.join(paths.channelsDir, SLUG, "data"); +const CONFIG = { dataDir: TARGET }; +const FINDMNT_MARK = path.join(ROOT, "findmnt-ran"); + +function seed(): void { + rmSync(CORPUS, { recursive: true, force: true }); + rmSync(DRIVE, { recursive: true, force: true }); + mkdirSync(path.join(TARGET, "vid1"), { recursive: true }); + writeFileSync( + path.join(TARGET, "vid1", "metadata.info.json"), + JSON.stringify({ id: "vid1", upload_date: "20260601" }), + ); + writeFileSync(path.join(TARGET, "vid1", "transcript.en.vtt"), "WEBVTT\n"); + mkdirSync(path.join(paths.channelsDir, SLUG), { recursive: true }); + writeFileSync( + path.join(paths.channelsDir, SLUG, "config.json"), + JSON.stringify({ + handling: "youtube", + name: SLUG, + url: `https://www.youtube.com/@${SLUG}/videos`, + dataDir: TARGET, + }), + ); + symlinkSync(TARGET, LINK); + writeFileSync(paths.findmntBin, `#!/bin/sh\ntouch ${FINDMNT_MARK}\nexit 1\n`, { + mode: 0o755, + }); + rmSync(FINDMNT_MARK, { force: true }); +} + +// Anything that reaches the drive: a path under its root, or through the +// channel's link (`data/` itself, followed, or anything below it). +function onDrive(c: Call): boolean { + const under = (p: string, base: string) => + p === base || p.startsWith(base + path.sep); + return under(c.path, DRIVE) || under(c.path, LINK); +} + +function stall(): void { + recordLocationHealth(LOC, "stalled", { now: Date.now() }); +} + +beforeEach(() => { + resetStorageHealth(); + forgetChannelMedia(); + resetStorageProbeMemo(); + clearRecencyCache(); + seed(); + calls = []; +}); + +// ── inspectChannelMedia: the gate and the memo ───────────────────────────── + +test("inspect with the config in hand: a stalled drive costs no call at all", async () => { + stall(); + calls = []; + const media = await inspectChannelMedia(paths, SLUG, CONFIG); + assert.deepEqual(calls, []); + assert.equal(media.status, "stalled"); + assert.equal(media.target, TARGET); + assert.match(String(media.detail), /^drive not answering \(location "USB drive", since /); + // Held by both pool-wide builds, with a reason that names no path. + assert.equal(isMediaHeld(media.status), true); + assert.doesNotMatch(HELD_REASON.stalled, /\//); +}); + +test("inspect without the config reads config.json and nothing on the drive", async () => { + stall(); + calls = []; + const media = await inspectChannelMedia(paths, SLUG); + assert.equal(media.status, "stalled"); + assert.deepEqual( + calls.map((c) => [c.fn, path.relative(ROOT, c.path)]), + [["readFile", path.join("corpus", "channels", SLUG, "config.json")]], + ); +}); + +test("the start-of-work guard refuses a stalled channel without a call", async () => { + stall(); + calls = []; + await assert.rejects( + () => assertChannelMediaReachable(paths, SLUG, CONFIG), + (err: unknown) => + err instanceof ChannelMediaUnreachableError && + err.status === "stalled" && + /drive not answering/.test(err.message), + ); + assert.deepEqual(calls, []); +}); + +test("the memo: two inspects within 5 s stat the drive once; fresh and age each ask again", async () => { + const t0 = 1_000_000; + const targetStats = () => + calls.filter((c) => c.fn === "stat" && c.path === TARGET).length; + const a = await inspectChannelMedia(paths, SLUG, CONFIG, { now: t0 }); + assert.equal(a.status, "ok"); + assert.equal(targetStats(), 1); + const b = await inspectChannelMedia(paths, SLUG, CONFIG, { now: t0 + 4_999 }); + assert.equal(b.status, "ok"); + assert.equal(targetStats(), 1, "a second inspect inside five seconds is remembered"); + await inspectChannelMedia(paths, SLUG, CONFIG, { now: t0 + 4_999, fresh: true }); + assert.equal(targetStats(), 2, "fresh bypasses the memo"); + await inspectChannelMedia(paths, SLUG, CONFIG, { now: t0 + CHANNEL_MEDIA_MEMO_MS }); + assert.equal(targetStats(), 3, "an answer five seconds old is asked again"); + // A fresh answer is not remembered: the memo still holds the one from t0+5000. + await inspectChannelMedia(paths, SLUG, CONFIG, { now: t0 + CHANNEL_MEDIA_MEMO_MS + 1 }); + assert.equal(targetStats(), 3); +}); + +test("the memo is keyed by the configured target, and a mover's forget clears it", async () => { + const now = 2_000_000; + const targetStats = () => + calls.filter((c) => c.fn === "stat" && c.path === TARGET).length; + await inspectChannelMedia(paths, SLUG, CONFIG, { now }); + // Another configured target is another key. + const other = await inspectChannelMedia(paths, SLUG, { dataDir: path.join(DRIVE, "x", "data") }, { now }); + assert.equal(other.status, "inconsistent"); + forgetChannelMedia(SLUG); + await inspectChannelMedia(paths, SLUG, CONFIG, { now }); + assert.equal(targetStats(), 2); + // clearRelocationMarker (the operator's last resort) forgets too. + await clearRelocationMarker(paths, SLUG); + await inspectChannelMedia(paths, SLUG, CONFIG, { now }); + assert.equal(targetStats(), 3); +}); + +test("the gate is asked before the memo: a stall is seen with an ok remembered", async () => { + const now = 3_000_000; + assert.equal((await inspectChannelMedia(paths, SLUG, CONFIG, { now })).status, "ok"); + stall(); + calls = []; + const media = await inspectChannelMedia(paths, SLUG, CONFIG, { now: now + 1 }); + assert.equal(media.status, "stalled"); + assert.deepEqual(calls, []); +}); + +// ── the other gated callers ──────────────────────────────────────────────── + +test("probeLocation and its memo: 'stalled', with no stat, no statfs and no findmnt", async () => { + // A remembered "available" first, so the memo's gate is the thing tested. + await probeLocationMemo(LOC, paths); + rmSync(FINDMNT_MARK, { force: true }); + stall(); + calls = []; + const direct = await probeLocation(LOC, paths); + const memo = await probeLocationMemo(LOC, paths); + assert.equal(direct.status, "stalled"); + assert.equal(memo.status, "stalled"); + assert.deepEqual(direct.identity, { known: false }); + assert.equal(direct.freeBytes, undefined); + assert.deepEqual(calls.filter(onDrive), []); + assert.equal(existsSync(FINDMNT_MARK), false, "findmnt was not run"); +}); + +test("volumeFreeBytes: a stalled location reads unknown, with no call on it", async () => { + stall(); + calls = []; + const out = await volumeFreeBytes({ paths, locations: [LOC] }); + assert.equal(out.usb, undefined); + assert.equal(typeof out.internal, "number"); + assert.deepEqual(calls.filter(onDrive), []); +}); + +test("readChannelStat: no walk of data/ on a stalled drive", async () => { + const before = await readChannelStat(paths, SLUG); + assert.equal(before?.videoCount, 1); + stall(); + calls = []; + assert.equal(await readChannelStat(paths, SLUG), null); + assert.deepEqual(calls.filter(onDrive), []); +}); + +test("recency: no tail read on a stalled drive, and no miss remembered for it", async () => { + const owner = new Map([["vid1", SLUG]]); + const meta = [{ slug: SLUG, config: CONFIG }]; + const args = { + paths, + meta, + candidateIds: new Set(["vid1"]), + owner, + interpolate: false, + fresh: true, + }; + stall(); + calls = []; + const stalledKeys = await buildRecencyKeys(args); + assert.deepEqual(calls.filter(onDrive), []); + // Layer 4: an undatable id sorts oldest. + assert.deepEqual(stalledKeys.get("vid1"), { key: "", estimated: false }); + // The drive answers again: the same id is read now, because the stall was + // not remembered as a miss. + resetStorageHealth(); + const keys = await buildRecencyKeys(args); + assert.deepEqual(keys.get("vid1"), { key: "20260601", estimated: false }); + assert.ok(calls.some((c) => c.fn === "open" && onDrive(c))); +}); + +test("a move onto a stalled location is refused without a stat of its root", async () => { + stall(); + calls = []; + const problem = await relocationRootPresenceProblem( + DRIVE, + { locations: [LOC], defaultLocationId: "" }, + paths, + ); + assert.match(String(problem), /drive not answering/); + assert.match(String(problem), /location "usb"/); + assert.deepEqual(calls.filter(onDrive), []); +}); + +test("a snapshot refresh of a stalled channel throws before its walk", async () => { + stall(); + calls = []; + await assert.rejects( + () => generateChannelSnapshot(paths, SLUG), + (err: unknown) => + err instanceof ChannelMediaUnreachableError && err.status === "stalled", + ); + assert.deepEqual(calls.filter(onDrive), []); +}); + +test("the saved-video store on a stalled drive reads unreachable, without a stat through its link", async () => { + const storeTarget = path.join(DRIVE, "saved-videos"); + mkdirSync(storeTarget, { recursive: true }); + symlinkSync(storeTarget, paths.savedVideosDir); + const settings = { + storage: { locations: [LOC], defaultLocationId: "", savedVideosLocationId: "usb" }, + } as unknown as SiteSettings; + assert.equal((await inspectSavedVideosStore(paths, settings)).status, "ok"); + stall(); + calls = []; + const store = await inspectSavedVideosStore(paths, settings); + assert.equal(store.status, "unreachable"); + assert.match(String(store.detail), /^drive not answering \(location "USB drive"/); + const followed = calls.filter( + (c) => + onDrive(c) || + (c.path.startsWith(paths.savedVideosDir) && !["lstat", "readlink"].includes(c.fn)), + ); + assert.deepEqual(followed, []); +}); diff --git a/common/controller/storageWatch.test.ts b/common/controller/storageWatch.test.ts @@ -12,13 +12,26 @@ import { import { LANES } from "../lib/autoQueueTypes"; import { resetStorageWatchSuspicion, + runStorageHealthPass, runStorageWatchPass, + startStorageWatch, + stopStorageWatch, } from "./storageWatch"; +import { inspectChannelMedia } from "../lib/channelMedia"; +import { + 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. -beforeEach(() => resetStorageWatchSuspicion()); +// 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(); +}); // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/storageWatch.test.ts @@ -375,3 +388,113 @@ test("the restore needs only one good pass", async () => { 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" }); + assert.deepEqual(r.transitions, [{ id: "cold", from: null, 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("arming the watch runs a health pass at once, and stopping it stops both timers", async () => { + await withTmp(async (h) => { + let asked = 0; + const armed = startStorageWatch({ + paths: h.paths, + io: h.io, + bins: h.paths, + write: false, + healthProbe: async () => { + asked += 1; + return "stalled"; + }, + log: () => {}, + }); + try { + assert.equal(armed, true); + // Armed once per process. + assert.equal(startStorageWatch({ 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 { + stopStorageWatch(); + } + assert.equal(startStorageWatch({ io: h.io, healthProbe: async () => "ok" }), true); + stopStorageWatch(); + }); +}); diff --git a/common/controller/storageWatch.ts b/common/controller/storageWatch.ts @@ -19,8 +19,24 @@ import { import { LANES } from "../lib/autoQueueTypes"; import { siteChannelIndex } from "../lib/site"; import { inspectChannelMedia } from "../lib/channelMedia"; -import { locationOfDataDir } from "../lib/storageLocations"; -import type { VolumeBins } from "../lib/storageVolumes"; +import { + locationOfDataDir, + type StorageLocation, +} from "../lib/storageLocations"; +import { + probeLocationHealth, + type LocationHealthProbe, + type VolumeBins, +} from "../lib/storageVolumes"; +import { + HEALTH_PROBE_INTERVAL_MS, + HEALTH_PROBE_TIMEOUT_MS, + NOT_ANSWERING, + pruneLocationHealth, + recordLocationHealth, + type HealthTransition, + type LocationHealthState, +} from "../lib/storageHealth"; import { listChannelConfigs } from "./channels"; import { maybeAutoRepoint, probeAllLocations } from "./storageLocations"; @@ -208,7 +224,12 @@ export async function runStorageWatchPass( const loc = locationOfDataDir(dataDir, locations); const probe = loc ? probes[loc.id] : undefined; const locationDown = Boolean(loc) && probe?.status !== "available"; - const media = await inspectChannelMedia(paths, slug, config); + // 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. @@ -308,26 +329,123 @@ export async function runStorageWatchPass( return out; } -// --------------------------------------------------------------------------- -// The cadence -// --------------------------------------------------------------------------- - // 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 15 s this asks each location's root, out of +// process (`probeLocationHealth`), and records the answer in +// `lib/storageHealth.ts`, which every page and poll consults before it touches +// a drive. One miss marks a location `stalled` at once; two clean answers in a +// row clear it (the rules are that module's). +// +// 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. Armed and stopped with the watch, so an +// idle boot (which does not arm the watch) has no health probe either. + +export type StorageHealthPassOpts = { + // Default: the configured locations, read from settings. + locations?: readonly StorageLocation[]; + io?: { read: () => SiteSettings }; + // Test seam. Default: `probeLocationHealth`, a child `stat` against a 3 s + // timer. + probe?: LocationHealthProbe; + now?: () => number; + log?: (line: string) => void; +}; + +export type StorageHealthPassResult = { + probed: number; + answers: Record<string, LocationHealthState>; + // Only the locations whose state changed. + transitions: HealthTransition[]; +}; + +export async function runStorageHealthPass( + opts: StorageHealthPassOpts = {}, +): Promise<StorageHealthPassResult> { + const log = opts.log ?? (() => {}); + const locations = + opts.locations ?? (opts.io ?? DEFAULT_IO).read().storage.locations; + pruneLocationHealth(locations.map((l) => l.id)); + const probe: LocationHealthProbe = + opts.probe ?? ((loc) => probeLocationHealth(loc)); + // 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 answer = await probe(loc).catch( + (): LocationHealthState => "ok", + ); + return [loc, answer] as const; + }), + ); + const out: StorageHealthPassResult = { + probed: answers.length, + answers: {}, + transitions: [], + }; + const now = opts.now?.() ?? Date.now(); + for (const [loc, answer] of answers) { + out.answers[loc.id] = answer; + const t = recordLocationHealth(loc, answer, { + now, + cause: `a stat of its root did not answer within ${HEALTH_PROBE_TIMEOUT_MS / 1000} s`, + }); + if (!t) continue; + out.transitions.push(t); + if (t.to === "stalled") { + log( + `[storage] "${loc.id}": ${NOT_ANSWERING} — a stat of its root did ` + + `not answer within ${HEALTH_PROBE_TIMEOUT_MS / 1000} s; pages and ` + + `polls skip it until two probes in a row answer`, + ); + } else if (t.from === "stalled") { + log(`[storage] "${loc.id}": answering again (${t.to})`); + } + } + return out; +} + +// --------------------------------------------------------------------------- +// The cadence +// --------------------------------------------------------------------------- + +// Fifteen seconds (lib/storageHealth.ts says why). Each pass is one short-lived +// `stat` subprocess per location. +export const STORAGE_HEALTH_INTERVAL_MS = HEALTH_PROBE_INTERVAL_MS; + // 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. +// 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<typeof setInterval> | null = null; +let healthTimer: ReturnType<typeof setInterval> | null = null; +let healthInFlight = false; // ARMED ONCE PER PROCESS. `unref()` so it never holds the event loop open — a -// CLI that imports a controller must still exit. +// CLI that imports a controller must still exit. Two timers: the five-minute +// watch pass and the fifteen-second health pass, which also runs once at once +// so a drive that is already stalled at boot is known before the first page. export function startStorageWatch( - opts: StorageWatchOpts & { intervalMs?: number } = {}, + opts: StorageWatchOpts & { + intervalMs?: number; + healthIntervalMs?: number; + healthProbe?: LocationHealthProbe; + } = {}, ): boolean { if (timer) return false; const every = opts.intervalMs ?? STORAGE_WATCH_INTERVAL_MS; @@ -339,11 +457,37 @@ export function startStorageWatch( }); }, every); timer.unref?.(); + const health = () => { + // One at a time: a pass is bounded by the probe's timer, 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, + probe: opts.healthProbe, + log: opts.log, + }) + .catch((err) => { + (opts.log ?? console.warn)( + `[storage] health pass failed: ${(err as Error).message}`, + ); + }) + .finally(() => { + healthInFlight = false; + }); + }; + healthTimer = setInterval( + health, + opts.healthIntervalMs ?? STORAGE_HEALTH_INTERVAL_MS, + ); + healthTimer.unref?.(); + health(); return true; } export function stopStorageWatch(): void { - if (!timer) return; - clearInterval(timer); + if (timer) clearInterval(timer); + if (healthTimer) clearInterval(healthTimer); timer = null; + healthTimer = null; } diff --git a/common/lib/channelMedia.ts b/common/lib/channelMedia.ts @@ -2,6 +2,12 @@ import path from "node:path"; import { lstat, readFile, readlink, rm, stat } from "node:fs/promises"; import type { Paths } from "./paths"; import type { ChannelConfig } from "./channelConfig"; +import { + NOT_ANSWERING, + sinceText, + stalledLocationForPath, + type LocationHealth, +} from "./storageHealth"; // WHERE A CHANNEL'S MEDIA ACTUALLY IS, and whether it can be reached. // @@ -60,7 +66,12 @@ export type ChannelMediaStatus = // A relocation is in flight (or was interrupted): the marker is present. | "in-transition" // Disk and config disagree, in either direction. Never guessed past. - | "inconsistent"; + | "inconsistent" + // Relocated onto a storage location whose drive is not answering + // (`lib/storageHealth.ts`). Answered from memory, WITHOUT a filesystem call: + // a call there would block one of the process's few I/O threads for as long + // as the drive takes to come back. Held and refused like `unreachable`. + | "stalled"; export type ChannelMediaLocation = { // Always channelDir/data — the path every reader uses, relocated or not. @@ -162,6 +173,7 @@ export async function clearRelocationMarker( slug: string, ): Promise<void> { await rm(relocationMarkerPath(paths, slug), { force: true }); + forgetChannelMedia(slug); } // The `dataDir` field alone, read straight off config.json. Deliberately NOT @@ -185,13 +197,111 @@ async function readConfiguredDataDir( } } +// The configured target's drive is not answering: the answer, from memory. The +// detail names the location and when it stopped, never a path — the badge's +// title shows it, and /storage has the paths. +export function stalledMediaLocation( + dataDir: string, + configured: string, + health: LocationHealth, +): ChannelMediaLocation { + return { + dataDir, + relocated: true, + target: configured, + status: "stalled", + detail: `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`, + }; +} + +// THE STALL, for a caller holding a parsed config: the stalled location the +// channel's media is on, or null. No I/O — the question every page and poll +// that reads a channel's `data/` asks before it does. +export function channelMediaStall( + config: Pick<ChannelConfig, "dataDir"> | null | undefined, +): LocationHealth | null { + const dir = config?.dataDir?.trim(); + return dir ? stalledLocationForPath(dir) : null; +} + +// --------------------------------------------------------------------------- +// The memo +// --------------------------------------------------------------------------- +// +// FIVE SECONDS, PER CHANNEL, KEYED BY SLUG AND THE CONFIGURED TARGET. The home +// page, /channels and the auto-queue status poll (every three seconds, four +// lanes) each inspect every channel; without this each of them costs three +// syscalls a relocated channel, every time, on a drive that may be the slow +// one. The key carries the configured `dataDir`, so a move that rewrites it is +// a new key at once, and the movers clear the memo outright +// (`forgetChannelMedia`) whenever a marker is written or removed. +// +// THE STALL GATE IS ASKED BEFORE THE MEMO, so a drive that stops answering is +// seen on the next call even when an `ok` from four seconds ago is remembered. +// +// `fresh: true` BYPASSES IT, and every caller that decides something from the +// answer passes it: the start-of-work guard (`assertChannelMediaReachable`), +// the movers, the index and stats builds, the storage watch, eviction and the +// re-point preflight. The memo is for pages and polls. +// +// ONE MAP PER PROCESS (on `globalThis`), for the reason `storageHealth.ts` +// gives: the movers that clear it run in one bundle layer and the pages that +// read it in another. + +export const CHANNEL_MEDIA_MEMO_MS = 5_000; + +export type InspectOptions = { + // Skip the memo: take a fresh answer, and do not remember it. + fresh?: boolean; + // Test seam for the clock. + now?: number; +}; + +type MediaMemo = Map<string, { at: number; location: ChannelMediaLocation }>; + +declare global { + // eslint-disable-next-line no-var + var __yttChannelMediaMemo__: MediaMemo | undefined; +} + +function mediaMemo(): MediaMemo { + if (!globalThis.__yttChannelMediaMemo__) { + globalThis.__yttChannelMediaMemo__ = new Map(); + } + return globalThis.__yttChannelMediaMemo__; +} + +// Forget what the memo holds: for one channel, or for every channel. The +// movers call it whenever the disk changes under a channel (a marker written or +// removed, a link swapped, a location re-pointed), so the next page sees it. +export function forgetChannelMedia(slug?: string): void { + const memo = mediaMemo(); + if (slug === undefined) { + memo.clear(); + return; + } + for (const key of [...memo.keys()]) { + if (key.split("\u0000")[1] === slug) memo.delete(key); + } +} + +function memoKey(paths: ChannelMediaPaths, slug: string, configured?: string): string { + return `${paths.channelsDir}\u0000${slug}\u0000${configured ?? ""}`; +} + +// Past this many entries, expired ones are swept on insert. A corpus has tens +// of channels; this only matters to a process that inspects many corpora. +const MEMO_SWEEP_AT = 512; + // Two stats and (at most) one small JSON read. Render-safe: nothing here walks a // directory, so calling it per channel on a listing page costs three syscalls a -// row. +// row — or none, for five seconds after the last answer (see the memo above), +// and none at all on a stalled location. export async function inspectChannelMedia( paths: ChannelMediaPaths, slug: string, config?: Pick<ChannelConfig, "dataDir"> | null, + opts: InspectOptions = {}, ): Promise<ChannelMediaLocation> { const dataDir = channelMediaDir(paths, slug); const configured = @@ -201,6 +311,41 @@ export async function inspectChannelMedia( ? config.dataDir.trim() : undefined; + // THE GATE, before any other call. A channel whose configured target is on a + // location the health probe found not answering is answered from memory: the + // marker read, the link and the target stat below would each hold an I/O + // thread (the stat on the drive itself for as long as the drive takes). With + // a config passed in, this costs no filesystem call at all. + if (configured) { + const stall = stalledLocationForPath(configured); + if (stall) return stalledMediaLocation(dataDir, configured, stall); + } + + const now = opts.now ?? Date.now(); + const key = memoKey(paths, slug, configured); + const memo = mediaMemo(); + if (!opts.fresh) { + const hit = memo.get(key); + if (hit && now - hit.at < CHANNEL_MEDIA_MEMO_MS) return { ...hit.location }; + } + const location = await inspectOnDisk(paths, slug, dataDir, configured); + if (!opts.fresh) { + if (memo.size >= MEMO_SWEEP_AT) { + for (const [k, v] of memo) { + if (now - v.at >= CHANNEL_MEDIA_MEMO_MS) memo.delete(k); + } + } + memo.set(key, { at: now, location }); + } + return { ...location }; +} + +async function inspectOnDisk( + paths: ChannelMediaPaths, + slug: string, + dataDir: string, + configured: string | undefined, +): Promise<ChannelMediaLocation> { const marker = await readRelocationMarker(paths, slug); if (marker) { return { @@ -316,6 +461,10 @@ export async function inspectChannelMedia( // "ok" and "in-place" pass; everything else throws. An in-transition or // inconsistent channel is refused for the same reason an unreachable one is: // the caller would otherwise read a half-populated or empty dir as the truth. +// A stalled one is refused because the work would block on the drive. +// +// ALWAYS FRESH: this is the start-of-work guard, and a remembered "ok" from a +// few seconds ago is not what a job about to read `data/` should be told. // // THIS CHECK HAS A TWIN. `checkChannelReachable` in // `umtool/report-to-video/cues.mjs` repeats the same statuses in plain `.mjs`, @@ -329,7 +478,9 @@ export async function assertChannelMediaReachable( slug: string, config?: Pick<ChannelConfig, "dataDir"> | null, ): Promise<ChannelMediaLocation> { - const location = await inspectChannelMedia(paths, slug, config); + const location = await inspectChannelMedia(paths, slug, config, { + fresh: true, + }); if (location.status === "ok" || location.status === "in-place") { return location; } diff --git a/common/lib/channelMediaHold.ts b/common/lib/channelMediaHold.ts @@ -33,6 +33,7 @@ export const HELD_REASON: Record<ChannelMediaStatus, string> = { unreachable: "its media is not reachable (drive not mounted?)", "in-transition": "a move of its media is in progress or was interrupted", inconsistent: "its data link and its config disagree", + stalled: "its drive is not answering (a stalled disk)", ok: "reachable", "in-place": "reachable", }; diff --git a/common/lib/storageHealth.test.ts b/common/lib/storageHealth.test.ts @@ -0,0 +1,129 @@ +import { beforeEach, test } from "node:test"; +import assert from "node:assert/strict"; +import { + HEALTH_CLEAN_TO_CLEAR, + allLocationHealth, + locationHealth, + notAnsweringText, + pruneLocationHealth, + recordLocationHealth, + resetStorageHealth, + sinceText, + stalledLocation, + stalledLocationForPath, +} from "./storageHealth"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/storageHealth.test.ts +// +// The rules the health state keeps (storageHealth.ts's header): one miss stalls +// at once, two clean answers in a row clear it, a miss in between starts the +// count again, and a new root starts the location over. + +beforeEach(() => resetStorageHealth()); + +const USB = { id: "usb", label: "USB drive", root: "/mnt/usb/media" }; + +test("one missed probe marks the location stalled at once", () => { + assert.equal(recordLocationHealth(USB, "ok", { now: 1_000 })?.to, "ok"); + const t = recordLocationHealth(USB, "stalled", { now: 2_000, cause: "no answer" }); + assert.deepEqual(t, { id: "usb", from: "ok", to: "stalled" }); + const h = locationHealth("usb"); + assert.equal(h?.state, "stalled"); + assert.equal(h?.since, 2_000); + assert.equal(h?.cause, "no answer"); +}); + +test("a location first seen stalled is stalled", () => { + const t = recordLocationHealth(USB, "stalled", { now: 5 }); + assert.deepEqual(t, { id: "usb", from: null, to: "stalled" }); + assert.equal(stalledLocation(USB)?.since, 5); +}); + +test("one clean probe after a stall does not clear it; two in a row do", () => { + assert.equal(HEALTH_CLEAN_TO_CLEAR, 2); + recordLocationHealth(USB, "ok", { now: 0 }); + recordLocationHealth(USB, "stalled", { now: 10 }); + assert.equal(recordLocationHealth(USB, "ok", { now: 20 }), null); + assert.equal(locationHealth("usb")?.state, "stalled"); + // `since` stays the stall's start while it lasts. + assert.equal(locationHealth("usb")?.since, 10); + const t = recordLocationHealth(USB, "ok", { now: 30 }); + assert.deepEqual(t, { id: "usb", from: "stalled", to: "ok" }); + const h = locationHealth("usb"); + assert.equal(h?.state, "ok"); + assert.equal(h?.since, 30); + assert.equal(h?.cause, undefined); +}); + +test("a miss between two clean probes starts the count again", () => { + recordLocationHealth(USB, "stalled", { now: 0 }); + recordLocationHealth(USB, "ok", { now: 1 }); + assert.equal(recordLocationHealth(USB, "stalled", { now: 2 }), null); + // The stall's start is not moved by a second miss. + assert.equal(locationHealth("usb")?.since, 0); + recordLocationHealth(USB, "ok", { now: 3 }); + assert.equal(locationHealth("usb")?.state, "stalled"); + recordLocationHealth(USB, "ok", { now: 4 }); + assert.equal(locationHealth("usb")?.state, "ok"); +}); + +test("an absent answer is a clean one: an unplugged drive answers at once", () => { + recordLocationHealth(USB, "stalled", { now: 0 }); + recordLocationHealth(USB, "absent", { now: 1 }); + const t = recordLocationHealth(USB, "absent", { now: 2 }); + assert.deepEqual(t, { id: "usb", from: "stalled", to: "absent" }); + assert.equal(stalledLocation(USB), null); +}); + +test("a re-pointed root starts the location over", () => { + recordLocationHealth(USB, "stalled", { now: 0 }); + const moved = { ...USB, root: "/mnt/elsewhere/media" }; + // The stall belonged to the old root: the new root's lookup is not stalled, + // and neither is the old one's (the entry now describes the new root). + assert.equal(stalledLocation(moved), null); + const t = recordLocationHealth(moved, "ok", { now: 1 }); + assert.deepEqual(t, { id: "usb", from: null, to: "ok" }); + assert.equal(locationHealth("usb")?.root, "/mnt/elsewhere/media"); + assert.equal(stalledLocation(USB), null); +}); + +test("the path gate matches like locationOfDataDir: longest root, strictly under", () => { + const parent = { id: "parent", label: "Parent", root: "/mnt/p" }; + const child = { id: "child", label: "Child", root: "/mnt/p/archive" }; + recordLocationHealth(parent, "stalled", { now: 0 }); + recordLocationHealth(child, "ok", { now: 0 }); + // A channel on the nested location answers for that location, not its parent. + assert.equal(stalledLocationForPath("/mnt/p/archive/chan/data"), null); + assert.equal(stalledLocationForPath("/mnt/p/chan/data")?.id, "parent"); + // The root itself is not "under" it, and an unrelated path is on nothing. + assert.equal(stalledLocationForPath("/mnt/p"), null); + assert.equal(stalledLocationForPath("/elsewhere/chan/data"), null); + assert.equal(stalledLocationForPath(""), null); +}); + +test("pruning drops locations no longer configured", () => { + recordLocationHealth(USB, "stalled", { now: 0 }); + recordLocationHealth({ id: "b", label: "B", root: "/b" }, "ok", { now: 0 }); + pruneLocationHealth(["b"]); + assert.deepEqual(Object.keys(allLocationHealth()), ["b"]); + assert.equal(stalledLocationForPath("/mnt/usb/media/x/data"), null); +}); + +test("the state is one map per process, on globalThis", () => { + recordLocationHealth(USB, "stalled", { now: 0 }); + // A second module copy reads the same object (the house pattern). + assert.equal( + globalThis.__yttStorageHealth__?.byId.get("usb")?.state, + "stalled", + ); +}); + +test("the words: since a time today, or a date and time", () => { + const now = new Date(2026, 8, 29, 12, 0).getTime(); + const today = new Date(2026, 8, 29, 11, 35).getTime(); + const yesterday = new Date(2026, 8, 28, 23, 5).getTime(); + assert.equal(sinceText(today, now), "since 11:35"); + assert.equal(sinceText(yesterday, now), "since 2026-09-28 23:05"); + assert.equal(notAnsweringText({ since: today }, now), "not answering since 11:35"); +}); diff --git a/common/lib/storageHealth.ts b/common/lib/storageHealth.ts @@ -0,0 +1,231 @@ +import { locationOfDataDir, type StorageLocation } from "./storageLocations"; + +// IS A STORAGE LOCATION'S DRIVE ANSWERING RIGHT NOW — the in-memory answer every +// page and poll asks before it touches the drive. +// +// A drive can be mounted and still not answer. An SMR disk in a USB enclosure +// under a long write stalls, the enclosure resets, and every filesystem call +// that has to reach the disk blocks for about 30 seconds. Node runs those calls +// on libuv's thread pool (four threads by default), so four of them block the +// whole editor: no page, no poll and no job log answers until the disk does. +// "Not mounted" does not describe that (a `stat` there does not fail, it hangs), +// and no in-process call can find it out without paying the hang itself. +// +// SO A CHILD PROCESS ASKS, AND THIS MODULE REMEMBERS WHAT IT SAID. +// `probeLocationHealth` (`storageVolumes.ts`) runs `stat` on the location's root +// as a subprocess raced against a 3 s timer, and the storage watch +// (`controller/storageWatch.ts`) asks every 15 s. A child stuck in the kernel +// blocks nothing of ours: the timer answers `stalled` and the child is left to +// finish on its own. Everything that would touch the drive in-process then asks +// `stalledLocationForPath` first and, on a stalled location, answers without the +// call: `inspectChannelMedia` reports `stalled`, the free-space column reads +// "—", the probe reads "Not answering", the recency layer skips the tail read. +// +// THE RULES, which are what keep a flaky drive from flapping the UI: +// - ONE failed probe marks the location `stalled` at once. A drive that did +// not answer in 3 s will not answer the next page either, and every page +// that asks costs a thread. +// - TWO consecutive clean probes clear it. A clean probe is one that answered +// in time, whether the root was there (`ok`) or not (`absent`: an unmounted +// drive answers ENOENT at once, and that is a different problem, which +// `inspectChannelMedia` already reports). One clean answer in the middle of +// a reset loop is not recovery. +// - A ROOT CHANGE (a re-point) starts the location over: the old root's stall +// says nothing about the new one. +// +// IN MEMORY ONLY. A stall is a fact about this minute, not about the corpus; +// persisting it would outlive the reset loop that caused it. A restarted +// process starts with no stall, and the first probe (at arm time) re-learns it. +// +// ONE MAP PER PROCESS, NOT PER MODULE COPY. The watch that probes is armed from +// `editor/instrumentation.ts`, and the pages that read are another bundle +// layer; Next can load this module once for each. The map lives on +// `globalThis`, the house pattern (`lib/jsonFile-server.ts`, `jobs/registry.ts`), +// so the writer and every reader see one map. +// +// PURE of I/O, and deliberately without execa, so `lib/channelMedia.ts` can ask +// it without pulling a subprocess module into everything that imports that. + +// One probe's answer, and a location's state. +export type LocationHealthState = + // The root answered in time and is a directory. + | "ok" + // The root did not answer within the probe's budget. + | "stalled" + // The root answered in time, and is not a directory (not mounted, or gone). + | "absent"; + +export type LocationHealth = { + id: string; + label: string; + root: string; + state: LocationHealthState; + // When the current state began, ms since epoch. + since: number; + // When the last probe (or observation) was recorded. + checkedAt: number; + // Consecutive clean answers since the location was marked stalled. It clears + // at HEALTH_CLEAN_TO_CLEAR. + cleanStreak: number; + // What did not answer, for a stalled location: the probe's own words. + cause?: string; +}; + +// The probe's budget. A `stat` of a directory on a healthy disk answers in +// microseconds; three seconds is a thousand times that, and short enough that +// a page asking during a stall has not waited long. +export const HEALTH_PROBE_TIMEOUT_MS = 3_000; +// How often the watch asks. Short, because the gate is only as current as the +// last answer: a stall that began just after a probe costs every page that +// touches the drive until the next one. +export const HEALTH_PROBE_INTERVAL_MS = 15_000; +// Clean answers in a row that clear a stall. +export const HEALTH_CLEAN_TO_CLEAR = 2; + +// The one wording of the state, for every surface that shows it. +export const NOT_ANSWERING = "drive not answering"; + +type HealthState = { byId: Map<string, LocationHealth> }; + +declare global { + // eslint-disable-next-line no-var + var __yttStorageHealth__: HealthState | undefined; +} + +function healthState(): HealthState { + if (!globalThis.__yttStorageHealth__) { + globalThis.__yttStorageHealth__ = { byId: new Map() }; + } + return globalThis.__yttStorageHealth__; +} + +// Test seam, and the escape hatch for a process that wants to forget. +export function resetStorageHealth(): void { + healthState().byId.clear(); +} + +export type HealthTransition = { + id: string; + from: LocationHealthState | null; + to: LocationHealthState; +}; + +// Record one answer about one location. Returns the transition when the state +// changed, null when it did not. +export function recordLocationHealth( + loc: Pick<StorageLocation, "id" | "label" | "root">, + answer: LocationHealthState, + opts: { now?: number; cause?: string } = {}, +): HealthTransition | null { + const now = opts.now ?? Date.now(); + const map = healthState().byId; + const prev = map.get(loc.id); + const label = loc.label || loc.id; + // First sighting, or the root moved under the same id: start over. + if (!prev || prev.root !== loc.root) { + map.set(loc.id, { + id: loc.id, + label, + root: loc.root, + state: answer, + since: now, + checkedAt: now, + cleanStreak: 0, + ...(answer === "stalled" && opts.cause ? { cause: opts.cause } : {}), + }); + return { id: loc.id, from: prev && prev.root === loc.root ? prev.state : null, to: answer }; + } + prev.label = label; + prev.checkedAt = now; + if (answer === "stalled") { + prev.cleanStreak = 0; + if (prev.state === "stalled") return null; + const from = prev.state; + prev.state = "stalled"; + prev.since = now; + if (opts.cause) prev.cause = opts.cause; + return { id: loc.id, from, to: "stalled" }; + } + if (prev.state === "stalled") { + prev.cleanStreak += 1; + if (prev.cleanStreak < HEALTH_CLEAN_TO_CLEAR) return null; + prev.state = answer; + prev.since = now; + prev.cleanStreak = 0; + delete prev.cause; + return { id: loc.id, from: "stalled", to: answer }; + } + if (prev.state === answer) return null; + const from = prev.state; + prev.state = answer; + prev.since = now; + return { id: loc.id, from, to: answer }; +} + +// Drop every location that is no longer configured. +export function pruneLocationHealth(liveIds: Iterable<string>): void { + const keep = new Set(liveIds); + const map = healthState().byId; + for (const id of [...map.keys()]) if (!keep.has(id)) map.delete(id); +} + +export function locationHealth(id: string): LocationHealth | undefined { + const h = healthState().byId.get(id); + return h ? { ...h } : undefined; +} + +// Every location's last answer, by id. A copy: the caller may keep it. +export function allLocationHealth(): Record<string, LocationHealth> { + const out: Record<string, LocationHealth> = {}; + for (const [id, h] of healthState().byId) out[id] = { ...h }; + return out; +} + +// THE GATE. The stalled location `p` is on, or null. `p` is a channel's +// `dataDir` (or anything under a location's root); the match is the one +// `locationOfDataDir` makes, longest root first, so a channel on a nested +// location answers for that location and not its parent. +export function stalledLocationForPath(p: string): LocationHealth | null { + const map = healthState().byId; + if (map.size === 0 || !p) return null; + const entries = [...map.values()]; + const loc = locationOfDataDir( + p, + entries.map((h) => ({ id: h.id, label: h.label, root: h.root, autoRepoint: false })), + ); + if (!loc) return null; + const h = map.get(loc.id); + return h && h.state === "stalled" ? { ...h } : null; +} + +// The location with this id, when it is stalled AND still at this root. The +// gate for a caller that holds a location rather than a channel path +// (`probeLocation`, `volumeFreeBytes`). +export function stalledLocation( + loc: Pick<StorageLocation, "id" | "root">, +): LocationHealth | null { + const h = healthState().byId.get(loc.id); + return h && h.state === "stalled" && h.root === loc.root ? { ...h } : null; +} + +// "since 11:35" in the server's clock, or "since 2026-09-28 11:35" when the +// stall began on another day than `now`. +export function sinceText(since: number, now: number = Date.now()): string { + const d = new Date(since); + const n = new Date(now); + const pad = (x: number) => String(x).padStart(2, "0"); + const hm = `${pad(d.getHours())}:${pad(d.getMinutes())}`; + const sameDay = + d.getFullYear() === n.getFullYear() && + d.getMonth() === n.getMonth() && + d.getDate() === n.getDate(); + return sameDay + ? `since ${hm}` + : `since ${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())} ${hm}`; +} + +// "not answering since 11:35" — the one line /storage, the /channels volume bar +// and the channel's Storage panel show for a stalled location. +export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number): string { + return `not answering ${sinceText(h.since, now)}`; +} diff --git a/common/lib/storageHealthProbe.test.ts b/common/lib/storageHealthProbe.test.ts @@ -0,0 +1,112 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import path from "node:path"; +import { tmpdir } from "node:os"; +import { chmod, mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { probeLocationHealth } from "./storageVolumes"; +import { HEALTH_PROBE_TIMEOUT_MS } from "./storageHealth"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/storageHealthProbe.test.ts +// +// THE HEALTH PROBE: `stat` on a location's root as a CHILD PROCESS, raced +// against a timer. No test stalls a real drive: a stalled `stat` is a fake +// binary that never answers (a child blocked in the kernel looks the same from +// here — it does not exit), and the case that matters is that the answer +// arrives on the timer while the child is still running. + +async function withDir(fn: (dir: string) => Promise<void>): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-health-")); + try { + await fn(dir); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +// A fake `stat`: sleeps `sleepMs` (a blocking sleep, so it is simply not +// answering), then prints `out` and exits `code`. +async function fakeStat( + dir: string, + opts: { sleepMs?: number; out?: string; code?: number }, +): Promise<string> { + const bin = path.join(dir, "fake-stat.mjs"); + await writeFile( + bin, + `#!/usr/bin/env node +if (${opts.sleepMs ?? 0} > 0) { + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, ${opts.sleepMs ?? 0}); +} +process.stdout.write(${JSON.stringify(opts.out ?? "")} + "\\n"); +process.exit(${opts.code ?? 0}); +`, + ); + await chmod(bin, 0o755); + return bin; +} + +test("the real stat: a directory is ok, a missing path and a file are absent", async () => { + await withDir(async (dir) => { + const root = path.join(dir, "media"); + await mkdir(root); + await writeFile(path.join(dir, "a-file"), "x"); + assert.equal(await probeLocationHealth({ root }), "ok"); + assert.equal(await probeLocationHealth({ root: path.join(dir, "gone") }), "absent"); + assert.equal(await probeLocationHealth({ root: path.join(dir, "a-file") }), "absent"); + assert.equal(await probeLocationHealth({ root: " " }), "absent"); + }); +}); + +test("a child that never answers is 'stalled' on the timer, without waiting for it", async () => { + await withDir(async (dir) => { + // Twenty seconds asleep: if the probe waited for the child, this test would + // take that long. + const statBin = await fakeStat(dir, { sleepMs: 20_000, out: "directory" }); + const started = Date.now(); + const answer = await probeLocationHealth( + { root: dir }, + { statBin, timeoutMs: 400 }, + ); + const took = Date.now() - started; + assert.equal(answer, "stalled"); + assert.ok(took >= 400, `answered after ${took} ms, before the timer`); + assert.ok(took < 3_000, `answered after ${took} ms`); + }); +}); + +test("the default budget is 3 s", async () => { + assert.equal(HEALTH_PROBE_TIMEOUT_MS, 3_000); + await withDir(async (dir) => { + const statBin = await fakeStat(dir, { sleepMs: 20_000, out: "directory" }); + const started = Date.now(); + assert.equal(await probeLocationHealth({ root: dir }, { statBin }), "stalled"); + const took = Date.now() - started; + assert.ok(took >= 3_000 && took < 6_000, `answered after ${took} ms`); + }); +}); + +test("an answer inside the budget is taken as given", async () => { + await withDir(async (dir) => { + const slowDir = await fakeStat(dir, { sleepMs: 100, out: "directory" }); + assert.equal( + await probeLocationHealth({ root: dir }, { statBin: slowDir, timeoutMs: 2_500 }), + "ok", + ); + const aFile = await fakeStat(dir, { out: "regular file" }); + assert.equal(await probeLocationHealth({ root: dir }, { statBin: aFile }), "absent"); + const missing = await fakeStat(dir, { code: 1 }); + assert.equal(await probeLocationHealth({ root: dir }, { statBin: missing }), "absent"); + }); +}); + +test("a stat that cannot be started fails open: ok, never stalled", async () => { + await withDir(async (dir) => { + assert.equal( + await probeLocationHealth( + { root: dir }, + { statBin: path.join(dir, "no-such-stat") }, + ), + "ok", + ); + }); +}); diff --git a/common/lib/storageVolumes.ts b/common/lib/storageVolumes.ts @@ -4,6 +4,11 @@ import { execa } from "execa"; import type { Paths } from "./paths"; import { getFreeBytes } from "./diskSpace"; import type { StorageLocation, StorageVolume } from "./storageLocations"; +import { + HEALTH_PROBE_TIMEOUT_MS, + stalledLocation, + type LocationHealthState, +} from "./storageHealth"; // STORAGE VOLUME PROBES — is this location's disk here, and if not, where? // @@ -33,7 +38,11 @@ export type StorageLocationStatus = | "absent" // The root is not there and we have no identity to look for — nothing to say // beyond "that path does not exist". - | "missing"; + | "missing" + // The drive did not answer the last health probe (`lib/storageHealth.ts`), so + // this probe did not ask: an in-process `stat` there would hold one of the + // process's I/O threads until the drive answered. No identity, no free space. + | "stalled"; // `known: false` is the fail-open answer and is NOT a problem report: it means // the probe could not ask (no findmnt, a container, a timeout), not that the @@ -226,6 +235,13 @@ export async function probeLocation( const timeoutMs = opts.findmntTimeoutMs ?? FINDMNT_TIMEOUT_MS; const root = loc.root.trim(); + // A DRIVE THAT IS NOT ANSWERING IS NOT ASKED. The `stat` and `statfs` below + // run in-process, and on a stalled disk each holds an I/O thread until the + // drive comes back; the health probe (below) already asked out of process. + if (stalledLocation(loc)) { + return { status: "stalled", identity: { known: false } }; + } + // AVAILABILITY IS `stat`, AND ONLY `stat`. A root that is a directory is // available even when every identity probe below fails — see the header. let isDir = false; @@ -304,6 +320,82 @@ export async function probeLocation( } } +// --------------------------------------------------------------------------- +// The health probe: is the drive ANSWERING, asked from a child process +// --------------------------------------------------------------------------- +// +// `probeLocation` answers "is the disk here" with an in-process `stat`, which is +// right until the disk is here and not answering: then that `stat` does not +// fail, it waits — for as long as the drive takes, on one of the four threads +// libuv runs every filesystem call on. A child process waiting in the kernel +// holds none of them. So this runs `stat` on the root as a SUBPROCESS and races +// it against a timer, and the timer's answer is `stalled`. +// +// THE TIMER WINS, AND NOTHING WAITS FOR THE CHILD. A process in uninterruptible +// I/O cannot be killed until the I/O returns, so awaiting its exit (which is +// what execa's own `timeout` does) would hand the wait straight back to us. +// The child is sent SIGKILL, which it takes when the drive lets it, and its +// promise settles into a handler nobody awaits. +// +// FAILS OPEN, like everything else here: a `stat` that could not be started +// (no binary) is "could not ask", which is `ok`, never `stalled`. Only a child +// that started and did not answer in time is a stall. + +export type HealthProbeOptions = { + timeoutMs?: number; + // Test seam: the binary to run (it is called as `<bin> -L -c %F -- <root>`). + statBin?: string; +}; + +export type LocationHealthProbe = ( + loc: Pick<StorageLocation, "id" | "root">, +) => Promise<LocationHealthState>; + +export async function probeLocationHealth( + loc: Pick<StorageLocation, "root">, + opts: HealthProbeOptions = {}, +): Promise<LocationHealthState> { + const root = loc.root.trim(); + if (root === "") return "absent"; + const timeoutMs = opts.timeoutMs ?? HEALTH_PROBE_TIMEOUT_MS; + let child: ReturnType<typeof execa>; + try { + child = execa(opts.statBin ?? "stat", ["-L", "-c", "%F", "--", root], { + buffer: true, + reject: false, + stdin: "ignore", + }); + } catch { + return "ok"; + } + const answered: Promise<LocationHealthState> = child.then( + (res) => { + // No exit code: the binary never ran (or was killed after the race was + // already decided). Could not ask. + if (typeof res.exitCode !== "number") return "ok"; + if (res.exitCode !== 0) return "absent"; + const out = typeof res.stdout === "string" ? res.stdout.trim() : ""; + return out === "directory" ? "ok" : "absent"; + }, + () => "ok", + ); + let timer: ReturnType<typeof setTimeout> | undefined; + const timedOut = new Promise<LocationHealthState>((resolve) => { + timer = setTimeout(() => resolve("stalled"), timeoutMs); + timer.unref?.(); + }); + const answer = await Promise.race([answered, timedOut]); + if (timer) clearTimeout(timer); + if (answer === "stalled") { + try { + child.kill("SIGKILL"); + } catch { + /* already gone */ + } + } + return answer; +} + // udisksctl availability, memoised per binary path. Same shape as a digest // app's probe (`digestApps.ts` claudeCode.probe): `--version`, reject:false, // short timeout. Memoised because /storage asks once per render and the answer @@ -401,6 +493,12 @@ export async function probeLocationMemo( opts: ProbeOptions & { refresh?: boolean; now?: number } = {}, ): Promise<MemoizedProbe> { const now = opts.now ?? Date.now(); + // The stall is asked before the memo, so a remembered "available" from a few + // seconds ago does not outlive the drive's answer. Not remembered either: + // the health probe's state is already the memory. + if (stalledLocation(loc)) { + return { status: "stalled", identity: { known: false }, probedAt: now }; + } const hit = probeMemo.get(loc.id); if ( !opts.refresh && diff --git a/common/views/pipeline/stageStatus.ts b/common/views/pipeline/stageStatus.ts @@ -601,7 +601,8 @@ export function computeStageStatuses( mediaStatus === "in-transition" ? "running" : mediaStatus === "unreachable" || - mediaStatus === "inconsistent" + mediaStatus === "inconsistent" || + mediaStatus === "stalled" ? "danger" : "neutral", }; diff --git a/common/views/storage.ts b/common/views/storage.ts @@ -217,6 +217,7 @@ export const STORAGE_STATUS_LABEL: Record<StorageLocationStatus, string> = { unmounted: "Not mounted", absent: "Not attached", missing: "Missing", + stalled: "Not answering", }; export const REPOINT_JOB_KIND = "repoint-storage-location"; diff --git a/editor/app/components/MediaLocationBadge.tsx b/editor/app/components/MediaLocationBadge.tsx @@ -10,18 +10,19 @@ import type { ChannelRowMedia } from "yt-dlp-transcript-common/views/channelRow" // the filesystem; the inspect() call that produces the location happens on the // server, once per row, and only its result travels. // -// WHAT THE FIVE STATUSES LOOK LIKE, and why there are only three appearances: +// WHAT THE SIX STATUSES LOOK LIKE, and why there are only three appearances: // // in-place → NOTHING. The overwhelming majority of channels are in place, // and a badge on every row saying "normal" is noise that makes // the two that matter harder to see, not easier. // ok → neutral. Relocated and reachable is a fact worth stating (the // bytes are not on the corpus disk) but it is not a problem. -// everything → red. unreachable, in-transition and inconsistent are all -// else "do not trust what this channel's dirs say right now": the +// everything → red. unreachable, in-transition, inconsistent and stalled are +// else all "do not trust what this channel's dirs say right now": the // first because the drive is not mounted, the second because a // move is half-done, the third because disk and config disagree -// and nothing here is willing to guess which one is right. +// and nothing here is willing to guess which one is right, the +// fourth because the drive is mounted and not answering. // // The `detail` string is the operator's prose from inspect() — the drive path, // the phase, the disagreement — and it goes on `title` so a row badge carries @@ -52,6 +53,7 @@ const LABELS: Record<ChannelMediaStatus, string | null> = { unreachable: "Media unreachable", "in-transition": "Media moving", inconsistent: "Media inconsistent", + stalled: "Media not answering", }; // The one-word state, for the compact rendering of a NAMED location: "on @@ -63,6 +65,7 @@ const SHORT_STATUS: Record<ChannelMediaStatus, string | null> = { unreachable: "unreachable", "in-transition": "moving", inconsistent: "inconsistent", + stalled: "not answering", }; // Null means "draw nothing" — an in-place channel, or no location at all (a