commit 7116122aec85ffbe5b13c2150c279eba8e0ecddf
parent 9babe0ee3ff820042c1b64655f8e298ea1a60877
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 1 Oct 2026 16:51:38 -0400
common: a report regeneration yields to the event loop, and runs one at a time
generateChannelSnapshot walks the video dirs in chunks of two waves (32 at 16
wide) with one setImmediate between chunks, so requests get a turn with no
snapshot work queued. Every refresh-report now runs on the one serial
"refresh-report" queue instead of the empty key, and startRefreshReport is the
one entry point (the debounced pass and Update all reports): a slug already
queued, running or being enqueued is not enqueued again.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
4 files changed, 218 insertions(+), 103 deletions(-)
diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts
@@ -674,6 +674,40 @@ export { SNAPSHOT_FILENAME, snapshotPath, readChannelSnapshot };
const SNAPSHOT_VIDEO_CONCURRENCY = 16;
+// THE WALK YIELDS TO THE EVENT LOOP every SNAPSHOT_YIELD_EVERY videos.
+//
+// It runs in the editor's own process, on the main thread, and a video's unit
+// parses its cues.json and (in the reconcile pass and for a non-YouTube
+// archive) its metadata.info.json — ~0.6 MB each on a long VOD, a few ms of
+// JSON.parse apiece. Sixteen units in flight keep every turn of the loop busy
+// with that work; on 2026-10-01 two regenerations of 2,000–3,000-video
+// channels ran for over an hour while `/` and `/jobs` did not answer. So the
+// ids are walked in chunks, each fanned out under the same limit, and between
+// chunks the walk waits one `setImmediate`: the loop gets a turn with no
+// snapshot work queued in it, and every request whose I/O completed meanwhile
+// runs before the next chunk starts.
+//
+// Two full waves of the limit (32 at 16 wide): a chunk is ~0.1–0.2 s of
+// parsing on such a channel, and a multiple of the width keeps every wave full
+// — 25 left the second wave nine wide, a fifth of the walk's throughput.
+// Measured in the slice D0 record (plans/release-17.md).
+const SNAPSHOT_YIELD_EVERY = 2 * SNAPSHOT_VIDEO_CONCURRENCY;
+
+// `fn` over `items` in chunks of `size`, in order, one setImmediate between
+// chunks. The concurrency inside a chunk is whatever `fn` imposes.
+async function mapInYieldingChunks<T, R>(
+ items: readonly T[],
+ size: number,
+ fn: (item: T) => Promise<R>,
+): Promise<R[]> {
+ const out: R[] = [];
+ for (let i = 0; i < items.length; i += size) {
+ if (i > 0) await new Promise<void>((resolve) => setImmediate(resolve));
+ out.push(...(await Promise.all(items.slice(i, i + size).map(fn))));
+ }
+ return out;
+}
+
// ONE LEVEL of a video dir's `clips/` — the files in it, and nothing deeper.
// Deliberately not recursive: `clipWindow-server.ts` writes `<from>-<to>.<ext>`
// and `<from>-<to>.json` flat into it and nothing else does, so a recursion
@@ -839,8 +873,10 @@ export async function generateChannelSnapshot(
}
const limit = pLimit(SNAPSHOT_VIDEO_CONCURRENCY);
- const perVideo = await Promise.all(
- videoDirNames.map((id) =>
+ const perVideo = await mapInYieldingChunks(
+ videoDirNames,
+ SNAPSHOT_YIELD_EVERY,
+ (id) =>
// One video directory's reads are one unit through the watchdog.
limit(() => through(async () => {
const dir = path.join(dataDir, id);
@@ -974,7 +1010,6 @@ export async function generateChannelSnapshot(
digest,
};
})),
- ),
);
const filesById = new Map<string, VideoFiles>();
diff --git a/common/jobs/snapshotScheduler.test.ts b/common/jobs/snapshotScheduler.test.ts
@@ -0,0 +1,74 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import path from "node:path";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/snapshotScheduler.test.ts
+//
+// ONE REGENERATION AT A TIME, ONCE PER CHANNEL (release 17 slice D0). Every
+// refresh-report used to run on the empty queue key — all at once — and two
+// passes a moment apart could both enqueue the same channel. Temp corpus only.
+
+const root = await mkdtemp(path.join(tmpdir(), "snapshot-scheduler-"));
+process.env.TRANSCRIPTS_DIR = path.join(root, "transcripts");
+process.env.SETTINGS_FILE = path.join(root, "settings.json");
+await writeFile(process.env.SETTINGS_FILE, "{}");
+for (const slug of ["alpha", "beta"]) {
+ const dir = path.join(process.env.TRANSCRIPTS_DIR, "channels", slug);
+ await mkdir(path.join(dir, "data"), { recursive: true });
+ await writeFile(path.join(dir, "config.json"), JSON.stringify({ handling: "youtube" }));
+}
+
+const { getPaths } = await import("../lib/paths");
+const { getRegistry } = await import("./registry");
+const {
+ REFRESH_REPORT_ACTIVE,
+ REFRESH_REPORT_QUEUE,
+ isRefreshReportActive,
+ startRefreshReport,
+} = await import("./snapshotScheduler");
+
+test.after(async () => {
+ await rm(root, { recursive: true, force: true });
+});
+
+test("two requests for one channel at once enqueue one regeneration", async () => {
+ const paths = getPaths();
+ const [a, b] = await Promise.all([
+ startRefreshReport(paths, "alpha"),
+ startRefreshReport(paths, "alpha"),
+ ]);
+ const started = [a, b].filter((r) => r.ok);
+ const refused = [a, b].filter((r) => !r.ok);
+ assert.equal(started.length, 1, "exactly one starts");
+ assert.equal(refused.length, 1);
+ const r = refused[0];
+ assert.ok(!r.ok && r.info === true && r.error === REFRESH_REPORT_ACTIVE);
+ const s = started[0];
+ assert.ok(s.ok);
+ assert.equal(isRefreshReportActive("alpha"), true);
+ assert.equal((await s.done).status, "done");
+ assert.equal(isRefreshReportActive("alpha"), false);
+ // Once it has finished, a new request starts a new one.
+ const again = await startRefreshReport(paths, "alpha");
+ assert.ok(again.ok);
+ await again.done;
+});
+
+test("every regeneration runs on the one serial refresh-report queue", async () => {
+ const paths = getPaths();
+ const a = await startRefreshReport(paths, "alpha");
+ const b = await startRefreshReport(paths, "beta");
+ assert.ok(a.ok && b.ok);
+ for (const id of [a.jobId, b.jobId]) {
+ assert.equal(getRegistry().get(id)?.queueKey, REFRESH_REPORT_QUEUE);
+ }
+ // Serial: the second is behind the first unless the first already finished.
+ const first = getRegistry().get(a.jobId);
+ const second = getRegistry().get(b.jobId);
+ if (first?.status === "running") assert.equal(second?.status, "queued");
+ assert.equal((await a.done).status, "done");
+ assert.equal((await b.done).status, "done");
+});
diff --git a/common/jobs/snapshotScheduler.ts b/common/jobs/snapshotScheduler.ts
@@ -10,6 +10,7 @@ import {
} from "../controller/channelSnapshot";
import { getRegistry } from "./registry";
import { drainStream } from "./drainStream";
+import type { StreamActionResult } from "./streamCommand";
// Global, debounced "channel report" (snapshot) regeneration scheduler.
//
@@ -93,6 +94,13 @@ type SchedulerState = {
// counts (/channels, the dashboard) would sit stale until the next unrelated
// change. One integer closes that gap at zero cost.
generation: number;
+ // Slugs whose refresh-report is being enqueued right now. Between the
+ // registry check and the record's registration runManagedFunction awaits
+ // (the media guard, the jobs dir), so two passes a moment apart both found
+ // nothing queued and both enqueued — two regenerations of the same
+ // 3,000-video channel, seen live on 2026-10-01. Held until the record is in
+ // the registry, which answers from then on.
+ starting: Set<string>;
};
declare global {
@@ -107,6 +115,7 @@ function getState(): SchedulerState {
timer: null,
firstDirtyAt: null,
generation: 0,
+ starting: new Set(),
};
}
return globalThis.__yttSnapshotScheduler__;
@@ -146,6 +155,85 @@ export function requestChannelSnapshot(paths: Paths, slug: string): void {
state.timer.unref?.();
}
+// THE ONE QUEUE EVERY REFRESH-REPORT RUNS ON: one regeneration at a time,
+// corpus-wide.
+//
+// It was the empty queueKey — "a local filesystem scan, so it runs in
+// parallel" — and every dirty channel's walk ran at once in the editor's own
+// process. On 2026-10-01 two of them (2,000 and 3,260 videos) ran side by side
+// for over an hour and `/` and `/jobs` did not answer. A non-empty key is how
+// the registry serializes (registry.ts `enqueue`: concurrency 1 per key), so
+// nothing new was needed; a report that waits its turn is a report a minute
+// late, and the pages keep answering meanwhile.
+export const REFRESH_REPORT_QUEUE = "refresh-report";
+
+// Why startRefreshReport started nothing for a slug that already has one.
+export const REFRESH_REPORT_ACTIVE = "already queued or running";
+
+// Is a regeneration of `slug` already queued, running, or being enqueued?
+export function isRefreshReportActive(slug: string): boolean {
+ if (getState().starting.has(slug)) return true;
+ return getRegistry()
+ .list()
+ .some(
+ (j) =>
+ j.kind === "refresh-report" &&
+ j.channelSlug === slug &&
+ (j.status === "queued" || j.status === "running"),
+ );
+}
+
+// Enqueue one channel's report regeneration on REFRESH_REPORT_QUEUE — unless
+// one is already queued or running for it, which is answered `info` with
+// REFRESH_REPORT_ACTIVE — the rule both callers had before this queue: a
+// queued one has not started reading, so it will see the change; a running one
+// may already have passed it, and the channel's next change marks it again.
+// The debounced pass and "Update all reports" both come through here, so they
+// dedup against each other too.
+export async function startRefreshReport(
+ paths: Paths,
+ slug: string,
+): Promise<StreamActionResult> {
+ if (isRefreshReportActive(slug)) {
+ return { ok: false, error: REFRESH_REPORT_ACTIVE, info: true };
+ }
+ const state = getState();
+ state.starting.add(slug);
+ try {
+ const { runManagedFunction } = await import("./streamCommand");
+ return await runManagedFunction({
+ kind: "refresh-report",
+ queueKey: REFRESH_REPORT_QUEUE,
+ paths,
+ channelSlug: slug,
+ fn: async (onLog) => {
+ onLog(`Regenerating report for ${slug}…`);
+ const snap = await generateChannelSnapshot(paths, slug);
+ getState().generation++;
+ const excluded = excludedDownloadIdSet(snap);
+ const awaitingTranscription = excluded.size
+ ? snap.buckets.downloadedNoTranscript.filter(
+ (id) => !excluded.has(id),
+ ).length
+ : snap.buckets.downloadedNoTranscript.length;
+ onLog(
+ `Done. ${snap.totals.videos} videos · ` +
+ `${snap.undownloadedIds.length} undownloaded · ` +
+ `${awaitingTranscription} awaiting transcription.`,
+ );
+ // No revalidatePath here — the caller revalidates once, after its
+ // streams drain.
+ },
+ });
+ } finally {
+ // The record is registered by now (or nothing was made): the registry
+ // answers for this slug from here on. A reset (resetSnapshotScheduler)
+ // may have swapped the state meanwhile; deleting from the old set is
+ // harmless.
+ state.starting.delete(slug);
+ }
+}
+
async function fire(): Promise<void> {
const state = getState();
// Snapshot and clear before the async work: any action that fires during
@@ -158,61 +246,18 @@ async function fire(): Promise<void> {
if (batch.length === 0) return;
try {
- const { runManagedFunction } = await import("./streamCommand");
-
- // Skip channels that already have a refresh-report job queued/running —
- // a fresh snapshot is already on its way (mirrors
- // refreshAllChannelSnapshotsAction's dedup).
- const active = new Set(
- getRegistry()
- .list()
- .filter(
- (j) =>
- j.kind === "refresh-report" &&
- (j.status === "queued" || j.status === "running") &&
- j.channelSlug,
- )
- .map((j) => j.channelSlug as string),
- );
-
const regenerated: string[] = [];
const streams: ReadableStream<string>[] = [];
for (const [slug, paths] of batch) {
- if (active.has(slug)) continue;
- const result = await runManagedFunction({
- kind: "refresh-report",
- // Empty queueKey: bypass queue serialization. Snapshot regen is a local
- // filesystem scan, so it runs in parallel rather than waiting behind
- // sync/download work (see registry.ts enqueue).
- queueKey: "",
- paths,
- channelSlug: slug,
- fn: async (onLog) => {
- onLog(`Regenerating report for ${slug}…`);
- const snap = await generateChannelSnapshot(paths, slug);
- getState().generation++;
- const excluded = excludedDownloadIdSet(snap);
- const awaitingTranscription = excluded.size
- ? snap.buckets.downloadedNoTranscript.filter(
- (id) => !excluded.has(id),
- ).length
- : snap.buckets.downloadedNoTranscript.length;
- onLog(
- `Done. ${snap.totals.videos} videos · ` +
- `${snap.undownloadedIds.length} undownloaded · ` +
- `${awaitingTranscription} awaiting transcription.`,
- );
- // No revalidatePath here — it's done once after all drains below.
- },
- });
+ const result = await startRefreshReport(paths, slug);
if (!result.ok) continue;
regenerated.push(slug);
streams.push(result.stream);
}
// Wait for every snapshot to finish writing before revalidating so the
- // re-rendered pages read fresh counts. queueKey "" runs them in parallel,
- // so this is ~the slowest snapshot, not the sum.
+ // re-rendered pages read fresh counts. They run one at a time on
+ // REFRESH_REPORT_QUEUE, so this is the sum of them.
await Promise.all(streams.map(drainStream));
if (regenerated.length > 0) {
diff --git a/editor/app/channels/actions.ts b/editor/app/channels/actions.ts
@@ -23,13 +23,12 @@ import {
import { renameChannel } from "yt-dlp-transcript-common/controller/renameChannel";
import { inspectChannelMedia } from "yt-dlp-transcript-common/lib/channelMedia";
import { channelMediaBusyReason } from "./lib/mediaBusy";
+import { generateChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot";
import {
- excludedDownloadIdSet,
- generateChannelSnapshot,
-} from "yt-dlp-transcript-common/controller/channelSnapshot";
-import { requestChannelSnapshot } from "yt-dlp-transcript-common/jobs/snapshotScheduler";
-import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
-import { runManagedFunction } from "yt-dlp-transcript-common/jobs/streamCommand";
+ REFRESH_REPORT_ACTIVE,
+ requestChannelSnapshot,
+ startRefreshReport,
+} from "yt-dlp-transcript-common/jobs/snapshotScheduler";
import { drainStream } from "yt-dlp-transcript-common/jobs/drainStream";
import {
siteChannelIndex,
@@ -397,57 +396,20 @@ export type RefreshAllResult = {
export async function refreshAllChannelSnapshotsAction(): Promise<RefreshAllResult> {
const paths = getPaths();
const channels = await listChannelConfigs(paths);
- const active = new Set(
- getRegistry()
- .list()
- .filter(
- (j) =>
- j.kind === "refresh-report" &&
- (j.status === "queued" || j.status === "running") &&
- j.channelSlug,
- )
- .map((j) => j.channelSlug as string),
- );
const queued: string[] = [];
const jobIds: string[] = [];
const skipped: { slug: string; reason: string }[] = [];
const streams: ReadableStream<string>[] = [];
for (const c of channels) {
- if (active.has(c.slug)) {
- skipped.push({ slug: c.slug, reason: "already running" });
- continue;
- }
- const result = await runManagedFunction({
- kind: "refresh-report",
- // Empty queueKey: bypass queue serialization. Snapshot regen is a
- // local filesystem scan that never touches the platform, so there's
- // no reason for it to wait behind sync/download work. See
- // registry.ts:69-72 for the documented escape hatch.
- queueKey: "",
- paths,
- channelSlug: c.slug,
- fn: async (onLog) => {
- onLog(`Regenerating report for ${c.slug}…`);
- const snap = await generateChannelSnapshot(paths, c.slug);
- const excluded = excludedDownloadIdSet(snap);
- const awaitingTranscription = excluded.size
- ? snap.buckets.downloadedNoTranscript.filter(
- (id) => !excluded.has(id),
- ).length
- : snap.buckets.downloadedNoTranscript.length;
- onLog(
- `Done. ${snap.totals.videos} videos · ` +
- `${snap.undownloadedIds.length} undownloaded · ` +
- `${awaitingTranscription} awaiting transcription.`,
- );
- // Deliberately no revalidatePath here — calling it from a
- // background fn races with the in-flight re-render that the action's
- // own revalidatePath triggers. The action's single revalidate at the
- // end picks up every fresh snapshot.
- },
- });
+ // The scheduler's one entry point: the refresh-report queue (one
+ // regeneration at a time — they used to run all at once, in parallel) and
+ // its per-slug dedup, shared with the debounced regeneration.
+ const result = await startRefreshReport(paths, c.slug);
if (!result.ok) {
- skipped.push({ slug: c.slug, reason: result.error });
+ skipped.push({
+ slug: c.slug,
+ reason: result.error === REFRESH_REPORT_ACTIVE ? "already running" : result.error,
+ });
continue;
}
queued.push(c.slug);
@@ -455,9 +417,8 @@ export async function refreshAllChannelSnapshotsAction(): Promise<RefreshAllResu
streams.push(result.stream);
}
// Wait for all snapshots to finish writing before revalidating so the
- // pages that read the snapshots read fresh counts. With queueKey === "" the
- // jobs all run in parallel, so this waits roughly the time of the
- // slowest snapshot, not the sum.
+ // pages that read the snapshots read fresh counts. They run one at a time,
+ // so this is the sum of them.
await Promise.all(streams.map(drainStream));
revalidatePath("/channels");
revalidatePath("/operations/[id]", "page");