commit 646b3a951d66ebe915c376f1814cd158130b60b1
parent 05e0cba726da093a9883b1e49a722a89b29ad99c
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 30 Sep 2026 09:00:08 -0400
common: the drive-health timings are settings — settings.storage.health (budgetMs, passIntervalMs, probeTimeoutMs, clearAfterCleanPasses, inFlightPerLocation), today's constants as the defaults
lib/storageHealthTimings.ts (pure) holds the defaults, the ranges, the
sanitizer (clamped on read, only a value that differs from its default is
kept) and the words. storageHealth.ts reads every number through one
accessor, healthTimings(); applyHealthTimings() sets them, the health pass
applies what it reads on every pass, and a changed interval re-arms the armed
pass through onPassIntervalChange. A raised cap admits waiting calls; a
lowered one is reached as calls return. The counters' sample spacing is
min(10 s, interval - 5 s), floored at half the interval. The index and stats
bins apply the settings once. SETTINGS.md regenerated. Tests.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
21 files changed, 855 insertions(+), 124 deletions(-)
diff --git a/SETTINGS.md b/SETTINGS.md
@@ -522,6 +522,7 @@ Where a channel's downloaded media goes when it is relocated off the corpus disk
| `locations` | `[]` | The named storage locations a channel's media may be relocated to — one entry per root, each with an id, label, root, `autoRepoint` and the learned volume identity. Order is display order. Managed on /storage. |
| `defaultLocationId` | `""` | The location prefilled as the destination of a move. "" = no default. |
| `savedVideosLocationId` | absent | WHERE THE SAVED-VIDEO STORE IS, by location id. "" = in place, under the corpus at `paths.savedVideosDir`.<br><br>A RECORD OF WHAT IS ON DISK, never an intention — the same contract as a channel's `config.dataDir`. It is written by the move, on success, after the copy has verified and the symlink is in place; nothing else writes it, and a reader that disagrees with the disk trusts the disk. Optional so an older settings.json parses (and an older binary that drops it leaves a store that still works, because the symlink is what every reader follows). |
+| `health` | absent | THE DRIVE-HEALTH TIMINGS: how long a read may take before a drive counts as not answering, how often the health pass looks, how long its look may take, how many clean looks clear a stall, and how many reads may be on one drive at once. Edited on /storage (Drive health timing). Absent = every default, and only a value that differs from its default is written, so an untuned install follows a default changed later. See `storage.health` below. |
#### `storage.locations[]`
@@ -547,6 +548,16 @@ Per entry — each entry spells its own values.
| `mountpoint` | Where the volume was mounted at the last successful probe, and the path of the location's root RELATIVE to that mountpoint. Invariant: `root === join(mountpoint, relPath)`. Keeping the two halves is what lets a probe compute a candidate root when the volume reappears elsewhere. |
| `relPath` | The location root's path RELATIVE to `mountpoint` (see there). Invariant: `root === join(mountpoint, relPath)`. |
+#### `storage.health`
+
+| Key | Default | Description |
+|---|---|---|
+| `budgetMs` | `3000` | How long one read may take before the drive counts as not answering, in ms (default 3000, 500–60000). The watchdog's budget (`onDrive`, lib/storageHealth.ts) for one unit of work — a video directory's reads, a page's reads of one video: a unit that has not answered by then is refused, and marks its location not answering unless the disk's request counters show it still completing others (slow, not stalled). A read waiting for a slot is refused when nothing on the drive has returned for this long plus a quarter of it (at most 250 ms). Takes effect on the next read after a save. |
+| `passIntervalMs` | `15000` | How often the health pass reads each location's disk counters, in ms (default 15000, 5000–300000). A stall that starts between two passes is seen by the next, or at once by a page's read. A save on /storage re-arms the pass's timer at once; a hand edit, at the next pass. Two counter samples are compared only when at least min(10 s, this − 5 s) apart, a spacing never less than half of this. |
+| `probeTimeoutMs` | `3000` | How long the health pass waits, in ms (default 3000, 500–30000), for the child `stat` of a root where no disk can be named (a timeout counts as not answering) and for the `findmnt` that names a root's disk (a timeout names none that pass). Takes effect on the next pass. |
+| `clearAfterCleanPasses` | `2` | How many clean answers in a row clear a location marked not answering (default 2, 1–10). Each health pass is one answer, and so is a Refresh on /storage; a miss in between starts the count again. Takes effect on the next answer. |
+| `inFlightPerLocation` | `4` | How many reads through the watchdog may be on one location's drive at once (default 4, 1–16); the rest wait in the editor's own queue, so a stall mid-walk holds this many of Node's threads, not all of them. Keep it under `UV_THREADPOOL_SIZE` (16 in the editor's start script). Takes effect on the next read. |
+
Default:
```json
diff --git a/common/bin/build-index.ts b/common/bin/build-index.ts
@@ -3,10 +3,16 @@
// as `tsx bin/build-index.ts` (export's build:index script).
import { getPaths, type Paths } from "../lib/paths";
import { buildIndex } from "../controller/buildIndex";
+import { settingsFromFile } from "../lib/settings";
+import { applyHealthTimings } from "../lib/storageHealth";
import { runIfEntryPoint } from "./_cli";
export async function main(opts: { paths?: Paths } = {}): Promise<void> {
- await buildIndex({ paths: opts.paths ?? getPaths() });
+ const paths = opts.paths ?? getPaths();
+ // A CLI process has no health pass: the drive-health timings the build's
+ // watchdog runs on (settings.storage.health) are applied here, once.
+ applyHealthTimings(settingsFromFile(paths.settingsFile).storage.health);
+ await buildIndex({ paths });
}
runIfEntryPoint(import.meta.url, () => main());
diff --git a/common/bin/build-stats.ts b/common/bin/build-stats.ts
@@ -3,10 +3,16 @@
// build:stats script); `main` is exported for the archilyzer CLI.
import { getPaths, type Paths } from "../lib/paths";
import { buildStats } from "../controller/buildStats";
+import { settingsFromFile } from "../lib/settings";
+import { applyHealthTimings } from "../lib/storageHealth";
import { runIfEntryPoint } from "./_cli";
export async function main(opts: { paths?: Paths } = {}): Promise<void> {
- await buildStats({ paths: opts.paths ?? getPaths() });
+ const paths = opts.paths ?? getPaths();
+ // A CLI process has no health pass: the drive-health timings the build's
+ // watchdog runs on (settings.storage.health) are applied here, once.
+ applyHealthTimings(settingsFromFile(paths.settingsFile).storage.health);
+ await buildStats({ paths });
}
runIfEntryPoint(import.meta.url, () => main());
diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts
@@ -743,7 +743,8 @@ export async function generateChannelSnapshot(
// each video directory's unit below. At most four of them wait on that drive
// at once — this walk runs in the editor's own process after every download
// or sync of the channel, sixteen wide, which is exactly while a long write is
- // stressing the drive — and one that does not answer in 3 s throws, so the
+ // stressing the drive — and one that does not answer within the budget
+ // (`storage.health.budgetMs`, 3 s by default) throws, so the
// scheduler keeps the last good snapshot.json, as on any failed refresh.
// The reconcile pass just below is sequential (one read at a time) and is
// not raced.
diff --git a/common/controller/channels.ts b/common/controller/channels.ts
@@ -90,8 +90,9 @@ async function hasDigestWithItems(videoDir: string): Promise<boolean> {
// `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
+// location is stalled, at most `inFlightPerLocation` (4) in flight on it, and
+// one that has not answered within the budget (`storage.health.budgetMs`, 3 s
+// by default) 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(
diff --git a/common/controller/recencyIndex.ts b/common/controller/recencyIndex.ts
@@ -199,8 +199,9 @@ async function readTailUploadDate(file: string): Promise<string | null> {
// `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
+// came back, 32 at a time — at most `inFlightPerLocation` (4) in flight on it,
+// and one that has not answered within the budget (3 s by default) 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");
diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts
@@ -300,8 +300,9 @@ 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.
+ // Through the watchdog: a stat that has not answered within the budget
+ // (`storage.health.budgetMs`, 3 s by default) marks the location stalled
+ // and refuses the same way.
let isDir = false;
try {
isDir = (await onDrive(named ?? r, () => stat(r))).isDirectory();
diff --git a/common/controller/storageLocations.ts b/common/controller/storageLocations.ts
@@ -259,7 +259,8 @@ export async function volumeFreeBytes(opts: {
// 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. 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.
+ // has not answered within the budget (3 s by default) marks the location
+ // stalled, and reads "—" too.
if (stalledLocation(loc)) {
out[loc.id] = undefined;
return;
diff --git a/common/controller/storageWatch.test.ts b/common/controller/storageWatch.test.ts
@@ -23,6 +23,8 @@ import {
} from "./storageWatch";
import { inspectChannelMedia } from "../lib/channelMedia";
import {
+ applyHealthTimings,
+ healthTimings,
locationHealth,
resetStorageHealth,
type LocationHealthState,
@@ -35,6 +37,7 @@ import {
beforeEach(() => {
resetStorageWatchSuspicion();
resetStorageHealth();
+ applyHealthTimings();
});
// Run with:
@@ -589,3 +592,70 @@ test("a stall auto-pauses after two passes, and says the drive is not answering
assert.match(String(autoPauseReasonOf(old, "a")), /drive that is not there/);
});
});
+
+// ── the timings are settings (release 15 slice DT) ─────────────────────────
+
+test("DT: every pass applies the timings it reads — the clear count and the log line follow storage.health", async () => {
+ await withTmp(async (h) => {
+ await seedRelocated(h, "slow", { targetExists: true });
+ h.io.read().storage.health = { clearAfterCleanPasses: 3, budgetMs: 5_000 };
+ const probe = scripted(["stalled", "ok", "ok", "ok"]);
+ const lines: string[] = [];
+ const pass = () => runStorageHealthPass({ io: h.io, probe, log: (l) => lines.push(l) });
+ await pass();
+ assert.equal(healthTimings().budgetMs, 5_000, "applied before anything was asked");
+ assert.match(lines.join("\n"), /until it answers 3 times in a row/);
+ await pass();
+ await pass();
+ assert.equal(locationHealth("cold")?.state, "stalled", "two clean passes are not three");
+ await pass();
+ assert.equal(locationHealth("cold")?.state, "ok");
+ // A pass handed its locations reads no settings: the timings stay.
+ await runStorageHealthPass({ locations: h.io.read().storage.locations, probe: scripted(["ok"]) });
+ assert.equal(healthTimings().clearAfterCleanPasses, 3);
+ });
+});
+
+test("DT: a changed pass interval re-arms the armed pass; an explicit interval follows nothing", async () => {
+ await withTmp(async (h) => {
+ let asked = 0;
+ const lines: string[] = [];
+ startStorageHealthWatch({
+ io: h.io,
+ probe: async () => {
+ asked += 1;
+ return "ok";
+ },
+ log: (l) => lines.push(l),
+ });
+ try {
+ assert.equal(asked, 1, "one pass at once");
+ // The save on /storage: written to settings, then applied at once.
+ h.io.read().storage.health = { passIntervalMs: 5_000 };
+ applyHealthTimings(h.io.read().storage.health);
+ assert.match(lines.join("\n"), /health pass re-armed: every 5 s/);
+ // At the default 15 s nothing would run for another 15 s; re-armed at
+ // 5 s, the next pass comes within about 5 s.
+ const started = Date.now();
+ while (asked < 2 && Date.now() - started < 7_000) {
+ await new Promise((r) => setTimeout(r, 50));
+ }
+ assert.equal(asked, 2, `a second pass after ${Date.now() - started} ms`);
+ assert.ok(Date.now() - started >= 4_500);
+ } finally {
+ stopStorageHealthWatch();
+ }
+ // Stopped: a later change re-arms nothing.
+ const before = lines.length;
+ applyHealthTimings({ passIntervalMs: 20_000 });
+ assert.equal(lines.length, before);
+ // Armed with an explicit interval, a change is not followed.
+ startStorageHealthWatch({ io: h.io, probe: async () => "ok", intervalMs: 60_000, log: (l) => lines.push(l) });
+ try {
+ applyHealthTimings({ passIntervalMs: 6_000 });
+ assert.equal(lines.some((l) => /re-armed/.test(l) && /6 s/.test(l)), false);
+ } finally {
+ stopStorageHealthWatch();
+ }
+ });
+});
diff --git a/common/controller/storageWatch.ts b/common/controller/storageWatch.ts
@@ -31,16 +31,18 @@ import {
type VolumeBins,
} from "../lib/storageVolumes";
import {
- HEALTH_PROBE_INTERVAL_MS,
- HEALTH_PROBE_TIMEOUT_MS,
NOT_ANSWERING,
+ applyHealthTimings,
+ healthTimings,
noteLocationDetector,
+ onPassIntervalChange,
pruneLocationHealth,
recordLocationHealth,
registerLocationHealth,
type HealthTransition,
type LocationHealthState,
} from "../lib/storageHealth";
+import { clearRuleText, secondsText } from "../lib/storageHealthTimings";
import { listChannelConfigs } from "./channels";
import { maybeAutoRepoint, probeAllLocations } from "./storageLocations";
@@ -361,12 +363,17 @@ export const STORAGE_WATCH_INTERVAL_MS = 5 * 60_000;
// 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 reads each location's block device
-// counters in /sys (`detectLocationHealth`; a child `stat` of the root only
-// where no device can be named) and records the answer in `lib/storageHealth.ts`,
-// which every page and poll consults before it touches a drive. One `stalled`
-// answer marks a location at once; two clean answers in a row clear it (the
-// rules are that module's).
+// answering at all. So every `storage.health.passIntervalMs` (15 s by default)
+// this reads each location's block device counters in /sys
+// (`detectLocationHealth`; a child `stat` of the root only where no device can
+// be named) and records the answer in `lib/storageHealth.ts`, which every page
+// and poll consults before it touches a drive. One `stalled` answer marks a
+// location at once; `clearAfterCleanPasses` clean answers in a row (two by
+// default) clear it (the rules are that module's).
+//
+// EVERY PASS APPLIES THE TIMINGS from the settings it reads (`applyHealthTimings`),
+// so a value changed by hand takes effect within one pass, and a changed
+// interval re-arms the pass's own timer (below).
//
// READ-ONLY AND IN MEMORY. It writes no settings and pauses nothing: the pass
// above sees a stalled location as down (its probe answers `stalled` without
@@ -382,7 +389,7 @@ export type StorageHealthPassOpts = {
// findmnt, for naming each root's block device. Default: getPaths().
bins?: Pick<VolumeBins, "findmntBin">;
// Test seam. Default: `detectLocationHealth` — the block device's counters,
- // or a child `stat` against a 3 s timer when no device can be named.
+ // or a child `stat` against `probeTimeoutMs` when no device can be named.
probe?: LocationHealthProbe;
now?: () => number;
log?: (line: string) => void;
@@ -401,7 +408,9 @@ function asVerdict(answer: LocationHealthState | HealthVerdict): HealthVerdict {
return typeof answer === "string" ? { answer } : answer;
}
-const STAT_CAUSE = `a stat of its root did not answer within ${HEALTH_PROBE_TIMEOUT_MS / 1000} s`;
+function statCause(): string {
+ return `a stat of its root did not answer within ${secondsText(healthTimings().probeTimeoutMs)}`;
+}
// Record one verdict. A verdict with no answer records nothing but the
// detector that gave it.
@@ -424,7 +433,7 @@ function recordVerdict(
}
return recordLocationHealth(loc, verdict.answer, {
now,
- cause: verdict.cause ?? STAT_CAUSE,
+ cause: verdict.cause ?? statCause(),
...(verdict.detector ? { detector: verdict.detector } : {}),
...(device !== undefined ? { device } : {}),
});
@@ -434,8 +443,12 @@ export async function runStorageHealthPass(
opts: StorageHealthPassOpts = {},
): Promise<StorageHealthPassResult> {
const log = opts.log ?? (() => {});
- const locations =
- opts.locations ?? (opts.io ?? DEFAULT_IO).read().storage.locations;
+ // The settings this pass runs on: its locations, and the drive-health
+ // timings, applied before anything is asked. A caller that hands in the
+ // locations reads no settings, and the timings stay as they were.
+ const storage = opts.locations ? null : (opts.io ?? DEFAULT_IO).read().storage;
+ if (storage) applyHealthTimings(storage.health);
+ const locations = opts.locations ?? storage?.locations ?? [];
pruneLocationHealth(locations.map((l) => l.id));
// Every configured location has an entry before anything is asked, so the
// watchdog (lib/storageHealth.ts `onDrive`) can find a channel's location
@@ -467,8 +480,9 @@ export async function runStorageHealthPass(
out.transitions.push(t);
if (t.to === "stalled") {
log(
- `[storage] "${loc.id}": ${NOT_ANSWERING} — ${verdict.cause ?? STAT_CAUSE}; ` +
- `pages and polls skip it until two passes in a row find it answering`,
+ `[storage] "${loc.id}": ${NOT_ANSWERING} — ${verdict.cause ?? statCause()}; ` +
+ `pages and polls skip it until it answers ` +
+ `${clearRuleText(healthTimings().clearAfterCleanPasses)}`,
);
} else if (t.from === "stalled") {
log(`[storage] "${loc.id}": answering again (${t.to})`);
@@ -481,7 +495,7 @@ export async function runStorageHealthPass(
// operator pressing it after doing something about the drive gets an answer
// taken afterwards. It counts as one answer like any other — a stalled location
// still needs two clean ones in a row — and the counters give none when their
-// last sample is under MIN_COUNTER_INTERVAL_MS old. Nothing is pruned.
+// last sample is under `minCounterIntervalMs()` old. Nothing is pruned.
export async function refreshLocationHealth(
loc: StorageLocation,
probe?: LocationHealthProbe,
@@ -500,11 +514,10 @@ export async function refreshLocationHealth(
// The cadence
// ---------------------------------------------------------------------------
-// Fifteen seconds (lib/storageHealth.ts says why). Each pass reads each
-// location's device counters in /sys (a findmnt only when the device is not
-// known yet or its /sys entry stopped reading; with no device, one short-lived
-// `stat`).
-export const STORAGE_HEALTH_INTERVAL_MS = HEALTH_PROBE_INTERVAL_MS;
+// `storage.health.passIntervalMs`, fifteen seconds by default
+// (lib/storageHealth.ts says why). Each pass reads each location's device
+// counters in /sys (a findmnt only when the device is not known yet or its /sys
+// entry stopped reading; with no device, one short-lived `stat`).
// A per-module-copy singleton, deliberately left so: it is not a temp-file
// name (slice W folded every tmp + rename onto lib/jsonFile-server.ts, whose
@@ -514,6 +527,7 @@ export const STORAGE_HEALTH_INTERVAL_MS = HEALTH_PROBE_INTERVAL_MS;
let timer: ReturnType<typeof setInterval> | null = null;
let healthTimer: ReturnType<typeof setInterval> | null = null;
let healthInFlight = false;
+let stopFollowingInterval: (() => void) | null = null;
// THE FIVE-MINUTE PASS, ARMED ONCE PER PROCESS, below the idle gate: it writes
// settings (an auto-pause). `unref()` so it never holds the event loop open — a
@@ -539,10 +553,16 @@ export function stopStorageWatch(): void {
timer = null;
}
-// THE FIFTEEN-SECOND HEALTH PASS, ARMED ONCE PER PROCESS, above the idle gate:
-// it is in memory and writes nothing (see the section header). It also runs
-// once at once, so the locations are registered and a drive that is already
-// stalled is on its way to being known before the first page.
+// THE HEALTH PASS, ARMED ONCE PER PROCESS, above the idle gate: it is in
+// memory and writes nothing (see the section header). It also runs once at
+// once, so the locations are registered and a drive that is already stalled is
+// on its way to being known before the first page.
+//
+// ITS INTERVAL FOLLOWS THE SETTING. Armed at `healthTimings().passIntervalMs`,
+// and RE-ARMED whenever an applied change moves it — a save on /storage (at
+// once, from the page's module copy: the subscription is on globalThis), or a
+// hand edit (at the next pass, which applies what it read). An explicit
+// `intervalMs` (the tests') is fixed and follows nothing.
export function startStorageHealthWatch(
opts: {
io?: { read: () => SiteSettings };
@@ -573,8 +593,19 @@ export function startStorageHealthWatch(
healthInFlight = false;
});
};
- healthTimer = setInterval(health, opts.intervalMs ?? STORAGE_HEALTH_INTERVAL_MS);
- healthTimer.unref?.();
+ const arm = (every: number) => {
+ if (healthTimer) clearInterval(healthTimer);
+ healthTimer = setInterval(health, every);
+ healthTimer.unref?.();
+ };
+ arm(opts.intervalMs ?? healthTimings().passIntervalMs);
+ if (opts.intervalMs === undefined) {
+ stopFollowingInterval = onPassIntervalChange((every) => {
+ if (!healthTimer) return;
+ (opts.log ?? console.log)(`[storage] health pass re-armed: every ${secondsText(every)}`);
+ arm(every);
+ });
+ }
health();
return true;
}
@@ -582,4 +613,6 @@ export function startStorageHealthWatch(
export function stopStorageHealthWatch(): void {
if (healthTimer) clearInterval(healthTimer);
healthTimer = null;
+ stopFollowingInterval?.();
+ stopFollowingInterval = null;
}
diff --git a/common/lib/channelMedia.ts b/common/lib/channelMedia.ts
@@ -3,14 +3,15 @@ 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,
+ healthTimings,
isDriveNotAnswering,
onDrive,
sinceText,
stalledLocationForPath,
type LocationHealth,
} from "./storageHealth";
+import { secondsText } from "./storageHealthTimings";
// WHERE A CHANNEL'S MEDIA ACTUALLY IS, and whether it can be reached.
//
@@ -221,7 +222,7 @@ export function stalledMediaLocation(
detail: health
? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`
: (detail ??
- `${NOT_ANSWERING} (a read did not answer within ${DRIVE_CALL_BUDGET_MS / 1000} s)`),
+ `${NOT_ANSWERING} (a read did not answer within ${secondsText(healthTimings().budgetMs)})`),
};
}
@@ -440,7 +441,8 @@ async function inspectOnDisk(
//
// 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.
+ // not answered within the budget (`storage.health.budgetMs`, 3 s by
+ // default) marks it stalled and answers `stalled` now.
try {
const st = await onDrive(configured, () => stat(configured));
if (!st.isDirectory()) {
diff --git a/common/lib/settingsDocs.ts b/common/lib/settingsDocs.ts
@@ -48,6 +48,10 @@ import {
STORAGE_SETTINGS_FIELD_DOCS,
STORAGE_VOLUME_FIELD_DOCS,
} from "./storageLocations";
+import {
+ HEALTH_TIMING_DEFAULTS,
+ STORAGE_HEALTH_SETTINGS_FIELD_DOCS,
+} from "./storageHealthTimings";
// `workers` IS LEFT OUT OF THE EXAMPLE, and that is the one place the example
// is not the literal default object. Its default is `[]`, and a settings.json
@@ -172,6 +176,13 @@ export function blockTables(d: SiteSettings): Partial<Record<keyof SiteSettings,
},
{ path: "storage.locations[]", docs: STORAGE_LOCATION_FIELD_DOCS },
{ path: "storage.locations[].volume", docs: STORAGE_VOLUME_FIELD_DOCS },
+ // Absent from the default block (only a tuned value is written), so the
+ // Default column is each timing's default, not the block's.
+ {
+ path: "storage.health",
+ docs: STORAGE_HEALTH_SETTINGS_FIELD_DOCS,
+ defaults: fromObject(HEALTH_TIMING_DEFAULTS),
+ },
],
buildPipeline: [
{
diff --git a/common/lib/settingsSchema.ts b/common/lib/settingsSchema.ts
@@ -59,6 +59,7 @@ import {
type StorageSettings,
type StorageVolume,
} from "./storageLocations";
+import { sanitizeStorageHealth } from "./storageHealthTimings";
import {
DEFAULT_DIARIZATION_ENGINE,
DEFAULT_DIARIZATION_THRESHOLD,
@@ -978,10 +979,15 @@ export function sanitizeStorage(value: unknown): StorageSettings {
const savedVideosLocationId = locations.some((l) => l.id === savedWanted)
? savedWanted
: "";
+ // THE DRIVE-HEALTH TIMINGS: each clamped into its range, and kept only where
+ // it differs from its default; a block with nothing left is not written
+ // (lib/storageHealthTimings.ts).
+ const health = sanitizeStorageHealth(r.health);
return {
locations,
defaultLocationId,
...(savedVideosLocationId ? { savedVideosLocationId } : {}),
+ ...(Object.keys(health).length > 0 ? { health } : {}),
};
}
diff --git a/common/lib/storageHealth.test.ts b/common/lib/storageHealth.test.ts
@@ -1,15 +1,15 @@
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,
+ applyHealthTimings,
driveCallsInFlight,
+ healthTimings,
isDriveNotAnswering,
noteLocationDetector,
onDrive,
+ onPassIntervalChange,
registerLocationHealth,
setCounterReader,
setDriveCallBudget,
@@ -22,6 +22,7 @@ import {
stalledLocation,
stalledLocationForPath,
} from "./storageHealth";
+import { HEALTH_TIMING_DEFAULTS } from "./storageHealthTimings";
// Run with:
// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/storageHealth.test.ts
@@ -32,6 +33,7 @@ import {
beforeEach(() => {
resetStorageHealth();
+ applyHealthTimings();
setDriveCallBudget();
setCounterReader(undefined);
});
@@ -55,7 +57,8 @@ test("a location first seen stalled is stalled", () => {
});
test("one clean probe after a stall does not clear it; two in a row do", () => {
- assert.equal(HEALTH_CLEAN_TO_CLEAR, 2);
+ assert.equal(HEALTH_TIMING_DEFAULTS.clearAfterCleanPasses, 2);
+ assert.equal(healthTimings().clearAfterCleanPasses, 2);
recordLocationHealth(USB, "ok", { now: 0 });
recordLocationHealth(USB, "stalled", { now: 10 });
assert.equal(recordLocationHealth(USB, "ok", { now: 20 }), null);
@@ -209,7 +212,7 @@ test("a stalled location is refused without the call being made", async () => {
});
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);
+ assert.equal(HEALTH_TIMING_DEFAULTS.inFlightPerLocation, 4);
registerLocationHealth([USB]);
setDriveCallBudget(80);
let made = 0;
@@ -380,7 +383,8 @@ test("L6: a call that times out while its disk is still completing requests is s
});
test("the default budget is 3 s", async () => {
- assert.equal(DRIVE_CALL_BUDGET_MS, 3_000);
+ assert.equal(HEALTH_TIMING_DEFAULTS.budgetMs, 3_000);
+ assert.equal(healthTimings().budgetMs, 3_000);
registerLocationHealth([USB]);
const started = Date.now();
await assert.rejects(() => onDrive(USB, never), DriveNotAnsweringError);
@@ -542,3 +546,125 @@ test("every slot held by an overdue call on a disk completing nothing: the next
assert.equal(locationHealth("usb")?.state, "stalled");
assert.match(String(locationHealth("usb")?.cause), /4 reads on it have not answered/);
});
+
+// ── the timings are settings (release 15 slice DT) ─────────────────────────
+// `applyHealthTimings` takes the stored `settings.storage.health` block, as the
+// health pass and /storage's save hand it over; `healthTimings()` is what every
+// number above is read through.
+
+test("DT: the applied budget feeds the watchdog — a unit slower than it is refused, and the same unit passes on the default", async () => {
+ registerLocationHealth([USB]);
+ // The settings' floor, 500 ms (a stored 200 clamps to it; storageHealthTimings.test.ts).
+ applyHealthTimings({ budgetMs: 200 });
+ assert.equal(healthTimings().budgetMs, 500);
+ const unit = () => sleep(700).then(() => "done");
+ await assert.rejects(() => onDrive(USB, unit), DriveNotAnsweringError);
+ assert.equal(locationHealth("usb")?.state, "stalled");
+ assert.match(String(locationHealth("usb")?.cause), /did not answer within 0.5 s/);
+ // Back on the defaults (3 s), with the location clear, the same unit answers.
+ await sleep(250);
+ applyHealthTimings();
+ resetStorageHealth();
+ registerLocationHealth([USB]);
+ assert.equal(await onDrive(USB, unit), "done");
+ assert.equal(locationHealth("usb")?.state, "ok");
+});
+
+test("DT: the applied cap — with inFlightPerLocation 2, the third call waits for a slot", async () => {
+ registerLocationHealth([USB]);
+ applyHealthTimings({ inFlightPerLocation: 2 });
+ const gates = Array.from({ length: 2 }, () => deferred<number>());
+ const first = gates.map((g) => onDrive(USB, () => g.promise));
+ let third = false;
+ const waiting = onDrive(USB, async () => {
+ third = true;
+ return 3;
+ });
+ await new Promise((r) => setImmediate(r));
+ assert.equal(driveCallsInFlight("usb"), 2);
+ assert.equal(third, false, "the third call is queued, not on the drive");
+ gates[0].resolve(1);
+ assert.equal(await waiting, 3);
+ gates[1].resolve(2);
+ await Promise.all(first);
+ assert.equal(driveCallsInFlight("usb"), 0);
+});
+
+test("DT: a raised cap admits the calls already waiting; a lowered one is reached as calls return", async () => {
+ registerLocationHealth([USB]);
+ applyHealthTimings({ inFlightPerLocation: 1 });
+ const gates = Array.from({ length: 3 }, () => deferred<number>());
+ let made = 0;
+ const calls = gates.map((g) =>
+ onDrive(USB, () => {
+ made += 1;
+ return g.promise;
+ }),
+ );
+ await new Promise((r) => setImmediate(r));
+ assert.equal(made, 1);
+ // Raised to 3 on a save: the two waiting take the new slots now.
+ applyHealthTimings({ inFlightPerLocation: 3 });
+ await new Promise((r) => setImmediate(r));
+ assert.equal(made, 3);
+ assert.equal(driveCallsInFlight("usb"), 3);
+ // Lowered to 1 with three in flight: a fourth call waits, and a return gives
+ // its slot back rather than handing it on while more than one is in flight.
+ applyHealthTimings({ inFlightPerLocation: 1 });
+ let fourth = false;
+ const late = onDrive(USB, async () => {
+ fourth = true;
+ return 4;
+ });
+ gates[0].resolve(0);
+ gates[1].resolve(0);
+ await new Promise((r) => setImmediate(r));
+ assert.equal(fourth, false, "two returns bring three down to one: no slot for the fourth yet");
+ assert.equal(driveCallsInFlight("usb"), 1);
+ gates[2].resolve(0);
+ assert.equal(await late, 4);
+ await Promise.all(calls);
+ assert.equal(driveCallsInFlight("usb"), 0);
+});
+
+test("DT: the applied clear count — three clean answers with clearAfterCleanPasses 3, one with 1", () => {
+ applyHealthTimings({ clearAfterCleanPasses: 3 });
+ recordLocationHealth(USB, "stalled", { now: 0 });
+ recordLocationHealth(USB, "ok", { now: 1 });
+ recordLocationHealth(USB, "ok", { now: 2 });
+ assert.equal(locationHealth("usb")?.state, "stalled", "two are not three");
+ assert.deepEqual(recordLocationHealth(USB, "ok", { now: 3 }), {
+ id: "usb",
+ from: "stalled",
+ to: "ok",
+ });
+ applyHealthTimings({ clearAfterCleanPasses: 1 });
+ recordLocationHealth(USB, "stalled", { now: 4 });
+ recordLocationHealth(USB, "ok", { now: 5 });
+ assert.equal(locationHealth("usb")?.state, "ok");
+});
+
+test("DT: a changed pass interval is told to the subscribers; the same one, or another timing, is not", () => {
+ const heard: number[] = [];
+ const stop = onPassIntervalChange((ms) => heard.push(ms));
+ try {
+ applyHealthTimings({ passIntervalMs: 30_000 });
+ applyHealthTimings({ passIntervalMs: 30_000, budgetMs: 5_000 });
+ applyHealthTimings();
+ assert.deepEqual(heard, [30_000, 15_000]);
+ } finally {
+ stop();
+ }
+ applyHealthTimings({ passIntervalMs: 60_000 });
+ assert.deepEqual(heard, [30_000, 15_000], "unsubscribed");
+});
+
+test("DT: the test seam's budget wins over the applied one, and the timings survive a reset", () => {
+ applyHealthTimings({ budgetMs: 8_000, inFlightPerLocation: 6 });
+ setDriveCallBudget(80);
+ assert.equal(healthTimings().budgetMs, 80);
+ setDriveCallBudget();
+ assert.equal(healthTimings().budgetMs, 8_000);
+ resetStorageHealth();
+ assert.equal(healthTimings().inFlightPerLocation, 6, "configuration, not health");
+});
diff --git a/common/lib/storageHealth.ts b/common/lib/storageHealth.ts
@@ -1,4 +1,11 @@
import { locationOfDataDir, type StorageLocation } from "./storageLocations";
+import {
+ HEALTH_TIMING_DEFAULTS,
+ resolveHealthTimings,
+ secondsText,
+ type HealthTimings,
+ type StorageHealthSettings,
+} from "./storageHealthTimings";
// IS A STORAGE LOCATION'S DRIVE ANSWERING RIGHT NOW — the in-memory answer every
// page and poll asks before it touches the drive.
@@ -13,11 +20,14 @@ import { locationOfDataDir, type StorageLocation } from "./storageLocations";
//
// SO SOMETHING THAT CANNOT HANG ASKS, AND THIS MODULE REMEMBERS WHAT IT SAID.
// Two detectors write here. The health pass (`controller/storageWatch.ts`,
-// every 15 s) reads each location's block device counters in /sys, which never
-// touch the drive (`detectLocationHealth` in `storageVolumes.ts`; a child `stat`
-// of the root raced against 3 s only where no device can be named). And
-// `onDrive` below races every in-process call the gate covers against 3 s and
-// marks the location the moment one does not answer. Everything that would
+// every `passIntervalMs`) reads each location's block device counters in /sys,
+// which never touch the drive (`detectLocationHealth` in `storageVolumes.ts`; a
+// child `stat` of the root raced against `probeTimeoutMs` only where no device
+// can be named). And `onDrive` below races every in-process call the gate
+// covers against `budgetMs` and marks the location the moment one does not
+// answer. The numbers are `settings.storage.health`, read through
+// `healthTimings()` below (by default a pass every 15 s, a 3 s probe and a 3 s
+// budget). Everything that would
// touch the drive in-process asks this state 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
@@ -27,12 +37,12 @@ import { locationOfDataDir, type StorageLocation } from "./storageLocations";
// - ONE `stalled` answer marks the location `stalled` at once. A drive that
// did not answer will not answer the next page either, and every page that
// asks costs a thread.
-// - TWO consecutive clean answers clear it. A clean answer is anything else:
-// the counters moving or idle, or a child `stat` answering 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.
+// - `clearAfterCleanPasses` (TWO by default) consecutive clean answers clear
+// it. A clean answer is anything else: the counters moving or idle, or a
+// child `stat` answering 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.
//
@@ -42,6 +52,13 @@ import { locationOfDataDir, type StorageLocation } from "./storageLocations";
// boot too) registers the locations; the watchdog re-learns a stall the moment
// a page reaches the drive.
//
+// THE TIMINGS ARE SETTINGS (`settings.storage.health`, release 15 slice DT),
+// held here as the process last applied them: the health pass applies them
+// from the settings it reads on every pass, and /storage's save applies them at
+// once. A process with no pass (a CLI) runs on the defaults unless its entry
+// applies them (the index and stats bins do). `healthTimings()` is the one
+// accessor. This module still reads no file.
+//
// 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
@@ -78,7 +95,7 @@ export type LocationHealth = {
// 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.
+ // at `clearAfterCleanPasses`.
cleanStreak: number;
// What did not answer, for a stalled location: the probe's own words.
cause?: string;
@@ -87,16 +104,16 @@ export type LocationHealth = {
device?: 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 DEFAULTS, and why they are what they are (lib/storageHealthTimings.ts
+// holds them, with their ranges):
+// - `probeTimeoutMs` 3 s: 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.
+// - `passIntervalMs` 15 s: short, because the gate is only as current as the
+// last answer — a stall that began just after a pass costs every page that
+// touches the drive until the next one.
+// - `clearAfterCleanPasses` 2: one clean answer in a reset loop is not
+// recovery.
// The one wording of the state, for every surface that shows it.
export const NOT_ANSWERING = "drive not answering";
@@ -123,7 +140,13 @@ type HealthState = {
inFlight?: Map<string, number>;
overdue?: Map<string, OverdueCall[]>;
waiters?: Map<string, Waiter[]>;
- // Test seam: the watchdog's budget.
+ // The timings as last applied (`applyHealthTimings`); absent = the defaults.
+ timings?: HealthTimings;
+ // Told when the pass interval changes, so the armed pass re-arms its timer
+ // (controller/storageWatch.ts). On globalThis like the rest: the save that
+ // changes it runs in a page's module copy, the pass in instrumentation's.
+ intervalListeners?: Set<(passIntervalMs: number) => void>;
+ // Test seam: the watchdog's budget, below the settings' 500 ms floor.
budgetMs?: number;
// Reads a block device's counters, synchronously and without touching the
// drive (set by `storageVolumes.ts`, which owns /sys). Absent: no check.
@@ -136,7 +159,7 @@ declare global {
}
type FilledState = HealthState &
- Required<Pick<HealthState, "inFlight" | "overdue" | "waiters">>;
+ Required<Pick<HealthState, "inFlight" | "overdue" | "waiters" | "intervalListeners">>;
function healthState(): FilledState {
if (!globalThis.__yttStorageHealth__) {
@@ -148,11 +171,57 @@ function healthState(): FilledState {
s.inFlight ??= new Map();
s.overdue ??= new Map();
s.waiters ??= new Map();
+ s.intervalListeners ??= new Set();
return s as FilledState;
}
+// THE ONE ACCESSOR of the drive-health timings: `settings.storage.health` as
+// this process last applied it, every absent key its default. Read at the
+// moment a number is needed (a call's timer, a slot, an answer's count, a
+// pass's probe), so an applied change takes effect on the next of each.
+export function healthTimings(): HealthTimings {
+ const s = healthState();
+ const t = s.timings ?? HEALTH_TIMING_DEFAULTS;
+ return s.budgetMs === undefined ? t : { ...t, budgetMs: s.budgetMs };
+}
+
+// Apply `settings.storage.health` (the stored block; absent = every default).
+// The health pass calls this with the settings it reads on every pass, and
+// /storage's save calls it at once. A changed pass interval is told to the
+// armed pass (`onPassIntervalChange`); a raised cap lets calls already waiting
+// for a slot take the new ones.
+export function applyHealthTimings(stored?: StorageHealthSettings): HealthTimings {
+ const s = healthState();
+ const before = healthTimings();
+ s.timings = resolveHealthTimings(stored);
+ const after = healthTimings();
+ if (after.inFlightPerLocation > before.inFlightPerLocation) {
+ for (const key of [...s.waiters.keys()]) admitWaiters(key);
+ }
+ if (after.passIntervalMs !== before.passIntervalMs) {
+ for (const fn of [...s.intervalListeners]) {
+ try {
+ fn(after.passIntervalMs);
+ } catch {
+ /* a listener that throws must not stop a save */
+ }
+ }
+ }
+ return after;
+}
+
+// Subscribe to pass-interval changes; returns the unsubscribe.
+export function onPassIntervalChange(fn: (passIntervalMs: number) => void): () => void {
+ const set = healthState().intervalListeners;
+ set.add(fn);
+ return () => {
+ set.delete(fn);
+ };
+}
+
// Test seam, and the escape hatch for a process that wants to forget. Calls
-// still waiting for a slot are released to run.
+// still waiting for a slot are released to run. The applied timings are
+// configuration, not health, and are kept.
export function resetStorageHealth(): void {
const s = healthState();
s.byId.clear();
@@ -246,7 +315,7 @@ function recordAnswer(
}
if (prev.state === "stalled") {
prev.cleanStreak += 1;
- if (prev.cleanStreak < HEALTH_CLEAN_TO_CLEAR) return null;
+ if (prev.cleanStreak < healthTimings().clearAfterCleanPasses) return null;
prev.state = answer;
prev.since = now;
prev.cleanStreak = 0;
@@ -362,17 +431,18 @@ export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number)
}
// ---------------------------------------------------------------------------
-// The watchdog: every call the gate covers, raced against 3 s
+// The watchdog: every call the gate covers, raced against the budget
// ---------------------------------------------------------------------------
//
-// 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 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.
+// THE DETECTOR THAT CANNOT BE FOOLED BY A CACHE. The pass reads the block
+// device's counters (or, with no device, a child `stat`), and a stall that
+// starts between two passes is not seen by it. A page or a poll that actually
+// reaches the drive is: `onDrive` runs the call against a timer of
+// `settings.storage.health.budgetMs` (3 s by default; `healthTimings()`), 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.
//
// THE BUDGET COVERS A WHOLE UNIT OF WORK. A caller sends a unit through as one
// call (a video directory's few reads, a page's reads of one video), so a slow
@@ -383,7 +453,8 @@ export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number)
// requests completed meanwhile the drive is slow, not stalled, and only this
// call is refused.
//
-// AT MOST DRIVE_CALLS_IN_FLIGHT CALLS PER SLOT KEY ARE IN FLIGHT. A walk of
+// AT MOST `inFlightPerLocation` (4 by default) CALLS PER SLOT KEY ARE IN
+// FLIGHT. 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. The rest wait in a queue of our own:
// - a transition to `stalled` (from any detector) refuses them at once;
@@ -397,9 +468,12 @@ export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number)
// began (slow, not stalled — the same test as a timeout's): none of those
// calls has returned, whatever the last pass said.
// A slot is released when its call really returns, not when the watchdog gave
-// up on it. So on one location at most four threads wait on its drive for the
-// calls that come through here — every page and poll path, and the snapshot
-// walk. A job's own reads that do not come through here are not capped.
+// up on it. So on one location at most `inFlightPerLocation` threads wait on
+// its drive for the calls that come through here — every page and poll path,
+// and the snapshot walk. A job's own reads that do not come through here are
+// not capped. A cap lowered while calls are in flight is reached as they
+// return (a returning call gives its slot back instead of handing it on); a
+// cap raised lets waiting calls take the new slots at once.
//
// THE SLOT KEY: a configured location's id; for a probe of another root under a
// location's id, that root; for a path on no configured location (a root typed
@@ -409,16 +483,14 @@ export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number)
// Do not nest `onDrive` for one key: the inner call would wait for a slot the
// outer one holds.
-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.
+// Test seam: shorten (or restore, with no argument) the watchdog's budget,
+// past the settings' 500 ms floor. It wins over the applied timings.
export function setDriveCallBudget(ms?: number): void {
healthState().budgetMs = ms;
}
function driveCallBudget(): number {
- return healthState().budgetMs ?? DRIVE_CALL_BUDGET_MS;
+ return healthTimings().budgetMs;
}
export class DriveNotAnsweringError extends Error {
@@ -429,7 +501,7 @@ export class DriveNotAnsweringError extends Error {
detail ??
(health
? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`
- : `${NOT_ANSWERING} (a read did not answer within ${driveCallBudget() / 1000} s)`),
+ : `${NOT_ANSWERING} (a read did not answer within ${secondsText(driveCallBudget())})`),
);
this.name = "DriveNotAnsweringError";
this.health = health;
@@ -500,13 +572,14 @@ function refuseWaiters(key: string, health: LocationHealth | null): void {
// refused at once when every slot is held by an overdue call.
async function acquireSlot(r: Resolved): Promise<void> {
const s = healthState();
+ const cap = healthTimings().inFlightPerLocation;
const n = s.inFlight.get(r.key) ?? 0;
- if (n < DRIVE_CALLS_IN_FLIGHT) {
+ if (n < cap) {
s.inFlight.set(r.key, n + 1);
return;
}
const overdue = s.overdue.get(r.key) ?? [];
- if (overdue.length >= DRIVE_CALLS_IN_FLIGHT) {
+ if (overdue.length >= cap) {
// EVERY SLOT IS HELD BY A CALL THE WATCHDOG GAVE UP ON: refused at once.
// Marked stalled only when the disk has not been completing requests
// since the oldest of them began — the watchdog's own slow-or-stalled test
@@ -516,14 +589,14 @@ async function acquireSlot(r: Resolved): Promise<void> {
if (oldest.before && now && now.completed > oldest.before.completed) {
throw new DriveNotAnsweringError(
null,
- `drive slow (${DRIVE_CALLS_IN_FLIGHT} reads on it are past ` +
- `${driveCallBudget() / 1000} s, while its disk is still completing others)`,
+ `drive slow (${overdue.length} reads on it are past ` +
+ `${secondsText(driveCallBudget())}, while its disk is still completing others)`,
);
}
let health: LocationHealth | null = null;
if (r.loc) {
recordLocationHealth(r.loc, "stalled", {
- cause: `${DRIVE_CALLS_IN_FLIGHT} reads on it have not answered`,
+ cause: `${overdue.length} reads on it have not answered`,
});
health = stalledLocation(r.loc);
}
@@ -559,7 +632,7 @@ async function acquireSlot(r: Resolved): Promise<void> {
reject(
new DriveNotAnsweringError(
r.loc ? stalledLocation(r.loc) : null,
- `${NOT_ANSWERING} (nothing on it answered for ${budget / 1000} s while a read waited)`,
+ `${NOT_ANSWERING} (nothing on it answered for ${secondsText(budget)} while a read waited)`,
),
);
};
@@ -595,16 +668,32 @@ async function acquireSlot(r: Resolved): Promise<void> {
// A slot is released when its call returns (or, on a refusal after a wait,
// without a call). A return is progress: every waiter on the key has its
-// deadline restarted, then the slot goes to the first waiter still waiting.
+// deadline restarted, then the slot goes to the first waiter still waiting —
+// unless the cap was lowered below what is in flight, when it is given back.
function releaseSlot(key: string): void {
const s = healthState();
const q = s.waiters.get(key);
if (q) for (const w of q) w.rearm();
- while (q && q.length > 0) {
+ const n = s.inFlight.get(key) ?? 1;
+ if (n <= healthTimings().inFlightPerLocation) {
+ while (q && q.length > 0) {
+ const next = q.shift() as Waiter;
+ if (next.resolve()) return;
+ }
+ }
+ s.inFlight.set(key, Math.max(0, n - 1));
+}
+
+// A raised cap: the calls already waiting on `key` take the new slots now,
+// rather than one at a time as calls return.
+function admitWaiters(key: string): void {
+ const s = healthState();
+ const q = s.waiters.get(key);
+ const cap = healthTimings().inFlightPerLocation;
+ while (q && q.length > 0 && (s.inFlight.get(key) ?? 0) < cap) {
const next = q.shift() as Waiter;
- if (next.resolve()) return;
+ if (next.resolve()) s.inFlight.set(key, (s.inFlight.get(key) ?? 0) + 1);
}
- s.inFlight.set(key, Math.max(0, (s.inFlight.get(key) ?? 1) - 1));
}
// How many calls are in flight on a location through `onDrive` (for tests and
@@ -626,8 +715,8 @@ function readDeviceCounters(loc: Resolved["loc"]): BlockStatSample | null {
}
// 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 its key, and raced against the budget. Throws
+// location is stalled, queued behind `inFlightPerLocation` calls already in
+// flight on its key, and raced against `budgetMs` (`healthTimings()`). 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> {
@@ -658,11 +747,13 @@ export async function onDrive<T>(where: Where, call: () => Promise<T>): Promise<
throw err;
}
const TIMED_OUT = Symbol("timed out");
+ // The budget this call runs against, read as it starts.
+ const budget = driveCallBudget();
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());
+ timer = setTimeout(() => resolve(TIMED_OUT), budget);
});
let answer: T | typeof TIMED_OUT;
try {
@@ -693,13 +784,13 @@ export async function onDrive<T>(where: Where, call: () => Promise<T>): Promise<
// SLOW, NOT STALLED: the disk completed requests while this call waited.
throw new DriveNotAnsweringError(
null,
- `drive slow (a read did not answer within ${driveCallBudget() / 1000} s, ` +
+ `drive slow (a read did not answer within ${secondsText(budget)}, ` +
`while its disk was still completing others)`,
);
}
if (r.loc) {
recordLocationHealth(r.loc, "stalled", {
- cause: `a read in the editor did not answer within ${driveCallBudget() / 1000} s`,
+ cause: `a read in the editor did not answer within ${secondsText(budget)}`,
});
throw new DriveNotAnsweringError(stalledLocation(r.loc));
}
diff --git a/common/lib/storageHealthCounters.test.ts b/common/lib/storageHealthCounters.test.ts
@@ -6,10 +6,10 @@ import { chmod, mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
import {
blockDeviceName,
detectLocationHealth,
- MIN_COUNTER_INTERVAL_MS,
+ minCounterIntervalMs,
resetHealthDetector,
} from "./storageVolumes";
-import { countersVerdict, parseBlockStat } from "./storageHealth";
+import { applyHealthTimings, countersVerdict, parseBlockStat } from "./storageHealth";
// Run with:
// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/storageHealthCounters.test.ts
@@ -19,7 +19,10 @@ import { countersVerdict, parseBlockStat } from "./storageHealth";
// the test writes — a stalled device is one whose in-flight count stays up
// while its completions stand still.
-beforeEach(() => resetHealthDetector());
+beforeEach(() => {
+ resetHealthDetector();
+ applyHealthTimings();
+});
// A real line from this machine's /sys/class/block/<dev>/stat (17 fields).
const LINE =
@@ -147,8 +150,9 @@ test("counters: a second sample sooner than the minimum interval gives no verdic
await detectLocationHealth(loc(h.root), h.bins, { sysBlockDir: h.sys, now: 0 });
const soon = await detectLocationHealth(loc(h.root), h.bins, {
sysBlockDir: h.sys,
- now: MIN_COUNTER_INTERVAL_MS - 1,
+ now: minCounterIntervalMs() - 1,
});
+ assert.equal(minCounterIntervalMs(), 10_000);
assert.equal(soon.answer, null);
// Compared with the FIRST sample, not the refused one.
const later = await detectLocationHealth(loc(h.root), h.bins, {
@@ -159,6 +163,31 @@ test("counters: a second sample sooner than the minimum interval gives no verdic
});
});
+test("counters: the minimum spacing follows the pass interval (storage.health.passIntervalMs)", async () => {
+ await withHarness(async (h) => {
+ await h.control({ source: "/dev/fakedisk1" });
+ await h.counters("fakedisk1", 100, 2);
+ // A pass every 8 s: samples 4 s apart are compared (8 − 5 = 3 s, floored
+ // at half the interval), which the default 15 s pass would refuse.
+ applyHealthTimings({ passIntervalMs: 8_000 });
+ assert.equal(minCounterIntervalMs(), 4_000);
+ await detectLocationHealth(loc(h.root), h.bins, { sysBlockDir: h.sys, now: 0 });
+ const early = await detectLocationHealth(loc(h.root), h.bins, {
+ sysBlockDir: h.sys,
+ now: 3_999,
+ });
+ assert.equal(early.answer, null);
+ const due = await detectLocationHealth(loc(h.root), h.bins, {
+ sysBlockDir: h.sys,
+ now: 4_000,
+ });
+ assert.equal(due.answer, "stalled");
+ // A pass every 5 minutes: 10 s, as at the default.
+ applyHealthTimings({ passIntervalMs: 300_000 });
+ assert.equal(minCounterIntervalMs(), 10_000);
+ });
+});
+
test("no device → the child stat, and the verdict says so", async () => {
await withHarness(async (h) => {
const opts = { sysBlockDir: h.sys, now: 0 };
diff --git a/common/lib/storageHealthProbe.test.ts b/common/lib/storageHealthProbe.test.ts
@@ -4,7 +4,8 @@ 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";
+import { HEALTH_TIMING_DEFAULTS } from "./storageHealthTimings";
+import { applyHealthTimings } from "./storageHealth";
// Run with:
// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/storageHealthProbe.test.ts
@@ -75,7 +76,7 @@ test("a child that never answers is 'stalled' on the timer, without waiting for
});
test("the default budget is 3 s", async () => {
- assert.equal(HEALTH_PROBE_TIMEOUT_MS, 3_000);
+ assert.equal(HEALTH_TIMING_DEFAULTS.probeTimeoutMs, 3_000);
await withDir(async (dir) => {
const statBin = await fakeStat(dir, { sleepMs: 20_000, out: "directory" });
const started = Date.now();
@@ -85,6 +86,21 @@ test("the default budget is 3 s", async () => {
});
});
+test("the budget is storage.health.probeTimeoutMs when the caller names none", async () => {
+ applyHealthTimings({ probeTimeoutMs: 500 });
+ try {
+ 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 >= 500 && took < 2_500, `answered after ${took} ms`);
+ });
+ } finally {
+ applyHealthTimings();
+ }
+});
+
test("an answer inside the budget is taken as given", async () => {
await withDir(async (dir) => {
const slowDir = await fakeStat(dir, { sleepMs: 100, out: "directory" });
diff --git a/common/lib/storageHealthTimings.test.ts b/common/lib/storageHealthTimings.test.ts
@@ -0,0 +1,138 @@
+// THE DRIVE-HEALTH TIMINGS AS SETTINGS: `settings.storage.health`.
+//
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/storageHealthTimings.test.ts
+//
+// The claims: absent is every default; a read clamps an out-of-range value into
+// its range (the schema's rule: a read never throws); only a value that differs
+// from its default is kept, so an untuned file has no `health` key and a save
+// writes none; the block survives a save and a read, and the mediaRoot
+// migration. Its own file for the SETTINGS SEAM (settingsWrite.test.ts says
+// why): SETTINGS_FILE is set before the settings module is imported.
+
+import { mkdtempSync, readFileSync, writeFileSync } from "node:fs";
+import { rm } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+import { after, test } from "node:test";
+import assert from "node:assert/strict";
+import {
+ HEALTH_TIMING_BOUNDS,
+ HEALTH_TIMING_DEFAULTS,
+ HEALTH_TIMING_KEYS,
+ clearRuleText,
+ counterSampleMinimumMs,
+ resolveHealthTimings,
+ sanitizeStorageHealth,
+ secondsText,
+} from "./storageHealthTimings";
+
+const ROOT = mkdtempSync(path.join(os.tmpdir(), "health-timings-"));
+process.env.TRANSCRIPTS_DIR = ROOT;
+process.env.SETTINGS_FILE = path.join(ROOT, "settings.json");
+
+const { getSettings, writeSettings, siteSettingsSchema, defaultSiteSettings } =
+ await import("./settings");
+
+after(() => rm(ROOT, { recursive: true, force: true }));
+
+const onDisk = () =>
+ JSON.parse(readFileSync(process.env.SETTINGS_FILE!, "utf8")) as {
+ storage: Record<string, unknown>;
+ };
+
+test("the defaults are the constants they replace", () => {
+ assert.deepEqual(HEALTH_TIMING_DEFAULTS, {
+ budgetMs: 3_000,
+ passIntervalMs: 15_000,
+ probeTimeoutMs: 3_000,
+ clearAfterCleanPasses: 2,
+ inFlightPerLocation: 4,
+ });
+ assert.deepEqual([...HEALTH_TIMING_KEYS].sort(), Object.keys(HEALTH_TIMING_DEFAULTS).sort());
+ for (const key of HEALTH_TIMING_KEYS) {
+ const { min, max } = HEALTH_TIMING_BOUNDS[key];
+ const d = HEALTH_TIMING_DEFAULTS[key];
+ assert.ok(min <= d && d <= max, `${key}'s default is in its range`);
+ }
+});
+
+test("absent is every default: no block, an empty block, junk", () => {
+ for (const raw of [undefined, null, {}, [], "3000", 7]) {
+ assert.deepEqual(sanitizeStorageHealth(raw), {});
+ assert.deepEqual(resolveHealthTimings(raw as never), HEALTH_TIMING_DEFAULTS);
+ }
+ assert.equal(getSettings().storage.health, undefined, "no file: no block");
+ assert.equal(defaultSiteSettings().storage.health, undefined);
+});
+
+test("a read clamps into the range, rounds, drops non-numbers and keeps only what differs from the default", () => {
+ assert.deepEqual(
+ sanitizeStorageHealth({
+ budgetMs: 200, // below 500
+ passIntervalMs: 900_000, // above 300000
+ probeTimeoutMs: "4000", // a string is not a number
+ clearAfterCleanPasses: 2, // the default
+ inFlightPerLocation: 7.6,
+ unknown: 1,
+ }),
+ { budgetMs: 500, passIntervalMs: 300_000, inFlightPerLocation: 8 },
+ );
+ // A value that clamps ONTO its default is the default, and is not kept.
+ assert.deepEqual(sanitizeStorageHealth({ clearAfterCleanPasses: 2.2 }), {});
+ assert.deepEqual(resolveHealthTimings({ budgetMs: 200 }), {
+ ...HEALTH_TIMING_DEFAULTS,
+ budgetMs: 500,
+ });
+});
+
+test("a settings.json round trip: a tuned value is written, read back, and a default is not written", async () => {
+ const base = defaultSiteSettings();
+ await writeSettings({
+ ...base,
+ storage: { ...base.storage, health: { budgetMs: 4_000, clearAfterCleanPasses: 2 } },
+ });
+ assert.deepEqual(onDisk().storage.health, { budgetMs: 4_000 });
+ assert.deepEqual(getSettings().storage.health, { budgetMs: 4_000 });
+ // Every value back to its default: the key is not written at all.
+ await writeSettings({ ...base, storage: { ...base.storage, health: { budgetMs: 3_000 } } });
+ assert.equal("health" in onDisk().storage, false);
+ // A hand-edited file out of range reads clamped.
+ writeFileSync(
+ process.env.SETTINGS_FILE!,
+ JSON.stringify({ storage: { locations: [], health: { inFlightPerLocation: 99 } } }),
+ );
+ assert.deepEqual(getSettings().storage.health, { inFlightPerLocation: 16 });
+});
+
+test("the timings survive the mediaRoot migration (a block with no locations)", () => {
+ const parsed = siteSettingsSchema.parse({
+ storage: { health: { passIntervalMs: 30_000 } },
+ });
+ assert.deepEqual(parsed.storage.health, { passIntervalMs: 30_000 });
+ writeFileSync(
+ process.env.SETTINGS_FILE!,
+ JSON.stringify({ storage: { mediaRoot: "/mnt/cold", health: { budgetMs: 6_000 } } }),
+ );
+ const s = getSettings();
+ assert.deepEqual(s.storage.locations.map((l) => l.root), ["/mnt/cold"]);
+ assert.deepEqual(s.storage.health, { budgetMs: 6_000 });
+});
+
+test("the counters' sample spacing: min(10 s, interval − 5 s), floored at half the interval", () => {
+ assert.equal(counterSampleMinimumMs(15_000), 10_000, "the default, as it was");
+ assert.equal(counterSampleMinimumMs(300_000), 10_000);
+ assert.equal(counterSampleMinimumMs(12_000), 7_000);
+ assert.equal(counterSampleMinimumMs(10_000), 5_000);
+ assert.equal(counterSampleMinimumMs(8_000), 4_000, "the floor: 3 s would be under half");
+ assert.equal(counterSampleMinimumMs(5_000), 2_500, "never 0 at the minimum interval");
+});
+
+test("the words", () => {
+ assert.equal(secondsText(3_000), "3 s");
+ assert.equal(secondsText(2_500), "2.5 s");
+ assert.equal(secondsText(80), "0.08 s");
+ assert.equal(clearRuleText(1), "once");
+ assert.equal(clearRuleText(2), "twice in a row");
+ assert.equal(clearRuleText(5), "5 times in a row");
+});
diff --git a/common/lib/storageHealthTimings.ts b/common/lib/storageHealthTimings.ts
@@ -0,0 +1,161 @@
+import type { FieldDocs } from "./fieldDocs";
+
+// THE DRIVE-HEALTH TIMINGS — `settings.storage.health` — and their defaults.
+//
+// The health gate (lib/storageHealth.ts) decides that a drive is not answering
+// from a handful of numbers: how long one read may take, how often the health
+// pass looks at the drive, how long that look may take, how many clean looks in
+// a row clear a stall, and how many reads may wait on one drive at once. They
+// were constants. Under heavy external-disk churn a stall can be misjudged, and
+// the ruling (release 15, slice DT) is that the operator then tunes the numbers
+// on /storage rather than the code. Today's constants are the defaults.
+//
+// THE STORED BLOCK HOLDS ONLY WHAT DIFFERS FROM A DEFAULT. Every key is
+// optional, an absent key is its default, and `sanitizeStorageHealth` drops a
+// value equal to its default — so a settings.json that never tuned anything has
+// no `health` key at all, and a default changed in a later release reaches it.
+//
+// READS CLAMP, THE FORM REFUSES. A hand-edited value outside its range is
+// clamped into it on read (the schema's rule: a read never throws); the /storage
+// form refuses one with a sentence instead, because a save that quietly stored
+// another number than the one typed reads as a form that did not listen.
+//
+// PURE, NO IMPORTS BUT A TYPE: the /storage form (a "use client" file) imports
+// the defaults, the ranges and the words from here.
+
+export type StorageHealthSettings = {
+ budgetMs?: number;
+ passIntervalMs?: number;
+ probeTimeoutMs?: number;
+ clearAfterCleanPasses?: number;
+ inFlightPerLocation?: number;
+};
+
+// Every timing, resolved: what the health module runs on.
+export type HealthTimings = Required<StorageHealthSettings>;
+
+export type HealthTimingKey = keyof HealthTimings;
+
+// Today's constants (release 15 slice DS), now the defaults.
+export const HEALTH_TIMING_DEFAULTS: Readonly<HealthTimings> = Object.freeze({
+ budgetMs: 3_000,
+ passIntervalMs: 15_000,
+ probeTimeoutMs: 3_000,
+ clearAfterCleanPasses: 2,
+ inFlightPerLocation: 4,
+});
+
+export const HEALTH_TIMING_BOUNDS: Readonly<
+ Record<HealthTimingKey, { min: number; max: number }>
+> = Object.freeze({
+ budgetMs: { min: 500, max: 60_000 },
+ passIntervalMs: { min: 5_000, max: 300_000 },
+ probeTimeoutMs: { min: 500, max: 30_000 },
+ clearAfterCleanPasses: { min: 1, max: 10 },
+ inFlightPerLocation: { min: 1, max: 16 },
+});
+
+// In the order the form shows them.
+export const HEALTH_TIMING_KEYS: readonly HealthTimingKey[] = [
+ "budgetMs",
+ "passIntervalMs",
+ "probeTimeoutMs",
+ "clearAfterCleanPasses",
+ "inFlightPerLocation",
+];
+
+// A whole number in range, or undefined (not a finite number).
+export function clampHealthTiming(key: HealthTimingKey, value: unknown): number | undefined {
+ if (typeof value !== "number" || !Number.isFinite(value)) return undefined;
+ const { min, max } = HEALTH_TIMING_BOUNDS[key];
+ return Math.min(max, Math.max(min, Math.round(value)));
+}
+
+// Coerce a raw `settings.storage.health` into the stored block: each value a
+// number clamped into its range, and only when it differs from its default.
+// Anything else — a string, a key the block does not name — is dropped.
+export function sanitizeStorageHealth(value: unknown): StorageHealthSettings {
+ if (!value || typeof value !== "object" || Array.isArray(value)) return {};
+ const r = value as Record<string, unknown>;
+ const out: StorageHealthSettings = {};
+ for (const key of HEALTH_TIMING_KEYS) {
+ const n = clampHealthTiming(key, r[key]);
+ if (n !== undefined && n !== HEALTH_TIMING_DEFAULTS[key]) out[key] = n;
+ }
+ return out;
+}
+
+// The stored block with every absent key filled from the defaults.
+export function resolveHealthTimings(stored?: StorageHealthSettings): HealthTimings {
+ const clean = sanitizeStorageHealth(stored);
+ return { ...HEALTH_TIMING_DEFAULTS, ...clean };
+}
+
+// THE COUNTERS' SAMPLE SPACING, derived from the pass interval. Two samples of a
+// disk's request counters closer than this are not compared: a healthy drive can
+// have a request in flight at two instants a moment apart without completing
+// one (release 15 DS, "kept at review"). The ruling: min(10 s, interval − 5 s),
+// so consecutive passes always compare — which is 10 s at the default 15 s, as
+// it was. Below a 10 s interval that formula falls toward nothing (0 at the 5 s
+// minimum, where a /storage Refresh just after a pass would compare two samples
+// taken milliseconds apart), so it is floored at half the interval.
+export function counterSampleMinimumMs(passIntervalMs: number): number {
+ return Math.min(10_000, Math.max(passIntervalMs - 5_000, Math.round(passIntervalMs / 2)));
+}
+
+// "3 s", "2.5 s", "0.08 s": a millisecond figure in the seconds every surface
+// uses.
+export function secondsText(ms: number): string {
+ return `${ms / 1000} s`;
+}
+
+// How many clean answers clear a stall, in words: "once", "twice in a row",
+// "3 times in a row".
+export function clearRuleText(n: number): string {
+ return n === 1 ? "once" : n === 2 ? "twice in a row" : `${n} times in a row`;
+}
+
+// What each field means, in the operator's words: the /storage form's hint line.
+export const HEALTH_TIMING_HINTS: Readonly<Record<HealthTimingKey, string>> = Object.freeze({
+ budgetMs:
+ "How long one read may take before the drive counts as not answering. Takes effect on the next read.",
+ passIntervalMs:
+ "How often each drive is checked, from its disk's own request counters. A save re-arms the check at once.",
+ probeTimeoutMs:
+ "How long one check may wait for the drive's root where no disk can be named (and for naming the disk).",
+ clearAfterCleanPasses:
+ "How many clean checks in a row it takes before a drive marked not answering is used again.",
+ inFlightPerLocation:
+ "How many reads may be on one drive at once; the rest wait their turn, and are refused if it stops answering.",
+});
+
+// SETTINGS.md's `storage.health` table.
+export const STORAGE_HEALTH_SETTINGS_FIELD_DOCS: FieldDocs<StorageHealthSettings> = {
+ budgetMs:
+ "How long one read may take before the drive counts as not answering, in ms (default 3000, " +
+ "500–60000). The watchdog's budget (`onDrive`, lib/storageHealth.ts) for one unit of work — a " +
+ "video directory's reads, a page's reads of one video: a unit that has not answered by then is " +
+ "refused, and marks its location not answering unless the disk's request counters show it still " +
+ "completing others (slow, not stalled). A read waiting for a slot is refused when nothing on the " +
+ "drive has returned for this long plus a quarter of it (at most 250 ms). Takes effect on the next " +
+ "read after a save.",
+ passIntervalMs:
+ "How often the health pass reads each location's disk counters, in ms (default 15000, " +
+ "5000–300000). A stall that starts between two passes is seen by the next, or at once by a page's " +
+ "read. A save on /storage re-arms the pass's timer at once; a hand edit, at the next pass. Two " +
+ "counter samples are compared only when at least min(10 s, this − 5 s) apart, a spacing never " +
+ "less than half of this.",
+ probeTimeoutMs:
+ "How long the health pass waits, in ms (default 3000, 500–30000), for the child `stat` of a root " +
+ "where no disk can be named (a timeout counts as not answering) and for the `findmnt` that names a " +
+ "root's disk (a timeout names none that pass). Takes effect on the next pass.",
+ clearAfterCleanPasses:
+ "How many clean answers in a row clear a location marked not answering (default 2, 1–10). Each " +
+ "health pass is one answer, and so is a Refresh on /storage; a miss in between starts the count " +
+ "again. Takes effect on the next answer.",
+ inFlightPerLocation:
+ "How many reads through the watchdog may be on one location's drive at once (default 4, 1–16); the " +
+ "rest wait in the editor's own queue, so a stall mid-walk holds this many of Node's threads, not " +
+ "all of them. Keep it under `UV_THREADPOOL_SIZE` (16 in the editor's start script). Takes effect on " +
+ "the next read.",
+};
diff --git a/common/lib/storageLocations.ts b/common/lib/storageLocations.ts
@@ -1,5 +1,6 @@
import path from "node:path";
import type { FieldDocs } from "./fieldDocs";
+import type { StorageHealthSettings } from "./storageHealthTimings";
// STORAGE LOCATIONS — the named places a channel's media may live.
//
@@ -95,6 +96,7 @@ export type StorageSettings = {
locations: StorageLocation[];
defaultLocationId: string;
savedVideosLocationId?: string;
+ health?: StorageHealthSettings;
};
export const STORAGE_SETTINGS_FIELD_DOCS: FieldDocs<StorageSettings> = {
@@ -112,6 +114,14 @@ export const STORAGE_SETTINGS_FIELD_DOCS: FieldDocs<StorageSettings> = {
"Optional so an older settings.json parses (and an older binary that " +
"drops it leaves a store that still works, because the symlink is what " +
"every reader follows).",
+ health:
+ "THE DRIVE-HEALTH TIMINGS: how long a read may take before a drive " +
+ "counts as not answering, how often the health pass looks, how long " +
+ "its look may take, how many clean looks clear a stall, and how many " +
+ "reads may be on one drive at once. Edited on /storage (Drive health " +
+ "timing). Absent = every default, and only a value that differs from " +
+ "its default is written, so an untuned install follows a default " +
+ "changed later. See `storage.health` below.",
};
// Strip trailing slashes so "/mnt/platter/" and "/mnt/platter" are one root.
@@ -214,9 +224,12 @@ export function migrateMediaRootToLocations(raw: unknown): unknown {
if (!raw || typeof raw !== "object" || Array.isArray(raw)) return raw;
const r = raw as Record<string, unknown>;
if (r.locations !== undefined) return raw;
+ // The drive-health timings are not about the locations: a block that spells
+ // them and no `locations` (a hand edit) keeps them through the migration.
+ const health = r.health !== undefined ? { health: r.health } : {};
const mediaRoot = typeof r.mediaRoot === "string" ? r.mediaRoot.trim() : "";
if (mediaRoot === "" || !path.isAbsolute(mediaRoot)) {
- return { locations: [], defaultLocationId: "" };
+ return { locations: [], defaultLocationId: "", ...health };
}
return {
locations: [
@@ -228,5 +241,6 @@ export function migrateMediaRootToLocations(raw: unknown): unknown {
},
],
defaultLocationId: "default",
+ ...health,
};
}
diff --git a/common/lib/storageVolumes.ts b/common/lib/storageVolumes.ts
@@ -6,8 +6,8 @@ import type { Paths } from "./paths";
import { getFreeBytes } from "./diskSpace";
import type { StorageLocation, StorageVolume } from "./storageLocations";
import {
- HEALTH_PROBE_TIMEOUT_MS,
countersVerdict,
+ healthTimings,
isDriveNotAnswering,
onDrive,
parseBlockStat,
@@ -17,6 +17,7 @@ import {
type HealthDetector,
type LocationHealthState,
} from "./storageHealth";
+import { counterSampleMinimumMs, secondsText } from "./storageHealthTimings";
// STORAGE VOLUME PROBES — is this location's disk here, and if not, where?
//
@@ -247,7 +248,8 @@ export async function probeLocation(
// run in-process, and on a stalled disk each holds an I/O thread until the
// 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.
+ // within the budget (`storage.health.budgetMs`, 3 s by default) marks the
+ // location stalled and this probe answers so.
const stalledProbe: StorageLocationProbe = {
status: "stalled",
identity: { known: false },
@@ -389,7 +391,7 @@ export async function probeLocationHealth(
): Promise<LocationHealthState> {
const root = loc.root.trim();
if (root === "") return "absent";
- const timeoutMs = opts.timeoutMs ?? HEALTH_PROBE_TIMEOUT_MS;
+ const timeoutMs = opts.timeoutMs ?? healthTimings().probeTimeoutMs;
let child: ReturnType<typeof execa>;
try {
child = execa(opts.statBin ?? "stat", ["-L", "-c", "%F", "--", root], {
@@ -440,8 +442,9 @@ export async function probeLocationHealth(
// reading them touches only /sys. So the health pass asks this first:
//
// 1. the root's device: `findmnt -J -T <root> -o SOURCE,UUID` as a child
-// raced against 3 s (a findmnt stuck resolving the root holds nothing of
-// ours), the `[subvolume]` suffix a bind or btrfs mount adds taken off,
+// raced against `probeTimeoutMs` (3 s by default; a findmnt stuck
+// resolving the root holds nothing of ours), the `[subvolume]` suffix a
+// bind or btrfs mount adds taken off,
// `/dev/mapper/<x>` resolved to its `dm-N`, then the basename. A partition
// and a mapper device both have `/sys/class/block/<name>/stat`. Asked only
// when the root has no device yet or its device's /sys entry cannot be read
@@ -450,7 +453,7 @@ export async function probeLocationHealth(
// none this pass; one whose UUID is not the location's recorded one names
// none (the root is then a directory on some other filesystem).
// 2. that device's `stat` line, compared with the previous pass's sample for
-// the location (same device, at least MIN_COUNTER_INTERVAL_MS earlier).
+// the location (same device, at least `minCounterIntervalMs()` earlier).
//
// NO DEVICE (a container, no findmnt, a network or tmpfs mount, no /sys entry)
// FALLS BACK TO THE CHILD `stat`, and the verdict says which detector answered.
@@ -464,8 +467,11 @@ export type DetectorOptions = HealthProbeOptions & {
export const SYS_BLOCK_DIR = "/sys/class/block";
// Two samples closer than this are not compared: a healthy drive can have a
// request in flight at two instants a moment apart without completing one.
-// The pass is 15 s apart; a /storage Refresh just after a pass gives no verdict.
-export const MIN_COUNTER_INTERVAL_MS = 10_000;
+// Derived from the pass interval (`counterSampleMinimumMs`): 10 s at the
+// default 15 s, so a /storage Refresh just after a pass gives no verdict.
+export function minCounterIntervalMs(): number {
+ return counterSampleMinimumMs(healthTimings().passIntervalMs);
+}
type CounterSample = BlockStatSample & { device: string; at: number };
@@ -604,7 +610,7 @@ export async function detectLocationHealth(
opts: DetectorOptions = {},
): Promise<HealthVerdict> {
const now = opts.now ?? Date.now();
- const timeoutMs = opts.timeoutMs ?? HEALTH_PROBE_TIMEOUT_MS;
+ const timeoutMs = opts.timeoutMs ?? healthTimings().probeTimeoutMs;
const sysBlockDir = opts.sysBlockDir ?? SYS_BLOCK_DIR;
const detector = detectorState();
// The device this root was last known on, if its /sys entry still reads;
@@ -623,7 +629,7 @@ export async function detectLocationHealth(
if (sample) {
detector.deviceByRoot.set(loc.root, device);
const prev = detector.samples.get(loc.id);
- if (prev && prev.device === device && now - prev.at < MIN_COUNTER_INTERVAL_MS) {
+ if (prev && prev.device === device && now - prev.at < minCounterIntervalMs()) {
return { answer: null, detector: "counters", device };
}
detector.samples.set(loc.id, { ...sample, device, at: now });
@@ -651,7 +657,7 @@ export async function detectLocationHealth(
answer,
detector: "stat",
...(answer === "stalled"
- ? { cause: `a stat of its root did not answer within ${timeoutMs / 1000} s` }
+ ? { cause: `a stat of its root did not answer within ${secondsText(timeoutMs)}` }
: {}),
};
}