commit 664a68d30365dfc22e8a37ba2bb5ec49f5b82f4f
parent c8f7e9d5e37cbadab7a9a64aca3c092e4b87a387
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 1 Oct 2026 20:59:33 -0400
common: a running regeneration gets one queued successor; REFRESH_REPORT_QUEUE in queueKeys (D0 review L2, N7)
startRefreshReport dedups only against a queued or starting regeneration of the
slug, so a change made while its walk runs queues one successor instead of being
dropped until the channel's next change; fire()'s trailing-edge comment is now
true. onError hands a failed walk's sentence to the caller.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
3 files changed, 82 insertions(+), 23 deletions(-)
diff --git a/common/jobs/snapshotScheduler.test.ts b/common/jobs/snapshotScheduler.test.ts
@@ -22,11 +22,11 @@ for (const slug of ["alpha", "beta"]) {
}
const { getPaths } = await import("../lib/paths");
-const { getRegistry } = await import("./registry");
+const { getRegistry, newJobId } = await import("./registry");
const {
REFRESH_REPORT_ACTIVE,
REFRESH_REPORT_QUEUE,
- isRefreshReportActive,
+ isRefreshReportPending,
startRefreshReport,
} = await import("./snapshotScheduler");
@@ -48,9 +48,8 @@ test("two requests for one channel at once enqueue one regeneration", async () =
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);
+ assert.equal(isRefreshReportPending("alpha"), false);
// Once it has finished, a new request starts a new one.
const again = await startRefreshReport(paths, "alpha");
assert.ok(again.ok);
@@ -72,3 +71,36 @@ test("every regeneration runs on the one serial refresh-report queue", async ()
assert.equal((await a.done).status, "done");
assert.equal((await b.done).status, "done");
});
+
+test("a running regeneration gets one queued successor, and no second", async () => {
+ const paths = getPaths();
+ // A regeneration of beta that is RUNNING and stays so until released: a
+ // record on the refresh-report queue whose start does nothing.
+ const registry = getRegistry();
+ const holderId = newJobId();
+ const holder = {
+ id: holderId,
+ kind: "refresh-report",
+ queueKey: REFRESH_REPORT_QUEUE,
+ channelSlug: "beta",
+ status: "queued" as const,
+ queuedAt: Date.now(),
+ logPath: "/dev/null",
+ };
+ registry.register(holder);
+ registry.enqueue(holder, { start: () => {}, onCancel: () => {} });
+ assert.equal(registry.get(holderId)?.status, "running");
+
+ // A change during the walk: it may already have passed it, so one successor
+ // is queued behind it…
+ const successor = await startRefreshReport(paths, "beta");
+ assert.ok(successor.ok, "queued behind the running one, not dropped");
+ assert.equal(registry.get(successor.jobId)?.status, "queued");
+ assert.equal(isRefreshReportPending("beta"), true);
+ // …and a further change finds it queued: it has not started reading yet.
+ const third = await startRefreshReport(paths, "beta");
+ assert.ok(!third.ok && third.info === true && third.error === REFRESH_REPORT_ACTIVE);
+
+ registry.finalize(holderId, "done");
+ assert.equal((await successor.done).status, "done");
+});
diff --git a/common/jobs/snapshotScheduler.ts b/common/jobs/snapshotScheduler.ts
@@ -11,6 +11,7 @@ import {
import { getRegistry } from "./registry";
import { drainStream } from "./drainStream";
import type { StreamActionResult } from "./streamCommand";
+import { REFRESH_REPORT_QUEUE } from "../lib/queueKeys";
// Global, debounced "channel report" (snapshot) regeneration scheduler.
//
@@ -155,8 +156,8 @@ 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.
+// THE ONE QUEUE EVERY REFRESH-REPORT RUNS ON (lib/queueKeys.ts): 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
@@ -165,36 +166,49 @@ export function requestChannelSnapshot(paths: Paths, slug: string): void {
// 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";
+export { REFRESH_REPORT_QUEUE };
// Why startRefreshReport started nothing for a slug that already has one.
-export const REFRESH_REPORT_ACTIVE = "already queued or running";
+export const REFRESH_REPORT_ACTIVE = "already queued";
-// 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()
+// The QUEUED regeneration of `slug` (not yet started), if there is one.
+export function queuedRefreshReportId(slug: string): string | null {
+ const queued = getRegistry()
.list()
- .some(
+ .find(
(j) =>
j.kind === "refresh-report" &&
j.channelSlug === slug &&
- (j.status === "queued" || j.status === "running"),
+ j.status === "queued",
);
+ return queued?.id ?? null;
+}
+
+// Is a regeneration of `slug` queued or being enqueued — one that has not
+// started reading yet, and so will see any change made before it does?
+//
+// A RUNNING one does not count. It may already have walked past the change
+// that asks for this regeneration, so one successor is queued behind it (the
+// queue is serial, so it starts when the running one ends). At most one: a
+// second request finds the successor queued.
+export function isRefreshReportPending(slug: string): boolean {
+ if (getState().starting.has(slug)) return true;
+ return queuedRefreshReportId(slug) !== null;
}
// 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.
+// one is already queued or being enqueued for it, which is answered `info`
+// with REFRESH_REPORT_ACTIVE (it has not started reading, so it will see the
+// change). Every walk goes through here — the debounced pass, "Update all
+// reports" and a channel's own Refresh report — so none runs outside the
+// queue, and they dedup against each other. `onError` hears the walk's error
+// sentence (an unmounted drive's) when the job fails.
export async function startRefreshReport(
paths: Paths,
slug: string,
+ opts: { onError?: (message: string) => void } = {},
): Promise<StreamActionResult> {
- if (isRefreshReportActive(slug)) {
+ if (isRefreshReportPending(slug)) {
return { ok: false, error: REFRESH_REPORT_ACTIVE, info: true };
}
const state = getState();
@@ -208,7 +222,13 @@ export async function startRefreshReport(
channelSlug: slug,
fn: async (onLog) => {
onLog(`Regenerating report for ${slug}…`);
- const snap = await generateChannelSnapshot(paths, slug);
+ let snap: Awaited<ReturnType<typeof generateChannelSnapshot>>;
+ try {
+ snap = await generateChannelSnapshot(paths, slug);
+ } catch (err) {
+ opts.onError?.((err as Error).message);
+ throw err;
+ }
getState().generation++;
const excluded = excludedDownloadIdSet(snap);
const awaitingTranscription = excluded.size
@@ -238,7 +258,10 @@ async function fire(): Promise<void> {
const state = getState();
// Snapshot and clear before the async work: any action that fires during
// regeneration re-marks the slug dirty and schedules a fresh pass (trailing
- // edge), instead of being swallowed by this in-flight batch.
+ // edge), instead of being swallowed by this in-flight batch. That pass is
+ // not dropped when the slug's regeneration is RUNNING: startRefreshReport
+ // queues one successor behind it (it is skipped only when one is already
+ // queued, which has not started reading and so will see the change).
const batch = [...state.dirty.entries()];
state.dirty.clear();
state.timer = null;
diff --git a/common/lib/queueKeys.ts b/common/lib/queueKeys.ts
@@ -30,6 +30,10 @@ export const DIGEST_REMOTE_QUEUE = "digest:remote";
// is a separate question, answered by backfillLimit() rather than by the queue.
export const BACKFILL_QUEUE = "backfill";
+// Every channel report regeneration (`refresh-report`), one at a time
+// corpus-wide — jobs/snapshotScheduler.ts says why (release 17 slice D0).
+export const REFRESH_REPORT_QUEUE = "refresh-report";
+
// Per-channel queue for channel-local bookkeeping jobs (clean/clear/verify).
export function channelQueueKey(slug: string): string {
return `channel:${slug}`;