commit 184777e4e6548bf98cece3cd8ebe80fc3e1a83d5
parent 6a469e4e7ac5e4d8ed03d0a6a39aae7fa4543f8b
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 30 Sep 2026 00:34:07 -0400
common: re-review M4 — a wait for a slot is refused only when nothing on the drive has answered for the budget, not for the queue's depth
A slot wait was timed from when the call queued (budget + grace), so on a
drive that was busy but answering, a deep queue was refused as "not answering"
for its depth alone: a waiter at depth d waits about d/4 units, so units of
about 1.1 s at 16 wide (the snapshot walk), 0.46 s at 32 (the recency reads)
and 0.2 s at 64 (readChannelStat, the inspect fan-outs) were enough.
Now every call that returns on a key (in time or late) restarts the deadline
of every waiter on that key, and a waiter is refused only when no call on its
key has returned for the budget plus the grace — "nothing on this drive
answered for 3 s". The overdue refusal and the refusal on a transition to
stalled are unchanged, so a queue behind four hung calls is still refused
within the budget. The keep-latest comment says what holds now.
Tests: a healthy 64-wide walk of units at half the budget (the last call waits
~15 units) has no refusals and marks nothing, on a configured location and on
a hand-typed root; a queue behind four hung calls is refused within the budget
(marked on a location, unmarked on a hand-typed root); a call that returns
late hands its slot to the calls waiting behind it. The M1 cases pass as they
were.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
3 files changed, 142 insertions(+), 20 deletions(-)
diff --git a/common/controller/keptVideos.ts b/common/controller/keptVideos.ts
@@ -73,9 +73,11 @@ async function keyedVideosNewestFirst(
if (ids.length === 0) return [];
const dataDir = path.join(paths.channelsDir, channelSlug, "data");
// BOUNDED, like every other corpus-shaped fan-out (lib/concurrency.ts): one
- // metadata read per video, and the largest channel has eleven thousand. It
- // also keeps a `through` queue short enough that a call waiting for a slot on
- // a healthy drive is never kept past the watchdog's budget.
+ // metadata read per video, and the largest channel has eleven thousand. With
+ // a `through` of `onDrive`, at most four of these are on the drive at once
+ // and the rest wait in its queue, which refuses a waiting read only when
+ // nothing on the drive has returned for the watchdog's budget — never for
+ // the queue's depth alone.
const keyed = await mapConcurrent(ids, KEY_READ_CONCURRENCY, async (id) => ({
id,
key: await through(() => uploadKey(path.join(dataDir, id), id)),
diff --git a/common/lib/storageHealth.test.ts b/common/lib/storageHealth.test.ts
@@ -365,8 +365,9 @@ test("L6: a call that times out while its disk is still completing requests is s
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/);
+ // Nothing on the drive returns: the waiter's wait runs out — refused, still
+ // without marking.
+ assert.match(await waiter, /nothing on it answered for 0.06 s while a read waited/);
assert.equal(locationHealth("usb")?.state, "ok");
// With the counters standing still, the same timeout marks it.
setCounterReader(() => ({ completed: 500, inFlight: 2 }));
@@ -397,3 +398,94 @@ test("registering locations creates entries without an answer, and a moved root
assert.equal(locationHealth("usb")?.state, "ok");
assert.equal(locationHealth("usb")?.root, "/mnt/new");
});
+
+// ── M4: the wait's deadline follows progress ───────────────────────────────
+
+const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));
+
+test("M4: a healthy 64-wide walk of units at half the budget — no refusals, nothing marked", async () => {
+ registerLocationHealth([USB]);
+ setDriveCallBudget(100);
+ let answered = 0;
+ const started = Date.now();
+ const outcomes = await Promise.all(
+ Array.from({ length: 64 }, () =>
+ onDrive("/mnt/usb/media/ch/data", async () => {
+ await sleep(50);
+ answered += 1;
+ return "ok";
+ }).catch((err: Error) => err.message),
+ ),
+ );
+ // Sixteen rounds of 50 ms: the last call waited about 750 ms, far past the
+ // 100 ms budget, and was not refused — the drive kept answering.
+ assert.ok(Date.now() - started >= 700, `took ${Date.now() - started} ms`);
+ assert.deepEqual(outcomes, Array(64).fill("ok"));
+ assert.equal(answered, 64);
+ assert.equal(locationHealth("usb")?.state, "ok");
+ assert.equal(driveCallsInFlight("usb"), 0);
+});
+
+test("M4: the same deep queue on a hand-typed root — no refusals either", async () => {
+ setDriveCallBudget(100);
+ const outcomes = await Promise.all(
+ Array.from({ length: 32 }, () =>
+ onDrive("/hand/typed/ch/data", async () => {
+ await sleep(50);
+ return "ok";
+ }).catch((err: Error) => err.message),
+ ),
+ );
+ assert.deepEqual(outcomes, Array(32).fill("ok"));
+});
+
+test("M4: a queue behind four hung calls is still refused within the budget", async () => {
+ setDriveCallBudget(100);
+ // A hand-typed root: nothing marks it, so the refusal is the timeouts'.
+ const hung = Array.from({ length: 4 }, () =>
+ onDrive("/hand/typed/ch/data", never).catch(() => "gave up"),
+ );
+ const started = Date.now();
+ const waiters = Array.from({ length: 6 }, () =>
+ onDrive("/hand/typed/ch/data", async () => "ran").catch((err: Error) => err.name),
+ );
+ assert.deepEqual(await Promise.all(waiters), Array(6).fill("DriveNotAnsweringError"));
+ assert.ok(Date.now() - started < 400, `refused after ${Date.now() - started} ms`);
+ await Promise.all(hung);
+ // And on a configured location, the timeouts mark it and refuse the queue.
+ registerLocationHealth([USB]);
+ const hungHere = Array.from({ length: 4 }, () => onDrive(USB, never).catch(() => "gave up"));
+ const t0 = Date.now();
+ const waitersHere = Array.from({ length: 6 }, () =>
+ onDrive(USB, async () => "ran").catch((err: Error) => err.name),
+ );
+ assert.deepEqual(await Promise.all(waitersHere), Array(6).fill("DriveNotAnsweringError"));
+ assert.ok(Date.now() - t0 < 400);
+ assert.equal(locationHealth("usb")?.state, "stalled");
+ await Promise.all(hungHere);
+});
+
+test("M4: a call that returns late frees its slot for the calls waiting behind it", async () => {
+ registerLocationHealth([USB]);
+ noteLocationDetector("usb", "counters", "sdz1");
+ // The counters keep completing, so the late calls are slow, not stalled.
+ let completed = 0;
+ setCounterReader(() => ({ completed: (completed += 5), inFlight: 3 }));
+ setDriveCallBudget(80);
+ // Four calls that take 150 ms: refused as slow at 80 ms, returning at 150.
+ const slow = Array.from({ length: 4 }, () =>
+ onDrive(USB, () => sleep(150)).then(
+ () => "answered",
+ (err: Error) => err.message,
+ ),
+ );
+ // Four more queue at 70 ms, deadline 70 + 80 + 20: the returns at 150 ms
+ // come first and hand them the slots.
+ await sleep(70);
+ const behind = Array.from({ length: 4 }, () =>
+ onDrive(USB, async () => "ran").catch((err: Error) => err.message),
+ );
+ for (const o of await Promise.all(slow)) assert.match(o, /^drive slow/);
+ assert.deepEqual(await Promise.all(behind), Array(4).fill("ran"));
+ assert.equal(locationHealth("usb")?.state, "ok");
+});
diff --git a/common/lib/storageHealth.ts b/common/lib/storageHealth.ts
@@ -102,8 +102,13 @@ export const HEALTH_CLEAN_TO_CLEAR = 2;
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).
-type Waiter = { resolve: () => boolean; reject: (err: Error) => void };
+// 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;
+};
type HealthState = {
byId: Map<string, LocationHealth>;
@@ -377,9 +382,10 @@ export function notAnsweringText(h: Pick<LocationHealth, "since">, now?: number)
// `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 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);
+// - 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: none of
// those calls has returned, whatever the last pass said.
@@ -503,15 +509,25 @@ async function acquireSlot(r: Resolved): Promise<void> {
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.
+ // 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<void>((resolve, reject) => {
let settled = false;
- const timer = setTimeout(() => {
+ let timer: ReturnType<typeof setTimeout> | undefined;
+ const giveUp = () => {
if (settled) return;
settled = true;
const q = s.waiters.get(r.key);
@@ -522,25 +538,33 @@ async function acquireSlot(r: Resolved): Promise<void> {
reject(
new DriveNotAnsweringError(
r.loc ? stalledLocation(r.loc) : null,
- `${NOT_ANSWERING} (a read waited ${budget / 1000} s for the drive and was not started)`,
+ `${NOT_ANSWERING} (nothing on it answered for ${budget / 1000} s while a read waited)`,
),
);
- }, budget + grace);
+ };
+ const arm = () => {
+ if (timer) clearTimeout(timer);
+ timer = setTimeout(giveUp, budget + grace);
+ };
const waiter: Waiter = {
resolve: () => {
if (settled) return false;
settled = true;
- clearTimeout(timer);
+ if (timer) clearTimeout(timer);
resolve();
return true;
},
reject: (err) => {
if (settled) return;
settled = true;
- clearTimeout(timer);
+ 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);
@@ -548,9 +572,13 @@ async function acquireSlot(r: Resolved): Promise<void> {
// 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.
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 next = q.shift() as Waiter;
if (next.resolve()) return;