commit 8c576dbf23b74da43352f45ad4e9416364877362
parent 74050956c077b8c02c3a1841de168790cf3f82e9
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 29 Sep 2026 23:36:45 -0400
common: review M3 — the snapshot walk of a relocated channel goes through the watchdog
generateChannelSnapshot runs in the editor's process after every download or
sync of a channel, sixteen video directories wide. For a channel whose media
is on another drive, its data/ listing, the keep-latest keys' metadata reads
and each video directory's unit now go through `onDrive(config.dataDir, …)`:
at most four on that drive at once, and one that does not answer in 3 s throws,
so the scheduler keeps the last good snapshot.json (its existing behaviour on
a failed refresh). A listing refused that way is rethrown, never read as an
empty channel. The reconcile pass before it is sequential and not raced.
keyedVideosNewestFirst (keep-latest) was an unbounded Promise.all over every
video; it is now mapConcurrent 16 wide, and takes the `through` wrapper.
Test: a video directory's read that never answers mid-walk throws
DriveNotAnsweringError, marks the location, writes no snapshot.json, and at
most four video directories reached the drive.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
3 files changed, 66 insertions(+), 11 deletions(-)
diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts
@@ -22,6 +22,7 @@ import {
} from "../lib/availability";
import { resolveCookiePolicy } from "../lib/cookiePolicy";
import { assertChannelMediaReachable } from "../lib/channelMedia";
+import { isDriveNotAnswering, onDrive } from "../lib/storageHealth";
import { CLIPS_DIR_NAME } from "../lib/clipWindow";
import { getSettings } from "../lib/settings";
import {
@@ -737,6 +738,19 @@ export async function generateChannelSnapshot(
const config = await readChannelConfig(paths, slug);
await assertChannelMediaReachable(paths, slug, config);
+ // A CHANNEL ON ANOTHER DRIVE IS WALKED THROUGH THE WATCHDOG
+ // (lib/storageHealth.ts `onDrive`): the data/ listing, the keep-latest keys and
+ // 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
+ // 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.
+ const drive = config?.dataDir?.trim() || undefined;
+ const through = <T>(read: () => Promise<T>): Promise<T> =>
+ drive ? onDrive(drive, read) : read();
+
// Heal any video dir that drifted from the canonical id layout before we read
// data/* (best-effort; never fail snapshot generation on a reconcile error).
try {
@@ -753,7 +767,11 @@ export async function generateChannelSnapshot(
maybeMissingRecord,
roster,
] = await Promise.all([
- readdir(dataDir, { withFileTypes: true }).catch(() => [] as Dirent[]),
+ through(() => readdir(dataDir, { withFileTypes: true })).catch((err) => {
+ // A drive that did not answer is not an empty channel: rethrown.
+ if (isDriveNotAnswering(err)) throw err;
+ return [] as Dirent[];
+ }),
readPlaylistUrls(playlistPath),
readArchive(archivePath),
loadFailedTranscriptions(paths, slug),
@@ -769,6 +787,7 @@ export async function generateChannelSnapshot(
paths,
channelSlug: slug,
keepLatest: config?.keepLatest ?? 0,
+ through,
});
const videoDirNames = dirEntries
.filter((d) => d.isDirectory())
@@ -821,7 +840,8 @@ export async function generateChannelSnapshot(
const limit = pLimit(SNAPSHOT_VIDEO_CONCURRENCY);
const perVideo = await Promise.all(
videoDirNames.map((id) =>
- limit(async () => {
+ // One video directory's reads are one unit through the watchdog.
+ limit(() => through(async () => {
const dir = path.join(dataDir, id);
const files = await readVideoFiles(dir, { checkUntranscribable: true });
// EVERY FILE IN THE DIR, STATTED ONCE, feeding two numbers.
@@ -952,7 +972,7 @@ export async function generateChannelSnapshot(
vttProvenance,
digest,
};
- }),
+ })),
),
);
diff --git a/common/controller/keptVideos.ts b/common/controller/keptVideos.ts
@@ -2,6 +2,10 @@ import path from "node:path";
import { readdir } from "node:fs/promises";
import type { Paths } from "../lib/paths";
import { loadRawMetadataFromDir } from "../lib/transcripts-server";
+import { mapConcurrent } from "../lib/concurrency";
+
+// Metadata reads in flight while keying a channel's videos by upload date.
+const KEY_READ_CONCURRENCY = 16;
// Rolling "keep-latest" window computation. Given a channel's keepLatest config,
// returns the ids of the newest N videos (by upload date). Used by both the
@@ -13,6 +17,11 @@ export type ComputeKeptOptions = {
paths: Paths;
channelSlug: string;
keepLatest: number;
+ // Runs each video's metadata read. The snapshot passes `onDrive` for a
+ // channel on another drive (lib/storageHealth.ts), which caps the reads in
+ // flight on that drive and gives up on one that does not answer. Default:
+ // the read itself.
+ through?: <T>(read: () => Promise<T>) => Promise<T>;
};
// List a channel's data-dir video ids (directories, skipping dotfiles). Mirrors
@@ -58,16 +67,19 @@ async function uploadKey(videoDir: string, id: string): Promise<string> {
async function keyedVideosNewestFirst(
paths: Paths,
channelSlug: string,
+ through: <T>(read: () => Promise<T>) => Promise<T> = (read) => read(),
): Promise<Array<{ id: string; key: string }>> {
- const ids = await listChannelVideoIds(paths, channelSlug);
+ const ids = await through(() => listChannelVideoIds(paths, channelSlug));
if (ids.length === 0) return [];
const dataDir = path.join(paths.channelsDir, channelSlug, "data");
- const keyed = await Promise.all(
- ids.map(async (id) => ({
- id,
- key: await uploadKey(path.join(dataDir, id), id),
- })),
- );
+ // 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.
+ const keyed = await mapConcurrent(ids, KEY_READ_CONCURRENCY, async (id) => ({
+ id,
+ key: await through(() => uploadKey(path.join(dataDir, id), id)),
+ }));
keyed.sort((a, b) =>
a.key === b.key ? b.id.localeCompare(a.id) : b.key.localeCompare(a.key),
);
@@ -82,9 +94,10 @@ export async function computeKeptVideoIds({
paths,
channelSlug,
keepLatest,
+ through,
}: ComputeKeptOptions): Promise<Set<string>> {
if (!Number.isFinite(keepLatest) || keepLatest <= 0) return new Set();
- const keyed = await keyedVideosNewestFirst(paths, channelSlug);
+ const keyed = await keyedVideosNewestFirst(paths, channelSlug, through);
return new Set(keyed.slice(0, Math.floor(keepLatest)).map((k) => k.id));
}
diff --git a/common/controller/storageStall.test.ts b/common/controller/storageStall.test.ts
@@ -520,3 +520,25 @@ test("watchdog: a move onto a root whose stat never answers is refused", async (
);
assert.match(String(problem), /drive not answering/);
});
+
+test("watchdog (M3): the snapshot walk's video unit never answers → the refresh throws and writes no snapshot", async () => {
+ for (const id of ["vid2", "vid3", "vid4", "vid5", "vid6", "vid7"]) {
+ mkdirSync(path.join(TARGET, id), { recursive: true });
+ }
+ const snapshotFile = path.join(paths.channelsDir, SLUG, "snapshot.json");
+ await watchdogCase(["readdir"], async () => {
+ // data/ itself answers (the listing, the reconcile pass); the video
+ // directories do not.
+ hang = (c) => c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid"));
+ await assert.rejects(
+ () => generateChannelSnapshot(paths, SLUG),
+ (err: unknown) => err instanceof Error && err.name === "DriveNotAnsweringError",
+ );
+ });
+ assert.equal(existsSync(snapshotFile), false, "the last snapshot.json stands");
+ // At most four video directories reached the drive.
+ const asked = calls.filter(
+ (c) => c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid")),
+ );
+ assert.ok(asked.length <= 4, `${asked.length} video dirs asked`);
+});