commit 78d9ef6751c8b918fff6b94f6326a880d6cb276f
parent 612855aad2d93e7baa7c0841813dac96b322c3a4
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 1 Oct 2026 23:15:39 -0400
common: requestRefreshReport, a bounded waitForRefreshReport and its notice (D0 re-review R1, R2, R3)
requestRefreshReport starts a channel's regeneration or finds the one queued
for it (waiting briefly for a caller mid-enqueue), cancelling the stream it
does not read. waitForRefreshReport waits up to REFRESH_REPORT_WAIT_MS (15 s)
and says done, failed (with the job log's [error] sentence, whoever started
it) or still waiting with the regenerations ahead of it; refreshReportWaitNotice
is the sentence. +2 tests.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
2 files changed, 189 insertions(+), 9 deletions(-)
diff --git a/common/jobs/snapshotScheduler.test.ts b/common/jobs/snapshotScheduler.test.ts
@@ -15,11 +15,21 @@ 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"]) {
+for (const slug of ["alpha", "beta", "gamma"]) {
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" }));
}
+// A channel relocated to a drive that is not there: its walk refuses (guard 3)
+// with a sentence naming the unreachable media.
+{
+ const dir = path.join(process.env.TRANSCRIPTS_DIR, "channels", "unmounted");
+ await mkdir(dir, { recursive: true });
+ await writeFile(
+ path.join(dir, "config.json"),
+ JSON.stringify({ handling: "youtube", dataDir: path.join(root, "no-such-drive", "unmounted", "data") }),
+ );
+}
const { getPaths } = await import("../lib/paths");
const { getRegistry, newJobId } = await import("./registry");
@@ -27,9 +37,31 @@ const {
REFRESH_REPORT_ACTIVE,
REFRESH_REPORT_QUEUE,
isRefreshReportPending,
+ refreshReportWaitNotice,
+ requestRefreshReport,
startRefreshReport,
+ waitForRefreshReport,
} = await import("./snapshotScheduler");
+// A refresh-report that is RUNNING until released: a record on the queue whose
+// start does nothing.
+function holdQueue(slug: string): () => void {
+ const registry = getRegistry();
+ const id = newJobId();
+ const holder = {
+ id,
+ kind: "refresh-report",
+ queueKey: REFRESH_REPORT_QUEUE,
+ channelSlug: slug,
+ status: "queued" as const,
+ queuedAt: Date.now(),
+ logPath: "/dev/null",
+ };
+ registry.register(holder);
+ registry.enqueue(holder, { start: () => {}, onCancel: () => {} });
+ return () => registry.finalize(id, "done");
+}
+
test.after(async () => {
await rm(root, { recursive: true, force: true });
});
@@ -104,3 +136,50 @@ test("a running regeneration gets one queued successor, and no second", async ()
registry.finalize(holderId, "done");
assert.equal((await successor.done).status, "done");
});
+
+// A PERSON'S REFRESH WAITS A BOUNDED TIME (re-review R1, R2).
+
+test("a refresh behind the queue answers with where it is, once the wait runs out", async () => {
+ const paths = getPaths();
+ const release = holdQueue("beta");
+ try {
+ const asked = await requestRefreshReport(paths, "gamma");
+ assert.ok(asked.ok && asked.started);
+ const wait = await waitForRefreshReport(paths, asked.jobId, { timeoutMs: 300, pollMs: 50 });
+ assert.equal(wait.state, "waiting");
+ assert.ok(wait.state === "waiting" && wait.ahead === 1, "one regeneration ahead of it");
+ assert.equal(
+ wait.state === "waiting" && refreshReportWaitNotice(wait),
+ `Queued behind 1 report regeneration — the report updates when it finishes (job ${asked.jobId}).`,
+ );
+ // A second person's click finds the same job, queued.
+ const again = await requestRefreshReport(paths, "gamma");
+ assert.deepEqual(again, { ok: true, jobId: asked.jobId, started: false });
+ release();
+ const finished = await waitForRefreshReport(paths, asked.jobId, { pollMs: 20 });
+ assert.deepEqual(finished, { state: "done", jobId: asked.jobId });
+ } finally {
+ release();
+ }
+});
+
+test("a failed walk is a failure with its sentence, for the click that found it queued too", async () => {
+ const paths = getPaths();
+ const release = holdQueue("beta");
+ let queuedId = "";
+ try {
+ const first = await requestRefreshReport(paths, "unmounted");
+ assert.ok(first.ok && first.started);
+ queuedId = first.jobId;
+ // Another click while it waits: it did not start this job, and still
+ // hears how it ended.
+ const second = await requestRefreshReport(paths, "unmounted");
+ assert.deepEqual(second, { ok: true, jobId: queuedId, started: false });
+ release();
+ const wait = await waitForRefreshReport(paths, queuedId, { pollMs: 20 });
+ assert.equal(wait.state, "failed");
+ assert.ok(wait.state === "failed" && /unmounted/.test(wait.error), "the walk's own sentence");
+ } finally {
+ release();
+ }
+});
diff --git a/common/jobs/snapshotScheduler.ts b/common/jobs/snapshotScheduler.ts
@@ -1,3 +1,5 @@
+import path from "node:path";
+import { readFile } from "node:fs/promises";
import type { Paths } from "../lib/paths";
import {
DEFAULT_REPORT_DEBOUNCE_PRESET,
@@ -12,6 +14,7 @@ import { getRegistry } from "./registry";
import { drainStream } from "./drainStream";
import type { StreamActionResult } from "./streamCommand";
import { REFRESH_REPORT_QUEUE } from "../lib/queueKeys";
+import { readJobMeta } from "./jobMeta";
// Global, debounced "channel report" (snapshot) regeneration scheduler.
//
@@ -206,7 +209,6 @@ export function isRefreshReportPending(slug: string): boolean {
export async function startRefreshReport(
paths: Paths,
slug: string,
- opts: { onError?: (message: string) => void } = {},
): Promise<StreamActionResult> {
if (isRefreshReportPending(slug)) {
return { ok: false, error: REFRESH_REPORT_ACTIVE, info: true };
@@ -222,13 +224,7 @@ export async function startRefreshReport(
channelSlug: slug,
fn: async (onLog) => {
onLog(`Regenerating report for ${slug}…`);
- let snap: Awaited<ReturnType<typeof generateChannelSnapshot>>;
- try {
- snap = await generateChannelSnapshot(paths, slug);
- } catch (err) {
- opts.onError?.((err as Error).message);
- throw err;
- }
+ const snap = await generateChannelSnapshot(paths, slug);
getState().generation++;
const excluded = excludedDownloadIdSet(snap);
const awaitingTranscription = excluded.size
@@ -254,6 +250,111 @@ export async function startRefreshReport(
}
}
+// ONE CHANNEL'S REPORT, ASKED FOR BY A PERSON OR A SCRIPT: the job that will
+// regenerate it — started now, or the one already queued for the channel
+// (which has not started reading, so it is as fresh). A caller mid-enqueue for
+// the same slug has no id yet; it is waited for, briefly. The stream is
+// cancelled: nothing reads it (the log is on disk).
+export async function requestRefreshReport(
+ paths: Paths,
+ slug: string,
+): Promise<{ ok: true; jobId: string; started: boolean } | { ok: false; error: string }> {
+ const started = await startRefreshReport(paths, slug);
+ if (started.ok) {
+ void started.stream.cancel();
+ return { ok: true, jobId: started.jobId, started: true };
+ }
+ if (!started.info) return { ok: false, error: started.error };
+ for (let i = 0; i < 40; i++) {
+ const queued = queuedRefreshReportId(slug);
+ if (queued) return { ok: true, jobId: queued, started: false };
+ if (!getState().starting.has(slug)) break;
+ await new Promise((resolve) => setTimeout(resolve, 50));
+ }
+ // The one we would have waited for started in the meantime (or another
+ // caller's enqueue failed): ask once more — now nothing is queued, so this
+ // queues one, or says why it cannot.
+ const again = await startRefreshReport(paths, slug);
+ if (again.ok) {
+ void again.stream.cancel();
+ return { ok: true, jobId: again.jobId, started: true };
+ }
+ const queued = queuedRefreshReportId(slug);
+ return queued
+ ? { ok: true, jobId: queued, started: false }
+ : { ok: false, error: again.error };
+}
+
+// HOW LONG A PERSON'S "Refresh report" WAITS before it answers with where its
+// job is instead. The queue is serial: behind a 3,000-video walk, or during
+// Update all reports, the report may be minutes away, and a button that says
+// "Refreshing…" for minutes says nothing.
+export const REFRESH_REPORT_WAIT_MS = 15_000;
+
+export type RefreshReportWait =
+ | { state: "done"; jobId: string }
+ | { state: "failed"; jobId: string; error: string }
+ // Still queued or running when the wait ran out. `ahead`: the jobs before it
+ // on the refresh-report queue — 0 means it is the one regenerating now.
+ | { state: "waiting"; jobId: string; ahead: number };
+
+// The `[error]` line a failed job's log ends with (streamCommand writes the
+// thrown sentence there — an unmounted drive's, for a walk), else null.
+async function lastErrorLine(paths: Paths, jobId: string): Promise<string | null> {
+ try {
+ const raw = await readFile(path.join(paths.jobsDir, `${jobId}.log`), "utf8");
+ const lines = raw.split("\n").filter((l) => l.startsWith("[error] "));
+ const last = lines[lines.length - 1];
+ return last ? last.slice("[error] ".length) : null;
+ } catch {
+ return null;
+ }
+}
+
+// Wait up to `timeoutMs` for a refresh-report job to end, and say how it ended
+// — or where it is. Polls the registry (it hands out no completion promise for
+// a job someone else started); a job evicted from it is read from its meta.
+export async function waitForRefreshReport(
+ paths: Paths,
+ jobId: string,
+ opts: { timeoutMs?: number; pollMs?: number } = {},
+): Promise<RefreshReportWait> {
+ const timeoutMs = opts.timeoutMs ?? REFRESH_REPORT_WAIT_MS;
+ const pollMs = opts.pollMs ?? 250;
+ const until = Date.now() + timeoutMs;
+ for (;;) {
+ const job = getRegistry().get(jobId);
+ let status = job?.status;
+ if (!job) {
+ status = (await readJobMeta(paths, jobId))?.status ?? "failed";
+ }
+ if (status === "done") return { state: "done", jobId };
+ if (status !== "queued" && status !== "running") {
+ const error =
+ (await lastErrorLine(paths, jobId)) ??
+ `Refresh report ${status} (job ${jobId})`;
+ return { state: "failed", jobId, error };
+ }
+ if (Date.now() >= until) {
+ return {
+ state: "waiting",
+ jobId,
+ ahead: Math.max(0, getRegistry().positionInQueue(jobId)),
+ };
+ }
+ await new Promise((resolve) => setTimeout(resolve, pollMs));
+ }
+}
+
+// The operator's sentence for a wait that ran out.
+export function refreshReportWaitNotice(w: Extract<RefreshReportWait, { state: "waiting" }>): string {
+ const where =
+ w.ahead === 0
+ ? "Regenerating now"
+ : `Queued behind ${w.ahead} report regeneration${w.ahead === 1 ? "" : "s"}`;
+ return `${where} — the report updates when it finishes (job ${w.jobId}).`;
+}
+
async function fire(): Promise<void> {
const state = getState();
// Snapshot and clear before the async work: any action that fires during