commit bd125ab0aea18df0e79644a8ea558f5d5a44a7e8
parent ae2c0f42c3c2e240af6d0efbc3605a8416595b52
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 29 Jul 2026 11:31:43 -0400
The sweep can now be started, survive a restart, pause, and yield the GPU
Four things stood between "the digest layer works" and "a corpus backfill is
running unattended". None of them were on the roadmap's gate list.
**A launcher.** Every entrypoint took a single channelSlug, so a corpus sweep
meant 63 manual clicks. `digestSweep.ts` walks the channels and launches the
existing per-channel job in sequence — deliberately NOT a new batch shape, since
one job per video would be 119,600 jobs against a 100-record registry and would
evict the history of the run it was recording. Heaviest channel first by
remaining AUDIO-HOURS, because cost is audio: HasanAbiVODs3 alone is 8,331 h and
leaving it for last is how a sweep spends 70 days looking nearly done. It stores
no cursor and re-plans every pass, so a finished transcription, a confirmed
duplicate cluster or a settings change that invalidates the corpus are all just
visible next time round.
**Boot resume.** instrumentation.ts restarted the sync heartbeat, the
transcription pause and the auto-runners, but not digests — so a restart
silently ended a sweep that had been running for days. Re-launching is safe, not
just convenient: nothing is frozen, so a resumed sweep re-does exactly zero work.
**A pause that something writes.** `digestBatch`'s limit() has always re-read
`digestsPaused` and returned 0 to make the pool idle-WAIT (a hold, not a stop, so
resuming re-derives nothing). Nothing ever set the flag —
settings/actions.ts passed it through with a comment deferring to
"the dashboard/channel controls", which did not exist. They do now, and they are
two controls, not one: arming a sweep and pausing one are different questions.
**GPU arbitration — the actual cost lever.** `queueKeys.ts` puts digest:local on
a different queue from TRANSCRIPTION_QUEUE so the lanes don't serialize, whose
consequence is that ollama and whisper share one 8 GB card. That is the measured
90 s/audio-hour against the 27 s an idle box projected: ~80 sweep days versus
~24. Bigger than every other lever in this plan by a factor of about fifty. The
fix reuses the pause: while whisper works, limit() returns 0 and the digest lane
steps aside. No scheduler, no priorities, no boot hook. It is a yield, not a
lock — a generation already in flight is not interrupted, because throwing away
partial work to save contention is a bad trade.
Also here, because a sweep you cannot watch is a sweep you cannot trust:
- `listChannels` emitted every count except the digest one it had already
computed. One line, and every cross-channel surface gets coverage.
- Coverage scalars on the widget poll (honouring its scalars-only rule), a
digest instrument on the dashboard, and the percentage shown to two decimals
because rounding 0.13% up to 1% flatters an 80-day backfill.
- **Progress could not move during a regeneration.** The bar re-counted
`current` from disk, and a regenerated digest is rewritten in place — so the
file count never moved and the bar sat at 0% for the whole job
({initial:39, current:39, target:72, pct:0}). Cosmetic once; over 81 days it
makes working work look wedged. The runner now reports its own `current`,
which also lets it raise its own target when sharing satisfies more videos
than were counted missing.
- That fix wires `setProgress`/`progressBaseline`, which were declared on the
options type, passed by the caller, and referenced nowhere in the body.
- ETAs for length-proportional work are computed in SECONDS PER AUDIO-HOUR.
Averaging tasks assumes every unit costs the same, which is true for downloads
and wildly false for digests — a VOD channel and a shorts channel have the
same task count and a 50x difference in cost, so a task average would be wrong
by an order of magnitude exactly when an operator most needs the number.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
15 files changed, 881 insertions(+), 32 deletions(-)
diff --git a/common/controller/channels.ts b/common/controller/channels.ts
@@ -197,6 +197,11 @@ export async function listChannels(paths: Paths): Promise<ChannelStat[]> {
videoCount: counts.videos,
transcriptCount: counts.transcripts,
downloadCount: counts.downloads,
+ // `countDataFiles` has always computed this and `listChannels` has always
+ // discarded it — unlike readChannelStat, which emits it. One line, and
+ // every cross-channel surface (the dashboard, /actionable, the widget)
+ // gets a coverage counter it was already paying the I/O for.
+ digestCount: counts.digests,
});
}
return out.sort((a, b) => a.slug.localeCompare(b.slug));
diff --git a/common/controller/digestBatch.ts b/common/controller/digestBatch.ts
@@ -32,6 +32,7 @@ import {
} from "./digestTarget";
import { CUES_JSON_FILENAME } from "../lib/videoStatus";
import { digestVideo } from "./digestVideo";
+import { transcriptionActivity } from "./digestYield";
import {
buildDigestClusterPlan,
planSlugForDir,
@@ -70,7 +71,14 @@ export type DigestBatchOptions = {
useClusters?: boolean;
clusterPlan?: DigestClusterPlan;
setProgress?: (snap: JobProgress) => void;
+ // Where the bar starts: how many videos already counted as done before this
+ // run. The batch reports `initial + completed` against it rather than letting
+ // the UI re-count files from disk, which cannot see a regeneration.
progressBaseline?: number;
+ // Where the bar ends. Supplied by the caller because it already resolves it
+ // (countMissingDigests) to size the job; the batch raises it if sharing turns
+ // out to satisfy more videos than were counted missing.
+ progressTarget?: number;
onLog?: (msg: string) => void;
signal?: AbortSignal;
drainSignal?: AbortSignal;
@@ -269,6 +277,7 @@ export async function runDigestBatch(
if (role?.kind === "mirror") {
attempted.add(candidate.id);
cursor++;
+ resolvedAudioSeconds += candidate.duration;
result.skipped++;
log(
`Skipping ${candidate.id}: duplicate of ${role.canonicalSlug}, which owns the digest for this cluster.`,
@@ -286,6 +295,7 @@ export async function runDigestBatch(
if (allFresh) {
attempted.add(candidate.id);
cursor++;
+ resolvedAudioSeconds += candidate.duration;
result.fresh++;
continue;
}
@@ -297,6 +307,47 @@ export async function runDigestBatch(
return null;
};
+ // WIRING THE PLUMBING THAT WAS DECLARED AND NEVER USED.
+ //
+ // `setProgress` and `progressBaseline` were on the options type and passed by
+ // digestActions.ts, but nothing in this file referenced either. Progress
+ // appeared to work only because buildActiveJobs re-counts `current` from disk
+ // — and that re-count is exactly what a REGENERATION defeats: a regenerated
+ // digest is rewritten in place, the file count never moves, and the bar sits
+ // at 0% for the whole job.
+ //
+ // The batch is the only thing that knows the truth, for two reasons the disk
+ // cannot express: it knows a regenerate happened, and it knows sharing moved
+ // the denominator — a canonical member's digest can satisfy a dozen mirrors
+ // at once, so a sweep genuinely CHANGES its own target as it runs.
+ const baseline = opts.progressBaseline ?? 0;
+ const progressTarget = opts.progressTarget ?? null;
+ // Audio-seconds retired from the worklist, by ANY route — generated, shared
+ // to, found fresh, or skipped as a mirror. All four remove work, and an ETA
+ // that only counted generations would keep quoting time for videos that are
+ // already done.
+ const totalAudioSeconds = candidates.reduce((n, c) => n + c.duration, 0);
+ let resolvedAudioSeconds = 0;
+ const reportProgress = (): void => {
+ if (!opts.setProgress) return;
+ const done = result.succeeded + result.shared + result.fresh;
+ const remainingAudioSeconds = Math.max(
+ 0,
+ totalAudioSeconds - resolvedAudioSeconds,
+ );
+ // Never let the bar exceed its target: sharing can satisfy more videos than
+ // the target was sized for, and a bar past 100% reads as a bug rather than
+ // as good news.
+ const current = baseline + done;
+ opts.setProgress({
+ metric: "digests",
+ initial: baseline,
+ target: Math.max(progressTarget ?? current, current),
+ current,
+ remainingAudioSeconds,
+ });
+ };
+
const runOne = async (
candidate: Candidate,
runSignal: AbortSignal,
@@ -305,6 +356,10 @@ export async function runDigestBatch(
id: candidate.id,
label: `digest ${opts.channelSlug}/${candidate.id}`,
kind: "digest",
+ // What makes seconds-per-audio-hour computable. A digest's cost is
+ // proportional to the transcript's LENGTH, not to it being one video, so
+ // a task-count average is the wrong denominator for a sweep ETA.
+ audioSeconds: candidate.duration,
});
try {
const outcome = await digestVideo({
@@ -356,6 +411,7 @@ export async function runDigestBatch(
else if (o.status === "misaligned") result.misaligned++;
}
}
+ reportProgress();
} catch (err) {
if (runSignal.aborted || opts.signal?.aborted) throw err;
result.attempted++;
@@ -364,10 +420,15 @@ export async function runDigestBatch(
`Failed ${candidate.id}: ${(err as Error)?.message ?? String(err)}`,
);
} finally {
+ resolvedAudioSeconds += candidate.duration;
task?.end();
}
};
+ // Edge-triggered, so the yield logs twice per contention window rather than
+ // once per poll over a multi-week sweep.
+ let yielding = false;
+
const NEVER = new AbortController().signal;
// Local lane: 1. The GPU is the bottleneck and a second concurrent generation
// just thrashes the same 8 GB of VRAM. The metered lane is network-bound, so it
@@ -375,6 +436,10 @@ export async function runDigestBatch(
const defaultConcurrency = app.lane === "local-gpu" ? 1 : 2;
const concurrency = Math.max(1, opts.concurrency ?? defaultConcurrency);
+ // Seed the bar before the first video finishes, so a job that spends its
+ // first minutes on a 3-hour VOD does not look like it never started.
+ reportProgress();
+
await runPool<Candidate>({
next,
run: runOne,
@@ -384,6 +449,33 @@ export async function runDigestBatch(
// idle-wait, which is a pause; returning null from next() would END the
// batch, which is not.
if (getSettings().digest.digestsPaused) return 0;
+ // Step aside for whisper. Same mechanism as the pause and for the same
+ // reason it works: a zero limit HOLDS the pool instead of ending the
+ // batch, so the sweep resumes the moment the card is free without
+ // re-deriving anything. Only the local (GPU) lane yields — the metered
+ // lane is network-bound and competes for nothing here.
+ if (
+ app.lane === "local-gpu" &&
+ getSettings().digest.yieldToTranscription
+ ) {
+ const activity = transcriptionActivity();
+ if (activity.busy) {
+ // Logged on the EDGE only. An operator watching a sweep sit at zero
+ // throughput has to be able to tell yielding from wedged, but a line
+ // per poll would bury the job log over a multi-week run.
+ if (!yielding) {
+ yielding = true;
+ log(
+ `Yielding the GPU to transcription (${activity.reason}); the digest lane will resume when it is free.`,
+ );
+ }
+ return 0;
+ }
+ if (yielding) {
+ yielding = false;
+ log("Transcription finished; resuming the digest lane.");
+ }
+ }
if (
app.metered &&
digestSettings.spendCapUsd > 0 &&
diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts
@@ -0,0 +1,325 @@
+// The corpus-wide digest sweep: one launcher for the thing that runs for weeks.
+//
+// Before this, "sweep the corpus" meant clicking Digest on 63 channel pages by
+// hand and losing the whole thing to a server restart. Everything underneath —
+// the per-channel batch, resume-by-re-deriving-from-disk, the pause, the GPU
+// yield — already existed and was already correct. What was missing was
+// something to walk the channels.
+//
+// Three shapes are deliberate:
+//
+// 1. IT LAUNCHES THE EXISTING PER-CHANNEL JOB, in sequence. It does not
+// introduce a batch over videos. digestBatch.ts is explicit that one job
+// per video would be 119,600 jobs against a 100-record registry and a
+// 500-record log, evicting the history of the very run it is recording.
+// 63 sequential channel jobs is the same contract, driven.
+//
+// 2. IT STORES NO CURSOR. Ordering is recomputed from
+// `buildDigestSweepPlan` each pass and eligibility is re-derived from disk
+// inside every batch, so a restart, a newly finished transcription, a
+// human confirming a duplicate cluster, and a settings change that
+// invalidates the corpus are all just visible on the next pass. A frozen
+// work-list would be wrong within hours of a multi-week run starting.
+//
+// 3. IT LOOPS UNTIL THE PLAN IS EMPTY, rather than passing over the channel
+// list once. A batch can legitimately leave work behind — the spend cap,
+// a drain, a video whose transcript landed mid-pass — and a one-pass sweep
+// would report itself finished with the corpus unfinished.
+//
+// Heaviest channel first, by remaining AUDIO-HOURS. Cost is audio, not videos:
+// HasanAbiVODs3 is 8,331 audio-hours and outweighs every duplicate mirror in
+// the corpus combined, so leaving it for last is how a sweep spends 70 days
+// looking nearly done.
+
+import type { Paths } from "../lib/paths";
+import { getPaths } from "../lib/paths";
+import { getSettings, writeSettings } from "../lib/settings";
+import { getRegistry } from "../jobs/registry";
+import { runManagedFunction } from "../jobs/streamCommand";
+import { drainStream } from "../jobs/drainStream";
+import { makeTaskTracker } from "../jobs/taskHooks";
+import { requestChannelSnapshot } from "../jobs/snapshotScheduler";
+import { DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE } from "../lib/queueKeys";
+import { countMissingDigests, runDigestBatch } from "./digestBatch";
+import {
+ audioHours,
+ buildDigestSweepPlan,
+ sweepDays,
+ type DigestSweepPlan,
+} from "./digestPlan";
+import type { DigestLaneChoice } from "./digestTarget";
+
+export const DIGEST_SWEEP_KIND = "digest-sweep";
+
+// How long the loop waits before recomputing the plan when a pass did no work.
+// Generous: nothing here is latency-sensitive and recomputing the plan reads
+// every channel's sidecars.
+const IDLE_POLL_MS = 60_000;
+
+type SweepLive = { jobId: string; startedAt: number };
+type SweepSingleton = { live: SweepLive | null };
+
+declare global {
+ // eslint-disable-next-line no-var
+ var __yttDigestSweep__: SweepSingleton | undefined;
+}
+
+function getSingleton(): SweepSingleton {
+ if (!globalThis.__yttDigestSweep__) {
+ globalThis.__yttDigestSweep__ = { live: null };
+ }
+ return globalThis.__yttDigestSweep__;
+}
+
+export function getDigestSweepJobId(): string | null {
+ const live = getSingleton().live;
+ if (!live) return null;
+ return getRegistry().get(live.jobId)?.status === "running"
+ ? live.jobId
+ : null;
+}
+
+export type DigestSweepOptions = {
+ paths?: Paths;
+ lane?: DigestLaneChoice;
+ // Restrict the sweep to these channels. Absent = the whole corpus.
+ channelSlugs?: string[];
+};
+
+// Run ONE channel through the existing per-channel batch, as its own managed
+// job — same kind, same queue and same progress metric a hand-clicked run
+// produces, so the sweep is inspectable with the tools that already exist
+// rather than being an opaque mega-job.
+export async function runDigestChannelJob(opts: {
+ paths: Paths;
+ channelSlug: string;
+ lane: DigestLaneChoice;
+ background?: boolean;
+ onDone?: () => void;
+}): Promise<{ ok: boolean; error?: string }> {
+ const { paths, channelSlug, lane } = opts;
+ const settings = getSettings();
+ const remote = lane === "remote";
+ const kind = remote ? "digest-channel-remote" : "digest-channel-local";
+
+ const result = await runManagedFunction({
+ kind,
+ queueKey: remote ? DIGEST_REMOTE_QUEUE : DIGEST_LOCAL_QUEUE,
+ paths,
+ channelSlug,
+ background: opts.background,
+ spec: { kind, slug: channelSlug, params: { lane } },
+ fn: async (onLog, signal, setProgress, ctx) => {
+ const missing = await countMissingDigests(
+ paths,
+ channelSlug,
+ undefined,
+ lane,
+ );
+ const batch = await runDigestBatch({
+ channelSlug,
+ paths,
+ lane,
+ ...(remote
+ ? { minDurationSeconds: settings.digest.longTailSeconds }
+ : {}),
+ setProgress,
+ // The bar measures THIS run, from zero. Seeding it with the count of
+ // digests already on disk was the old shape and it could not represent
+ // a regeneration, where the file count never moves.
+ progressBaseline: 0,
+ progressTarget: missing,
+ onLog,
+ signal,
+ drainSignal: ctx.drainSignal,
+ tracker: makeTaskTracker(ctx, onLog),
+ });
+ onLog(
+ `Digest batch: ${batch.succeeded} generated, ${batch.fresh} already current, ` +
+ `${batch.shared} shared to mirrors, ${batch.misaligned} mirror(s) refused by the alignment gate, ` +
+ `${batch.skipped} skipped, ${batch.failed} failed; ` +
+ `${batch.engineCalls} model call(s), ${batch.warnings} warning(s)` +
+ (batch.costUsd > 0 ? `, $${batch.costUsd.toFixed(4)}` : "") +
+ (batch.spendCapped ? " (stopped at the spend cap)" : "") +
+ ".",
+ );
+ requestChannelSnapshot(paths, channelSlug);
+ opts.onDone?.();
+ },
+ });
+ if (!result.ok) return { ok: false, error: result.error };
+ // Wait for the channel to finish before returning: the sweep is sequential by
+ // design (one GPU), and the queue would serialize these anyway — awaiting
+ // makes that explicit and lets the loop re-plan against real results.
+ await drainStream(result.stream);
+ return { ok: true };
+}
+
+async function runSweepLoop(
+ paths: Paths,
+ lane: DigestLaneChoice,
+ channelSlugs: string[] | undefined,
+ onLog: (msg: string) => void,
+ signal: AbortSignal,
+ drainSignal: AbortSignal,
+): Promise<void> {
+ let pass = 0;
+ for (;;) {
+ if (signal.aborted || drainSignal.aborted) return;
+ // The operator turned the sweep off: stop cleanly rather than being
+ // cancelled, so the job ends "done" and the queue is released.
+ if (!getSettings().digest.sweepEnabled) {
+ onLog("Sweep disarmed in settings — stopping.");
+ return;
+ }
+
+ pass++;
+ const plan: DigestSweepPlan = await buildDigestSweepPlan({
+ paths,
+ lane,
+ channelSlugs,
+ onLog,
+ });
+ const work = plan.channels.filter((c) => c.generateSeconds > 0);
+ if (work.length === 0) {
+ onLog(
+ `Pass ${pass}: nothing left to generate — ${plan.fresh.videos.toLocaleString()} video(s) already digested at the current identity.`,
+ );
+ return;
+ }
+ onLog(
+ `Pass ${pass}: ${work.length} channel(s), ` +
+ `${audioHours(plan.generateSeconds).toFixed(0)} audio-hours to generate ` +
+ `(~${sweepDays(plan.generateSeconds).toFixed(1)} days at the measured rate), ` +
+ `${audioHours(plan.sharedSeconds).toFixed(0)} audio-hours covered by cluster sharing.`,
+ );
+
+ let didWork = false;
+ for (const channel of work) {
+ if (signal.aborted || drainSignal.aborted) return;
+ if (!getSettings().digest.sweepEnabled) {
+ onLog("Sweep disarmed in settings — stopping.");
+ return;
+ }
+ onLog(
+ `→ ${channel.channelSlug}: ${audioHours(channel.generateSeconds).toFixed(0)} audio-hours ` +
+ `(~${sweepDays(channel.generateSeconds).toFixed(1)} days).`,
+ );
+ const outcome = await runDigestChannelJob({
+ paths,
+ channelSlug: channel.channelSlug,
+ lane,
+ // Behind anything an operator clicks by hand: a sweep is weeks long and
+ // must never make a deliberate single-channel run wait for it.
+ background: true,
+ });
+ if (!outcome.ok) {
+ // One channel failing is not the sweep failing. Log it and carry on —
+ // an unreachable engine will fail every channel and the loop's own
+ // no-progress guard below is what stops the spin.
+ onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`);
+ continue;
+ }
+ didWork = true;
+ }
+
+ if (!didWork) {
+ // Every channel refused to start (engine down, queue wedged). Back off
+ // rather than spinning through 63 failures a second.
+ onLog(
+ `Pass ${pass} started nothing; waiting ${IDLE_POLL_MS / 1000}s before re-planning.`,
+ );
+ await sleep(IDLE_POLL_MS, signal);
+ }
+ }
+}
+
+function sleep(ms: number, signal: AbortSignal): Promise<void> {
+ return new Promise((resolve) => {
+ const t = setTimeout(resolve, ms);
+ signal.addEventListener(
+ "abort",
+ () => {
+ clearTimeout(t);
+ resolve();
+ },
+ { once: true },
+ );
+ });
+}
+
+// Arm and start the corpus-wide sweep. Persists `sweepEnabled` so a restart
+// resumes it (see resumeDigestSweepIfEnabled).
+export async function startDigestSweep(
+ opts: DigestSweepOptions = {},
+): Promise<string | null> {
+ const paths = opts.paths ?? getPaths();
+ const running = getDigestSweepJobId();
+ if (running) return running;
+
+ const settings = getSettings();
+ if (!settings.digest.sweepEnabled) {
+ writeSettings({
+ ...settings,
+ digest: { ...settings.digest, sweepEnabled: true },
+ });
+ }
+
+ const live: SweepLive = { jobId: "", startedAt: Date.now() };
+ const lane = opts.lane ?? "local";
+ const result = await runManagedFunction({
+ kind: DIGEST_SWEEP_KIND,
+ // Empty key: the orchestrator itself does no work, it waits on per-channel
+ // jobs that serialize on the digest queue. Putting it ON that queue would
+ // deadlock — it would hold the only slot while waiting for a job that needs
+ // the same slot.
+ queueKey: "",
+ paths,
+ fn: async (onLog, signal, _setProgress, ctx) => {
+ live.jobId = ctx.jobId;
+ try {
+ await runSweepLoop(
+ paths,
+ lane,
+ opts.channelSlugs,
+ onLog,
+ signal,
+ ctx.drainSignal,
+ );
+ } finally {
+ if (getSingleton().live === live) getSingleton().live = null;
+ }
+ },
+ });
+ if (!result.ok) return null;
+ live.jobId = result.jobId;
+ getSingleton().live = live;
+ return result.jobId;
+}
+
+// Boot hook. The auto-transcribe/auto-download runners are restarted at server
+// start the same way; digests were not, so a restart silently ended a sweep
+// that had been running for days and nothing said so.
+export async function resumeDigestSweepIfEnabled(
+ paths: Paths = getPaths(),
+): Promise<void> {
+ if (!getSettings().digest.sweepEnabled) return;
+ await startDigestSweep({ paths });
+}
+
+// Disarm and stop. Persisting the flag is the point: without it a restart would
+// resurrect a sweep the operator had just stopped.
+export function stopDigestSweep(): boolean {
+ const settings = getSettings();
+ if (settings.digest.sweepEnabled) {
+ writeSettings({
+ ...settings,
+ digest: { ...settings.digest, sweepEnabled: false },
+ });
+ }
+ const live = getSingleton().live;
+ if (!live) return false;
+ // Drain, not cancel: the in-flight channel finishes rather than losing a
+ // part-generated video, and the loop exits at its next check.
+ return getRegistry().requestDrain(live.jobId);
+}
diff --git a/common/controller/digestYield.ts b/common/controller/digestYield.ts
@@ -0,0 +1,68 @@
+// GPU arbitration between the digest lane and the transcription lane.
+//
+// `queueKeys.ts` deliberately puts `digest:local` on a DIFFERENT registry queue
+// from TRANSCRIPTION_QUEUE, so the two never serialize against each other. That
+// reasoning is right for digest-local vs digest-remote (GPU-bound vs
+// network-bound), but its consequence for digest-local vs transcription is that
+// **ollama and whisper run concurrently on the same 8 GB card**.
+//
+// That is not hypothetical: the bake-off projected 27 s per audio-hour on an
+// idle box and the 102-video validation run measured **90 s** on a box also
+// running auto-transcribe. Over a 77,000-audio-hour sweep the difference is
+// roughly 80 days versus 24 — far and away the largest cost lever in the
+// backfill, and about fifty times bigger than everything duplicate sharing can
+// save (see bin/digest-plan.ts).
+//
+// The fix is deliberately NOT a scheduler. `digestBatch`'s limit() already
+// returns 0 to idle-wait on a pause, and runPool treats a zero limit as "hold,
+// don't finish" — so the digest lane can step aside using machinery that is
+// already proven, with no priority system, no new queue and no boot hook.
+//
+// This is a YIELD, not a lock. A digest already in flight is not interrupted, so
+// the two can still overlap for the length of one generation if transcription
+// starts just after a dispatch. Interrupting would waste the partial work; the
+// next dispatch sees the busy lane and holds.
+
+import { getWorkerPool } from "../jobs/workerPool";
+import { getRegistry } from "../jobs/registry";
+import { TRANSCRIPTION_QUEUE } from "../lib/queueKeys";
+
+export type TranscriptionActivity = {
+ busy: boolean;
+ // Which signal fired, for the log line. An operator watching a sweep sit at
+ // zero throughput needs to be told it is yielding rather than wedged.
+ reason: "local-worker" | "queued-job" | null;
+};
+
+// Is the local transcription lane using the GPU right now?
+//
+// Two signals, because one alone leaves a hole:
+// - a busy LOCAL worker is whisper actually running on this box's card
+// (`remote` workers delegate to another machine and compete for nothing
+// here, so they are deliberately excluded);
+// - a RUNNING job on TRANSCRIPTION_QUEUE covers the gaps between worker
+// acquisitions — audio extraction, model load, the moment between two
+// videos in a batch. Those gaps are exactly where a multi-minute digest
+// generation would otherwise slip in and hold the card.
+export function transcriptionActivity(): TranscriptionActivity {
+ try {
+ for (const w of getWorkerPool().summary()) {
+ if (w.busy && w.kind === "local") {
+ return { busy: true, reason: "local-worker" };
+ }
+ }
+ } catch {
+ // A pool that cannot be read must not wedge the digest lane: fail OPEN, i.e.
+ // keep generating. The cost of being wrong here is contention, not deadlock.
+ }
+ try {
+ for (const job of getRegistry().list()) {
+ if (job.queueKey === TRANSCRIPTION_QUEUE && job.status === "running") {
+ return { busy: true, reason: "queued-job" };
+ }
+ }
+ } catch {
+ /* same: fail open */
+ }
+ return { busy: false, reason: null };
+}
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -20,6 +20,25 @@ export type JobProgress = {
metric: JobProgressMetric;
initial: number;
target: number;
+ // Progress as counted BY THE RUNNER, when the runner knows better than the
+ // disk does.
+ //
+ // The default is to re-count `current` from on-disk channel stats, which is
+ // right for downloads and transcripts: a file appears, the count goes up.
+ // It is WRONG for digests under regeneration — a regenerated digest is
+ // rewritten in place, so the file count never moves and the bar sits at 0%
+ // for the whole job (observed: {initial:39, current:39, target:72, pct:0}).
+ // Cosmetic for a five-minute job; over an 81-day sweep it makes working work
+ // look wedged.
+ //
+ // Only set it where the runner has a genuinely better number. Absent → the
+ // disk re-count, unchanged.
+ current?: number;
+ // Audio-seconds still to process. Paired with completedTaskAudioSeconds it
+ // gives an ETA in the unit the work is actually priced in; without it the
+ // estimate falls back to averaging TASKS, which for digests is wrong by more
+ // than an order of magnitude between a VOD channel and a shorts channel.
+ remainingAudioSeconds?: number;
};
export type JobTaskKind = "download" | "transcribe" | "digest";
@@ -83,6 +102,16 @@ export type JobRecord = {
// reflects the whole batch. Both undefined until the first task completes.
completedTaskCount?: number;
completedTaskMs?: number;
+ // Audio-seconds covered by completed sub-operations. Only digest/transcribe
+ // tasks report it, because only they consume work proportional to a video's
+ // LENGTH — a digest of a 3-hour VOD is not one task's worth of anything.
+ //
+ // It is what makes a sweep ETA expressible: `completedTaskMs` over this gives
+ // seconds-per-audio-hour, the unit the whole backfill is estimated in, and
+ // the only unit in which "how long is this going to take" has an answer. A
+ // task-count average cannot express it — the corpus is 77k videos and 77k
+ // audio-hours, and those distribute completely differently.
+ completedTaskAudioSeconds?: number;
};
export type QueueSnapshot = {
@@ -305,11 +334,19 @@ class JobRegistry {
// Fold one finished sub-operation's wall-clock duration into the job's
// running totals (see JobRecord.completedTaskCount/Ms). No-op if the job is
// gone or the duration is nonsensical.
- recordTaskDuration(jobId: string, durationMs: number): void {
+ recordTaskDuration(
+ jobId: string,
+ durationMs: number,
+ audioSeconds?: number,
+ ): void {
const job = this.jobs.get(jobId);
if (!job || durationMs < 0) return;
job.completedTaskCount = (job.completedTaskCount ?? 0) + 1;
job.completedTaskMs = (job.completedTaskMs ?? 0) + durationMs;
+ if (typeof audioSeconds === "number" && audioSeconds > 0) {
+ job.completedTaskAudioSeconds =
+ (job.completedTaskAudioSeconds ?? 0) + audioSeconds;
+ }
}
// Snapshot of queues for UI, resolved from the scheduler's id views back to
diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts
@@ -91,7 +91,7 @@ export type JobRunContext = {
},
) => void;
removeTask: (taskId: string) => void;
- recordTaskDone: (durationMs: number) => void;
+ recordTaskDone: (durationMs: number, audioSeconds?: number) => void;
};
export type RunManagedFunctionOpts = CommonOpts & {
@@ -324,8 +324,8 @@ export async function runManagedFunction(
addTask: (task) => registry.addTask(id, task),
updateTask: (taskId, patch) => registry.updateTask(id, taskId, patch),
removeTask: (taskId) => registry.removeTask(id, taskId),
- recordTaskDone: (ms) => {
- registry.recordTaskDuration(id, ms);
+ recordTaskDone: (ms, audioSeconds) => {
+ registry.recordTaskDuration(id, ms, audioSeconds);
// Refresh the report after EACH completed sub-operation (each video
// downloaded/transcribed in a batch), not only when the whole batch
// finishes — so a long batch updates incrementally. The global debounce
diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts
@@ -28,6 +28,11 @@ export type TaskTracker = {
appId?: string;
// The worker that owns this task, surfaced on the Workers page.
workerId?: string;
+ // The video's duration. Recorded against the task's measured wall clock so
+ // the job can report SECONDS PER AUDIO-HOUR — the unit a digest sweep is
+ // estimated in, and one a task-count average cannot express. Omit where the
+ // work is not proportional to length (a download is bytes, not minutes).
+ audioSeconds?: number;
}) => TaskHandle;
};
@@ -44,7 +49,7 @@ export function makeTaskTracker(
forwardLog: (line: string) => void,
): TaskTracker {
return {
- start({ id, label, kind, appId, workerId }) {
+ start({ id, label, kind, appId, workerId, audioSeconds }) {
if (!ctx) {
return { onLog: forwardLog, update: () => {}, end: () => {} };
}
@@ -97,7 +102,7 @@ export function makeTaskTracker(
const end = () => {
if (ended) return;
ended = true;
- ctx.recordTaskDone(Date.now() - startedAt);
+ ctx.recordTaskDone(Date.now() - startedAt, audioSeconds);
ctx.removeTask(id);
};
return { onLog, update, end };
diff --git a/common/lib/settings.ts b/common/lib/settings.ts
@@ -216,6 +216,19 @@ export type DigestSettings = {
// pattern), so a pause survives a restart with no boot hook — unlike
// transcriptionsPaused, which needs editor/instrumentation.ts to re-apply it.
digestsPaused: boolean;
+ // Yield the GPU to the transcription lane: while whisper is working, the
+ // digest batch's limit() returns 0 and the pool idle-waits. ON by default,
+ // because `digest:local` is deliberately on a different queue from
+ // TRANSCRIPTION_QUEUE and so would otherwise run ollama and whisper on the
+ // same 8 GB card — measured at 90 s per audio-hour against the 27 s an idle
+ // box projected. See controller/digestYield.ts.
+ yieldToTranscription: boolean;
+ // The corpus-wide sweep is armed. Read at boot by the editor's instrumentation
+ // hook, the same way the auto-transcribe/auto-download runners are, so a sweep
+ // survives a server restart. It is persisted INTENT, not a cursor: the batch
+ // re-derives eligibility from disk on every pull, so a resumed sweep does zero
+ // rework and needs nothing else remembered.
+ sweepEnabled: boolean;
// Hard ceiling on cumulative metered spend per job, USD. 0 = no cap. Only ever
// consulted for a metered app.
spendCapUsd: number;
@@ -502,6 +515,13 @@ export function defaultDigest(): DigestSettings {
// here would give the same number two homes and let them drift.
apps: {},
digestsPaused: false,
+ // ON. Contention with whisper is the single largest cost in the backfill
+ // (90 s/audio-hour measured vs 27 projected on an idle box), so the safe
+ // default is to step aside; turning it off is the deliberate choice.
+ yieldToTranscription: true,
+ // OFF. A corpus-wide sweep is GPU-weeks of work and is never armed by
+ // default — an operator starts it.
+ sweepEnabled: false,
spendCapUsd: 0,
sections: ["chapters"],
timestampMode: DEFAULT_DIGEST_TIMESTAMP_MODE,
@@ -569,6 +589,11 @@ export function sanitizeDigest(value: unknown): DigestSettings {
: d.remoteAppId,
apps: sanitizeDigestApps(r.apps),
digestsPaused: r.digestsPaused === true,
+ // Defaults to ON when absent — `=== false` rather than `!== true`, so a
+ // settings file written before this field existed keeps the GPU-safe
+ // behaviour instead of silently opting into contention.
+ yieldToTranscription: r.yieldToTranscription !== false,
+ sweepEnabled: r.sweepEnabled === true,
spendCapUsd:
typeof r.spendCapUsd === "number" && r.spendCapUsd > 0
? Math.round(r.spendCapUsd * 100) / 100
diff --git a/editor/app/api/widget/sync/route.ts b/editor/app/api/widget/sync/route.ts
@@ -23,6 +23,20 @@ export type WidgetSyncPayload = {
nextRunAt: number | null; // soonest eligible channel's nextDueAt
overdue: boolean; // some eligible channel is due now
};
+ // Corpus-wide digest coverage. SCALARS ONLY, honouring this payload's stated
+ // constraint — three numbers and two booleans, not a per-channel breakdown.
+ //
+ // It belongs on the poll rather than a page load because the number it
+ // reports moves over WEEKS: an operator watching an 80-day backfill needs to
+ // see it move at all, and "0.13% of 77,207" is not a figure any single
+ // channel page can show.
+ digest: {
+ digested: number; // videos carrying a non-empty ai-digest.json
+ videos: number; // videos in the corpus (the denominator)
+ channelsWithAny: number; // channels the layer has reached at all
+ paused: boolean; // settings.digest.digestsPaused
+ sweeping: boolean; // a corpus-wide sweep is armed
+ };
};
// Build the widget sync payload. Exported so the dashboard cockpit can seed its
@@ -61,6 +75,18 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> {
}
}
+ // Free: listChannels already counted these, it just used to throw the digest
+ // half away (see controller/channels.ts).
+ let digested = 0;
+ let videos = 0;
+ let channelsWithAny = 0;
+ for (const c of channels) {
+ videos += c.videoCount;
+ const n = c.digestCount ?? 0;
+ digested += n;
+ if (n > 0) channelsWithAny++;
+ }
+
return {
lastSyncAllAt: state.lastSyncAllAt,
lastIndividualSyncAt,
@@ -70,6 +96,13 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> {
nextRunAt,
overdue,
},
+ digest: {
+ digested,
+ videos,
+ channelsWithAny,
+ paused: settings.digest.digestsPaused,
+ sweeping: settings.digest.sweepEnabled,
+ },
} satisfies WidgetSyncPayload;
}
diff --git a/editor/app/channels/[slug]/digestActions.ts b/editor/app/channels/[slug]/digestActions.ts
@@ -8,7 +8,6 @@ import {
DIGEST_REMOTE_QUEUE,
resolveQueueKey,
} from "yt-dlp-transcript-common/lib/queueKeys";
-import { readChannelStat } from "yt-dlp-transcript-common/controller/channels";
import {
countMissingDigests,
runDigestBatch,
@@ -80,18 +79,15 @@ export async function digestChannelAction(
params: { queueKey, lane, order, limitCount, force },
},
fn: async (onLog, signal, setProgress, ctx) => {
- const stat = await readChannelStat(paths, slug);
- if (stat) {
- // The LANE matters: the progress target must be computed against the
- // engine that is actually about to run, not whichever one the local
- // setting names.
- const missing = await countMissingDigests(paths, slug, undefined, lane);
- setProgress({
- metric: "digests",
- initial: stat.digestCount ?? 0,
- target: (stat.digestCount ?? 0) + missing,
- });
- }
+ // The LANE matters: the progress target must be computed against the
+ // engine that is actually about to run, not whichever one the local
+ // setting names.
+ const missing = await countMissingDigests(paths, slug, undefined, lane);
+ // The bar measures THIS RUN, from zero, and the batch reports its own
+ // `current`. It used to be seeded with the count of digests already on
+ // disk and left to the UI's disk re-count — which cannot move during a
+ // REGENERATION, because a regenerated digest is rewritten in place. The
+ // bar sat at 0% for whole jobs ({initial:39, current:39, target:72}).
const result = await runDigestBatch({
channelSlug: slug,
paths,
@@ -107,7 +103,8 @@ export async function digestChannelAction(
typeof limitCount === "number" && limitCount > 0 ? limitCount : undefined,
force: force === true,
setProgress,
- progressBaseline: stat?.digestCount ?? 0,
+ progressBaseline: 0,
+ progressTarget: missing,
onLog,
signal,
drainSignal: ctx.drainSignal,
@@ -158,20 +155,18 @@ export async function digestBucketAction(
paths,
channelSlug: slug,
fn: async (onLog, signal, setProgress, ctx) => {
- const stat = await readChannelStat(paths, slug);
- if (stat) {
- const missing = await countMissingDigests(paths, slug, cleaned, lane);
- setProgress({
- metric: "digests",
- initial: stat.digestCount ?? 0,
- target: (stat.digestCount ?? 0) + missing,
- });
- }
+ // Same shape as the channel action: the bar measures this run from zero
+ // and the batch reports its own `current`, because the disk re-count it
+ // would otherwise use cannot see a digest rewritten in place.
+ const missing = await countMissingDigests(paths, slug, cleaned, lane);
const result = await runDigestBatch({
channelSlug: slug,
paths,
lane,
ids: cleaned,
+ setProgress,
+ progressBaseline: 0,
+ progressTarget: missing,
onLog,
signal,
drainSignal: ctx.drainSignal,
diff --git a/editor/app/components/dashboard/PipelineBand.tsx b/editor/app/components/dashboard/PipelineBand.tsx
@@ -8,6 +8,7 @@ import type { WidgetSyncPayload } from "../../api/widget/sync/route";
import { ActiveJobsLive } from "../../jobs/components/ActiveJobsLive";
import { PauseTranscriptionsButton } from "../../jobs/components/PauseTranscriptionsButton";
import { PauseDownloadsButton } from "../../jobs/components/PauseDownloadsButton";
+import { DigestSweepControls } from "../../jobs/components/DigestSweepControls";
import { syncAllChannelsAction, type SyncAllResult } from "../../channels/actions";
import { fmtTime } from "../../widget/lib/relativeTime";
@@ -37,6 +38,13 @@ export function PipelineBand({
const busy = workerList.filter((w) => w.busy).length;
const paused = workers?.paused ?? false;
const downloadsPaused = workers?.downloadsPaused ?? false;
+ const digest = sync?.digest ?? null;
+ // Deliberately not rounded up. At 0.13% a "1%" would be a lie of the kind
+ // that makes an 80-day backfill look nearly begun.
+ const digestPct =
+ digest && digest.videos > 0
+ ? (digest.digested / digest.videos) * 100
+ : null;
const lastSyncText =
sync == null
@@ -112,6 +120,36 @@ export function PipelineBand({
label={<span className="text-destructive">downloads paused</span>}
/>
)}
+ {digest && (
+ <Instrument
+ dotClass={
+ digest.paused
+ ? "bg-destructive"
+ : digest.sweeping
+ ? "bg-success"
+ : "bg-muted-foreground/40"
+ }
+ label={
+ <>
+ ◆ digests{" "}
+ <span className="font-medium">
+ {digest.digested.toLocaleString()}
+ </span>
+ {digestPct !== null && (
+ <span className="text-muted-foreground">
+ {" "}
+ ({digestPct < 1 ? digestPct.toFixed(2) : digestPct.toFixed(1)}%
+ of {digest.videos.toLocaleString()})
+ </span>
+ )}
+ {digest.paused && (
+ <span className="text-destructive"> · paused</span>
+ )}
+ {!digest.paused && digest.sweeping && " · sweeping"}
+ </>
+ }
+ />
+ )}
</div>
</div>
@@ -121,6 +159,11 @@ export function PipelineBand({
paused={downloadsPaused}
onChange={onWorkersChange}
/>
+ <DigestSweepControls
+ sweeping={digest?.sweeping ?? false}
+ paused={digest?.paused ?? false}
+ onChange={onSynced}
+ />
<SyncAllButton onSynced={onSynced} />
<Link
href="/channels/new"
diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts
@@ -8,6 +8,10 @@ import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta";
import type { JobSpec } from "yt-dlp-transcript-common/jobs/jobSpec";
import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand";
+import {
+ startDigestSweep,
+ stopDigestSweep,
+} from "yt-dlp-transcript-common/controller/digestSweep";
import { runJobSpec } from "./runJobSpec";
import { buildQueueView } from "./queue/buildQueueView";
@@ -169,6 +173,73 @@ export async function resumeDownloadsAction(): Promise<DownloadsPauseResult> {
return setDownloadsPaused(false);
}
+// Global, restart-surviving digest pause — the counterpart of the downloads one
+// and, until now, a flag with no writer. `digestBatch`'s limit() has always
+// re-read `digestsPaused` at dispatch time and returned 0 to make the pool
+// idle-wait (a real pause: the job HOLDS rather than ending, so nothing has to
+// be re-derived on resume), and settings/actions.ts passed the field through
+// untouched with a comment saying "the dashboard/channel controls own the
+// pause". Those controls did not exist, so the pause could not be set from
+// anywhere. This is them.
+//
+// It needs no boot hook, for the same reason downloadsPaused doesn't: the flag
+// is consulted at dispatch, not applied to a live pool.
+export type DigestPauseResult = { ok: boolean; error?: string };
+
+async function setDigestsPaused(paused: boolean): Promise<DigestPauseResult> {
+ try {
+ const current = getSettings();
+ if (current.digest.digestsPaused !== paused) {
+ await writeSettings({
+ ...current,
+ digest: { ...current.digest, digestsPaused: paused },
+ });
+ }
+ } catch (e) {
+ return { ok: false, error: (e as Error).message };
+ }
+ revalidatePath("/");
+ revalidatePath("/jobs");
+ return { ok: true };
+}
+
+export async function pauseDigestsAction(): Promise<DigestPauseResult> {
+ return setDigestsPaused(true);
+}
+
+export async function resumeDigestsAction(): Promise<DigestPauseResult> {
+ return setDigestsPaused(false);
+}
+
+// Arm / disarm the corpus-wide sweep. Distinct from the pause: a pause holds a
+// running sweep at zero throughput, this decides whether there is a sweep at
+// all — and it persists, so the boot hook resumes it.
+export type DigestSweepResult = { ok: boolean; jobId?: string; error?: string };
+
+export async function startDigestSweepAction(): Promise<DigestSweepResult> {
+ try {
+ const jobId = await startDigestSweep();
+ revalidatePath("/");
+ revalidatePath("/jobs");
+ return jobId
+ ? { ok: true, jobId }
+ : { ok: false, error: "The sweep could not be started (see job logs)." };
+ } catch (e) {
+ return { ok: false, error: (e as Error).message };
+ }
+}
+
+export async function stopDigestSweepAction(): Promise<DigestSweepResult> {
+ try {
+ stopDigestSweep();
+ revalidatePath("/");
+ revalidatePath("/jobs");
+ return { ok: true };
+ } catch (e) {
+ return { ok: false, error: (e as Error).message };
+ }
+}
+
// Retention scopes offered by the ClearLogsMenu. "all" clears every finished
// job's log; the day-scopes clear anything older than that. Running/queued jobs
// are never deleted (see pruneJobLogs).
diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts
@@ -39,12 +39,27 @@ function computeEtaSeconds(
job: JobRecord,
remaining: number,
now: number,
+ remainingAudioSeconds?: number,
): number | undefined {
const count = job.completedTaskCount ?? 0;
const totalMs = job.completedTaskMs ?? 0;
if (count < 1 || remaining <= 0 || job.startedAt === undefined) {
return undefined;
}
+ // Prefer an AUDIO-HOUR estimate where the work is proportional to length.
+ // Averaging tasks assumes every unit costs about the same, which is true for
+ // downloads and wildly false for digests: this corpus is ~77k videos and ~77k
+ // audio-hours, and a channel of 9-hour VODs and a channel of 10-minute clips
+ // have the same task count and a 50x difference in cost. A task average would
+ // therefore quote an ETA that is wrong by more than an order of magnitude at
+ // exactly the moment an operator most needs it — the start of an 80-day run.
+ const doneAudio = job.completedTaskAudioSeconds ?? 0;
+ if (doneAudio > 0 && remainingAudioSeconds && remainingAudioSeconds > 0) {
+ const secondsPerAudioSecond = totalMs / 1000 / doneAudio;
+ const elapsedMs = Math.max(1, now - job.startedAt);
+ const concurrency = Math.max(1, totalMs / elapsedMs);
+ return (remainingAudioSeconds * secondsPerAudioSecond) / concurrency;
+ }
const avgProcMs = totalMs / count; // measured average per task
const elapsedMs = Math.max(1, now - job.startedAt);
const concurrency = Math.max(1, totalMs / elapsedMs); // effective parallelism
@@ -65,12 +80,18 @@ function computeJobProgressView(
// `current` is RE-COUNTED from disk (readChannelStat), never reported by the
// runner — which is why the digest metric needed its own on-disk counter
// (digestCount) rather than a number the batch could have just told us.
+ // A runner-reported `current` WINS. The disk re-count below cannot see a
+ // regeneration — a regenerated digest is rewritten in place, so the file
+ // count never moves and the bar sits at 0% for the whole job. Only the runner
+ // knows it did the work. Downloads and transcripts report nothing and keep
+ // the disk re-count, unchanged.
const current =
- snap.metric === "downloads"
+ snap.current ??
+ (snap.metric === "downloads"
? stat.downloadCount
: snap.metric === "digests"
? (stat.digestCount ?? 0)
- : stat.transcriptCount;
+ : stat.transcriptCount);
const range = Math.max(0, snap.target - snap.initial);
const advance = Math.max(0, current - snap.initial);
const pct =
@@ -82,7 +103,12 @@ function computeJobProgressView(
current,
target: snap.target,
pct,
- etaSeconds: computeEtaSeconds(job, remaining, now),
+ etaSeconds: computeEtaSeconds(
+ job,
+ remaining,
+ now,
+ snap.remainingAudioSeconds,
+ ),
};
}
diff --git a/editor/app/jobs/components/DigestSweepControls.tsx b/editor/app/jobs/components/DigestSweepControls.tsx
@@ -0,0 +1,104 @@
+"use client";
+
+import { useEffect, useState, useTransition } from "react";
+import { useRouter } from "next/navigation";
+import {
+ pauseDigestsAction,
+ resumeDigestsAction,
+ startDigestSweepAction,
+ stopDigestSweepAction,
+} from "../actions";
+
+// The two digest controls, side by side, because they are NOT the same control
+// and conflating them is how an operator loses a week of GPU time:
+//
+// Sweep on/off — is there a corpus-wide backfill at all. Persisted, so a
+// server restart resumes it (editor/instrumentation.ts).
+// Pause/resume — hold a running sweep at zero throughput without ending it.
+// The batch's limit() returns 0, which makes the pool idle-WAIT rather than
+// finish, so resuming costs nothing and re-derives nothing.
+//
+// Stopping the sweep drains rather than cancels: the channel in flight finishes
+// instead of losing a part-generated video.
+export function DigestSweepControls({
+ sweeping,
+ paused,
+ onChange,
+}: {
+ sweeping: boolean;
+ paused: boolean;
+ onChange?: () => void | Promise<void>;
+}) {
+ const [pending, startTransition] = useTransition();
+ const [error, setError] = useState<string | null>(null);
+ const router = useRouter();
+ // Disabled until hydrated. A click on a server-rendered button before React
+ // attaches fires nothing at all — no request, no job, no error — which is the
+ // recorded root cause of the digest pilot's "un-created job".
+ const [mounted, setMounted] = useState(false);
+ useEffect(() => setMounted(true), []);
+
+ function run(fn: () => Promise<{ ok: boolean; error?: string }>) {
+ setError(null);
+ startTransition(async () => {
+ const result = await fn();
+ if (!result.ok) setError(result.error ?? "Failed.");
+ if (onChange) await onChange();
+ else router.refresh();
+ });
+ }
+
+ const disabled = pending || !mounted;
+
+ return (
+ <div className="flex items-center gap-2">
+ <button
+ type="button"
+ disabled={disabled}
+ aria-label={sweeping ? "stop digest sweep" : "start digest sweep"}
+ onClick={() =>
+ run(sweeping ? stopDigestSweepAction : startDigestSweepAction)
+ }
+ title={
+ sweeping
+ ? "Stop the corpus-wide digest backfill. The channel in flight finishes first; nothing already generated is lost."
+ : "Start the corpus-wide digest backfill: every channel in turn, heaviest first by remaining audio-hours. Survives a restart."
+ }
+ className={
+ sweeping
+ ? "px-3 py-1.5 rounded-md bg-warning text-warning-foreground text-sm font-medium hover:bg-warning/90 disabled:opacity-50"
+ : "px-3 py-1.5 rounded-md border border-border text-sm hover:bg-muted disabled:opacity-50"
+ }
+ >
+ {sweeping ? "Stop Digest Sweep" : "Start Digest Sweep"}
+ </button>
+ {(sweeping || paused) && (
+ <button
+ type="button"
+ disabled={disabled}
+ aria-label={paused ? "resume digests" : "pause digests"}
+ onClick={() =>
+ run(paused ? resumeDigestsAction : pauseDigestsAction)
+ }
+ title={
+ paused
+ ? "Resume digest generation. The running job picks up where it left off — it was holding, not stopped."
+ : "Hold digest generation without ending the sweep. The running job idles at zero and resumes instantly."
+ }
+ className={
+ paused
+ ? "px-3 py-1.5 rounded-md bg-warning text-warning-foreground text-sm font-medium hover:bg-warning/90 disabled:opacity-50"
+ : "px-3 py-1.5 rounded-md border border-border text-sm hover:bg-muted disabled:opacity-50"
+ }
+ >
+ {paused ? "Resume Digests" : "Pause Digests"}
+ </button>
+ )}
+ {error && (
+ <span role="alert" className="text-xs text-destructive">
+ {error}
+ </span>
+ )}
+ </div>
+ );
+}
diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts
@@ -46,4 +46,24 @@ export async function register() {
} catch {
/* a runner that fails to start must not block server readiness */
}
+
+ // Resume the corpus-wide digest sweep, if one is armed. A sweep is GPU-WEEKS
+ // long, so it will outlive several restarts by construction — and before this
+ // hook a restart silently ended one that had been running for days, with
+ // nothing to say so.
+ //
+ // Re-launching is safe rather than merely convenient: the sweep stores no
+ // cursor and the batch re-derives eligibility from disk on every pull, so a
+ // resumed sweep re-does exactly zero work. The digest PAUSE needs no hook
+ // here for the same reason downloadsPaused doesn't — it is read at dispatch
+ // time — but "is a sweep running" is process state, and process state is what
+ // a restart destroys.
+ try {
+ const { resumeDigestSweepIfEnabled } = await import(
+ "yt-dlp-transcript-common/controller/digestSweep"
+ );
+ await resumeDigestSweepIfEnabled();
+ } catch {
+ /* a sweep that fails to resume must not block server readiness */
+ }
}