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. // // A drive can be mounted and still not answer. An SMR disk in a USB enclosure // under a long write stalls, the enclosure resets, and every filesystem call // that has to reach the disk blocks for about 30 seconds. Node runs those calls // on libuv's thread pool (four threads by default), so four of them block the // whole editor: no page, no poll and no job log answers until the disk does. // "Not mounted" does not describe that (a `stat` there does not fail, it hangs), // and no in-process call can find it out without paying the hang itself. // // SO SOMETHING THAT CANNOT HANG ASKS, AND THIS MODULE REMEMBERS WHAT IT SAID. // Two detectors write here. The health pass (`controller/storageWatch.ts`, // 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 // layer skips the tail read. // // THE RULES, which are what keep a flaky drive from flapping the UI: // - 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. // - `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. // // 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 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. // // 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 // `globalThis`, the house pattern (`lib/jsonFile-server.ts`, `jobs/registry.ts`), // so the writer and every reader see one map. // // PURE of I/O, and deliberately without execa, so `lib/channelMedia.ts` can ask // it without pulling a subprocess module into everything that imports that. // One probe's answer, and a location's state. export type LocationHealthState = // The root answered in time and is a directory. | "ok" // The root did not answer within the probe's budget. | "stalled" // The root answered in time, and is not a directory (not mounted, or gone). | "absent"; // What decided a location's state. `counters`: the block device's own request // counters (`/sys/class/block//stat`), which never touch the drive; // `stat`: a child `stat` of the root, when no device could be named (a // container, no findmnt, no /sys entry). A stall marked by `onDrive`'s watchdog // keeps the detector the last pass used; its `cause` says what did not answer. export type HealthDetector = "counters" | "stat"; export type LocationHealth = { id: string; label: string; root: string; state: LocationHealthState; detector?: HealthDetector; // When the current state began, ms since epoch. since: number; // When the last probe (or observation) was recorded. checkedAt: number; // Consecutive clean answers since the location was marked stalled. It clears // at `clearAfterCleanPasses`. 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 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"; // A call waiting for a slot. `resolve` answers whether it took the slot (a // waiter whose own wait already timed out does not); `rearm` restarts its // deadline, which a call returning on its key does (see `acquireSlot`). type Waiter = { resolve: () => boolean; reject: (err: Error) => void; rearm: () => void; }; // A call the watchdog gave up on that has not returned yet: when it began, // and the device's counters then (null when the counters detector has named // no device for its location). type OverdueCall = { startedAt: number; before: BlockStatSample | null }; type HealthState = { byId: Map; // `onDrive`'s bookkeeping, by slot key (see `resolveWhere`): calls in // flight, the ones among them the watchdog has already given up on, and // the calls waiting for a slot. inFlight?: Map; overdue?: Map; waiters?: Map; // 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. readCounters?: (device: string) => BlockStatSample | null; }; declare global { // eslint-disable-next-line no-var var __yttStorageHealth__: HealthState | undefined; } type FilledState = HealthState & Required>; function healthState(): FilledState { if (!globalThis.__yttStorageHealth__) { globalThis.__yttStorageHealth__ = { byId: new Map() }; } const s = globalThis.__yttStorageHealth__; // Filled lazily: a dev server's hot reload keeps an object made by an older // copy of this module. s.inFlight ??= new Map(); s.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. The applied timings are // configuration, not health, and are kept. 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; to: LocationHealthState; }; // Record one answer about one location. Returns the transition when the state // changed, null when it did not. export function recordLocationHealth( loc: Pick, answer: LocationHealthState, 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, answer: LocationHealthState, opts: { now?: number; cause?: string; detector?: HealthDetector; device?: string | null; }, ): HealthTransition | null { const now = opts.now ?? Date.now(); const map = healthState().byId; const prev = map.get(loc.id); const label = loc.label || loc.id; // First sighting, or the root moved under the same id: start over. if (!prev || prev.root !== loc.root) { map.set(loc.id, { id: loc.id, label, root: loc.root, state: answer, ...(opts.detector ? { detector: opts.detector } : {}), ...(opts.device ? { device: opts.device } : {}), since: now, checkedAt: now, cleanStreak: 0, ...(answer === "stalled" && opts.cause ? { cause: opts.cause } : {}), }); return { id: loc.id, from: null, to: answer }; } 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; const from = prev.state; prev.state = "stalled"; prev.since = now; if (opts.cause) prev.cause = opts.cause; return { id: loc.id, from, to: "stalled" }; } if (prev.state === "stalled") { prev.cleanStreak += 1; if (prev.cleanStreak < healthTimings().clearAfterCleanPasses) return null; prev.state = answer; prev.since = now; prev.cleanStreak = 0; delete prev.cause; return { id: loc.id, from: "stalled", to: answer }; } if (prev.state === answer) return null; const from = prev.state; prev.state = answer; prev.since = now; return { id: loc.id, from, to: answer }; } // Make sure every configured location has an entry, without recording an // answer: a new location starts `ok`, and a known one whose root moved starts // over. The watchdog below can only mark a location it can find, and it finds // a channel's location among these entries. export function registerLocationHealth( locations: ReadonlyArray>, now: number = Date.now(), ): void { const map = healthState().byId; for (const loc of locations) { const prev = map.get(loc.id); if (prev && prev.root === loc.root) { prev.label = loc.label || loc.id; continue; } recordLocationHealth(loc, "ok", { now }); } } // 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, device?: string | null, ): void { const h = healthState().byId.get(id); 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. export function pruneLocationHealth(liveIds: Iterable): void { const keep = new Set(liveIds); const map = healthState().byId; for (const id of [...map.keys()]) if (!keep.has(id)) map.delete(id); } export function locationHealth(id: string): LocationHealth | undefined { const h = healthState().byId.get(id); return h ? { ...h } : undefined; } // Every location's last answer, by id. A copy: the caller may keep it. export function allLocationHealth(): Record { const out: Record = {}; for (const [id, h] of healthState().byId) out[id] = { ...h }; return out; } // THE GATE. The stalled location `p` is on, or null. `p` is a channel's // `dataDir` (or anything under a location's root); the match is the one // `locationOfDataDir` makes, longest root first, so a channel on a nested // location answers for that location and not its parent. export function stalledLocationForPath(p: string): LocationHealth | null { const map = healthState().byId; if (map.size === 0 || !p) return null; const entries = [...map.values()]; const loc = locationOfDataDir( p, entries.map((h) => ({ id: h.id, label: h.label, root: h.root, autoRepoint: false })), ); if (!loc) return null; const h = map.get(loc.id); return h && h.state === "stalled" ? { ...h } : null; } // The location with this id, when it is stalled AND still at this root. The // gate for a caller that holds a location rather than a channel path // (`probeLocation`, `volumeFreeBytes`). export function stalledLocation( loc: Pick, ): LocationHealth | null { const h = healthState().byId.get(loc.id); return h && h.state === "stalled" && h.root === loc.root ? { ...h } : null; } // "since 11:35" in the server's clock, or "since 2026-09-28 11:35" when the // stall began on another day than `now`. export function sinceText(since: number, now: number = Date.now()): string { const d = new Date(since); const n = new Date(now); const pad = (x: number) => String(x).padStart(2, "0"); const hm = `${pad(d.getHours())}:${pad(d.getMinutes())}`; const sameDay = d.getFullYear() === n.getFullYear() && d.getMonth() === n.getMonth() && d.getDate() === n.getDate(); return sameDay ? `since ${hm}` : `since ${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())} ${hm}`; } // "not answering since 11:35" — the one line /storage, the /channels volume bar // and the channel's Storage panel show for a stalled location. export function notAnsweringText(h: Pick, now?: number): string { return `not answering ${sinceText(h.since, now)}`; } // --------------------------------------------------------------------------- // The watchdog: every call the gate covers, raced against the budget // --------------------------------------------------------------------------- // // 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 // 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. // // 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; // - a wait is refused (without marking anything) only when no call on its key // has returned for the budget plus a small grace: every call that returns // restarts every waiter's deadline, so a deep queue on a drive that is busy // but answering waits as long as it takes; // - 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 unless // the disk has been completing requests since the oldest of those calls // 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 `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 // by hand), the root the path is under (`//data` → ``). 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. // 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 healthTimings().budgetMs; } export class DriveNotAnsweringError extends Error { // The stalled location, when a known one was marked. readonly health: LocationHealth | null; constructor(health: LocationHealth | null, detail?: string) { super( detail ?? (health ? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})` : `${NOT_ANSWERING} (a read did not answer within ${secondsText(driveCallBudget())})`), ); this.name = "DriveNotAnsweringError"; this.health = health; } } export function isDriveNotAnswering(err: unknown): err is DriveNotAnsweringError { return err instanceof Error && err.name === "DriveNotAnsweringError"; } type Where = string | Pick; type Resolved = { key: string; // The location to mark, when marking is right. loc: Pick | null; }; function slotKeyOfLocation(id: string): string { return `loc:${id}`; } // The root a path on no configured location is under: `//media` // (relocatedMediaDir's shape, release 17) or the retired `//data` // gives ``; anything else is its own key. export function rootOfUnknownPath(p: string): string { const clean = p.replace(/\/+$/, ""); const parts = clean.split("/"); const last = parts[parts.length - 1]; return parts.length > 2 && (last === "media" || last === "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); // 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 }; } 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)); } // 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 { const s = healthState(); const cap = healthTimings().inFlightPerLocation; const n = s.inFlight.get(r.key) ?? 0; if (n < cap) { s.inFlight.set(r.key, n + 1); return; } const overdue = s.overdue.get(r.key) ?? []; 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 // (see `onDrive`): four slow reads on a busy disk are not a stall. const oldest = overdue.reduce((a, b) => (b.startedAt < a.startedAt ? b : a)); const now = readDeviceCounters(r.loc); if (oldest.before && now && now.completed > oldest.before.completed) { throw new DriveNotAnsweringError( null, `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: `${overdue.length} reads on it have not answered`, }); health = stalledLocation(r.loc); } throw new DriveNotAnsweringError(health); } const budget = driveCallBudget(); // THE WAIT'S DEADLINE FOLLOWS PROGRESS, NOT THE QUEUE. A waiter is refused // only when NO call on its key has returned for a full budget — "nothing on // this drive answered for 3 s", the condition the watchdog exists for. // Every call that returns on the key (in time or late) re-arms every waiter // behind it, so a deep queue on a drive that is busy but answering (a // 64-wide walk of units that each take a second) is never refused for its // depth alone. Timed from when the call queued, it was: a waiter at depth d // waits about d/4 units, whatever the drive is doing. // // 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((resolve, reject) => { let settled = false; let timer: ReturnType | undefined; const giveUp = () => { 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} (nothing on it answered for ${secondsText(budget)} while a read waited)`, ), ); }; const arm = () => { if (timer) clearTimeout(timer); timer = setTimeout(giveUp, budget + grace); }; const waiter: Waiter = { resolve: () => { if (settled) return false; settled = true; if (timer) clearTimeout(timer); resolve(); return true; }, reject: (err) => { if (settled) return; settled = true; if (timer) clearTimeout(timer); reject(err); }, rearm: () => { if (!settled) arm(); }, }; arm(); 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. } // 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 — // 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(); 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()) s.inFlight.set(key, (s.inFlight.get(key) ?? 0) + 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(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 `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(where: Where, call: () => Promise): Promise { const r = resolveWhere(where); const refused = r.loc ? stalledLocation(r.loc) : null; if (refused) throw new DriveNotAnsweringError(refused); 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) return; released = true; releaseSlot(r.key); }; const startedAt = Date.now(); const before = readDeviceCounters(r.loc); let pending: Promise; try { pending = call(); } catch (err) { release(); throw err; } const TIMED_OUT = Symbol("timed out"); // The budget this call runs against, read as it starts. const budget = driveCallBudget(); let timer: ReturnType | 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((resolve) => { timer = setTimeout(() => resolve(TIMED_OUT), budget); }); let answer: T | typeof TIMED_OUT; try { answer = await Promise.race([pending, timeout]); } catch (err) { release(); throw err; } finally { if (timer) clearTimeout(timer); } if (answer !== TIMED_OUT) { release(); return answer; } // The call is still waiting on the drive. Its slot stays held, and counted // overdue, until it returns; nothing awaits it. const entry: OverdueCall = { startedAt, before }; s.overdue.set(r.key, [...(s.overdue.get(r.key) ?? []), entry]); const done = () => { const left = (s.overdue.get(r.key) ?? []).filter((e) => e !== entry); if (left.length > 0) s.overdue.set(r.key, left); else s.overdue.delete(r.key); 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 ${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 ${secondsText(budget)}`, }); throw new DriveNotAnsweringError(stalledLocation(r.loc)); } refuseWaiters(r.key, null); // The budget this call ran against, not the one in force now: a save during // the call must not rewrite what it was given. throw new DriveNotAnsweringError( null, `${NOT_ANSWERING} (a read did not answer within ${secondsText(budget)})`, ); } // --------------------------------------------------------------------------- // The block device's counters // --------------------------------------------------------------------------- // // `/sys/class/block//stat` is the kernel's own count of a device's // requests (Documentation/block/stat.rst): field 1 reads completed, 5 writes // completed, 9 requests in flight now. Reading it never touches the drive. A // drive that is merely slow, even one grinding through a long write, keeps // 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 nothing completed between (reads, // writes, discards, flushes) // ok ⇔ anything else (nothing in flight at one of them, or completions // moved) // // and the health rules above turn one `stalled` into a stall and two `ok`s in // a row into its end. 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; 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( prev: BlockStatSample, cur: BlockStatSample, ): "ok" | "stalled" { return prev.inFlight > 0 && cur.inFlight > 0 && cur.completed === prev.completed ? "stalled" : "ok"; }