import path from "node:path"; import { readFile } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { DEFAULT_REPORT_DEBOUNCE_PRESET, REPORT_DEBOUNCE_PRESETS, getSettings, } from "../lib/settings"; import { excludedDownloadIdSet, generateChannelSnapshot, } from "../controller/channelSnapshot"; 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. // // Any action that changes a channel's report-relevant data calls // `requestChannelSnapshot(paths, slug)` when it finishes. That marks the slug // dirty and (re)arms a single shared timer. When the timer fires, every dirty // channel's snapshot is regenerated in parallel via the existing // `refresh-report` managed-job pattern, then the affected pages are revalidated. // // The debounce coalesces bursts of actions into one regen pass instead of // rewriting snapshot.json after every individual operation. The window is read // from global settings (REPORT_DEBOUNCE_PRESETS) at arm time. // Job kinds that must NOT trigger a regen. "refresh-report" is the regen itself // — including it would make the central runManagedFunction hook re-arm the timer // from inside the regen job, an infinite loop. "detect-duplicates" is a global // read-only scan with no per-channel snapshot impact. // // The digest kinds are here for a different, load-bearing reason: SCALE. // ctx.recordTaskDone arms a channel-snapshot regen after EVERY completed // sub-operation, and each regen is a full 16-way per-video fan-out over the whole // channel. A digest sweep is ~119,600 sub-operations corpus-wide, so leaving them // in would spend more machine time regenerating snapshots than generating // digests. The digest actions instead regenerate ONCE per channel, at job end. const NO_REGEN_KINDS = new Set([ "refresh-report", "detect-duplicates", "digest-channel-local", "digest-channel-remote", "digest-share-cluster", // The backfill kinds are here for exactly the digest reason above: a corpus // sweep is tens of thousands of sub-operations, ctx.recordTaskDone arms a // regen after every one, and each regen is a full 16-way per-video walk of the // channel. Leaving them in would spend more machine time regenerating reports // than doing the backfill. The per-channel backfill job regenerates ONCE, at // job end. "backfill-channel", // A RELOCATION MOVES BYTES, IT DOES NOT CHANGE THEM. Every video dir, every // transcript and every mtime is identical afterwards — `rsync -a` preserves // them, which is the same property that makes the LMDB index a no-op — so a // regen would walk the whole channel to write a byte-identical snapshot with // a newer generatedAt, claiming a measurement it did not take. On the 130 GB // channel this exists for that is a 16-way walk of 11,000 video dirs for // nothing, immediately after a job that just moved 130 GB. // // It also removes a race the operator would feel: the regen is a queued job // for this channel, and the Storage panel refuses a move while the channel // has one — so "move out, then move back" would be blocked by a report // nobody needed. "relocate-channel-media", // Moving the saved-video store changes where bytes are, not what any channel // has — every count in every report is identical afterwards. "relocate-saved-videos", // A RE-POINT DOES NOT EVEN MOVE THE BYTES. It rewrites n symlinks and n // `dataDir` fields so the media that came back on a different mountpoint is // reachable again; every count in every report is what it was before. It // also carries no channelSlug on its record — it is a job about a LOCATION, // not a channel — so the central hook has nothing to regenerate anyway, and // listing it here says that on purpose rather than by omission. "repoint-storage-location", ]); export function shouldRequestSnapshot(kind: string): boolean { return !NO_REGEN_KINDS.has(kind); } type SchedulerState = { // slug -> Paths to regen with. getPaths() is process-stable, so last-writer // -wins is fine; the Map dedups repeat requests for the same channel. dirty: Map; timer: ReturnType | null; // When the current (uncleared) dirty batch started, for the max-wait cap. firstDirtyAt: number | null; // Bumped every time a report is actually rewritten. // // The editor's /api/pulse change token is deliberately built from in-memory // state only — it must never walk the corpus — so it cannot notice that a // snapshot FILE changed on disk. Without this counter, a report regenerated // on the debounce AFTER the triggering job already reached its final status // would move nothing the pulse can see, and pages showing snapshot-derived // 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; }; declare global { // eslint-disable-next-line no-var var __yttSnapshotScheduler__: SchedulerState | undefined; } function getState(): SchedulerState { if (!globalThis.__yttSnapshotScheduler__) { globalThis.__yttSnapshotScheduler__ = { dirty: new Map(), timer: null, firstDirtyAt: null, generation: 0, starting: new Set(), }; } return globalThis.__yttSnapshotScheduler__; } function resolveWindow(): { debounceMs: number; maxWaitMs: number | null } { try { const preset = getSettings().reportDebouncePreset; return REPORT_DEBOUNCE_PRESETS[preset] ?? REPORT_DEBOUNCE_PRESETS.fast; } catch { // getSettings reads the filesystem; fall back to the default window if it // is unavailable for any reason. return REPORT_DEBOUNCE_PRESETS[DEFAULT_REPORT_DEBOUNCE_PRESET]; } } // Mark a channel's report dirty and (re)arm the global debounce timer. // Synchronous and non-throwing — safe to call from job completion handlers and // request-scoped server actions without awaiting or try/catch. export function requestChannelSnapshot(paths: Paths, slug: string): void { if (!slug) return; const state = getState(); state.dirty.set(slug, paths); if (state.firstDirtyAt === null) state.firstDirtyAt = Date.now(); const { debounceMs, maxWaitMs } = resolveWindow(); if (state.timer) clearTimeout(state.timer); const delay = maxWaitMs === null ? debounceMs : Math.max( 0, Math.min(debounceMs, maxWaitMs - (Date.now() - state.firstDirtyAt)), ); state.timer = setTimeout(() => void fire(), delay); // Never let a pending regen keep the process (or the e2e test runner) alive. state.timer.unref?.(); } // 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 // 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 { REFRESH_REPORT_QUEUE }; // Why startRefreshReport started nothing for a slug that already has one. export const REFRESH_REPORT_ACTIVE = "already queued"; // The QUEUED regeneration of `slug` (not yet started), if there is one. export function queuedRefreshReportId(slug: string): string | null { const queued = getRegistry() .list() .find( (j) => j.kind === "refresh-report" && j.channelSlug === slug && 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 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, ): Promise { if (isRefreshReportPending(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); } } // 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 { 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 { 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): 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 { 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. 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; state.firstDirtyAt = null; if (batch.length === 0) return; try { const regenerated: string[] = []; const streams: ReadableStream[] = []; for (const [slug, paths] of batch) { 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. 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) { try { // Imported lazily so non-Next consumers of common/jobs never resolve // next/cache at module load. fire() only runs in the editor runtime. const { revalidatePath } = await import("next/cache"); for (const slug of regenerated) revalidatePath(`/channels/${slug}`); revalidatePath("/channels"); revalidatePath("/operations/[id]", "page"); revalidatePath("/cleanup"); revalidatePath("/"); } catch { // revalidatePath outside a request/Next runtime — the snapshots are // still written; the next render picks them up. } } } catch { // Backstop: a thrown fire() must never become an unhandledRejection. The // per-channel refresh-report jobs handle their own errors/logging. } } // Fire any pending regen immediately and await it. For tests/diagnostics. export async function flushChannelSnapshotsNow(): Promise { const state = getState(); if (state.timer) { clearTimeout(state.timer); state.timer = null; } await fire(); } // Clear the timer and drop all scheduler state. Called by the e2e cache-reset // route so a pending regen can't fire against the next spec. export function resetSnapshotScheduler(): void { const existing = globalThis.__yttSnapshotScheduler__; if (existing?.timer) clearTimeout(existing.timer); globalThis.__yttSnapshotScheduler__ = undefined; }