commit 7497bfbaa8898ad9f8de5823b13694025716812d
parent e77ad9eb02faf162f24c95bb14f702187e308cac
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 29 Sep 2026 22:52:03 -0400
common: the relocation marker is read before the stall gate; a 3 s watchdog on every gated call, and at most four calls in flight per location
Rulings Q2 and Q1(b).
- inspectChannelMedia reads the channel's `.relocating.json` first (it is in
the channel dir, on the corpus disk), so a channel mid-move on a stalled
drive reads `in-transition`; the gate then covers the link and the target.
A remembered `in-transition` is given as it is; any other remembered answer
is gated first. A stall is still not remembered.
- `onDrive(where, call)` (lib/storageHealth.ts): refused at once on a stalled
location; otherwise raced against DRIVE_CALL_BUDGET_MS (3 s). A call that
has not answered marks its location `stalled` (since now, cause "a read in
the editor did not answer within 3 s") and throws DriveNotAnsweringError;
the call is left to settle. At most DRIVE_CALLS_IN_FLIGHT (4) calls per
location are in flight through it; the rest wait in a queue of its own and
are refused without a call when the location stalls, so a stalled drive
holds at most four of the pool's threads. A slot is released when its call
really returns. A path on no known location is raced but marks nothing; a
probe of another root under a location's id does not rewrite that location.
- Through it: inspect's target stat, probeLocation's stat and statfs,
volumeFreeBytes' root and mountpoint stats and statfs, readChannelStat's
walk (the data/ readdir, then each video directory as one call; null when
the drive stops answering), the recency tail reads (not remembered as
misses), the move-root stat, the saved-video store's stat, and
listSavedVideos' reads when a page asks. registerLocationHealth gives every
configured location an entry without an answer, so the watchdog can find a
channel's location. listChannelStatsFromDisk (the batch jobs' walk) passes no
drive and is unchanged.
Tests: onDrive (pass-through, a never-settling call stalls on the timer and
marks the location, refusal without a call, the four-slot cap and its queue
refused on a stall, a freed slot runs a waiter, an unknown path, a candidate
root, the 3 s default); storageStall (marker first, and every gated caller
against a spy whose drive call never settles: stalled within the race, the
location marked, no further call on the drive until cleared).
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
11 files changed, 860 insertions(+), 161 deletions(-)
diff --git a/common/controller/channels.ts b/common/controller/channels.ts
@@ -17,6 +17,7 @@ import {
} from "../lib/videoStatus";
import { loadDigest } from "../lib/digest-server";
import { channelMediaStall, readRelocationMarker } from "../lib/channelMedia";
+import { isDriveNotAnswering, onDrive } from "../lib/storageHealth";
// 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
@@ -87,39 +88,56 @@ async function hasDigestWithItems(videoDir: string): Promise<boolean> {
);
}
-async function countDataFiles(dataDir: string): Promise<{
+// `drive` is the channel's configured target (`config.dataDir`) when its media
+// is on another drive. Then every read goes through `onDrive`: none while that
+// location is stalled, at most four in flight on it, and one that has not
+// answered in 3 s marks it stalled — and the walk answers null ("the drive did
+// not answer") instead of counts. The rest of the walk is refused without a
+// call. An in-place channel's walk is on the corpus disk and is not wrapped.
+async function countDataFiles(
+ dataDir: string,
+ drive?: string,
+): Promise<{
videos: number;
transcripts: number;
downloads: number;
digests: number;
-}> {
+} | null> {
+ const through = <T>(call: () => Promise<T>): Promise<T> =>
+ drive ? onDrive(drive, call) : call();
let dirs: Dirent[];
try {
- dirs = await readdir(dataDir, { withFileTypes: true });
- } catch {
+ dirs = await through(() => readdir(dataDir, { withFileTypes: true }));
+ } catch (err) {
+ if (isDriveNotAnswering(err)) return null;
return { videos: 0, transcripts: 0, downloads: 0, digests: 0 };
}
const videoDirs = dirs.filter((d) => d.isDirectory());
// Bounded: the largest channel has 11,224 video dirs and this used to open
// them all at once.
- const flags = await mapConcurrent(
- videoDirs,
- VIDEO_READ_CONCURRENCY,
- async (d) => {
- const dir = path.join(dataDir, d.name);
- const files = await readVideoFiles(dir);
- return {
- transcript: isVideoTranscribed(files),
- download: isVideoDownloaded(files),
- // Only transcribed videos can carry a digest, so the sidecar read is
- // skipped for the rest — the same conditional per-video sidecar-read
- // pattern channelSnapshot.ts uses for coverage and VTT provenance.
- digest: isVideoTranscribed(files)
- ? await hasDigestWithItems(dir)
- : false,
- };
- },
- );
+ let flags: Array<{ transcript: boolean; download: boolean; digest: boolean }>;
+ try {
+ flags = await mapConcurrent(videoDirs, VIDEO_READ_CONCURRENCY, (d) =>
+ // One video directory's few reads are one call through the watchdog.
+ through(async () => {
+ const dir = path.join(dataDir, d.name);
+ const files = await readVideoFiles(dir);
+ return {
+ transcript: isVideoTranscribed(files),
+ download: isVideoDownloaded(files),
+ // Only transcribed videos can carry a digest, so the sidecar read is
+ // skipped for the rest — the same conditional per-video sidecar-read
+ // pattern channelSnapshot.ts uses for coverage and VTT provenance.
+ digest: isVideoTranscribed(files)
+ ? await hasDigestWithItems(dir)
+ : false,
+ };
+ }),
+ );
+ } catch (err) {
+ if (isDriveNotAnswering(err)) return null;
+ throw err;
+ }
let transcripts = 0;
let downloads = 0;
let digests = 0;
@@ -196,7 +214,11 @@ export async function readChannelStat(
// 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"));
+ const counts = await countDataFiles(
+ path.join(channelDir, "data"),
+ config.dataDir?.trim() || undefined,
+ );
+ if (!counts) return null;
return {
slug,
config,
@@ -271,7 +293,14 @@ export async function listChannelStatsFromDisk(
if (!config) continue;
const channelDir = path.join(paths.channelsDir, slug);
const dataDir = path.join(channelDir, "data");
- const counts = await countDataFiles(dataDir);
+ // A batch job's ground truth: no drive passed, so nothing is raced and the
+ // answer is never null.
+ const counts = (await countDataFiles(dataDir)) ?? {
+ videos: 0,
+ transcripts: 0,
+ downloads: 0,
+ digests: 0,
+ };
out.push({
slug,
config,
diff --git a/common/controller/recencyIndex.ts b/common/controller/recencyIndex.ts
@@ -8,7 +8,11 @@ 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";
+import {
+ isDriveNotAnswering,
+ onDrive,
+ stalledLocationForPath,
+} from "../lib/storageHealth";
// Upload-date lookup for the auto-queue's "newest first" ordering.
//
@@ -192,17 +196,21 @@ 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.
+// `drives` maps a channel whose media is on another drive to its configured
+// target. Its reads go through `onDrive`: none while that drive's location is
+// stalled (lib/storageHealth.ts) — each would hold an I/O thread until the drive
+// came back, 32 at a time — at most four in flight on it, and one that has not
+// answered in 3 s marks it stalled. An id not read for that reason is NOT
+// memoized as a miss: it falls through to layers 3 and 4 for now, and a later
+// refresh with the drive answering reads it.
+const NOT_READ = Symbol("not read: the drive is not answering");
+
async function datesFromMetadata(
paths: Paths,
owner: ReadonlyMap<string, string>,
wanted: Set<string>,
out: Map<string, RecencyKey>,
- stalled: ReadonlySet<string> = new Set(),
+ drives: ReadonlyMap<string, string> = new Map(),
): Promise<void> {
const todo: string[] = [];
for (const id of wanted) {
@@ -215,7 +223,9 @@ async function datesFromMetadata(
continue;
}
const slug = owner.get(id);
- if (slug === undefined || stalled.has(slug)) continue;
+ if (slug === undefined) continue;
+ const drive = drives.get(slug);
+ if (drive && stalledLocationForPath(drive)) continue;
todo.push(id);
if (todo.length >= TAIL_READS_PER_BUILD) {
if (!tailCapLogged) {
@@ -229,19 +239,31 @@ async function datesFromMetadata(
}
}
if (todo.length === 0) return;
- const dates = await mapConcurrent(todo, TAIL_READ_CONCURRENCY, (id) =>
- readTailUploadDate(
- path.join(
+ const dates = await mapConcurrent(
+ todo,
+ TAIL_READ_CONCURRENCY,
+ async (id): Promise<string | null | typeof NOT_READ> => {
+ const slug = owner.get(id) as string;
+ const file = path.join(
paths.channelsDir,
- owner.get(id) as string,
+ slug,
"data",
id,
"metadata.info.json",
- ),
- ),
+ );
+ const drive = drives.get(slug);
+ if (!drive) return readTailUploadDate(file);
+ try {
+ return await onDrive(drive, () => readTailUploadDate(file));
+ } catch (err) {
+ if (isDriveNotAnswering(err)) return NOT_READ;
+ throw err;
+ }
+ },
);
for (const [i, id] of todo.entries()) {
const date = dates[i];
+ if (date === NOT_READ) continue;
if (tailMemo.size < TAIL_MEMO_CAP) tailMemo.set(id, date);
if (!date) continue;
out.set(id, { key: date, estimated: false });
@@ -323,7 +345,7 @@ export type BuildRecencyKeysArgs = {
paths: Paths;
// The channels whose work is in play — the runner's own channel-meta list.
// 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).
+ // are on another drive (their reads go through the stall watchdog).
meta: ReadonlyArray<{
slug: string;
config?: Pick<ChannelConfig, "dataDir"> | null;
@@ -483,10 +505,12 @@ export async function buildRecencyKeys({
// auto-transcribe, whose entire candidate set is by definition absent from the
// transcript index.
if (owner && missing.size > 0) {
- const stalled = new Set(
- meta.filter((m) => channelMediaStall(m.config)).map((m) => m.slug),
- );
- await datesFromMetadata(paths, owner, missing, out, stalled);
+ const drives = new Map<string, string>();
+ for (const m of meta) {
+ const dir = m.config?.dataDir?.trim();
+ if (dir) drives.set(m.slug, dir);
+ }
+ await datesFromMetadata(paths, owner, missing, out, drives);
}
// Layer 3: playlist interpolation, for videos with nothing on disk at all.
diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts
@@ -40,6 +40,8 @@ import {
import { getFreeBytes } from "../lib/diskSpace";
import {
NOT_ANSWERING,
+ isDriveNotAnswering,
+ onDrive,
sinceText,
stalledLocation,
} from "../lib/storageHealth";
@@ -298,10 +300,19 @@ export async function relocationRootPresenceProblem(
);
}
+ // Through the watchdog: a stat that has not answered in 3 s marks the
+ // location stalled and refuses the same way.
let isDir = false;
try {
- isDir = (await stat(r)).isDirectory();
- } catch {
+ isDir = (await onDrive(named ?? r, () => stat(r))).isDirectory();
+ } catch (err) {
+ if (isDriveNotAnswering(err)) {
+ return (
+ `The destination root ${r}${where}: ${NOT_ANSWERING} ` +
+ `(${err.health ? sinceText(err.health.since) : "a stat of it did not answer"}). ` +
+ `Wait for it to answer, or check the drive.`
+ );
+ }
isDir = false;
}
if (!isDir) {
diff --git a/common/controller/relocateSavedVideos.ts b/common/controller/relocateSavedVideos.ts
@@ -21,6 +21,8 @@ import {
import type { RelocationMarker, RelocationPhase } from "../lib/channelMedia";
import {
NOT_ANSWERING,
+ isDriveNotAnswering,
+ onDrive,
sinceText,
stalledLocation,
} from "../lib/storageHealth";
@@ -174,7 +176,20 @@ export async function inspectSavedVideosStore(
detail: `${NOT_ANSWERING} (location "${stall.label}", ${sinceText(stall.since)})`,
};
}
- if (!(await isDirectory(target))) {
+ let reachable: boolean;
+ try {
+ reachable = await onDrive(loc ?? target, () => isDirectory(target));
+ } catch (err) {
+ if (!isDriveNotAnswering(err)) throw err;
+ return {
+ dir,
+ locationId,
+ target,
+ status: "unreachable",
+ detail: err.message,
+ };
+ }
+ if (!reachable) {
return {
dir,
locationId,
diff --git a/common/controller/savedVideoInventory.ts b/common/controller/savedVideoInventory.ts
@@ -4,7 +4,11 @@ import type { Paths } from "../lib/paths";
import { mapConcurrent } from "../lib/concurrency";
import { loadSavedVideo } from "../lib/savedVideo-server";
import { savedVideoPath, type SavedVideoPointer } from "../lib/savedVideo";
-import { stalledLocationForPath } from "../lib/storageHealth";
+import {
+ isDriveNotAnswering,
+ onDrive,
+ stalledLocationForPath,
+} from "../lib/storageHealth";
const { pathExists, readdir, readlink } = fs;
@@ -44,7 +48,9 @@ async function listChannelSlugs(paths: Paths): Promise<string[]> {
// (lib/storageHealth.ts) is skipped — its pointers are one read per video dir,
// each of which would wait on the drive — and its slug is pushed there, so the
// page can say which channels it did not read. The link is read, not followed:
-// it is on the corpus disk. Omitted (the backup job), nothing is skipped.
+// it is on the corpus disk. A relocated channel's reads then go through
+// `onDrive`'s watchdog, so a drive that stops answering mid-list is skipped and
+// named the same way. Omitted (the backup job), nothing is skipped or raced.
export async function listSavedVideos({
paths,
channelSlug,
@@ -63,32 +69,43 @@ export async function listSavedVideos({
CHANNEL_CONCURRENCY,
async (slug): Promise<SavedVideoEntry[]> => {
const dataDir = path.join(paths.channelsDir, slug, "data");
+ let drive = "";
if (notAnswering) {
- const target = await readlink(dataDir).catch(() => "");
- if (target && stalledLocationForPath(target)) {
+ drive = await readlink(dataDir).catch(() => "");
+ if (drive && stalledLocationForPath(drive)) {
notAnswering.push(slug);
return [];
}
}
- if (!(await pathExists(dataDir))) return [];
- const ids = await readdir(dataDir).catch(() => [] as string[]);
- const entries = await mapConcurrent(
- ids,
- VIDEO_CONCURRENCY,
- async (videoId): Promise<SavedVideoEntry | null> => {
- const videoDir = path.join(dataDir, videoId);
- const pointer = await loadSavedVideo(videoDir);
- if (!pointer) return null;
- return {
- slug,
- videoId,
- videoDir,
- storedPath: savedVideoPath(pointer),
- pointer,
- };
- },
- );
- return entries.filter((e): e is SavedVideoEntry => e !== null);
+ const through = <T>(call: () => Promise<T>): Promise<T> =>
+ drive ? onDrive(drive, call) : call();
+ try {
+ if (!(await through(() => pathExists(dataDir)))) return [];
+ const ids = await through(() =>
+ readdir(dataDir).catch(() => [] as string[]),
+ );
+ const entries = await mapConcurrent(
+ ids,
+ VIDEO_CONCURRENCY,
+ async (videoId): Promise<SavedVideoEntry | null> => {
+ const videoDir = path.join(dataDir, videoId);
+ const pointer = await through(() => loadSavedVideo(videoDir));
+ if (!pointer) return null;
+ return {
+ slug,
+ videoId,
+ videoDir,
+ storedPath: savedVideoPath(pointer),
+ pointer,
+ };
+ },
+ );
+ return entries.filter((e): e is SavedVideoEntry => e !== null);
+ } catch (err) {
+ if (!notAnswering || !isDriveNotAnswering(err)) throw err;
+ notAnswering.push(slug);
+ return [];
+ }
},
);
const out = perChannel.flat();
diff --git a/common/controller/storageLocations.ts b/common/controller/storageLocations.ts
@@ -29,7 +29,11 @@ import {
inspectChannelMedia,
relocatedDataDir,
} from "../lib/channelMedia";
-import { stalledLocation } from "../lib/storageHealth";
+import {
+ isDriveNotAnswering,
+ onDrive,
+ stalledLocation,
+} from "../lib/storageHealth";
import { relocationQueueKey } from "../lib/queueKeys";
import {
runManagedFunction,
@@ -253,66 +257,84 @@ export async function volumeFreeBytes(opts: {
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.
+ // below would hold an I/O thread until it did. "—", like unmounted. The
+ // calls that reach the drive go through `onDrive`'s watchdog; one that
+ // has not answered in 3 s marks the location stalled, and reads "—" too.
if (stalledLocation(loc)) {
out[loc.id] = undefined;
return;
}
- const st = await stat(loc.root).catch(() => null);
- if (!st?.isDirectory()) {
+ try {
+ out[loc.id] = await freeOnLocation(loc);
+ } catch (err) {
+ if (!isDriveNotAnswering(err)) throw err;
out[loc.id] = undefined;
- return;
}
- // THE DIRECTORY EXISTING IS NOT THE DRIVE BEING THERE, and this is the
- // half the stat above could not catch. An unmounted mountpoint is a real,
- // empty directory ON ITS PARENT'S FILESYSTEM — so `getFreeBytes` succeeds
- // and reports the parent volume's free space, which on this machine is
- // the disk the operator is trying to empty. The location would then read
- // "233 GB free" about a platter that is not plugged in.
- //
- // THE BOUNDARY IS TESTED AT THE MOUNTPOINT, NEVER AT THE ROOT, and the
- // first version of this got that wrong in the one way that matters on
- // this machine. A location's root is `join(mountpoint, relPath)` — the
- // production `platter` is `/run/media/<user>/<uuid>/archilyzer-media`, a
- // SUBDIRECTORY of the mountpoint — so the root and its parent are on the
- // same filesystem BY CONSTRUCTION whenever `relPath` is non-empty, and
- // comparing those two devices reported "free space unknown" for a
- // correctly mounted drive. A bind mount reads the same way.
- //
- // Crossing a mount changes the device number, so `stat(mountpoint).dev`
- // against `stat(dirname(mountpoint)).dev` is the honest question, and it
- // is two syscalls — TABLES NEVER PROBE (see the header) rules out asking
- // `probeLocation`, which is up to three subprocesses.
- //
- // Only asked of a location that has learned a `volume.uuid`: that field
- // is the assertion that the root is supposed to be on its own volume. A
- // location on a plain directory (never probed, a container, a
- // subdirectory of the system disk by design) shares its parent's device
- // legitimately, and withholding its free space would be wrong.
- const mountpoint = loc.volume?.uuid
- ? (loc.volume.mountpoint ?? "").trim()
- : "";
- // `/` is its own parent, so a volume mounted at the root has no boundary
- // to test and is trivially there — the process is reading from it.
- if (mountpoint && mountpoint !== path.dirname(mountpoint)) {
- const [atMount, aboveMount] = await Promise.all([
- stat(mountpoint).catch(() => null),
- stat(path.dirname(mountpoint)).catch(() => null),
- ]);
- // Gone entirely, or present as an ordinary directory on the parent
- // filesystem: either way nothing is mounted there.
- if (!atMount || (aboveMount && aboveMount.dev === atMount.dev)) {
- out[loc.id] = undefined;
- return;
- }
- }
- const free = await getFreeBytes(loc.root);
- out[loc.id] = Number.isFinite(free) ? free : undefined;
}),
);
return out;
}
+// One location's free space, for volumeFreeBytes. Throws DriveNotAnsweringError
+// when its drive does not answer.
+async function freeOnLocation(
+ loc: StorageLocation,
+): Promise<number | undefined> {
+ const orNull = <T>(p: Promise<T>): Promise<T | null> =>
+ p.catch((err) => {
+ if (isDriveNotAnswering(err)) throw err;
+ return null;
+ });
+ const st = await orNull(onDrive(loc, () => stat(loc.root)));
+ if (!st?.isDirectory()) return undefined;
+ // THE DIRECTORY EXISTING IS NOT THE DRIVE BEING THERE, and this is the
+ // half the stat above could not catch. An unmounted mountpoint is a real,
+ // empty directory ON ITS PARENT'S FILESYSTEM — so `getFreeBytes` succeeds
+ // and reports the parent volume's free space, which on this machine is
+ // the disk the operator is trying to empty. The location would then read
+ // "233 GB free" about a platter that is not plugged in.
+ //
+ // THE BOUNDARY IS TESTED AT THE MOUNTPOINT, NEVER AT THE ROOT, and the
+ // first version of this got that wrong in the one way that matters on
+ // this machine. A location's root is `join(mountpoint, relPath)` — the
+ // production `platter` is `/run/media/<user>/<uuid>/archilyzer-media`, a
+ // SUBDIRECTORY of the mountpoint — so the root and its parent are on the
+ // same filesystem BY CONSTRUCTION whenever `relPath` is non-empty, and
+ // comparing those two devices reported "free space unknown" for a
+ // correctly mounted drive. A bind mount reads the same way.
+ //
+ // Crossing a mount changes the device number, so `stat(mountpoint).dev`
+ // against `stat(dirname(mountpoint)).dev` is the honest question, and it
+ // is two syscalls — TABLES NEVER PROBE (see the header) rules out asking
+ // `probeLocation`, which is up to three subprocesses.
+ //
+ // Only asked of a location that has learned a `volume.uuid`: that field
+ // is the assertion that the root is supposed to be on its own volume. A
+ // location on a plain directory (never probed, a container, a
+ // subdirectory of the system disk by design) shares its parent's device
+ // legitimately, and withholding its free space would be wrong.
+ const mountpoint = loc.volume?.uuid
+ ? (loc.volume.mountpoint ?? "").trim()
+ : "";
+ // `/` is its own parent, so a volume mounted at the root has no boundary
+ // to test and is trivially there — the process is reading from it.
+ if (mountpoint && mountpoint !== path.dirname(mountpoint)) {
+ // The mountpoint is the drive's own top directory; its parent is on the
+ // filesystem above it, so only the first goes through the watchdog.
+ const [atMount, aboveMount] = await Promise.all([
+ orNull(onDrive(loc, () => stat(mountpoint))),
+ stat(path.dirname(mountpoint)).catch(() => null),
+ ]);
+ // Gone entirely, or present as an ordinary directory on the parent
+ // filesystem: either way nothing is mounted there.
+ if (!atMount || (aboveMount && aboveMount.dev === atMount.dev)) {
+ return undefined;
+ }
+ }
+ const free = await onDrive(loc, () => getFreeBytes(loc.root));
+ return Number.isFinite(free) ? free : undefined;
+}
+
// ---------------------------------------------------------------------------
// The probe memo
// ---------------------------------------------------------------------------
diff --git a/common/controller/storageStall.test.ts b/common/controller/storageStall.test.ts
@@ -38,8 +38,11 @@ import {
} from "../lib/channelMedia";
import { HELD_REASON, isMediaHeld } from "../lib/channelMediaHold";
import {
+ locationHealth,
recordLocationHealth,
+ registerLocationHealth,
resetStorageHealth,
+ setDriveCallBudget,
} from "../lib/storageHealth";
import {
probeLocation,
@@ -58,6 +61,10 @@ import type { SiteSettings } from "../lib/settings";
// ── the spy ────────────────────────────────────────────────────────────────
type Call = { fn: string; path: string };
let calls: Call[] = [];
+// THE HANG: a promise-API call whose (function, path) matches never settles —
+// what a read blocked on a stalled drive looks like from here. Recorded like
+// any other call. No drive is involved.
+let hang: ((c: Call) => boolean) | null = null;
{
const req = createRequire(import.meta.url);
const fsCjs = req("node:fs") as Record<string, unknown>;
@@ -75,17 +82,21 @@ let calls: Call[] = [];
: Buffer.isBuffer(v)
? v.toString()
: null;
- const wrap = (mod: Record<string, unknown>, name: string) => {
+ const wrap = (mod: Record<string, unknown>, name: string, promises = false) => {
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) });
+ if (p !== null) {
+ const call = { fn: name, path: path.resolve(p) };
+ calls.push(call);
+ if (promises && hang?.(call)) return new Promise(() => {});
+ }
return (fn as (...a: unknown[]) => unknown).apply(this, args);
};
};
for (const n of NAMES) {
- wrap(fspCjs, n);
+ wrap(fspCjs, n, true);
wrap(fsCjs, n);
wrap(fsCjs, `${n}Sync`);
}
@@ -161,6 +172,8 @@ function stall(): void {
}
beforeEach(() => {
+ hang = null;
+ setDriveCallBudget();
resetStorageHealth();
forgetChannelMedia();
resetStorageProbeMemo();
@@ -171,11 +184,14 @@ beforeEach(() => {
// ── inspectChannelMedia: the gate and the memo ─────────────────────────────
-test("inspect with the config in hand: a stalled drive costs no call at all", async () => {
+const MARKER = path.join(paths.channelsDir, SLUG, ".relocating.json");
+
+test("inspect with the config in hand: on a stalled drive only the marker is read, on the corpus disk", async () => {
stall();
calls = [];
const media = await inspectChannelMedia(paths, SLUG, CONFIG);
- assert.deepEqual(calls, []);
+ assert.deepEqual(calls, [{ fn: "readFile", path: MARKER }]);
+ assert.deepEqual(calls.filter(onDrive), []);
assert.equal(media.status, "stalled");
assert.equal(media.target, TARGET);
assert.match(String(media.detail), /^drive not answering \(location "USB drive", since /);
@@ -191,10 +207,29 @@ test("inspect without the config reads config.json and nothing on the drive", as
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")]],
+ [
+ ["readFile", path.join("corpus", "channels", SLUG, "config.json")],
+ ["readFile", path.join("corpus", "channels", SLUG, ".relocating.json")],
+ ],
);
});
+test("the marker first: a channel mid-move on a stalled drive reads in-transition", async () => {
+ writeFileSync(
+ MARKER,
+ JSON.stringify({ target: TARGET, direction: "back", startedAt: "", phase: "copy" }),
+ );
+ stall();
+ calls = [];
+ const media = await inspectChannelMedia(paths, SLUG, CONFIG);
+ assert.equal(media.status, "in-transition");
+ assert.deepEqual(calls.filter(onDrive), []);
+ // Remembered as a move: a second look within five seconds gives the move
+ // back, not the stall.
+ const again = await inspectChannelMedia(paths, SLUG, CONFIG);
+ assert.equal(again.status, "in-transition");
+});
+
test("the start-of-work guard refuses a stalled channel without a call", async () => {
stall();
calls = [];
@@ -205,7 +240,7 @@ test("the start-of-work guard refuses a stalled channel without a call", async (
err.status === "stalled" &&
/drive not answering/.test(err.message),
);
- assert.deepEqual(calls, []);
+ assert.deepEqual(calls.filter(onDrive), []);
});
test("the memo: two inspects within 5 s stat the drive once; fresh and age each ask again", async () => {
@@ -377,3 +412,111 @@ test("the saved-video inventory, for a page: a stalled channel is named, not rea
// Without the array (the backup job) nothing is skipped: it reads the drive.
assert.equal((await listSavedVideos({ paths })).length, 1);
});
+
+// ── the watchdog: a call that does not answer in the budget ────────────────
+// The location is registered (the health pass does that in the editor) but
+// answering; the hang makes one call on the drive never settle. The budget is
+// shortened to 100 ms; lib/storageHealth.test.ts holds the 3 s default.
+
+function hangOnDrive(fns: string[]): void {
+ hang = (c) => fns.includes(c.fn) && onDrive(c);
+}
+
+async function watchdogCase(
+ fns: string[],
+ run: () => Promise<unknown>,
+): Promise<unknown> {
+ registerLocationHealth([LOC]);
+ setDriveCallBudget(100);
+ hangOnDrive(fns);
+ calls = [];
+ const started = Date.now();
+ const out = await run();
+ assert.ok(Date.now() - started < 2_000, "answered on the watchdog, not on the drive");
+ assert.equal(locationHealth("usb")?.state, "stalled", "the location is marked at once");
+ assert.match(String(locationHealth("usb")?.cause), /a read in the editor did not answer/);
+ return out;
+}
+
+test("watchdog: inspect's target stat never answers → stalled, marked, and nothing more is asked of the drive", async () => {
+ const media = (await watchdogCase(["stat"], () =>
+ inspectChannelMedia(paths, SLUG, CONFIG),
+ )) as { status: string; detail?: string };
+ assert.equal(media.status, "stalled");
+ assert.match(String(media.detail), /^drive not answering \(location "USB drive"/);
+ // The next inspect, the guard and the walk make no call on the drive.
+ calls = [];
+ assert.equal((await inspectChannelMedia(paths, SLUG, CONFIG)).status, "stalled");
+ await assert.rejects(() => assertChannelMediaReachable(paths, SLUG, CONFIG));
+ assert.equal(await readChannelStat(paths, SLUG), null);
+ assert.deepEqual(calls.filter(onDrive), []);
+ // Cleared (two clean answers), the drive is asked again.
+ recordLocationHealth(LOC, "ok");
+ recordLocationHealth(LOC, "ok");
+ hang = null;
+ assert.equal(
+ (await inspectChannelMedia(paths, SLUG, CONFIG, { fresh: true })).status,
+ "ok",
+ );
+});
+
+test("watchdog: a video directory's read never answers mid-walk → readChannelStat answers null", async () => {
+ for (const id of ["vid2", "vid3", "vid4", "vid5", "vid6"]) {
+ mkdirSync(path.join(TARGET, id), { recursive: true });
+ }
+ const out = await watchdogCase(["readdir"], async () => {
+ // The walk's first readdir (of data/ itself, through the link) answers;
+ // the video directories' do not.
+ hang = (c) =>
+ c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid"));
+ return readChannelStat(paths, SLUG);
+ });
+ assert.equal(out, null);
+ // At most four video directories were asked before the stall refused the rest.
+ const asked = calls.filter((c) => c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid")));
+ assert.ok(asked.length <= 4, `${asked.length} video dirs asked`);
+});
+
+test("watchdog: probeLocation's stat never answers → 'stalled'", async () => {
+ const probe = (await watchdogCase(["stat"], () => probeLocation(LOC, paths))) as {
+ status: string;
+ };
+ assert.equal(probe.status, "stalled");
+});
+
+test("watchdog: volumeFreeBytes' stat never answers → unknown", async () => {
+ const out = (await watchdogCase(["stat"], () =>
+ volumeFreeBytes({ paths, locations: [LOC] }),
+ )) as Record<string, number | undefined>;
+ assert.equal(out.usb, undefined);
+});
+
+test("watchdog: a recency tail read never answers → not dated, not remembered as a miss", async () => {
+ const args = {
+ paths,
+ meta: [{ slug: SLUG, config: CONFIG }],
+ candidateIds: new Set(["vid1"]),
+ owner: new Map([["vid1", SLUG]]),
+ interpolate: false,
+ fresh: true,
+ };
+ const keys = (await watchdogCase(["open"], () => buildRecencyKeys(args))) as Map<
+ string,
+ { key: string }
+ >;
+ assert.equal(keys.get("vid1")?.key, "");
+ resetStorageHealth();
+ hang = null;
+ assert.equal((await buildRecencyKeys(args)).get("vid1")?.key, "20260601");
+});
+
+test("watchdog: a move onto a root whose stat never answers is refused", async () => {
+ const problem = await watchdogCase(["stat"], () =>
+ relocationRootPresenceProblem(
+ DRIVE,
+ { locations: [LOC], defaultLocationId: "" },
+ paths,
+ ),
+ );
+ assert.match(String(problem), /drive not answering/);
+});
diff --git a/common/lib/channelMedia.ts b/common/lib/channelMedia.ts
@@ -3,7 +3,10 @@ import { lstat, readFile, readlink, rm, stat } from "node:fs/promises";
import type { Paths } from "./paths";
import type { ChannelConfig } from "./channelConfig";
import {
+ DRIVE_CALL_BUDGET_MS,
NOT_ANSWERING,
+ isDriveNotAnswering,
+ onDrive,
sinceText,
stalledLocationForPath,
type LocationHealth,
@@ -203,14 +206,18 @@ async function readConfiguredDataDir(
export function stalledMediaLocation(
dataDir: string,
configured: string,
- health: LocationHealth,
+ // Null when the drive is on no location the health state knows (a root typed
+ // by hand) and the watchdog found it not answering.
+ health: LocationHealth | null,
): ChannelMediaLocation {
return {
dataDir,
relocated: true,
target: configured,
status: "stalled",
- detail: `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`,
+ detail: health
+ ? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`
+ : `${NOT_ANSWERING} (a read did not answer within ${DRIVE_CALL_BUDGET_MS / 1000} s)`,
};
}
@@ -236,8 +243,11 @@ export function channelMediaStall(
// 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.
+// THE STALL GATE IS ASKED BEFORE A REMEMBERED ANSWER IS GIVEN, so a drive that
+// stops answering is seen on the next call even when an `ok` from four seconds
+// ago is remembered. A remembered `in-transition` is given as it is: the
+// marker is asked before the gate (below), and the movers forget the channel
+// whenever they write or remove it.
//
// `fresh: true` BYPASSES IT, and every caller that decides something from the
// answer passes it: the start-of-work guard (`assertChannelMediaReachable`),
@@ -295,8 +305,10 @@ 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 — or none, for five seconds after the last answer (see the memo above),
-// and none at all on a stalled location.
+// row — or none, for five seconds after the last answer (see the memo above).
+// On a stalled location the one call that reaches the drive (the target's
+// `stat`) is not made; the marker, the link and config.json are on the corpus
+// disk and are read as usual.
export async function inspectChannelMedia(
paths: ChannelMediaPaths,
slug: string,
@@ -311,25 +323,22 @@ 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 };
+ if (hit && now - hit.at < CHANNEL_MEDIA_MEMO_MS) {
+ if (hit.location.status !== "in-transition" && configured) {
+ const stall = stalledLocationForPath(configured);
+ if (stall) return stalledMediaLocation(dataDir, configured, stall);
+ }
+ return { ...hit.location };
+ }
}
const location = await inspectOnDisk(paths, slug, dataDir, configured);
- if (!opts.fresh) {
+ // A stall is not remembered: the health state is already its memory.
+ if (!opts.fresh && location.status !== "stalled") {
if (memo.size >= MEMO_SWEEP_AT) {
for (const [k, v] of memo) {
if (now - v.at >= CHANNEL_MEDIA_MEMO_MS) memo.delete(k);
@@ -346,6 +355,10 @@ async function inspectOnDisk(
dataDir: string,
configured: string | undefined,
): Promise<ChannelMediaLocation> {
+ // THE MARKER FIRST. It is in the channel dir, on the corpus disk, so reading
+ // it costs the drive nothing — and a channel mid-move reads `in-transition`
+ // whatever its drive is doing, which is what a resumed move and every guard
+ // key off.
const marker = await readRelocationMarker(paths, slug);
if (marker) {
return {
@@ -360,6 +373,15 @@ async function inspectOnDisk(
};
}
+ // THE GATE: a channel whose configured target is on a location whose drive
+ // is not answering is answered from memory, before the link is looked at and
+ // before the target's stat, which would hold an I/O thread for as long as the
+ // drive takes.
+ if (configured) {
+ const stall = stalledLocationForPath(configured);
+ if (stall) return stalledMediaLocation(dataDir, configured, stall);
+ }
+
let link: Awaited<ReturnType<typeof lstat>> | null = null;
try {
link = await lstat(dataDir);
@@ -411,8 +433,12 @@ async function inspectOnDisk(
// The link points at a DEEP path (<root>/<slug>/data), so an unmounted root
// gives ENOENT here. An empty mountpoint can never be mistaken for the
// media, which is the whole reason the suffix is fixed.
+ //
+ // THE ONE CALL HERE THAT REACHES THE DRIVE, so it goes through the
+ // watchdog: not made while the location is stalled, and a stat that has
+ // not answered in 3 s marks it stalled and answers `stalled` now.
try {
- const st = await stat(configured);
+ const st = await onDrive(configured, () => stat(configured));
if (!st.isDirectory()) {
return {
dataDir,
@@ -422,7 +448,10 @@ async function inspectOnDisk(
detail: `${configured} exists but is not a directory`,
};
}
- } catch {
+ } catch (err) {
+ if (isDriveNotAnswering(err)) {
+ return stalledMediaLocation(dataDir, configured, err.health);
+ }
return {
dataDir,
relocated: true,
diff --git a/common/lib/storageHealth.test.ts b/common/lib/storageHealth.test.ts
@@ -1,8 +1,16 @@
import { beforeEach, test } from "node:test";
import assert from "node:assert/strict";
import {
+ DRIVE_CALLS_IN_FLIGHT,
+ DRIVE_CALL_BUDGET_MS,
+ DriveNotAnsweringError,
HEALTH_CLEAN_TO_CLEAR,
allLocationHealth,
+ driveCallsInFlight,
+ isDriveNotAnswering,
+ onDrive,
+ registerLocationHealth,
+ setDriveCallBudget,
locationHealth,
notAnsweringText,
pruneLocationHealth,
@@ -20,7 +28,10 @@ import {
// 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());
+beforeEach(() => {
+ resetStorageHealth();
+ setDriveCallBudget();
+});
const USB = { id: "usb", label: "USB drive", root: "/mnt/usb/media" };
@@ -127,3 +138,149 @@ test("the words: since a time today, or a date and time", () => {
assert.equal(sinceText(yesterday, now), "since 2026-09-28 23:05");
assert.equal(notAnsweringText({ since: today }, now), "not answering since 11:35");
});
+
+// ── the watchdog (`onDrive`) ────────────────────────────────────────────────
+// A call that never answers is a promise that never settles: exactly what a
+// read blocked on a stalled drive looks like from here. No drive is involved.
+
+const never = () => new Promise<never>(() => {});
+
+function deferred<T>() {
+ let resolve!: (v: T) => void;
+ let reject!: (e: Error) => void;
+ const promise = new Promise<T>((res, rej) => {
+ resolve = res;
+ reject = rej;
+ });
+ return { promise, resolve, reject };
+}
+
+test("a call that answers in time passes through, value or error, and frees its slot", async () => {
+ registerLocationHealth([USB]);
+ assert.equal(await onDrive("/mnt/usb/media/ch/data", async () => 42), 42);
+ await assert.rejects(
+ () => onDrive(USB, async () => {
+ throw new Error("ENOENT");
+ }),
+ /ENOENT/,
+ );
+ assert.equal(driveCallsInFlight("usb"), 0);
+ assert.equal(locationHealth("usb")?.state, "ok");
+});
+
+test("a call that never answers: stalled on the timer, the location marked at once, the call left to settle", async () => {
+ registerLocationHealth([USB], 0);
+ setDriveCallBudget(80);
+ const late = deferred<string>();
+ const started = Date.now();
+ await assert.rejects(
+ () => onDrive("/mnt/usb/media/ch/data", () => late.promise),
+ (err: unknown) =>
+ err instanceof DriveNotAnsweringError &&
+ isDriveNotAnswering(err) &&
+ err.health?.id === "usb" &&
+ /^drive not answering \(location "USB drive", since /.test(err.message),
+ );
+ assert.ok(Date.now() - started < 1_000);
+ const h = locationHealth("usb");
+ assert.equal(h?.state, "stalled");
+ assert.ok((h?.since ?? 0) >= started, "since is now, not the entry's first sighting");
+ assert.match(String(h?.cause), /a read in the editor did not answer within 0.08 s/);
+ // The slot is held until the call really returns.
+ assert.equal(driveCallsInFlight("usb"), 1);
+ late.resolve("finally");
+ await new Promise((r) => setImmediate(r));
+ assert.equal(driveCallsInFlight("usb"), 0);
+});
+
+test("a stalled location is refused without the call being made", async () => {
+ registerLocationHealth([USB]);
+ recordLocationHealth(USB, "stalled");
+ let made = 0;
+ await assert.rejects(
+ () => onDrive("/mnt/usb/media/ch/data", async () => ++made),
+ DriveNotAnsweringError,
+ );
+ await assert.rejects(() => onDrive(USB, async () => ++made), DriveNotAnsweringError);
+ assert.equal(made, 0);
+});
+
+test("at most four calls in flight on a location; the rest wait, and are refused without a call when it stalls", async () => {
+ assert.equal(DRIVE_CALLS_IN_FLIGHT, 4);
+ registerLocationHealth([USB]);
+ setDriveCallBudget(80);
+ let made = 0;
+ const calls = Array.from({ length: 7 }, () =>
+ onDrive(USB, () => {
+ made += 1;
+ return never();
+ }).then(
+ () => "answered",
+ (err: Error) => err.name,
+ ),
+ );
+ await new Promise((r) => setImmediate(r));
+ assert.equal(made, 4, "four in flight, three waiting");
+ const outcomes = await Promise.all(calls);
+ assert.deepEqual(outcomes, Array(7).fill("DriveNotAnsweringError"));
+ assert.equal(made, 4, "the three that waited were refused without a call");
+ assert.equal(locationHealth("usb")?.state, "stalled");
+});
+
+test("a waiting call runs when a slot frees, if the location is still answering", async () => {
+ registerLocationHealth([USB]);
+ const gates = Array.from({ length: 4 }, () => deferred<number>());
+ const first = gates.map((g) => onDrive(USB, () => g.promise));
+ let fifth = false;
+ const waiting = onDrive(USB, async () => {
+ fifth = true;
+ return 5;
+ });
+ await new Promise((r) => setImmediate(r));
+ assert.equal(fifth, false);
+ gates[0].resolve(1);
+ assert.equal(await waiting, 5);
+ for (const g of gates.slice(1)) g.resolve(0);
+ await Promise.all(first);
+ assert.equal(driveCallsInFlight("usb"), 0);
+});
+
+test("a path on no known location is still raced, and names no location", async () => {
+ setDriveCallBudget(50);
+ await assert.rejects(
+ () => onDrive("/elsewhere/ch/data", never),
+ (err: unknown) => err instanceof DriveNotAnsweringError && err.health === null,
+ );
+ assert.deepEqual(allLocationHealth(), {});
+});
+
+test("a probe of another root under a location's id is raced but does not rewrite that location", async () => {
+ registerLocationHealth([USB]);
+ setDriveCallBudget(50);
+ await assert.rejects(
+ () => onDrive({ ...USB, root: "/mnt/candidate" }, never),
+ DriveNotAnsweringError,
+ );
+ assert.equal(locationHealth("usb")?.root, "/mnt/usb/media");
+ assert.equal(locationHealth("usb")?.state, "ok");
+});
+
+test("the default budget is 3 s", async () => {
+ assert.equal(DRIVE_CALL_BUDGET_MS, 3_000);
+ registerLocationHealth([USB]);
+ const started = Date.now();
+ await assert.rejects(() => onDrive(USB, never), DriveNotAnsweringError);
+ const took = Date.now() - started;
+ assert.ok(took >= 3_000 && took < 4_500, `answered after ${took} ms`);
+});
+
+test("registering locations creates entries without an answer, and a moved root starts over", () => {
+ registerLocationHealth([USB], 5);
+ assert.equal(locationHealth("usb")?.state, "ok");
+ recordLocationHealth(USB, "stalled", { now: 6 });
+ registerLocationHealth([USB], 7);
+ assert.equal(locationHealth("usb")?.state, "stalled", "registering is not an answer");
+ registerLocationHealth([{ ...USB, root: "/mnt/new" }], 8);
+ assert.equal(locationHealth("usb")?.state, "ok");
+ assert.equal(locationHealth("usb")?.root, "/mnt/new");
+});
diff --git a/common/lib/storageHealth.ts b/common/lib/storageHealth.ts
@@ -55,11 +55,19 @@ export type LocationHealthState =
// The root answered in time, and is not a directory (not mounted, or gone).
| "absent";
+// What decided a location's state. `counters`: the block device's own request
+// counters (`/sys/class/block/<dev>/stat`), which never touch the drive;
+// `stat`: a child `stat` of the root, when no device could be named (a
+// container, no findmnt, no /sys entry). A stall marked by `onDrive`'s watchdog
+// keeps the detector the last pass used; its `cause` says what did not answer.
+export type HealthDetector = "counters" | "stat";
+
export type LocationHealth = {
id: string;
label: string;
root: string;
state: LocationHealthState;
+ detector?: HealthDetector;
// When the current state began, ms since epoch.
since: number;
// When the last probe (or observation) was recorded.
@@ -85,23 +93,43 @@ 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> };
+type Waiter = { resolve: () => void; reject: (err: Error) => void };
+
+type HealthState = {
+ byId: Map<string, LocationHealth>;
+ // `onDrive`'s bookkeeping, by location id: calls in flight, and the calls
+ // waiting for one of them to finish.
+ inFlight?: Map<string, number>;
+ waiters?: Map<string, Waiter[]>;
+ // Test seam: the watchdog's budget.
+ budgetMs?: number;
+};
declare global {
// eslint-disable-next-line no-var
var __yttStorageHealth__: HealthState | undefined;
}
-function healthState(): HealthState {
+function healthState(): Required<Omit<HealthState, "budgetMs">> & HealthState {
if (!globalThis.__yttStorageHealth__) {
globalThis.__yttStorageHealth__ = { byId: new Map() };
}
- return globalThis.__yttStorageHealth__;
+ const s = globalThis.__yttStorageHealth__;
+ // Filled lazily: a dev server's hot reload keeps an object made by an older
+ // copy of this module.
+ s.inFlight ??= new Map();
+ s.waiters ??= new Map();
+ return s as Required<Omit<HealthState, "budgetMs">> & HealthState;
}
-// Test seam, and the escape hatch for a process that wants to forget.
+// Test seam, and the escape hatch for a process that wants to forget. Calls
+// still waiting for a slot are released to run.
export function resetStorageHealth(): void {
- healthState().byId.clear();
+ const s = healthState();
+ s.byId.clear();
+ s.inFlight.clear();
+ for (const q of s.waiters.values()) for (const w of q) w.resolve();
+ s.waiters.clear();
}
export type HealthTransition = {
@@ -115,7 +143,7 @@ export type HealthTransition = {
export function recordLocationHealth(
loc: Pick<StorageLocation, "id" | "label" | "root">,
answer: LocationHealthState,
- opts: { now?: number; cause?: string } = {},
+ opts: { now?: number; cause?: string; detector?: HealthDetector } = {},
): HealthTransition | null {
const now = opts.now ?? Date.now();
const map = healthState().byId;
@@ -128,15 +156,17 @@ export function recordLocationHealth(
label,
root: loc.root,
state: answer,
+ ...(opts.detector ? { detector: opts.detector } : {}),
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 };
+ return { id: loc.id, from: null, to: answer };
}
prev.label = label;
prev.checkedAt = now;
+ if (opts.detector) prev.detector = opts.detector;
if (answer === "stalled") {
prev.cleanStreak = 0;
if (prev.state === "stalled") return null;
@@ -162,6 +192,25 @@ export function recordLocationHealth(
return { id: loc.id, from, to: answer };
}
+// Make sure every configured location has an entry, without recording an
+// answer: a new location starts `ok`, and a known one whose root moved starts
+// over. The watchdog below can only mark a location it can find, and it finds
+// a channel's location among these entries.
+export function registerLocationHealth(
+ locations: ReadonlyArray<Pick<StorageLocation, "id" | "label" | "root">>,
+ now: number = Date.now(),
+): void {
+ const map = healthState().byId;
+ for (const loc of locations) {
+ const prev = map.get(loc.id);
+ if (prev && prev.root === loc.root) {
+ prev.label = loc.label || loc.id;
+ continue;
+ }
+ recordLocationHealth(loc, "ok", { now });
+ }
+}
+
// Drop every location that is no longer configured.
export function pruneLocationHealth(liveIds: Iterable<string>): void {
const keep = new Set(liveIds);
@@ -229,3 +278,193 @@ export function sinceText(since: number, now: number = Date.now()): string {
export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number): string {
return `not answering ${sinceText(h.since, now)}`;
}
+
+// ---------------------------------------------------------------------------
+// The watchdog: every call the gate covers, raced against 3 s
+// ---------------------------------------------------------------------------
+//
+// THE DETECTOR THAT CANNOT BE FOOLED BY A CACHE. The 15 s pass reads the
+// block device's counters (or, with no device, a child `stat`), and a stall
+// that starts between two passes, or that a cached answer hides, is not seen
+// by it. A page or a poll that actually reaches the drive is: `onDrive` runs
+// the call against a 3 s timer, and a call that has not answered by then marks
+// its location `stalled` at once (since now) and throws
+// `DriveNotAnsweringError`. The call itself is left to settle on its own: its
+// thread is held until the drive answers, which is the stated limit.
+//
+// AND NO LOCATION HOLDS MORE THAN FOUR OF THE POOL'S THREADS. A walk of
+// `data/` fans out 64 wide, and a stall mid-walk would otherwise put all 64 in
+// libuv's queue before the watchdog fired — every thread held for as long as
+// the drive takes, and every other call in the process queued behind them. So
+// at most DRIVE_CALLS_IN_FLIGHT calls per location are in flight through here;
+// the rest wait in a queue of our own, and the moment the location is marked
+// stalled they are refused without a call. A slot is released when its call
+// really returns, not when the watchdog gave up on it.
+//
+// Do not nest `onDrive` for one location: the inner call would wait for a slot
+// the outer one holds. A unit of work (a video directory's few reads) goes
+// through as one call.
+
+export const DRIVE_CALL_BUDGET_MS = 3_000;
+export const DRIVE_CALLS_IN_FLIGHT = 4;
+
+// Test seam: shorten (or restore, with no argument) the watchdog's budget.
+export function setDriveCallBudget(ms?: number): void {
+ healthState().budgetMs = ms;
+}
+
+function driveCallBudget(): number {
+ return healthState().budgetMs ?? DRIVE_CALL_BUDGET_MS;
+}
+
+export class DriveNotAnsweringError extends Error {
+ // The location the drive is on, when it is a known one.
+ readonly health: LocationHealth | null;
+ constructor(health: LocationHealth | null) {
+ super(
+ health
+ ? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`
+ : `${NOT_ANSWERING} (a read did not answer within ${driveCallBudget() / 1000} s)`,
+ );
+ this.name = "DriveNotAnsweringError";
+ this.health = health;
+ }
+}
+
+export function isDriveNotAnswering(err: unknown): err is DriveNotAnsweringError {
+ return err instanceof Error && err.name === "DriveNotAnsweringError";
+}
+
+type Where = string | Pick<StorageLocation, "id" | "label" | "root">;
+
+// The location `where` is on, as far as the health state knows: a path is
+// matched among the entries (`locationOfDataDir`'s rule); a location is itself.
+// `markable` is false for a location object whose root is not the one its
+// entry records (a probe of a candidate root under the same id): racing it is
+// right, rewriting that location's entry for another root is not.
+function resolveWhere(
+ where: Where,
+): { loc: Pick<StorageLocation, "id" | "label" | "root">; markable: boolean } | null {
+ const map = healthState().byId;
+ if (typeof where !== "string") {
+ const entry = map.get(where.id);
+ return { loc: where, markable: !entry || entry.root === where.root };
+ }
+ if (map.size === 0 || !where) return null;
+ const found = locationOfDataDir(
+ where,
+ [...map.values()].map((h) => ({ id: h.id, label: h.label, root: h.root, autoRepoint: false })),
+ );
+ return found ? { loc: found, markable: true } : null;
+}
+
+async function acquireSlot(id: string): Promise<void> {
+ const s = healthState();
+ const n = s.inFlight.get(id) ?? 0;
+ if (n < DRIVE_CALLS_IN_FLIGHT) {
+ s.inFlight.set(id, n + 1);
+ return;
+ }
+ await new Promise<void>((resolve, reject) => {
+ const q = s.waiters.get(id) ?? [];
+ q.push({ resolve, reject });
+ s.waiters.set(id, q);
+ });
+ // Resolved by a release that handed its slot over: the count is unchanged.
+}
+
+function releaseSlot(id: string): void {
+ const s = healthState();
+ const next = s.waiters.get(id)?.shift();
+ if (next) {
+ next.resolve();
+ return;
+ }
+ s.inFlight.set(id, Math.max(0, (s.inFlight.get(id) ?? 1) - 1));
+}
+
+function refuseWaiters(id: string, health: LocationHealth | null): void {
+ const s = healthState();
+ const q = s.waiters.get(id) ?? [];
+ s.waiters.delete(id);
+ for (const w of q) w.reject(new DriveNotAnsweringError(health));
+}
+
+// How many calls are in flight on a location through `onDrive` (for tests and
+// for a reader that wants to say so).
+export function driveCallsInFlight(id: string): number {
+ return healthState().inFlight.get(id) ?? 0;
+}
+
+// Run `call` against the drive `where` is on: refused at once when that
+// location is stalled, queued behind DRIVE_CALLS_IN_FLIGHT calls already in
+// flight on it, and raced against the budget. Throws DriveNotAnsweringError
+// for a refusal or a timeout; any other error is the call's own.
+export async function onDrive<T>(where: Where, call: () => Promise<T>): Promise<T> {
+ const resolved = resolveWhere(where);
+ const stalledNow = () =>
+ resolved
+ ? stalledLocation(resolved.loc)
+ : typeof where === "string"
+ ? stalledLocationForPath(where)
+ : null;
+ const refused = stalledNow();
+ if (refused) throw new DriveNotAnsweringError(refused);
+ const id = resolved?.loc.id;
+ if (id !== undefined) {
+ await acquireSlot(id);
+ // Stalled while this call waited: refused, and the slot passed on.
+ const late = stalledNow();
+ if (late) {
+ releaseSlot(id);
+ throw new DriveNotAnsweringError(late);
+ }
+ }
+ let released = false;
+ const release = () => {
+ if (released || id === undefined) return;
+ released = true;
+ releaseSlot(id);
+ };
+ let pending: Promise<T>;
+ try {
+ pending = call();
+ } catch (err) {
+ release();
+ throw err;
+ }
+ const TIMED_OUT = Symbol("timed out");
+ let timer: ReturnType<typeof setTimeout> | undefined;
+ // Not unref'd: it is cleared the moment the call answers, and while the call
+ // is outstanding the timer is what must fire.
+ const timeout = new Promise<typeof TIMED_OUT>((resolve) => {
+ timer = setTimeout(() => resolve(TIMED_OUT), driveCallBudget());
+ });
+ let answer: T | typeof TIMED_OUT;
+ try {
+ answer = await Promise.race([pending, timeout]);
+ } catch (err) {
+ release();
+ throw err;
+ } finally {
+ if (timer) clearTimeout(timer);
+ }
+ if (answer !== TIMED_OUT) {
+ release();
+ return answer;
+ }
+ // The call is still waiting on the drive. Its slot stays held until it
+ // returns; nothing awaits it.
+ pending.then(release, release);
+ let health: LocationHealth | null = null;
+ if (resolved) {
+ if (resolved.markable) {
+ recordLocationHealth(resolved.loc, "stalled", {
+ cause: `a read in the editor did not answer within ${driveCallBudget() / 1000} s`,
+ });
+ health = stalledLocation(resolved.loc);
+ }
+ refuseWaiters(resolved.loc.id, health);
+ }
+ throw new DriveNotAnsweringError(health);
+}
diff --git a/common/lib/storageVolumes.ts b/common/lib/storageVolumes.ts
@@ -6,6 +6,8 @@ import { getFreeBytes } from "./diskSpace";
import type { StorageLocation, StorageVolume } from "./storageLocations";
import {
HEALTH_PROBE_TIMEOUT_MS,
+ isDriveNotAnswering,
+ onDrive,
stalledLocation,
type LocationHealthState,
} from "./storageHealth";
@@ -237,18 +239,23 @@ export async function probeLocation(
// 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 } };
- }
+ // drive comes back; the health pass (below) already asked without touching
+ // it. Both go through `onDrive`'s watchdog, too: one that has not answered
+ // in 3 s marks the location stalled and this probe answers so.
+ const stalledProbe: StorageLocationProbe = {
+ status: "stalled",
+ identity: { known: false },
+ };
+ if (stalledLocation(loc)) return stalledProbe;
// 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;
if (root !== "") {
try {
- isDir = (await stat(root)).isDirectory();
- } catch {
+ isDir = (await onDrive(loc, () => stat(root))).isDirectory();
+ } catch (err) {
+ if (isDriveNotAnswering(err)) return stalledProbe;
isDir = false;
}
}
@@ -261,7 +268,13 @@ export async function probeLocation(
// `null` and every byte formatter downstream gets a surprise. An
// unmeasurable root reports no free space at all, which is the honest
// answer and the one the field is already optional for.
- const free = await getFreeBytes(root);
+ let free: number;
+ try {
+ free = await onDrive(loc, () => getFreeBytes(root));
+ } catch (err) {
+ if (isDriveNotAnswering(err)) return stalledProbe;
+ throw err;
+ }
// Two ways a mount will not be there after a reboot: udisks put it under
// /run/media (or /media) because a human plugged it in, or there is no
// fstab entry naming its UUID. The fstab call is skipped when the