commit 30ff28db2d184133fafa743bbeb35f0bb11af533
parent d9205b09248b933219751c519a04f5633f32a320
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 29 Sep 2026 23:36:45 -0400
common: review M1, M2 (keys), L2, L5, L6, L7, L8, L10 — onDrive's queue has a deadline and a stall refuses it; the counters count every completion and ask findmnt only when they must
onDrive (lib/storageHealth.ts):
- M1: calls past their budget are counted per slot key (`overdue`). When every
slot is held by one, a new call is refused at once and the location marked
stalled again (none of those calls has returned, whatever the last pass
said). Otherwise a wait for a slot is raced against the budget (plus a
quarter of it, at most 250 ms, so the calls it waits behind time out and mark
first) and refused without marking. Every transition to `stalled`, the
pass's and a Refresh's included, refuses the waiting calls.
- M2: a path on no configured location is capped by the root it is under
(`<root>/<slug>/data` → `<root>`), so a hand-typed root holds at most four
threads too; nothing can mark it.
- L7: a probe of another root under a location's id has its own slots
(keyed by that root) and does not refuse the location's waiters.
- L6: on a timeout, when the counters detector has named the location's
device, its counters are read (synchronously, from /sys, never the drive)
and compared with a reading taken when the call began: requests completed
meanwhile means slow, not stalled — the call is refused, nothing is marked.
storageVolumes.ts registers the reader (`setCounterReader`); the device is
kept on the health entry.
- L8: the header says the budget covers a whole unit of work; L3: the header
no longer calls the child stat the detector.
The counters (lib/storageVolumes.ts, storageHealth.ts):
- L2: completed counts discards (field 12) and flushes (16) when present,
since in_flight counts them.
- L10: a root's device is read from the last pass; findmnt only when there is
none or its /sys entry stops reading.
- L5: the detector's samples are on globalThis (`__yttHealthDetector__`), so
the pass and /storage's Refresh compare against one previous sample.
inspectChannelMedia carries the watchdog's own words for a refusal that
marked nothing.
Tests: M1's case (four hung calls, two clean answers, a fifth refused at once
and the location marked again); waiters refused on the pass's transition; a
hand-typed root's four slots shared by its channels; a candidate root's own
slots; L6 (slow: refused unmarked, a waiter's wait runs out unmarked; still
counters: marked); the parser with discards and flushes; a known device read
without findmnt, findmnt again when /sys stops reading.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
5 files changed, 507 insertions(+), 152 deletions(-)
diff --git a/common/lib/channelMedia.ts b/common/lib/channelMedia.ts
@@ -207,8 +207,11 @@ export function stalledMediaLocation(
dataDir: string,
configured: string,
// Null when the drive is on no location the health state knows (a root typed
- // by hand) and the watchdog found it not answering.
+ // by hand) and the watchdog found it not answering, or when the watchdog
+ // found it slow rather than stalled.
health: LocationHealth | null,
+ // The watchdog's own words, for a refusal that marked no location.
+ detail?: string,
): ChannelMediaLocation {
return {
dataDir,
@@ -217,7 +220,8 @@ export function stalledMediaLocation(
status: "stalled",
detail: health
? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`
- : `${NOT_ANSWERING} (a read did not answer within ${DRIVE_CALL_BUDGET_MS / 1000} s)`,
+ : (detail ??
+ `${NOT_ANSWERING} (a read did not answer within ${DRIVE_CALL_BUDGET_MS / 1000} s)`),
};
}
@@ -450,7 +454,7 @@ async function inspectOnDisk(
}
} catch (err) {
if (isDriveNotAnswering(err)) {
- return stalledMediaLocation(dataDir, configured, err.health);
+ return stalledMediaLocation(dataDir, configured, err.health, err.message);
}
return {
dataDir,
diff --git a/common/lib/storageHealth.test.ts b/common/lib/storageHealth.test.ts
@@ -8,8 +8,10 @@ import {
allLocationHealth,
driveCallsInFlight,
isDriveNotAnswering,
+ noteLocationDetector,
onDrive,
registerLocationHealth,
+ setCounterReader,
setDriveCallBudget,
locationHealth,
notAnsweringText,
@@ -31,6 +33,7 @@ import {
beforeEach(() => {
resetStorageHealth();
setDriveCallBudget();
+ setCounterReader(undefined);
});
const USB = { id: "usb", label: "USB drive", root: "/mnt/usb/media" };
@@ -245,24 +248,134 @@ test("a waiting call runs when a slot frees, if the location is still answering"
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,
+test("a path on no known location is raced and capped by its root, and names no location", async () => {
+ setDriveCallBudget(80);
+ let made = 0;
+ // Two channels under one hand-typed root share its four slots.
+ const calls = [
+ ...Array.from({ length: 3 }, () => "/hand/typed/chan-a/data"),
+ ...Array.from({ length: 3 }, () => "/hand/typed/chan-b/data"),
+ ].map((p) =>
+ onDrive(p, () => {
+ made += 1;
+ return never();
+ }).then(
+ () => "answered",
+ (err: DriveNotAnsweringError) => (err.health === null ? "refused" : "marked"),
+ ),
);
+ await new Promise((r) => setImmediate(r));
+ assert.equal(made, 4, "one root, four slots");
+ assert.deepEqual(await Promise.all(calls), Array(6).fill("refused"));
+ assert.equal(made, 4);
assert.deepEqual(allLocationHealth(), {});
+ // Every slot is held by a call given up on: the next is refused at once.
+ const started = Date.now();
+ await assert.rejects(() => onDrive("/hand/typed/chan-c/data", async () => ++made));
+ assert.ok(Date.now() - started < 50);
+ assert.equal(made, 4);
});
-test("a probe of another root under a location's id is raced but does not rewrite that location", async () => {
+test("a probe of another root under a location's id has its own slots and does not rewrite that location", async () => {
registerLocationHealth([USB]);
setDriveCallBudget(50);
+ // Four candidate probes that never answer: they hold the CANDIDATE root's
+ // four slots, not the location's.
+ for (let i = 0; i < 4; i++) {
+ await assert.rejects(
+ () => onDrive({ ...USB, root: "/mnt/candidate" }, never),
+ DriveNotAnsweringError,
+ );
+ }
+ assert.equal(locationHealth("usb")?.root, "/mnt/usb/media");
+ assert.equal(locationHealth("usb")?.state, "ok");
+ assert.equal(driveCallsInFlight("usb"), 0);
+ // The real location's calls run as if nothing happened.
+ assert.deepEqual(
+ await Promise.all([1, 2, 3, 4, 5].map((n) => onDrive(USB, async () => n))),
+ [1, 2, 3, 4, 5],
+ );
+});
+
+test("M1: every slot held by a call given up on — a new call is refused within the budget, after the pass cleared the location", async () => {
+ registerLocationHealth([USB]);
+ setDriveCallBudget(80);
+ const hung = Array.from({ length: 4 }, () => onDrive(USB, never).catch(() => "gave up"));
+ assert.deepEqual(await Promise.all(hung), Array(4).fill("gave up"));
+ assert.equal(locationHealth("usb")?.state, "stalled");
+ // Two clean answers from the pass (in a reset loop's good moment).
+ recordLocationHealth(USB, "ok");
+ recordLocationHealth(USB, "ok");
+ assert.equal(locationHealth("usb")?.state, "ok");
+ let made = false;
+ const started = Date.now();
await assert.rejects(
- () => onDrive({ ...USB, root: "/mnt/candidate" }, never),
- DriveNotAnsweringError,
+ () => onDrive(USB, async () => {
+ made = true;
+ }),
+ (err: unknown) => err instanceof DriveNotAnsweringError && err.health?.id === "usb",
);
- assert.equal(locationHealth("usb")?.root, "/mnt/usb/media");
+ assert.ok(Date.now() - started < 80, "refused at once, not after a wait");
+ assert.equal(made, false);
+ // And the location is marked again: none of the four has returned.
+ assert.equal(locationHealth("usb")?.state, "stalled");
+ assert.match(String(locationHealth("usb")?.cause), /4 reads on it have not answered/);
+});
+
+test("M1: every transition to stalled refuses the waiting calls at once, whoever decided it", async () => {
+ registerLocationHealth([USB]);
+ setDriveCallBudget(5_000);
+ const gates = Array.from({ length: 4 }, () => deferred<number>());
+ const inFlight = gates.map((g) => onDrive(USB, () => g.promise));
+ const waiting = onDrive(USB, async () => 5).then(
+ () => "ran",
+ (err: Error) => err.name,
+ );
+ await new Promise((r) => setImmediate(r));
+ const started = Date.now();
+ // The pass (not the watchdog) finds the drive stalled.
+ recordLocationHealth(USB, "stalled");
+ assert.equal(await waiting, "DriveNotAnsweringError");
+ assert.ok(Date.now() - started < 1_000);
+ for (const g of gates) g.resolve(0);
+ await Promise.all(inFlight);
+});
+
+test("L6: a call that times out while its disk is still completing requests is slow — refused, the location not marked", async () => {
+ registerLocationHealth([USB]);
+ noteLocationDetector("usb", "counters", "sdz1");
+ let completed = 100;
+ setCounterReader((device) => {
+ assert.equal(device, "sdz1");
+ completed += 7;
+ return { completed, inFlight: 2 };
+ });
+ setDriveCallBudget(60);
+ // Four slow calls and one waiting behind them.
+ const slow = Array.from({ length: 4 }, () =>
+ onDrive(USB, never).then(
+ () => "answered",
+ (err: Error) => err.message,
+ ),
+ );
+ const waiter = onDrive(USB, async () => 1).then(
+ () => "ran",
+ (err: Error) => err.message,
+ );
+ const outcomes = await Promise.all(slow);
+ for (const o of outcomes) assert.match(o, /^drive slow/);
+ assert.equal(locationHealth("usb")?.state, "ok", "slow is not stalled");
+ // The waiter's own wait runs out: refused, still without marking.
+ assert.match(await waiter, /waited 0.06 s for the drive and was not started/);
assert.equal(locationHealth("usb")?.state, "ok");
+ // With the counters standing still, the same timeout marks it.
+ setCounterReader(() => ({ completed: 500, inFlight: 2 }));
+ resetStorageHealth();
+ registerLocationHealth([USB]);
+ noteLocationDetector("usb", "counters", "sdz1");
+ setCounterReader(() => ({ completed: 500, inFlight: 2 }));
+ await assert.rejects(() => onDrive(USB, never), DriveNotAnsweringError);
+ assert.equal(locationHealth("usb")?.state, "stalled");
});
test("the default budget is 3 s", async () => {
diff --git a/common/lib/storageHealth.ts b/common/lib/storageHealth.ts
@@ -11,23 +11,26 @@ import { locationOfDataDir, type StorageLocation } from "./storageLocations";
// "Not mounted" does not describe that (a `stat` there does not fail, it hangs),
// and no in-process call can find it out without paying the hang itself.
//
-// SO A CHILD PROCESS ASKS, AND THIS MODULE REMEMBERS WHAT IT SAID.
-// `probeLocationHealth` (`storageVolumes.ts`) runs `stat` on the location's root
-// as a subprocess raced against a 3 s timer, and the storage watch
-// (`controller/storageWatch.ts`) asks every 15 s. A child stuck in the kernel
-// blocks nothing of ours: the timer answers `stalled` and the child is left to
-// finish on its own. Everything that would touch the drive in-process then asks
-// `stalledLocationForPath` first and, on a stalled location, answers without the
-// call: `inspectChannelMedia` reports `stalled`, the free-space column reads
-// "—", the probe reads "Not answering", the recency layer skips the tail read.
+// 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
+// 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
+// layer skips the tail read.
//
// THE RULES, which are what keep a flaky drive from flapping the UI:
-// - ONE failed probe marks the location `stalled` at once. A drive that did
-// not answer in 3 s will not answer the next page either, and every page
-// that asks costs a thread.
-// - TWO consecutive clean probes clear it. A clean probe is one that answered
-// in time, whether the root was there (`ok`) or not (`absent`: an unmounted
-// drive answers ENOENT at once, and that is a different problem, which
+// - 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.
// - A ROOT CHANGE (a re-point) starts the location over: the old root's stall
@@ -35,7 +38,9 @@ import { locationOfDataDir, type StorageLocation } from "./storageLocations";
//
// IN MEMORY ONLY. A stall is a fact about this minute, not about the corpus;
// persisting it would outlive the reset loop that caused it. A restarted
-// process starts with no stall, and the first probe (at arm time) re-learns it.
+// process starts with no stall, and the first pass (at arm time, on an idle
+// boot too) registers the locations; the watchdog re-learns a stall the moment
+// a page reaches the drive.
//
// 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
@@ -77,6 +82,9 @@ export type LocationHealth = {
cleanStreak: number;
// What did not answer, for a stalled location: the probe's own words.
cause?: string;
+ // The block device the counters detector reads for this location, when it
+ // could name one. The watchdog reads its counters too (see `onDrive`).
+ device?: string;
};
// The probe's budget. A `stat` of a directory on a healthy disk answers in
@@ -93,16 +101,23 @@ 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 Waiter = { resolve: () => void; reject: (err: Error) => void };
+// A call waiting for a slot. `resolve` answers whether it took the slot (a
+// waiter whose own wait already timed out does not).
+type Waiter = { resolve: () => boolean; 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.
+ // `onDrive`'s bookkeeping, by slot key (see `resolveWhere`): calls in
+ // flight, how many of those the watchdog has already given up on, and the
+ // calls waiting for a slot.
inFlight?: Map<string, number>;
+ overdue?: Map<string, number>;
waiters?: Map<string, Waiter[]>;
// Test seam: the watchdog's budget.
budgetMs?: number;
+ // Reads a block device's counters, synchronously and without touching the
+ // drive (set by `storageVolumes.ts`, which owns /sys). Absent: no check.
+ readCounters?: (device: string) => BlockStatSample | null;
};
declare global {
@@ -110,7 +125,10 @@ declare global {
var __yttStorageHealth__: HealthState | undefined;
}
-function healthState(): Required<Omit<HealthState, "budgetMs">> & HealthState {
+type FilledState = HealthState &
+ Required<Pick<HealthState, "inFlight" | "overdue" | "waiters">>;
+
+function healthState(): FilledState {
if (!globalThis.__yttStorageHealth__) {
globalThis.__yttStorageHealth__ = { byId: new Map() };
}
@@ -118,8 +136,9 @@ function healthState(): Required<Omit<HealthState, "budgetMs">> & HealthState {
// Filled lazily: a dev server's hot reload keeps an object made by an older
// copy of this module.
s.inFlight ??= new Map();
+ s.overdue ??= new Map();
s.waiters ??= new Map();
- return s as Required<Omit<HealthState, "budgetMs">> & HealthState;
+ return s as FilledState;
}
// Test seam, and the escape hatch for a process that wants to forget. Calls
@@ -128,10 +147,20 @@ export function resetStorageHealth(): void {
const s = healthState();
s.byId.clear();
s.inFlight.clear();
+ s.overdue.clear();
for (const q of s.waiters.values()) for (const w of q) w.resolve();
s.waiters.clear();
}
+// `storageVolumes.ts` registers its /sys reader here (see `onDrive`, L6 in
+// the record: a call that times out while its disk is still completing
+// requests is slow, not stalled). A test passes its own, or undefined.
+export function setCounterReader(
+ read: ((device: string) => BlockStatSample | null) | undefined,
+): void {
+ healthState().readCounters = read;
+}
+
export type HealthTransition = {
id: string;
from: LocationHealthState | null;
@@ -143,7 +172,33 @@ export type HealthTransition = {
export function recordLocationHealth(
loc: Pick<StorageLocation, "id" | "label" | "root">,
answer: LocationHealthState,
- opts: { now?: number; cause?: string; detector?: HealthDetector } = {},
+ opts: {
+ now?: number;
+ cause?: string;
+ detector?: HealthDetector;
+ // The counters' device; `null` forgets it (the stat detector answered).
+ device?: string | null;
+ } = {},
+): HealthTransition | null {
+ const t = recordAnswer(loc, answer, opts);
+ // EVERY TRANSITION TO STALLED refuses the calls waiting for a slot on the
+ // location, whoever decided it (the pass, the watchdog, a Refresh): they
+ // would otherwise wait for a slot held by a call the drive is not answering.
+ if (t?.to === "stalled") {
+ refuseWaiters(slotKeyOfLocation(loc.id), stalledLocation(loc));
+ }
+ return t;
+}
+
+function recordAnswer(
+ loc: Pick<StorageLocation, "id" | "label" | "root">,
+ answer: LocationHealthState,
+ opts: {
+ now?: number;
+ cause?: string;
+ detector?: HealthDetector;
+ device?: string | null;
+ },
): HealthTransition | null {
const now = opts.now ?? Date.now();
const map = healthState().byId;
@@ -157,6 +212,7 @@ export function recordLocationHealth(
root: loc.root,
state: answer,
...(opts.detector ? { detector: opts.detector } : {}),
+ ...(opts.device ? { device: opts.device } : {}),
since: now,
checkedAt: now,
cleanStreak: 0,
@@ -167,6 +223,8 @@ export function recordLocationHealth(
prev.label = label;
prev.checkedAt = now;
if (opts.detector) prev.detector = opts.detector;
+ if (opts.device === null) delete prev.device;
+ else if (opts.device) prev.device = opts.device;
if (answer === "stalled") {
prev.cleanStreak = 0;
if (prev.state === "stalled") return null;
@@ -213,9 +271,16 @@ export function registerLocationHealth(
// Which detector decided a location's last answer, when it gave none (the
// counters' first sample has nothing to compare with).
-export function noteLocationDetector(id: string, detector: HealthDetector): void {
+export function noteLocationDetector(
+ id: string,
+ detector: HealthDetector,
+ device?: string | null,
+): void {
const h = healthState().byId.get(id);
- if (h) h.detector = detector;
+ if (!h) return;
+ h.detector = detector;
+ if (device === null) delete h.device;
+ else if (device) h.device = device;
}
// Drop every location that is no longer configured.
@@ -292,25 +357,44 @@ export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number)
//
// 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.
+// 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 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
+// drive that is still answering can be marked by one long unit. Unless its
+// disk is visibly still completing requests: on a timeout, when the counters
+// detector has named the location's device, its counters are read (from /sys,
+// never the drive) and compared with a reading taken when the call began; if
+// requests completed meanwhile the drive is slow, not stalled, and only this
+// call is refused.
//
-// AND NO LOCATION HOLDS MORE THAN FOUR OF THE POOL'S THREADS. A walk of
+// AT MOST DRIVE_CALLS_IN_FLIGHT 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 — 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.
+// 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;
+// - a wait is raced against the same budget (plus a small grace, below), and
+// one that runs out is refused without marking anything (the calls holding
+// the slots may be answering slowly);
+// - when every slot is held by a call the watchdog already gave up on, a new
+// call is refused at once and the location marked stalled again: 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.
//
-// 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.
+// 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
+// by hand), the root the path is under (`<root>/<slug>/data` → `<root>`). Only a
+// configured location can be marked.
+//
+// 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;
@@ -325,13 +409,14 @@ function driveCallBudget(): number {
}
export class DriveNotAnsweringError extends Error {
- // The location the drive is on, when it is a known one.
+ // The stalled location, when a known one was marked.
readonly health: LocationHealth | null;
- constructor(health: LocationHealth | null) {
+ constructor(health: LocationHealth | null, detail?: string) {
super(
- health
- ? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})`
- : `${NOT_ANSWERING} (a read did not answer within ${driveCallBudget() / 1000} s)`,
+ detail ??
+ (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;
@@ -344,95 +429,177 @@ export function isDriveNotAnswering(err: unknown): err is 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 {
+type Resolved = {
+ key: string;
+ // The location to mark, when marking is right.
+ loc: Pick<StorageLocation, "id" | "label" | "root"> | null;
+};
+
+function slotKeyOfLocation(id: string): string {
+ return `loc:${id}`;
+}
+
+// The root a path on no configured location is under: `<root>/<slug>/data`
+// (relocatedDataDir's shape) gives `<root>`; anything else is its own key.
+function rootOfUnknownPath(p: string): string {
+ const clean = p.replace(/\/+$/, "");
+ const parts = clean.split("/");
+ return parts.length > 2 && parts[parts.length - 1] === "data"
+ ? parts.slice(0, -2).join("/") || "/"
+ : clean;
+}
+
+function resolveWhere(where: Where): Resolved {
const map = healthState().byId;
if (typeof where !== "string") {
const entry = map.get(where.id);
- return { loc: where, markable: !entry || entry.root === where.root };
+ // A probe of a candidate root under a location's id is its own key, and
+ // marks nothing: racing it is right, rewriting that location is not.
+ if (entry && entry.root !== where.root) {
+ return { key: `root:${where.root}`, loc: null };
+ }
+ return { key: slotKeyOfLocation(where.id), loc: where };
}
- 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;
+ const found =
+ map.size > 0 && where
+ ? locationOfDataDir(
+ where,
+ [...map.values()].map((h) => ({
+ id: h.id,
+ label: h.label,
+ root: h.root,
+ autoRepoint: false,
+ })),
+ )
+ : null;
+ if (found) return { key: slotKeyOfLocation(found.id), loc: found };
+ return { key: `root:${rootOfUnknownPath(where)}`, loc: null };
+}
+
+function refuseWaiters(key: string, health: LocationHealth | null): void {
+ const s = healthState();
+ const q = s.waiters.get(key) ?? [];
+ s.waiters.delete(key);
+ for (const w of q) w.reject(new DriveNotAnsweringError(health));
}
-async function acquireSlot(id: string): Promise<void> {
+// Take a slot on `key`, or wait for one (raced against the budget), or be
+// refused at once when every slot is held by an overdue call.
+async function acquireSlot(r: Resolved): Promise<void> {
const s = healthState();
- const n = s.inFlight.get(id) ?? 0;
+ const n = s.inFlight.get(r.key) ?? 0;
if (n < DRIVE_CALLS_IN_FLIGHT) {
- s.inFlight.set(id, n + 1);
+ s.inFlight.set(r.key, n + 1);
return;
}
+ if ((s.overdue.get(r.key) ?? 0) >= DRIVE_CALLS_IN_FLIGHT) {
+ let health: LocationHealth | null = null;
+ if (r.loc) {
+ recordLocationHealth(r.loc, "stalled", {
+ cause: `${DRIVE_CALLS_IN_FLIGHT} reads on it have not answered`,
+ });
+ health = stalledLocation(r.loc);
+ }
+ throw new DriveNotAnsweringError(health);
+ }
+ const budget = driveCallBudget();
+ // THE WAIT'S DEADLINE IS THE BUDGET PLUS A SMALL GRACE. A waiter's timer is
+ // created when it queues — before the race timers of calls that took their
+ // slots in the same tick — so with equal deadlines the waiters would give up
+ // a moment before the calls they wait behind, and be refused without the
+ // mark those calls are about to make. The grace lets them time out first.
+ const grace = Math.min(250, Math.round(budget / 4));
await new Promise<void>((resolve, reject) => {
- const q = s.waiters.get(id) ?? [];
- q.push({ resolve, reject });
- s.waiters.set(id, q);
+ let settled = false;
+ const timer = setTimeout(() => {
+ if (settled) return;
+ settled = true;
+ const q = s.waiters.get(r.key);
+ if (q) {
+ const i = q.indexOf(waiter);
+ if (i >= 0) q.splice(i, 1);
+ }
+ reject(
+ new DriveNotAnsweringError(
+ r.loc ? stalledLocation(r.loc) : null,
+ `${NOT_ANSWERING} (a read waited ${budget / 1000} s for the drive and was not started)`,
+ ),
+ );
+ }, budget + grace);
+ const waiter: Waiter = {
+ resolve: () => {
+ if (settled) return false;
+ settled = true;
+ clearTimeout(timer);
+ resolve();
+ return true;
+ },
+ reject: (err) => {
+ if (settled) return;
+ settled = true;
+ clearTimeout(timer);
+ reject(err);
+ },
+ };
+ const q = s.waiters.get(r.key) ?? [];
+ q.push(waiter);
+ s.waiters.set(r.key, q);
});
// Resolved by a release that handed its slot over: the count is unchanged.
}
-function releaseSlot(id: string): void {
+function releaseSlot(key: string): void {
const s = healthState();
- const next = s.waiters.get(id)?.shift();
- if (next) {
- next.resolve();
- return;
+ const q = s.waiters.get(key);
+ while (q && q.length > 0) {
+ const next = q.shift() as Waiter;
+ if (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));
+ 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
// for a reader that wants to say so).
export function driveCallsInFlight(id: string): number {
- return healthState().inFlight.get(id) ?? 0;
+ return healthState().inFlight.get(slotKeyOfLocation(id)) ?? 0;
+}
+
+function readDeviceCounters(loc: Resolved["loc"]): BlockStatSample | null {
+ if (!loc) return null;
+ const s = healthState();
+ const device = s.byId.get(loc.id)?.device;
+ if (!device || !s.readCounters) return null;
+ try {
+ return s.readCounters(device);
+ } catch {
+ return 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 it, and raced against the budget. Throws DriveNotAnsweringError
-// for a refusal or a timeout; any other error is the call's own.
+// flight on its key, 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();
+ const r = resolveWhere(where);
+ const refused = r.loc ? stalledLocation(r.loc) : null;
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);
- }
+ await acquireSlot(r);
+ // Stalled while this call waited: refused, and the slot passed on.
+ const late = r.loc ? stalledLocation(r.loc) : null;
+ if (late) {
+ releaseSlot(r.key);
+ throw new DriveNotAnsweringError(late);
}
+ const s = healthState();
let released = false;
const release = () => {
- if (released || id === undefined) return;
+ if (released) return;
released = true;
- releaseSlot(id);
+ releaseSlot(r.key);
};
+ const before = readDeviceCounters(r.loc);
let pending: Promise<T>;
try {
pending = call();
@@ -460,20 +627,31 @@ export async function onDrive<T>(where: Where, call: () => Promise<T>): Promise<
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);
+ // The call is still waiting on the drive. Its slot stays held, and counted
+ // overdue, until it returns; nothing awaits it.
+ s.overdue.set(r.key, (s.overdue.get(r.key) ?? 0) + 1);
+ const done = () => {
+ s.overdue.set(r.key, Math.max(0, (s.overdue.get(r.key) ?? 1) - 1));
+ release();
+ };
+ pending.then(done, done);
+ const after = readDeviceCounters(r.loc);
+ if (before && after && after.completed > before.completed) {
+ // 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, ` +
+ `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`,
+ });
+ throw new DriveNotAnsweringError(stalledLocation(r.loc));
}
- throw new DriveNotAnsweringError(health);
+ refuseWaiters(r.key, null);
+ throw new DriveNotAnsweringError(null);
}
// ---------------------------------------------------------------------------
@@ -487,7 +665,8 @@ export async function onDrive<T>(where: Where, call: () => Promise<T>): Promise<
// completing requests; one in a reset loop has requests in flight and
// completes none. So, between two samples a pass apart:
//
-// stalled ⇔ in flight at both samples AND no read or write completed between
+// stalled ⇔ in flight at both samples AND nothing completed between (reads,
+// writes, discards, flushes)
// ok ⇔ anything else (nothing in flight at one of them, or completions
// moved)
//
@@ -496,10 +675,15 @@ export async function onDrive<T>(where: Where, call: () => Promise<T>): Promise<
export type BlockStatSample = { completed: number; inFlight: number };
+// Completed: reads (field 1) + writes (5), and, on kernels that count them,
+// discards (12) and flushes (16) — in_flight counts those too, so a long flush
+// alone (an SMR drive emptying its media cache) must not read as nothing
+// completing.
export function parseBlockStat(line: string): BlockStatSample | null {
const f = line.trim().split(/\s+/).map(Number);
if (f.length < 9 || f.slice(0, 9).some((n) => !Number.isFinite(n))) return null;
- return { completed: f[0] + f[4], inFlight: f[8] };
+ const opt = (i: number) => (Number.isFinite(f[i]) ? f[i] : 0);
+ return { completed: f[0] + f[4] + opt(11) + opt(15), inFlight: f[8] };
}
export function countersVerdict(
diff --git a/common/lib/storageHealthCounters.test.ts b/common/lib/storageHealthCounters.test.ts
@@ -25,11 +25,16 @@ beforeEach(() => resetHealthDetector());
const LINE =
"368126412 135952867 36463191157 80435430 29670281 46781260 4001255365 167469727 0 44151394 253218297 877589 0 861936568 4100450 578669 1212688";
-test("the stat line: reads + writes completed, and requests in flight", () => {
+test("the stat line: reads + writes (+ discards + flushes) completed, and requests in flight", () => {
+ // Fields 1, 5, 12 and 16 completed; field 9 in flight.
assert.deepEqual(parseBlockStat(LINE), {
- completed: 368126412 + 29670281,
+ completed: 368126412 + 29670281 + 877589 + 578669,
inFlight: 0,
});
+ // A long flush alone moves completions: not a stall.
+ const before = parseBlockStat("10 0 0 0 5 0 0 0 1 0 0 0 0 0 0 7 0");
+ const after = parseBlockStat("10 0 0 0 5 0 0 0 1 0 0 0 0 0 0 8 0");
+ assert.equal(countersVerdict(before!, after!), "ok");
// An 11-field line from an older kernel parses the same way.
assert.deepEqual(parseBlockStat("10 0 0 0 5 0 0 0 3 0 0\n"), { completed: 15, inFlight: 3 });
assert.equal(parseBlockStat(""), null);
@@ -194,11 +199,12 @@ test("no device → the child stat, and the verdict says so", async () => {
});
});
-test("a findmnt that does not answer: the last device named for that root is read", async () => {
+test("a device already named is read without findmnt; findmnt again only when its /sys entry stops reading", async () => {
await withHarness(async (h) => {
await h.control({ source: "/dev/fakedisk1" });
await h.counters("fakedisk1", 100, 2);
await detectLocationHealth(loc(h.root), h.bins, { sysBlockDir: h.sys, now: 0 });
+ // findmnt now hangs: it is not asked, the known device is read.
await h.control({ sleepMs: 5_000, source: "/dev/fakedisk1" });
const started = Date.now();
const v = await detectLocationHealth(loc(h.root), h.bins, {
@@ -206,8 +212,20 @@ test("a findmnt that does not answer: the last device named for that root is rea
now: 15_000,
timeoutMs: 300,
});
- assert.ok(Date.now() - started < 2_000, "not waited for");
+ assert.ok(Date.now() - started < 250, "no findmnt was run");
assert.equal(v.detector, "counters");
assert.equal(v.answer, "stalled");
+ // The device is gone from /sys (replugged under another name): findmnt is
+ // asked again — here it does not answer, so no device, and the child stat.
+ await rm(path.join(h.sys, "fakedisk1"), { recursive: true });
+ const again = await detectLocationHealth(loc(h.root), h.bins, {
+ sysBlockDir: h.sys,
+ now: 30_000,
+ timeoutMs: 300,
+ });
+ assert.equal(again.detector, "stat");
+ // The detector's samples are on globalThis (the pass and a Refresh run in
+ // different module copies and compare against one previous sample).
+ assert.ok(globalThis.__yttHealthDetector__);
});
});
diff --git a/common/lib/storageVolumes.ts b/common/lib/storageVolumes.ts
@@ -1,4 +1,5 @@
import path from "node:path";
+import { readFileSync } from "node:fs";
import { lstat, readFile, realpath, stat } from "node:fs/promises";
import { execa } from "execa";
import type { Paths } from "./paths";
@@ -10,6 +11,7 @@ import {
isDriveNotAnswering,
onDrive,
parseBlockStat,
+ setCounterReader,
stalledLocation,
type BlockStatSample,
type HealthDetector,
@@ -437,14 +439,16 @@ export async function probeLocationHealth(
// `countersVerdict`) are the device's own account of what it has done, and
// reading them touches only /sys. So the health pass asks this first:
//
-// 1. the root's device, once per pass: `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, `/dev/mapper/<x>` resolved to its `dm-N`, then the basename.
-// A partition and a mapper device both have `/sys/class/block/<name>/stat`.
-// A findmnt that timed out reuses the device the last pass found for that
-// root; one whose UUID is not the location's recorded one names none (the
-// root is then a directory on some other filesystem, not the drive).
+// 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,
+// `/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
+// (a drive replugged under another name): a findmnt per pass would leave
+// one child stuck per pass during a long stall. One that times out names
+// 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).
//
@@ -470,16 +474,40 @@ type DetectorState = {
deviceByRoot: Map<string, string>;
};
-// Module state, not globalThis: only the health pass (one module copy, armed
-// from instrumentation) and /storage's Refresh read it, and a Refresh reaching
-// another copy just starts that copy's comparison over.
-const detector: DetectorState = { samples: new Map(), deviceByRoot: new Map() };
+declare global {
+ // eslint-disable-next-line no-var
+ var __yttHealthDetector__: DetectorState | undefined;
+}
+
+// ON globalThis, the house pattern: the health pass runs in instrumentation's
+// module copy and /storage's Refresh in a page's, and they must compare
+// against the same previous sample.
+function detectorState(): DetectorState {
+ globalThis.__yttHealthDetector__ ??= {
+ samples: new Map(),
+ deviceByRoot: new Map(),
+ };
+ return globalThis.__yttHealthDetector__;
+}
export function resetHealthDetector(): void {
- detector.samples.clear();
- detector.deviceByRoot.clear();
+ detectorState().samples.clear();
+ detectorState().deviceByRoot.clear();
}
+// THE WATCHDOG'S READING of a device's counters (lib/storageHealth.ts,
+// `onDrive`): synchronous, so it does not wait behind the thread pool it is
+// judging, and cheap — a /sys read is answered by the kernel from memory and
+// never reaches the drive.
+function readCountersNow(device: string): BlockStatSample | null {
+ try {
+ return parseBlockStat(readFileSync(path.join(SYS_BLOCK_DIR, device, "stat"), "utf8"));
+ } catch {
+ return null;
+ }
+}
+setCounterReader(readCountersNow);
+
type Raced = { answered: false } | ({ answered: true } & Run);
// `run`, but raced against a timer that nobody waits past: a child stuck in
@@ -577,13 +605,21 @@ export async function detectLocationHealth(
): Promise<HealthVerdict> {
const now = opts.now ?? Date.now();
const timeoutMs = opts.timeoutMs ?? HEALTH_PROBE_TIMEOUT_MS;
- const mapped = loc.root.trim()
- ? await blockDeviceOfRoot(loc, bins, timeoutMs)
- : { device: null, timedOut: false };
- const device =
- mapped.device ?? (mapped.timedOut ? (detector.deviceByRoot.get(loc.root) ?? null) : null);
+ const sysBlockDir = opts.sysBlockDir ?? SYS_BLOCK_DIR;
+ const detector = detectorState();
+ // The device this root was last known on, if its /sys entry still reads;
+ // findmnt only when there is none (see step 1 above).
+ let device: string | null = detector.deviceByRoot.get(loc.root) ?? null;
+ let sample = device ? await readBlockStat(device, sysBlockDir) : null;
+ if (!sample) {
+ detector.deviceByRoot.delete(loc.root);
+ const mapped = loc.root.trim()
+ ? await blockDeviceOfRoot(loc, bins, timeoutMs)
+ : { device: null, timedOut: false };
+ device = mapped.device;
+ sample = device ? await readBlockStat(device, sysBlockDir) : null;
+ }
if (device) {
- const sample = await readBlockStat(device, opts.sysBlockDir ?? SYS_BLOCK_DIR);
if (sample) {
detector.deviceByRoot.set(loc.root, device);
const prev = detector.samples.get(loc.id);