Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 534fd1c322e3595cc733f44de23d577671effad3
parent 07705a4c73333d6c578fe002a9aa22c20ae2ce20
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  7 Aug 2026 21:46:57 -0400

Make backfill a first-class thing the system knows about

Every derived-data feature that lands hits the same wall: the corpus that
already exists does not have what it needs. Diarization hit it first — its
input is AUDIO, which cleanAudioFromTranscribed deletes once a video is
transcribed — and the only catch-up was a per-channel button on a controller
written for that one feature. There was no way to ask "how much of the corpus
is missing this?", and no way to run catch-up alongside new-video work without
one starving the other.

So, in the shape jobKinds.ts already uses: a feature declares ONE entry in
lib/backfillKinds.ts, and the system supplies the lane, the resource share and
the indicator. Diarization is the first entry.

FOUR STATES, NOT TWO, and the split is the load-bearing part. Measured here:
77,106 videos, 836 with media still on disk, 1 diarized. A single "remaining"
number would read 77,105, and 91x of that is unreachable without a
re-download. The repo has already been burned by exactly this once —
api/widget/actionable/route.ts refuses to filter on `noDigest` because during
the backfill that is 99.87% of the corpus. So `missing` (reachable now) and
`missing-input` (needs re-acquiring) are separate numbers everywhere, and no
surface sums them.

STALENESS IS PROVENANCE, NOT AGE. isDiarizationFresh compares the recorded
models and threshold against what settings WOULD produce now, with an absent
field reading as today's default so adding one does not invalidate the corpus
(lib/digest.ts's isSectionFresh trick). Before this, diarizeOne short-circuited
on mere existence, so a threshold or model change — the two most likely reasons
to re-run — left everything looking done. `version` is deliberately excluded: an
engine point-release must not invalidate ~500-680 s/audio-hour of captured work.

The lane gets its own queueKey, so it is concurrent with transcription rather
than behind it, and backfill.weight is a SHARE not a priority: 0 (the default)
is idle-only, standing aside the moment transcription has work, and >0 is a
guaranteed share floored at 1 so a small weight is a slow lane, not a stopped
one. The yield deliberately watches only the transcription lane — two yielding
lanes can park each other.

Re-acquiring deleted media (backfillReacquire.ts) is OFF by default and is the
part to read carefully: each file is fetched, used, and removed in a `finally`
whether the backfill succeeded, failed or crashed, unless the video is marked
do-not-clean, and nothing starts under the disk floor. On a disk at 97% a leak
here fills it, which is why those four behaviours have e2e tests.

backfillSweep.ts is cloned from digestSweep.ts, keeping its four load-bearing
shapes verbatim: launch the per-channel job (not 77,000 jobs the registry would
evict), store no cursor, persist intent WITH scope, and run the orchestrator on
queueKey "" so it cannot deadlock waiting on the queue it holds.

Indicators are generic and iterate snapshot.backfill[kindId]: the channel
Backfill stage card, an /actionable section, a dashboard instrument and an
optional widget strip. Registering a second kind lights all four with no UI
work — which is the claim the registry has to earn next.

Verification: common 530/531, tsc clean in common and editor, editor e2e 442
passed. The one common failure is the pre-existing digest corpus census, which
asserts an absolute video count that moves whenever a video is downloaded
(191,116 -> 194,053); it reproduces on a pristine HEAD.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

Diffstat:
Acommon/controller/backfillBatch.test.ts | 86+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/backfillBatch.ts | 0
Acommon/controller/backfillReacquire.ts | 203+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/backfillSweep.ts | 416+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/channelSnapshot.ts | 65+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/diarizeOne.ts | 25++++++++++++++++++++-----
Mcommon/jobs/jobKinds.ts | 23+++++++++++++++++++++++
Mcommon/jobs/registry.ts | 8++++++--
Mcommon/jobs/snapshotScheduler.ts | 8++++++++
Acommon/lib/backfillKinds.test.ts | 327+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/backfillKinds.ts | 288+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/diarization.ts | 112+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/queueKeys.ts | 12++++++++++++
Mcommon/lib/settings.ts | 106+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/lib/videoStatus.ts | 11+++++++++++
Meditor/CHANGELOG.md | 1+
Meditor/app/actionable/components/InlineActionButton.tsx | 6++++++
Meditor/app/actionable/lib/loadActionable.ts | 27+++++++++++++++++++++++++++
Meditor/app/actionable/page.tsx | 33++++++++++++++++++++++++++++++++-
Aeditor/app/api/test/resume-backfill-sweep/route.ts | 35+++++++++++++++++++++++++++++++++++
Meditor/app/api/widget/sync/route.ts | 37+++++++++++++++++++++++++++++++++++++
Aeditor/app/channels/[slug]/backfillActions.ts | 65+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/channels/[slug]/components/stages/BackfillStage.tsx | 133+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/channels/[slug]/lib/stageStatus.ts | 66++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/channels/[slug]/page.tsx | 29+++++++++++++++++++++++++++++
Meditor/app/components/dashboard/PipelineBand.tsx | 39+++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/actions.ts | 38++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/BackfillSweepControls.tsx | 96+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/components/RunningJobsList.tsx | 4++++
Meditor/app/jobs/jobReplayRegistry.ts | 12++++++++++++
Meditor/app/settings/actions.ts | 19+++++++++++++++++++
Meditor/app/settings/components/SettingsForm.tsx | 67+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/widget/builder/components/SectionMini.tsx | 12++++++++++++
Meditor/app/widget/components/MonitorWidget.tsx | 40++++++++++++++++++++++++++++++++++++++++
Meditor/app/widget/lib/config.ts | 9+++++++++
Meditor/app/widget/lib/sections.ts | 17++++++++++++++++-
Aeditor/e2e/backfill.spec.ts | 586+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/instrumentation.ts | 14++++++++++++++
38 files changed, 3064 insertions(+), 11 deletions(-)

diff --git a/common/controller/backfillBatch.test.ts b/common/controller/backfillBatch.test.ts @@ -0,0 +1,86 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { backfillLimit } from "./backfillBatch"; +import { defaultBackfill, sanitizeBackfill } from "../lib/settings"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/backfillBatch.test.ts +// +// The PURE function, not a live pool — same discipline as digestYield.test.ts, +// and for the same reason: this decides whether a multi-day lane runs at all, +// and a test that needs a GPU, a registry and a worker pool is a test nobody +// runs before shipping. + +test("weight 0 is idle-only: full slots when the primary lane is quiet", () => { + assert.equal( + backfillLimit({ weight: 0, slots: 4, primaryBusy: false }), + 4, + ); + assert.equal(backfillLimit({ weight: 0, slots: 1, primaryBusy: false }), 1); +}); + +test("weight 0 stands aside completely while the primary lane works", () => { + // Zero is a HOLD, not a stop: runPool idle-waits at a zero limit rather than + // finishing, so the lane resumes the moment transcription is free without + // re-deriving anything. That is the whole reason this is a limit and not a + // scheduler. + assert.equal(backfillLimit({ weight: 0, slots: 4, primaryBusy: true }), 0); + assert.equal(backfillLimit({ weight: 0, slots: 16, primaryBusy: true }), 0); +}); + +test("a positive weight is a guaranteed share, busy or not", () => { + // The point of a non-zero weight: the primary lane being busy no longer parks + // the backfill. + assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: true }), 2); + assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: false }), 2); + assert.equal(backfillLimit({ weight: 1, slots: 4, primaryBusy: true }), 4); +}); + +test("a small weight is a SLOW lane, never a stopped one", () => { + // floor(4 * 0.25) is 1, and floor(1 * 0.25) is 0 — which without the floor + // would silently turn "a quarter of the machine" into "never runs", and look + // exactly like a wedge. + assert.equal(backfillLimit({ weight: 0.25, slots: 4, primaryBusy: true }), 1); + assert.equal(backfillLimit({ weight: 0.25, slots: 1, primaryBusy: true }), 1); + assert.equal(backfillLimit({ weight: 0.01, slots: 8, primaryBusy: true }), 1); +}); + +test("no slots means no work, whatever the weight", () => { + assert.equal(backfillLimit({ weight: 0, slots: 0, primaryBusy: false }), 0); + assert.equal(backfillLimit({ weight: 1, slots: 0, primaryBusy: false }), 0); +}); + +test("the shipped default is idle-only", () => { + // Stated as a test because it is a promise the feature makes: catch-up work on + // a corpus that already exists must never slow down new arrivals. + const d = defaultBackfill(); + assert.equal(d.weight, 0); + assert.equal(d.enabled, false); + assert.equal(d.allowRedownload, false); + assert.equal(backfillLimit({ ...d, slots: d.concurrency, primaryBusy: true }), 0); +}); + +test("a hand-edited weight is clamped rather than rejected", () => { + // settings.json is hand-editable. A 5 means "as much as possible", and reading + // it as the idle-only 0 would be the opposite of the intent. + assert.equal(sanitizeBackfill({ weight: 5 }).weight, 1); + assert.equal(sanitizeBackfill({ weight: -3 }).weight, 0); + assert.equal(sanitizeBackfill({ weight: "half" }).weight, 0); + assert.equal(sanitizeBackfill({ weight: 0.25 }).weight, 0.25); + // Concurrency shares clampPositiveInt's floor of 1, so a 0 cannot silently + // park the lane either. + assert.equal(sanitizeBackfill({ concurrency: 0 }).concurrency, 1); +}); + +test("a stale sweep scope survives sanitization as a list of slugs", () => { + // The scope is persisted WITH the flag, so it has to round-trip: junk entries + // are dropped, real ones kept. + assert.deepEqual( + sanitizeBackfill({ sweepChannels: ["a", "", 7, "b"], sweepKinds: ["diarization"] }), + { + ...defaultBackfill(), + sweepChannels: ["a", "b"], + sweepKinds: ["diarization"], + }, + ); +}); diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts Binary files differ. diff --git a/common/controller/backfillReacquire.ts b/common/controller/backfillReacquire.ts @@ -0,0 +1,203 @@ +// Re-acquire the media a backfill needs for a video whose input is GONE. +// +// This is the piece with no precedent anywhere in the repo, and it is the +// highest-risk code in the feature: it writes media onto a disk that is 97% full +// with ~45 GB free, for a population measured at ~76,270 videos. Everything +// below is arranged around one property — A RE-FETCHED FILE DOES NOT SURVIVE THE +// ITEM THAT FETCHED IT. +// +// Four guards, in the order they fire: +// +// 1. OPT-IN. The caller only reaches here when settings.backfill.allowRedownload +// (or an explicit per-run flag) is set. Default off. +// 2. DISK FLOOR, per item. checkDiskSpace against settings.minFreeDiskGB before +// each fetch, not once at the start: a long run's twentieth video must not +// inherit the first one's headroom. Under the floor is a refusal to START, +// reported as its own outcome so the batch can stop rather than fail 76,000 +// times. +// 3. NOTHING PERMANENTLY GONE. isPermanentlyGone / EXCLUDED_FROM_DOWNLOAD, +// verbatim, so a deleted or members-only video is not re-attempted forever. +// 4. REMOVAL IN A `finally`, ALWAYS. cleanup() removes exactly the files this +// fetch created — diffed against a listing taken BEFORE it ran, so nothing +// that was already on disk can ever be deleted by it — and the caller runs +// it whether the backfill succeeded, failed or threw. +// +// The one exception to (4) is a video marked do-not-clean. That marker is the +// operator saying "this media is archived, keep it", and it is honoured here for +// the same reason every cleanup controller honours it. cleanup() reports that it +// kept the file so the batch can say so in its log rather than leaving it to be +// discovered as a mystery on a full disk. + +import path from "node:path"; +import { readdir, rm } from "node:fs/promises"; +import type { Paths } from "../lib/paths"; +import { getSettings } from "../lib/settings"; +import { checkDiskSpace } from "../lib/diskSpace"; +import { formatBytes } from "../lib/format"; +import { isPermanentlyGone } from "../lib/availability"; +import { resolveEffectiveAvailability } from "../lib/availability-server"; +import { isDoNotClean } from "../lib/doNotClean-server"; +import { resolveCookiePolicy } from "../lib/cookiePolicy"; +import { isRealAudioFile, isSourceMediaFile } from "../lib/videoStatus"; +import { readChannelConfig } from "./channels"; +import { findVideoSourceUrl } from "./undownloadedVideos"; +import { downloadOneManaged } from "../ytdlp/downloadOneManaged"; +import { resolveDiarizableMedia } from "./diarizeOne"; + +export type ReacquireOutcome = + // Media was already there — nothing fetched, nothing to clean up. + | { status: "present"; cleanup: () => Promise<boolean> } + // Fetched. `cleanup()` removes it again and returns whether it did (false = + // deliberately kept because the video is marked do-not-clean). + | { status: "fetched"; cleanup: () => Promise<boolean> } + // Free disk is under the configured floor. A refusal to start, not a failure. + | { status: "disk-floor"; cleanup: () => Promise<boolean> } + // Deleted / members-only / private, or no resolvable source URL. Re-trying + // this video will not help. + | { status: "gone"; cleanup: () => Promise<boolean> } + | { status: "failed"; cleanup: () => Promise<boolean> }; + +const NOTHING_TO_CLEAN = async (): Promise<boolean> => false; + +export async function reacquireMediaFor(opts: { + paths: Paths; + channelSlug: string; + videoId: string; + videoDir: string; + onLog?: (msg: string) => void; + signal?: AbortSignal; +}): Promise<ReacquireOutcome> { + const log = opts.onLog ?? ((m: string) => console.log(m)); + + // Re-checked HERE, not trusted from the caller's pull: the classification that + // sent us here may be minutes old (the pool holds items at a zero limit), and + // a video that got its audio back in the meantime must not be re-fetched. + if (await resolveDiarizableMedia(opts.videoDir)) { + return { status: "present", cleanup: NOTHING_TO_CLEAN }; + } + + const settings = getSettings(); + const disk = await checkDiskSpace(opts.paths, settings); + if (!disk.ok) { + log( + `Not re-acquiring ${opts.videoId}: only ${formatBytes(disk.freeBytes)} free, ` + + `below the ${formatBytes(disk.thresholdBytes)} floor.`, + ); + return { status: "disk-floor", cleanup: NOTHING_TO_CLEAN }; + } + + const availability = await resolveEffectiveAvailability(opts.videoDir); + if (isPermanentlyGone(availability)) { + log(`Not re-acquiring ${opts.videoId}: ${availability}.`); + return { status: "gone", cleanup: NOTHING_TO_CLEAN }; + } + + const config = await readChannelConfig(opts.paths, opts.channelSlug); + if (!config) { + log(`Not re-acquiring ${opts.videoId}: channel config unreadable.`); + return { status: "failed", cleanup: NOTHING_TO_CLEAN }; + } + const url = await findVideoSourceUrl( + opts.paths, + opts.channelSlug, + opts.videoId, + config, + ); + if (!url) { + log(`Not re-acquiring ${opts.videoId}: no resolvable source URL.`); + return { status: "gone", cleanup: NOTHING_TO_CLEAN }; + } + + // The listing BEFORE the fetch. cleanup() removes only what is not in it, so + // this function can never delete media that was already on disk — including + // anything a concurrent lane wrote. + const before = new Set( + await readdir(opts.videoDir).catch(() => [] as string[]), + ); + + log(`Re-acquiring media for ${opts.videoId} (${url})…`); + try { + await downloadOneManaged({ + channelSlug: opts.channelSlug, + channelConfig: config, + paths: opts.paths, + videoUrl: url, + onLog: (s) => log(s.trimEnd()), + signal: opts.signal ?? new AbortController().signal, + cookiePolicy: resolveCookiePolicy(settings, config), + inlineTranscribeOnFallback: false, + globalSkipLiveDownloads: settings.skipLiveDownloads, + // The archive already lists this video — it was downloaded once. A second + // line would just be noise. + appendArchive: false, + // Audio now, container discarded: a backfill wants the smallest thing that + // can be read, and this is the flag the save-disk path already uses. + extractImmediately: true, + keepSourceVideoOverride: false, + }); + } catch (err) { + log(`Re-acquire ${opts.videoId} failed: ${(err as Error).message}`); + // Fall through to the cleanup builder anyway: a failed download can still + // have left a partial file, and that is exactly the leak this exists to stop. + return { + status: "failed", + cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log), + }; + } + + // Trust the disk, not the exit code. + if (!(await resolveDiarizableMedia(opts.videoDir))) { + log(`Re-acquire ${opts.videoId}: nothing usable landed.`); + return { + status: "failed", + cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log), + }; + } + return { + status: "fetched", + cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log), + }; +} + +// Remove exactly the media this fetch added. Returns true when it removed +// something, false when it deliberately kept it (do-not-clean) or there was +// nothing new. +// +// Restricted to REAL AUDIO and SOURCE-MEDIA files plus their partials, never +// "everything new in the dir": a download also writes metadata, a download log +// and an outcome sidecar, and those are records of what happened that outlive +// the media on purpose. +function buildCleanup( + videoDir: string, + before: Set<string>, + videoId: string, + log: (msg: string) => void, +): () => Promise<boolean> { + return async () => { + const after = await readdir(videoDir).catch(() => [] as string[]); + const added = after.filter( + (name) => + !before.has(name) && + (isRealAudioFile(name) || + isSourceMediaFile(name) || + name.startsWith("audio.") || + name.startsWith("source-media.")), + ); + if (added.length === 0) return false; + if (await isDoNotClean(videoDir)) { + log( + `Keeping re-acquired media for ${videoId}: marked "do not clean" (${added.join(", ")}).`, + ); + return false; + } + for (const name of added) { + await rm(path.join(videoDir, name), { force: true }).catch((err) => { + // Reported, never thrown: this runs in a `finally`, and throwing here + // would replace the real outcome of the item with a cleanup error. + log(`Could not remove ${name} for ${videoId}: ${(err as Error).message}`); + }); + } + log(`Removed re-acquired media for ${videoId} (${added.join(", ")}).`); + return true; + }; +} diff --git a/common/controller/backfillSweep.ts b/common/controller/backfillSweep.ts @@ -0,0 +1,416 @@ +// The corpus-wide backfill sweep: one launcher for the catch-up that runs for +// days, and survives the restarts it will certainly outlive. +// +// Cloned from digestSweep.ts, which is the complete worked example of +// armed-and-resumable in this repo. Four of its shapes are load-bearing and are +// kept verbatim, each because getting it wrong has already cost something here: +// +// 1. IT LAUNCHES THE EXISTING PER-CHANNEL JOB, in sequence. Not one job per +// video: the registry keeps 100 records and the log 500, so 77,000 jobs +// would evict the history of the very run recording them. +// 2. IT STORES NO CURSOR. The work-list is recomputed each pass and +// eligibility is re-derived from disk inside every batch, so a restart, a +// newly-finished transcription and a settings change that invalidates the +// corpus are all just visible on the next pass. +// 3. IT PERSISTS INTENT AND SCOPE TOGETHER. The boot hook re-launches from +// settings alone, so arming without recording the scope resurrects a +// deliberately-bounded run as a corpus-wide one — and a stop that does not +// CLEAR the scope leaves a stale one to silently narrow the next sweep. +// Both were real bugs in the digest sweep. +// 4. THE ORCHESTRATOR RUNS ON queueKey "". It waits on the per-channel jobs +// that serialize on BACKFILL_QUEUE; putting it on that queue would deadlock +// — it would hold the only slot while waiting for a job that needs it. +// +// The termination guard is measured against REMAINING WORK, not against jobs +// starting. An engine that is broken makes every channel job start normally and +// then do nothing, which from out here looks exactly like progress. + +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 { BACKFILL_QUEUE } from "../lib/queueKeys"; +import { resolveBackfillKinds } from "../lib/backfillKinds"; +import { listChannelStatsFromDisk } from "./channels"; +import { countBackfillWork, runBackfillBatch } from "./backfillBatch"; + +export const BACKFILL_SWEEP_KIND = "backfill-sweep"; +export const BACKFILL_CHANNEL_KIND = "backfill-channel"; + +// How long the loop waits before re-planning when a pass did no work. Generous: +// nothing here is latency-sensitive and planning reads every channel's videos. +const IDLE_POLL_MS = 60_000; + +// Give up after this many consecutive passes that move no work. Three, not one: +// a single barren pass is legitimate (everything remaining failed a guard this +// time and may not next time), but three in a row is a broken engine or a corpus +// the sweep cannot make progress on, and both want a human. +const MAX_BARREN_PASSES = 3; + +type SweepLive = { + jobId: string; + startedAt: number; + // The per-channel job the sweep is currently waiting on. Tracked so a stop can + // drain THAT too — without it, "stop" means "after the current channel", which + // on a large channel is hours. + channelJobId: string | null; +}; +type SweepSingleton = { live: SweepLive | null }; + +declare global { + // eslint-disable-next-line no-var + var __yttBackfillSweep__: SweepSingleton | undefined; +} + +function getSingleton(): SweepSingleton { + if (!globalThis.__yttBackfillSweep__) { + globalThis.__yttBackfillSweep__ = { live: null }; + } + return globalThis.__yttBackfillSweep__; +} + +export function getBackfillSweepJobId(): string | null { + const live = getSingleton().live; + if (!live) return null; + return getRegistry().get(live.jobId)?.status === "running" ? live.jobId : null; +} + +export type BackfillSweepOptions = { + paths?: Paths; + // Restrict to these kinds / channels. Absent = every enabled lane kind, whole + // corpus. + kindIds?: string[]; + channelSlugs?: string[]; +}; + +export type BackfillPlanEntry = { + channelSlug: string; + reachable: number; + missingInput: number; +}; + +// What is left, per channel, ordered heaviest-first by REACHABLE work. +// +// Reachable only, deliberately: ordering by total remaining would sort the +// corpus by how much of it is unreachable, which on the measured numbers (835 +// reachable vs ~76,270 needing a re-download) means sorting by noise. +export async function buildBackfillSweepPlan(opts: { + paths: Paths; + kindIds?: string[]; + channelSlugs?: string[]; + onLog?: (msg: string) => void; +}): Promise<BackfillPlanEntry[]> { + const wanted = new Set(opts.channelSlugs ?? []); + const channels = (await listChannelStatsFromDisk(opts.paths)).filter( + (ch) => wanted.size === 0 || wanted.has(ch.slug), + ); + const plan: BackfillPlanEntry[] = []; + for (const ch of channels) { + const counts = await countBackfillWork(opts.paths, ch.slug, opts.kindIds); + plan.push({ + channelSlug: ch.slug, + reachable: counts.reachable, + missingInput: counts.missingInput, + }); + } + return plan.sort( + (a, b) => b.reachable - a.reachable || a.channelSlug.localeCompare(b.channelSlug), + ); +} + +// Run ONE channel through the batch, as its own managed job — same kind, same +// queue and same progress metric a hand-clicked run produces, so a sweep is +// inspectable with the tools that already exist rather than being an opaque +// mega-job. +export async function runBackfillChannelJob(opts: { + paths: Paths; + channelSlug: string; + kindIds?: string[]; + queueKey?: string; + background?: boolean; + onStarted?: (jobId: string) => void; +}): Promise<{ ok: boolean; error?: string }> { + const { paths, channelSlug } = opts; + const result = await runManagedFunction({ + kind: BACKFILL_CHANNEL_KIND, + queueKey: opts.queueKey ?? BACKFILL_QUEUE, + paths, + channelSlug, + background: opts.background, + spec: { + kind: BACKFILL_CHANNEL_KIND, + slug: channelSlug, + params: { kindIds: opts.kindIds, queueKey: opts.queueKey }, + }, + fn: async (onLog, signal, setProgress, ctx) => { + const before = await countBackfillWork(paths, channelSlug, opts.kindIds); + const batch = await runBackfillBatch({ + channelSlug, + paths, + kindIds: opts.kindIds, + setProgress, + // The bar measures THIS run, from zero. Seeding it with what is already + // on disk cannot represent a regeneration, where no file count moves. + progressBaseline: 0, + progressTarget: before.reachable, + onLog, + signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), + }); + onLog( + `Backfill ${channelSlug}: ${batch.succeeded} done, ${batch.fresh} already current, ` + + `${batch.failed} failed; ${batch.missingInput} still need their media re-acquired` + + (batch.reacquired > 0 + ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` + : "") + + (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + + ".", + ); + requestChannelSnapshot(paths, channelSlug); + }, + }); + if (!result.ok) return { ok: false, error: result.error }; + opts.onStarted?.(result.jobId); + // Wait for the channel to finish: the sweep is sequential by design, 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, + kindIds: string[] | undefined, + channelSlugs: string[] | undefined, + live: SweepLive, + onLog: (msg: string) => void, + signal: AbortSignal, + drainSignal: AbortSignal, +): Promise<void> { + let pass = 0; + // Consecutive passes that left the remaining work UNCHANGED — see the header + // for why this measures the plan and not whether jobs started. + let barrenPasses = 0; + let lastReachable: number | null = null; + for (;;) { + if (signal.aborted || drainSignal.aborted) return; + // The operator turned it off: stop CLEANLY rather than being cancelled, so + // the job ends "done" and the queue is released. + if (!getSettings().backfill.sweepEnabled) { + onLog("Backfill sweep disarmed in settings — stopping."); + return; + } + // The feature itself may have been switched off underneath us. Same clean + // stop: a sweep for a disabled backfill has nothing to do and should not + // sit there looking wedged. + if (resolveBackfillKinds(getSettings(), kindIds).length === 0) { + onLog( + "No enabled backfill kind is in scope — stopping. (Enable the feature in Settings, then re-arm.)", + ); + return; + } + + pass++; + const plan = await buildBackfillSweepPlan({ + paths, + kindIds, + channelSlugs, + onLog, + }); + const work = plan.filter((c) => c.reachable > 0); + const reachable = plan.reduce((n, c) => n + c.reachable, 0); + const missingInput = plan.reduce((n, c) => n + c.missingInput, 0); + if (work.length === 0) { + onLog( + `Pass ${pass}: nothing reachable left to backfill. ` + + `${missingInput.toLocaleString()} video(s) would need their media re-acquired first.`, + ); + return; + } + // BOTH NUMBERS, never their sum. See lib/backfillKinds.ts's header. + onLog( + `Pass ${pass}: ${work.length} channel(s), ${reachable.toLocaleString()} video(s) ` + + `reachable now, ${missingInput.toLocaleString()} needing their media re-acquired.`, + ); + + let didWork = false; + for (const channel of work) { + if (signal.aborted || drainSignal.aborted) return; + if (!getSettings().backfill.sweepEnabled) { + onLog("Backfill sweep disarmed in settings — stopping."); + return; + } + onLog( + `→ ${channel.channelSlug}: ${channel.reachable.toLocaleString()} reachable.`, + ); + const outcome = await runBackfillChannelJob({ + paths, + channelSlug: channel.channelSlug, + kindIds, + // Behind anything an operator clicks by hand: a sweep is days long and + // must never make a deliberate single-channel run wait for it. + background: true, + onStarted: (jobId) => { + live.channelJobId = jobId; + }, + }); + live.channelJobId = null; + if (!outcome.ok) { + // One channel failing to START is not the sweep failing. + onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`); + continue; + } + didWork = true; + } + + const moved = lastReachable === null || lastReachable > reachable; + barrenPasses = moved ? 0 : barrenPasses + 1; + lastReachable = reachable; + + if (barrenPasses >= MAX_BARREN_PASSES) { + onLog( + `Pass ${pass} made no progress for ${MAX_BARREN_PASSES} passes in a row ` + + `(${reachable.toLocaleString()} still outstanding). Stopping rather than ` + + `spinning — check the engine and the job logs, then restart the sweep.`, + ); + return; + } + + if (!didWork || !moved) { + onLog( + `Pass ${pass} made no progress; 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` WITH its scope so +// a restart resumes exactly the run that was armed (see +// resumeBackfillSweepIfEnabled). +export async function startBackfillSweep( + opts: BackfillSweepOptions = {}, +): Promise<string | null> { + const paths = opts.paths ?? getPaths(); + const running = getBackfillSweepJobId(); + if (running) return running; + + const settings = getSettings(); + const channelScope = opts.channelSlugs ?? settings.backfill.sweepChannels; + const kindScope = opts.kindIds ?? settings.backfill.sweepKinds; + const same = (a: string[], b: string[]) => + a.join(" + if ( + !settings.backfill.sweepEnabled || + !same(settings.backfill.sweepChannels, channelScope) || + !same(settings.backfill.sweepKinds, kindScope) + ) { + // AWAITED, and it matters: the loop reads `sweepEnabled` at the top of its + // first pass, so an un-awaited write races it and the sweep quits + // immediately with "disarmed in settings" — refusing to start at all. + await writeSettings({ + ...settings, + backfill: { + ...settings.backfill, + sweepEnabled: true, + sweepChannels: channelScope, + sweepKinds: kindScope, + }, + }); + } + + const live: SweepLive = { + jobId: "", + startedAt: Date.now(), + channelJobId: null, + }; + const result = await runManagedFunction({ + kind: BACKFILL_SWEEP_KIND, + // Empty key — see the header. The orchestrator does no work itself; it waits + // on per-channel jobs that serialize on BACKFILL_QUEUE. + queueKey: "", + paths, + fn: async (onLog, signal, _setProgress, ctx) => { + live.jobId = ctx.jobId; + try { + await runSweepLoop( + paths, + kindScope.length > 0 ? kindScope : undefined, + channelScope.length > 0 ? channelScope : undefined, + live, + 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. 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. "Is a sweep running" is process +// state, and process state is what a restart destroys. +export async function resumeBackfillSweepIfEnabled( + paths: Paths = getPaths(), +): Promise<void> { + if (!getSettings().backfill.sweepEnabled) return; + await startBackfillSweep({ paths }); +} + +// Disarm and stop. +export async function stopBackfillSweep(): Promise<boolean> { + const settings = getSettings(); + // Note the second and third clauses: guarding on `sweepEnabled` alone means a + // stop on an already-disarmed sweep skips the write entirely and leaves the + // SCOPE behind — which is how a stale scope survives to silently narrow the + // next sweep. That was measured on the digest sweep, not theorised. + if ( + settings.backfill.sweepEnabled || + settings.backfill.sweepChannels.length > 0 || + settings.backfill.sweepKinds.length > 0 + ) { + await writeSettings({ + ...settings, + backfill: { + ...settings.backfill, + sweepEnabled: false, + sweepChannels: [], + sweepKinds: [], + }, + }); + } + const live = getSingleton().live; + if (!live) return false; + // Drain BOTH: the orchestrator, and the per-channel job it is waiting on. + // Draining only the orchestrator makes "stop" mean "after the current + // channel", while the button sits there looking wedged. A drain is still + // graceful — the video in flight completes and nothing part-done is discarded. + if (live.channelJobId) getRegistry().requestDrain(live.channelJobId); + return getRegistry().requestDrain(live.jobId); +} diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -25,6 +25,13 @@ import { import { isDoNotClean } from "../lib/doNotClean-server"; import { loadDigest } from "../lib/digest-server"; import { isSectionFresh } from "../lib/digest"; +import { + addBackfillState, + emptyBackfillCounts, + laneBackfillKinds, + type BackfillClassification, + type BackfillSnapshotEntry, +} from "../lib/backfillKinds"; import { resolveDigestTarget } from "./digestTarget"; import { isExcludedFromTruncatedCheck } from "../lib/excludeTruncatedCheck-server"; import { loadDownloadOutcome } from "../lib/downloadOutcome-server"; @@ -71,6 +78,20 @@ export type ChannelSnapshot = { // local lane is actually carrying the corpus. Optional: older snapshots lack // it; readers default to {}. digestEngines?: Record<string, number>; + // Per-backfill-kind work counts, keyed by kind id (see lib/backfillKinds.ts). + // Beside `totals` and NOT in `buckets`, following the digestEngines precedent + // above for the same reason: buckets is a closed literal of `string[]` id + // lists, and this is per-kind counts. + // + // THE THREE NUMBERS ARE NEVER SUMMED. `missing` + `stale` is work the lane can + // do today; `missingInput` needs the media re-acquired and, measured on this + // corpus, is 91x larger. A single "remaining" figure here would put every + // channel permanently at the top of every list — which is the documented + // reason /api/widget/actionable refuses to filter on `noDigest`. + // + // Optional: snapshots written before this existed lack it, and readers default + // to {}. + backfill?: Record<string, BackfillSnapshotEntry>; buckets: { noTranscript: string[]; downloadedNoTranscript: string[]; @@ -355,6 +376,17 @@ export async function generateChannelSnapshot( (e) => !e.toLowerCase().startsWith("youtube"), ); + // Which backfills are live, and what identity each would produce right now. + // Resolved ONCE per channel — a settings read and some string work — so the + // per-video probe below is a comparison rather than a derivation, exactly as + // digestTarget is above. An empty list is the default (every backfill feature + // ships off), and it makes the whole indicator free. + const backfillKinds = laneBackfillKinds(getSettings()); + const backfillTargets: Record<string, unknown> = {}; + for (const kind of backfillKinds) { + backfillTargets[kind.id] = kind.resolveTarget(getSettings()); + } + const limit = pLimit(SNAPSHOT_VIDEO_CONCURRENCY); const perVideo = await Promise.all( videoDirNames.map((id) => @@ -414,9 +446,26 @@ export async function generateChannelSnapshot( isVideoTranscribed(files) && !files.isUntranscribable ? await loadDigest(dir) : null; + // Backfill state, per registered kind. Costs NOTHING when no backfill + // feature is enabled — `backfillKinds` is empty then, and this loop does + // not run — which is the normal case and why this can live in the + // snapshot's hot path at all. When one IS enabled it is one small + // sidecar read per transcribed video, the same shape and cost as the + // digest read directly above; the classification itself reuses `files` + // rather than re-listing the directory. + const backfill: Record<string, BackfillClassification> = {}; + for (const kind of backfillKinds) { + backfill[kind.id] = await kind.state({ + videoDir: dir, + videoId: id, + files, + target: backfillTargets[kind.id], + }); + } return { id, files, + backfill, audioSizes, nativeId, availability, @@ -795,6 +844,21 @@ export async function generateChannelSnapshot( } : undefined; + // Fold the per-video classifications into per-kind counts. Only the reachable + // half carries ids — see BackfillSnapshotEntry for why missing-input does not. + const backfillCounts: Record<string, BackfillSnapshotEntry> = {}; + for (const kind of backfillKinds) { + const counts = emptyBackfillCounts(); + const ids: string[] = []; + for (const v of perVideo) { + const state = v.backfill[kind.id]; + if (!state) continue; + addBackfillState(counts, state); + if (state === "missing" || state === "stale") ids.push(v.id); + } + backfillCounts[kind.id] = { ...counts, ids: ids.sort() }; + } + const corruptSourceSet = new Set(corruptSource); const corruptFullSourceSet = new Set(corruptFullSource); const snapshot: ChannelSnapshot = { @@ -835,6 +899,7 @@ export async function generateChannelSnapshot( digestWarnings: digestWarnings.sort(), }, digestEngines, + backfill: backfillCounts, undownloadedIds, excludedFromDownload, keptCount: keptIds.size, diff --git a/common/controller/diarizeOne.ts b/common/controller/diarizeOne.ts @@ -16,8 +16,12 @@ import { execa } from "execa"; import type { Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; import type { DiarizationSettings } from "../lib/settings"; -import { DIARIZATION_FILENAME } from "../lib/diarization"; -import { hasDiarization } from "../lib/diarization-server"; +import { + DIARIZATION_FILENAME, + diarizationTarget, + isDiarizationFresh, +} from "../lib/diarization"; +import { hasDiarization, loadDiarization } from "../lib/diarization-server"; import { findSourceMedia, isRealAudioFile } from "../lib/videoStatus"; import { resolveSavedVideo } from "../lib/savedVideo-server"; @@ -41,7 +45,9 @@ export type DiarizeOneOptions = { export type DiarizeOneOutcome = // Wrote (or rewrote) diarization.json. | "diarized" - // A valid sidecar already existed and force was not set. + // A valid sidecar already existed AT THE CURRENT IDENTITY (engine, models, + // threshold) and force was not set. A sidecar from a different identity is + // stale, not this — it gets regenerated. | "already-exists" // The lane is switched off in settings. | "disabled" @@ -87,8 +93,17 @@ export async function diarizeOneVideo( ); return "not-configured"; } - if (!opts.force && (await hasDiarization(opts.videoDir))) { - return "already-exists"; + // FRESHNESS, not existence. This used to short-circuit on `hasDiarization` + // alone, which meant the two most likely reasons to re-run — a different + // clustering threshold, a different model — left the entire corpus looking + // done, and the only way to redo anything was `force` over everything. The + // record carries its provenance precisely so this comparison is possible; see + // isDiarizationFresh for what is compared and why `version` is not. + if (!opts.force) { + const existing = await loadDiarization(opts.videoDir); + if (existing && isDiarizationFresh(existing, diarizationTarget(cfg))) { + return "already-exists"; + } } const media = await resolveDiarizableMedia(opts.videoDir); diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -200,6 +200,29 @@ const JOB_KINDS: Record<string, JobKindMeta> = { bookmarkable: true, queueKeyStrategy: "custom", }, + // The corpus-wide backfill sweep's orchestrator. Drainable, because a "stop" + // has to reach the channel job it is waiting on rather than meaning "after the + // current channel", and NOT bookmarkable: the sweep is armed through a + // persisted settings flag, so replaying a record would be a second way to + // start the same singleton. + "backfill-sweep": { + kind: "backfill-sweep", + label: "Backfill sweep", + drainable: true, + bookmarkable: false, + queueKeyStrategy: "parallel", + }, + // One channel through the backfill lane. Drainable (the batch stops taking new + // videos and lets the in-flight one finish) and bookmarkable, because it + // re-derives its work-list from disk on every run and so replays correctly + // with nothing remembered — the same contract diarize-channel has. + "backfill-channel": { + kind: "backfill-channel", + label: "Backfill channel", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "custom", + }, "check-kept-deleted": { kind: "check-kept-deleted", label: "Check kept videos", diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -14,7 +14,11 @@ export type JobStatus = // Extended for the digest sweep. NOTE: this union has one re-spelled copy in // RunningJobsList.tsx's RunningJobsListItem — kept as an IMPORT there now, so // TypeScript actually flags the next member added here. -export type JobProgressMetric = "downloads" | "transcripts" | "digests"; +export type JobProgressMetric = + | "downloads" + | "transcripts" + | "digests" + | "backfills"; export type JobProgress = { metric: JobProgressMetric; @@ -41,7 +45,7 @@ export type JobProgress = { remainingAudioSeconds?: number; }; -export type JobTaskKind = "download" | "transcribe" | "digest"; +export type JobTaskKind = "download" | "transcribe" | "digest" | "backfill"; // A single in-flight sub-operation within a job (one video download or one // transcription). Only currently-running tasks are kept on the record — they diff --git a/common/jobs/snapshotScheduler.ts b/common/jobs/snapshotScheduler.ts @@ -40,6 +40,14 @@ const NO_REGEN_KINDS = new Set<string>([ "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. backfillSweep's per-channel job regenerates ONCE, + // at job end. + "backfill-sweep", + "backfill-channel", ]); export function shouldRequestSnapshot(kind: string): boolean { diff --git a/common/lib/backfillKinds.test.ts b/common/lib/backfillKinds.test.ts @@ -0,0 +1,327 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import path from "node:path"; +import os from "node:os"; +import { mkdtemp, mkdir, writeFile, rm } from "node:fs/promises"; +import { + getBackfillKind, + laneBackfillKinds, + resolveBackfillKinds, + addBackfillState, + emptyBackfillCounts, + reachableBackfillWork, + type BackfillClassification, +} from "./backfillKinds"; +import { readVideoFiles, SOURCE_MEDIA_BASENAME } from "./videoStatus"; +import { + DEFAULT_DIARIZATION_THRESHOLD, + DIARIZATION_FILENAME, + isDiarizationFresh, + type DiarizationRecord, +} from "./diarization"; +import { defaultSiteSettings, type SiteSettings } from "./settings"; +import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test common/lib/backfillKinds.test.ts +// +// Every case here is a FIXTURE DIRECTORY, on purpose: `state()` is defined as +// "read from disk, never stored", and a test that hands it a hand-built record +// would not be testing the thing the indicators actually call. + +const diarization = getBackfillKind("diarization")!; + +function settingsWithDiarization( + over: Partial<SiteSettings["diarization"]> = {}, +): SiteSettings { + const s = defaultSiteSettings(); + return { + ...s, + diarization: { + ...s.diarization, + enabled: true, + segModel: "/opt/models/seg-1.onnx", + embModel: "/opt/models/emb-1.onnx", + threshold: DEFAULT_DIARIZATION_THRESHOLD, + ...over, + }, + }; +} + +// A video dir built from a description of what is on disk. `sidecar` is written +// verbatim so a MALFORMED file can be tested — that case is not hypothetical, +// it is what a crash mid-write leaves behind. +async function fixture(opts: { + transcript?: boolean; + audio?: boolean; + container?: boolean; + savedPointer?: boolean; + sidecar?: string; +}): Promise<{ dir: string; cleanup: () => Promise<void> }> { + const root = await mkdtemp(path.join(os.tmpdir(), "backfill-kinds-")); + const dir = path.join(root, "vid1"); + await mkdir(dir, { recursive: true }); + if (opts.transcript !== false) { + await writeFile( + path.join(dir, "transcript.json"), + JSON.stringify({ transcription: [{ text: "hi" }] }), + ); + } + if (opts.audio) await writeFile(path.join(dir, "audio.mp3"), "x"); + // `source-media.<ext>` specifically — isSourceMediaFile keys off the basename, + // and an arbitrarily-named mp4 in a video dir is not a persisted container. + if (opts.container) { + await writeFile(path.join(dir, `${SOURCE_MEDIA_BASENAME}.mp4`), "x"); + } + if (opts.savedPointer) { + await writeFile( + path.join(dir, SAVED_VIDEO_POINTER_FILENAME), + JSON.stringify({ storedAt: "now", dir: "/nowhere", file: "video.mp4" }), + ); + } + if (opts.sidecar !== undefined) { + await writeFile(path.join(dir, DIARIZATION_FILENAME), opts.sidecar); + } + return { dir, cleanup: () => rm(root, { recursive: true, force: true }) }; +} + +function sidecar(engine: DiarizationRecord["engine"]): string { + return JSON.stringify({ + videoId: "vid1", + generatedAt: "2026-08-07T00:00:00.000Z", + speakers: 2, + turns: [{ start: 0, end: 4, speaker: 0 }], + engine, + } satisfies DiarizationRecord); +} + +// The current identity, as scripts/diarize.mjs would record it: BASENAMES, not +// the configured full paths. +const CURRENT = { + engine: "sherpa-onnx", + segmentationModel: "seg-1.onnx", + embeddingModel: "emb-1.onnx", + threshold: DEFAULT_DIARIZATION_THRESHOLD, +}; + +async function classify( + dirOpts: Parameters<typeof fixture>[0], + settings: SiteSettings = settingsWithDiarization(), +): Promise<BackfillClassification> { + const { dir, cleanup } = await fixture(dirOpts); + try { + const files = await readVideoFiles(dir, { checkUntranscribable: true }); + return await diarization.state({ + videoDir: dir, + videoId: "vid1", + files, + target: diarization.resolveTarget(settings), + }); + } finally { + await cleanup(); + } +} + +test("present: the sidecar matches the identity we would produce now", async () => { + assert.equal( + await classify({ audio: true, sidecar: sidecar(CURRENT) }), + "present", + ); + // Still present with the audio already cleaned away — the whole point of + // capturing while the audio exists is that the result outlives it. + assert.equal(await classify({ sidecar: sidecar(CURRENT) }), "present"); +}); + +test("stale: a different threshold or model is work, not coverage", async () => { + // The threshold is the single most consequential knob (it decides how many + // speakers come out), so a change to it MUST show as work. Before the + // comparator, this read as done. + assert.equal( + await classify({ + audio: true, + sidecar: sidecar({ ...CURRENT, threshold: 0.5 }), + }), + "stale", + ); + assert.equal( + await classify({ + audio: true, + sidecar: sidecar({ ...CURRENT, segmentationModel: "seg-OLD.onnx" }), + }), + "stale", + ); + assert.equal( + await classify({ + audio: true, + sidecar: sidecar({ ...CURRENT, embeddingModel: "emb-OLD.onnx" }), + }), + "stale", + ); +}); + +// THE ENGINE IS NOT COMPARED, and that is a fix rather than an omission. +// scripts/diarize.mjs records whichever binary actually ran, and nothing in +// DiarizationSettings can predict that — so a hardcoded engine in the target +// would mark every sidecar from any other wrapper permanently stale, which at +// ~500-680 s/audio-hour is an infinite regeneration loop. (The e2e fake engine +// records "fake-diarize" and would have tripped it on the first run.) +test("a different engine binary is not, by itself, stale", async () => { + assert.equal( + await classify({ + audio: true, + sidecar: sidecar({ ...CURRENT, engine: "some-other-engine" }), + }), + "present", + ); + // The comparison is still WRITTEN, so adding an engine setting later needs no + // new logic: a target that does declare one still rejects a mismatch. + assert.equal( + isDiarizationFresh( + { + videoId: "v", + generatedAt: "now", + speakers: 1, + turns: [], + engine: { ...CURRENT, engine: "some-other-engine" }, + }, + { ...CURRENT, engine: "sherpa-onnx" }, + ), + false, + ); +}); + +// THE COMPATIBILITY RULE, and the reason it is written down: without it, adding +// a field to the provenance would mark all 77,000 videos stale at once. +test("an absent recorded field compares equal to today's default", async () => { + assert.equal( + await classify({ + audio: true, + // A sidecar from before the threshold was recorded at all. + sidecar: sidecar({ + engine: "sherpa-onnx", + segmentationModel: "seg-1.onnx", + embeddingModel: "emb-1.onnx", + }), + }), + "present", + ); + // ...and it is equal to the DEFAULT specifically, not to anything: with a + // non-default threshold configured, the same record is stale. + assert.equal( + await classify( + { + audio: true, + sidecar: sidecar({ + engine: "sherpa-onnx", + segmentationModel: "seg-1.onnx", + embeddingModel: "emb-1.onnx", + }), + }, + settingsWithDiarization({ threshold: 0.5 }), + ), + "stale", + ); +}); + +test("a version bump alone does not invalidate captured work", async () => { + // ~500-680 s/audio-hour of CPU says an engine point-release is not a reason to + // redo everything when the models and the threshold are unchanged. + assert.equal( + await classify({ + audio: true, + sidecar: sidecar({ ...CURRENT, version: "9.9.9" }), + }), + "present", + ); +}); + +test("missing: no sidecar, and the input is still here", async () => { + assert.equal(await classify({ audio: true }), "missing"); + // A persisted source container counts — ffmpeg reads it directly, which is + // what resolveDiarizableMedia does. + assert.equal(await classify({ container: true }), "missing"); +}); + +test("missing-input: no sidecar and nothing to diarize from", async () => { + assert.equal(await classify({}), "missing-input"); + // A saved-video POINTER whose stored file has gone (an unmounted backup disk) + // is not an input either — the pointer is not the media. + assert.equal(await classify({ savedPointer: true }), "missing-input"); +}); + +test("a malformed sidecar reads as absent, never as done", async () => { + // The same rule diarization-server.ts's hasDiarization() encodes: a + // half-written file must not be what convinces anything the work is captured + // — that is what would let the cleanup sweep delete the only copy of the audio. + assert.equal( + await classify({ audio: true, sidecar: "{ not json" }), + "missing", + ); + assert.equal( + await classify({ audio: true, sidecar: JSON.stringify({ videoId: "x" }) }), + "missing", + ); +}); + +test("not-applicable: an untranscribed video is not this backfill's business", async () => { + assert.equal( + await classify({ transcript: false, audio: true }), + "not-applicable", + ); +}); + +test("a disabled or unconfigured feature reports no backfill at all", () => { + const s = defaultSiteSettings(); + // Off by default, so nothing advertises catch-up work for it. + assert.equal(diarization.enabled(s), false); + assert.equal(laneBackfillKinds(s).length, 0); + // Enabled but with no models is "not set up", which must also report nothing + // rather than a corpus-sized work list nobody can act on. + assert.equal( + diarization.enabled(settingsWithDiarization({ segModel: "" })), + false, + ); + assert.equal(diarization.enabled(settingsWithDiarization()), true); + assert.equal(laneBackfillKinds(settingsWithDiarization()).length, 1); +}); + +test("an unknown or stale kind id in the sweep scope is dropped, not fatal", () => { + const s = settingsWithDiarization(); + // Empty scope = every enabled lane kind. + assert.deepEqual( + resolveBackfillKinds(s, []).map((k) => k.id), + ["diarization"], + ); + assert.deepEqual( + resolveBackfillKinds(s, undefined).map((k) => k.id), + ["diarization"], + ); + // A settings file naming a kind from another build must not wedge the lane. + assert.deepEqual(resolveBackfillKinds(s, ["from-the-future"]), []); + assert.deepEqual( + resolveBackfillKinds(s, ["diarization", "from-the-future"]).map((k) => k.id), + ["diarization"], + ); +}); + +// The accounting rule the whole feature turns on: reachable work and +// needs-re-acquiring are never added together. +test("counts keep reachable work and needs-re-acquiring apart", () => { + const counts = emptyBackfillCounts(); + for (const state of [ + "missing", + "missing", + "stale", + "missing-input", + "missing-input", + "missing-input", + "present", + "not-applicable", + ] as BackfillClassification[]) { + addBackfillState(counts, state); + } + assert.deepEqual(counts, { missing: 2, stale: 1, missingInput: 3 }); + // 3, not 6. Measured on the real corpus the difference is 835 vs 77,105, and + // reporting the larger number is what would make every surface useless. + assert.equal(reachableBackfillWork(counts), 3); +}); diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts @@ -0,0 +1,288 @@ +// The single source of truth for what a BACKFILL is. +// +// The problem this exists for repeats: a derived-data feature lands, and the +// corpus that already exists does not have what it needs. For diarization that +// input is AUDIO, which cleanAudioFromTranscribed deletes once a video is +// transcribed. Before this table the only catch-up was a per-channel button on +// a controller written for that one feature (controller/diarizeAll.ts) — there +// was no way to ask "how much of the corpus is missing this?", and no way to run +// catch-up alongside new-video work without one starving the other. +// +// So, in the shape jobKinds.ts already uses: adding a backfill should mean +// adding ONE entry here, and the system supplies the lane, the resource share +// and the indicator. +// +// FOUR STATES, NOT TWO, and the split is the load-bearing part. Measured on this +// corpus at the time of writing: 77,106 videos, 836 with media still on disk, 1 +// diarized. A single "remaining" number would therefore read 77,105 — and 91x of +// that is unreachable without re-downloading. The repo has already been burned by +// exactly this once: editor/app/api/widget/actionable/route.ts deliberately +// refuses to filter on `noDigest` because during the backfill that is 99.87% of +// the corpus and counting it would put every channel in the list forever. So +// `missing` (reachable now) and `missing-input` (needs re-acquiring) are +// SEPARATE numbers, everywhere, and no surface is allowed to add them together. +// +// STALENESS IS PROVENANCE, NOT AGE. `state` is derived from disk on every read +// and never stored, and a kind that records what produced its output compares +// that against what we would produce now — lib/digest.ts's isSectionFresh, whose +// absent-field-equals-today's-default trick is what stops adding a field from +// invalidating the whole corpus. Diarization writes that provenance and, before +// this, had no comparator at all: diarizeOne short-circuited on mere existence, +// so the two most likely reasons to re-run (a threshold or model change) left +// everything looking done. +// +// SERVER-ONLY, despite living in lib/. It reads the filesystem and calls a +// controller, so it is `-server.ts` in everything but name; the path is the one +// the plan named. No client component imports it — the UI is handed plain +// numbers off the channel snapshot, and labels as props. +// +// WHY ONLY ONE ENTRY TODAY. The other catch-up mechanisms in the repo do not fit +// this per-video probe, and forcing them in would make the table lie: +// controller/backfillAvailability.ts is CHANNEL-scoped (one JSON map, folded +// into sync at runYtdlp.ts), and controller/normalizeAll.ts has no recorded +// provenance to compare, so its "stale" is undefined. `tier` still exists +// because it is what keeps them apart if they are ever added: `inline` folds +// into an existing pass and `lane` gets the concurrent queue and the share. + +import type { Paths } from "./paths"; +import type { SiteSettings } from "./settings"; +import { + diarizationTarget, + isDiarizationFresh, + type DiarizationFreshnessTarget, +} from "./diarization"; +import { loadDiarization } from "./diarization-server"; +import { + findSourceMedia, + isVideoTranscribed, + type VideoFiles, +} from "./videoStatus"; +import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo"; +import { resolveSavedVideo } from "./savedVideo-server"; +import { diarizeOneVideo } from "../controller/diarizeOne"; + +// What a video's relationship to a backfill is, right now, read from disk. +// +// present — has it, at the identity we would produce now. +// stale — has it, but from a different engine/model/threshold. +// missing — does not have it, and the input to produce it is HERE. +// missing-input — does not have it, and the input is gone. Reachable only by +// re-acquiring the media, which is opt-in and bounded. +export type BackfillState = "present" | "stale" | "missing" | "missing-input"; + +// Videos this backfill has no opinion about (not transcribed, marked +// untranscribable). Kept out of BackfillState so it can never be counted. +export type BackfillClassification = BackfillState | "not-applicable"; + +// How expensive one video is, which decides where the work runs. +// +// inline — microseconds to cheap I/O; folds into a pass that already walks the +// corpus, and never gets a lane of its own. +// lane — expensive enough to need its own queue and a resource share. +// Diarization is ~500-680 s/audio-hour of CPU. +export type BackfillCostTier = "inline" | "lane"; + +export type BackfillProbe = { + videoDir: string; + videoId: string; + // Already read by the caller. Taking it rather than re-reading is what makes + // the channel snapshot's per-video classification free — see VideoFiles.entries. + files: VideoFiles; + target: unknown; +}; + +export type BackfillRunOptions = { + paths: Paths; + videoDir: string; + videoId: string; + channelSlug: string; + target: unknown; + // Redo a `present` video anyway (an operator forcing a regeneration). + force?: boolean; + onLog?: (msg: string) => void; + signal?: AbortSignal; +}; + +// Deliberately mirrors DiarizeOneOutcome's discipline: an expected condition is +// an outcome, never a throw. The batch counts these and keeps going. +export type BackfillRunOutcome = + | "done" + | "already-present" + | "missing-input" + | "not-configured" + | "disabled" + | "failed"; + +export type BackfillKind = { + id: string; + label: string; + // One line of UI copy: what this backfill is, in the operator's terms. + hint: string; + tier: BackfillCostTier; + // The feature's OWN gate. A disabled feature reports no backfill at all — + // otherwise every surface would advertise catch-up work for something the + // operator has switched off. + enabled(settings: SiteSettings): boolean; + // The identity we would produce now, resolved ONCE per run rather than per + // video. `unknown` here is the one erasure point in the table: each entry + // narrows it back to its own type on the line below. The alternative — making + // the whole registry generic — infects every consumer with a type parameter + // for no gain, since none of them look inside a target. + resolveTarget(settings: SiteSettings): unknown; + state(probe: BackfillProbe): Promise<BackfillClassification>; + run(opts: BackfillRunOptions): Promise<BackfillRunOutcome>; +}; + +// Diarization: the first entry, and the reason the table exists. +const diarization: BackfillKind = { + id: "diarization", + label: "Speaker diarization", + hint: "Speaker turns captured from the audio, written to diarization.json beside the transcript.", + tier: "lane", + enabled: (settings) => + settings.diarization.enabled && + !!settings.diarization.segModel && + !!settings.diarization.embModel, + // The SAME target diarizeOne's own short-circuit uses. Two derivations would + // let the counter and the runner disagree about what is stale. + resolveTarget: (settings): DiarizationFreshnessTarget => + diarizationTarget(settings.diarization), + async state({ videoDir, files, target }) { + // Same eligibility as diarizeAll's transcribedOnly default: the capture lane + // exists to pair speaker turns with a transcript, and an untranscribed + // video's audio is not at risk from the cleanup sweep yet. + if (!isVideoTranscribed(files) || files.isUntranscribable) { + return "not-applicable"; + } + // Cheap negative first: no sidecar in the listing means no read at all. + if (files.hasDiarization) { + const record = await loadDiarization(videoDir); + // A malformed file reads as ABSENT here, exactly as hasDiarization() in + // diarization-server.ts treats it: a half-written sidecar must never be + // what convinces anything the work is done. + if (record) { + return isDiarizationFresh(record, target as DiarizationFreshnessTarget) + ? "present" + : "stale"; + } + } + return (await hasDiarizableInput(videoDir, files)) + ? "missing" + : "missing-input"; + }, + async run(opts) { + // diarizeOneVideo re-reads settings when none is passed, which is what we + // want: the batch may run for hours and a model change mid-run should be + // picked up. `force` is how a `stale` video gets redone at all — the + // existence short-circuit inside is now a freshness check, but an operator + // forcing a regeneration still needs to win. + const outcome = await diarizeOneVideo({ + paths: opts.paths, + videoDir: opts.videoDir, + videoId: opts.videoId, + force: opts.force, + onLog: opts.onLog, + signal: opts.signal, + }); + if (outcome === "diarized") return "done"; + if (outcome === "already-exists") return "already-present"; + if (outcome === "no-audio") return "missing-input"; + if (outcome === "not-configured") return "not-configured"; + if (outcome === "disabled") return "disabled"; + return "failed"; + }, +}; + +// Is there anything on disk ffmpeg could read for this video? Mirrors +// controller/diarizeOne.ts's resolveDiarizableMedia, but answered from the +// listing the caller already has so the common cases cost no I/O: +// extracted audio, then a persisted source container, then — only when the +// pointer file is actually present — the saved-video store. +async function hasDiarizableInput( + videoDir: string, + files: VideoFiles, +): Promise<boolean> { + if (files.audioFiles.length > 0) return true; + if (findSourceMedia(files.entries)) return true; + if (!files.entries.includes(SAVED_VIDEO_POINTER_FILENAME)) return false; + // The pointer exists but the stored file may not (an unmounted backup disk), + // so this last step really does have to touch the filesystem. + return (await resolveSavedVideo(videoDir)) !== null; +} + +// One entry per backfill known to the system. +export const BACKFILL_KINDS: readonly BackfillKind[] = [diarization]; + +export const BACKFILL_KIND_BY_ID: Record<string, BackfillKind> = + Object.fromEntries(BACKFILL_KINDS.map((k) => [k.id, k])); + +export function getBackfillKind(id: string): BackfillKind | undefined { + return BACKFILL_KIND_BY_ID[id]; +} + +// The kinds that get the concurrent lane, filtered by their own feature gate. +// Both filters matter: `inline` kinds must never take a lane slot, and a +// disabled feature must report and run nothing. +export function laneBackfillKinds(settings: SiteSettings): BackfillKind[] { + return BACKFILL_KINDS.filter( + (k) => k.tier === "lane" && k.enabled(settings), + ); +} + +// Resolve a caller-supplied list of kind ids against the registry. An empty or +// absent list means "every enabled lane kind" — the sweep's scope default. +// Unknown ids are dropped rather than throwing: a settings file may name a kind +// from a newer build, and a stale scope must not wedge the lane. +export function resolveBackfillKinds( + settings: SiteSettings, + ids: readonly string[] | undefined, +): BackfillKind[] { + const lane = laneBackfillKinds(settings); + if (!ids || ids.length === 0) return lane; + const wanted = new Set(ids); + return lane.filter((k) => wanted.has(k.id)); +} + +// Per-kind counts, the shape every indicator reads. `missing` and `missingInput` +// are never summed — see the header. +export type BackfillCounts = { + missing: number; + stale: number; + missingInput: number; +}; + +export function emptyBackfillCounts(): BackfillCounts { + return { missing: 0, stale: 0, missingInput: 0 }; +} + +// Fold one classification into a counts record. Central so no surface invents +// its own accounting: `present` and `not-applicable` add to nothing, which is +// what makes these counts a WORK LIST rather than a coverage measure. +export function addBackfillState( + counts: BackfillCounts, + state: BackfillClassification, +): void { + if (state === "missing") counts.missing++; + else if (state === "stale") counts.stale++; + else if (state === "missing-input") counts.missingInput++; +} + +// What the lane can act on WITHOUT re-acquiring media. The number every "how +// much is left?" surface should lead with. +export function reachableBackfillWork(counts: BackfillCounts): number { + return counts.missing + counts.stale; +} + +// What a channel snapshot stores per kind: the three counts, plus the ids of the +// REACHABLE work only. +// +// The asymmetry is deliberate. A stage card has to list what it would act on, so +// those ids have to be somewhere the render path can read without walking the +// corpus (there is a guard test forbidding exactly that). But `missingInput` is +// ~76,000 videos corpus-wide, and writing that list into all 66 snapshots would +// put tens of megabytes of ids on disk to say a number we already have. So: ids +// for the actionable half, a count for the other. +export type BackfillSnapshotEntry = BackfillCounts & { + // missing + stale, sorted. Never includes missing-input. + ids: string[]; +}; diff --git a/common/lib/diarization.ts b/common/lib/diarization.ts @@ -60,6 +60,118 @@ export type DiarizationRecord = { export const DIARIZATION_FILENAME = "diarization.json"; +// The engine id scripts/diarize.mjs records when it runs its own default +// (sherpa-onnx) rather than a `--engine` replacement. Named here so the +// freshness comparator below and the wrapper agree on one spelling. +export const DEFAULT_DIARIZATION_ENGINE = "sherpa-onnx"; + +// The clustering-threshold default, measured on this corpus (see +// DiarizationSettings.threshold for the sweep that produced it). It lives here +// rather than only in settings.ts because isDiarizationFresh needs it to +// normalize an ABSENT recorded threshold, and this module must stay importable +// from anywhere. settings.ts imports it, so there is still one source of truth. +export const DEFAULT_DIARIZATION_THRESHOLD = 0.9; + +// What we WOULD produce for this video now, as an identity. The mirror of +// DigestFreshnessTarget in lib/digest.ts, and deliberately the same shape of +// idea: a sidecar is stale when its recorded provenance differs from this, not +// when it is old. +export type DiarizationFreshnessTarget = { + // OPTIONAL, and compared only when present — which today means never. + // + // The engine is not a setting: scripts/diarize.mjs records whatever binary + // actually ran (its own default, or the basename of a `--engine` / + // DIARIZE_ENGINE_CMD replacement), and nothing in DiarizationSettings can say + // which that will be. Asserting a hardcoded "sherpa-onnx" here would mark + // every sidecar produced by any other wrapper permanently stale — an infinite + // regeneration loop at ~500-680 s/audio-hour, and one the e2e fake engine + // would trip on its first run. + // + // So the comparison is restricted to what the configuration can genuinely + // predict: the models and the threshold, which are also the knobs that + // actually change the output. The field stays here, and the comparison stays + // written, so that adding an engine setting later needs no new logic. + engine?: string; + // Basenames, as DiarizationEngine records them — full paths are + // machine-specific and would make every record stale on another shard. + segmentationModel?: string; + embeddingModel?: string; + threshold?: number; +}; + +// Basename without importing node:path — this module is client-safe and has no +// dependencies, which is what lets both the engine wrapper's consumers and the +// backfill registry share the comparator below. +function baseName(p: string): string { + const parts = p.split(/[/\\]/); + return parts[parts.length - 1] ?? p; +} + +// The identity the CURRENT configuration would produce. Structurally typed +// rather than taking DiarizationSettings, so this module stays free of a +// settings import (settings.ts imports the threshold default FROM here). +// +// One definition, two callers — controller/diarizeOne.ts's short-circuit and +// lib/backfillKinds.ts's state() — because a comparator and the writer it +// guards disagreeing about the identity is how a corpus ends up either +// regenerating forever or never. +export function diarizationTarget(cfg: { + segModel?: string; + embModel?: string; + threshold?: number; +}): DiarizationFreshnessTarget { + return { + // `engine` is deliberately NOT set — see DiarizationFreshnessTarget. Settings + // cannot know which binary will run, so claiming to compare it would mark + // every sidecar from any other wrapper permanently stale. + ...(cfg.segModel ? { segmentationModel: baseName(cfg.segModel) } : {}), + ...(cfg.embModel ? { embeddingModel: baseName(cfg.embModel) } : {}), + threshold: cfg.threshold ?? DEFAULT_DIARIZATION_THRESHOLD, + }; +} + +// Absent and "" are the same thing (nothing configured), so a record written by +// an engine that reports no model name compares equal to a target that has none +// either. Copied from digest.ts's sameVariant for the same reason it exists +// there: the compatibility rule has to be explicit or adding a field silently +// invalidates the whole corpus. +function sameModel(a: string | undefined, b: string | undefined): boolean { + return (a ?? "") === (b ?? ""); +} + +// An absent recorded threshold is read as TODAY'S DEFAULT. That is the trick +// that stops a newly-recorded field from invalidating everything written before +// it existed: a sidecar from before the field was written compares equal as long +// as the current setting is still the default. +function sameThreshold(a: number | undefined, b: number | undefined): boolean { + return ( + (a ?? DEFAULT_DIARIZATION_THRESHOLD) === (b ?? DEFAULT_DIARIZATION_THRESHOLD) + ); +} + +// Is this sidecar what the current configuration would produce? +// +// Before this, diarizeOne short-circuited on mere EXISTENCE, which meant a +// threshold or model change — the two most likely reasons to re-run at all — +// left the whole corpus looking done. Three fields are compared and one is +// deliberately not: `version` is excluded, because an engine point-release must +// not invalidate ~500-680 s/audio-hour of captured work when the models and the +// threshold it was clustered at are unchanged. +export function isDiarizationFresh( + record: DiarizationRecord | null, + target: DiarizationFreshnessTarget, +): boolean { + if (!record) return false; + const e = record.engine; + if (!e) return false; + return ( + (target.engine === undefined || e.engine === target.engine) && + sameModel(e.segmentationModel, target.segmentationModel) && + sameModel(e.embeddingModel, target.embeddingModel) && + sameThreshold(e.threshold, target.threshold) + ); +} + // Total attributed speech, in seconds. Turns may overlap (two people talking at // once), so this can exceed the audio duration — that is a real signal, not a // bug, and callers that want wall-time coverage should merge ranges first. diff --git a/common/lib/queueKeys.ts b/common/lib/queueKeys.ts @@ -18,6 +18,18 @@ export { TRANSCRIPTION_QUEUE }; export const DIGEST_LOCAL_QUEUE = "digest:local"; export const DIGEST_REMOTE_QUEUE = "digest:remote"; +// The backfill lane (controller/backfillBatch.ts). Its own key for exactly the +// reason the two digest keys have theirs: registry.ts submits every non-empty +// queueKey at concurrency 1, so a distinct key is the ONLY way to get a lane +// that runs CONCURRENTLY with transcription and the digest lanes rather than +// behind them — and that concurrency is the entire point of a catch-up lane. +// +// Serialized against ITSELF, which is also deliberate: catch-up work is +// per-video CPU (diarization is ~500-680 s/audio-hour) and two channels' worth +// at once would just thrash. How much of the machine it may take while running +// is a separate question, answered by backfillLimit() rather than by the queue. +export const BACKFILL_QUEUE = "backfill"; + // Per-channel queue for channel-local bookkeeping jobs (clean/clear/verify). export function channelQueueKey(slug: string): string { return `channel:${slug}`; diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -24,6 +24,7 @@ import { defaultAutoQueue, sanitizeAutoQueue, } from "../jobs/autoQueuePolicy"; +import { DEFAULT_DIARIZATION_THRESHOLD } from "./diarization"; import { DEFAULT_COOKIE_MODE, isCookieMode, @@ -202,6 +203,56 @@ export type SiteSettings = { // Speaker diarization captured right after transcription, while the audio is // still on disk. OFF by default. See DiarizationSettings. diarization: DiarizationSettings; + // The generic catch-up lane for derived data the existing corpus predates. + // OFF by default, and idle-only when on. See BackfillSettings. + backfill: BackfillSettings; +}; + +// Configuration for the backfill lane — the generic answer to "a derived-data +// feature landed and 77,000 existing videos do not have it". +// +// The one knob that matters is `weight`, and it is a SHARE, not a priority: the +// registry submits every named queue at concurrency 1 and SchedulerTier only +// orders work within a single key, so there is no priority system to join. What +// the lane actually gets is its own queueKey (concurrency with transcription) +// plus a limit() that returns 0 to stand aside — the same mechanism the digest +// yield uses, which fails OPEN so a bad read costs contention rather than a +// deadlock. +export type BackfillSettings = { + // Master switch for the lane. Off means the registry still REPORTS what is + // missing (that is the indicator's whole job) but nothing runs. + enabled: boolean; + // The resource share, 0..1. + // + // 0 (default) — idle-only: run only while the primary transcription lane is + // doing nothing. A backfill is by definition catch-up work on a corpus + // that already exists, so it must never slow down new arrivals. + // >0 — a guaranteed share of the lane's slots, floored at 1 so a small + // weight is a slow lane rather than a stopped one. + // + // Clamped to [0, 1]. See backfillLimit() in controller/backfillBatch.ts. + weight: number; + // Slots the lane may use when it is not standing aside. Kept at 1 by default + // for the same reason diarization.concurrency is: this is CPU-bound work + // competing with GPU feeding and the digest sweep for the same 8 threads. + concurrency: number; + // A corpus-wide sweep is armed. Persisted so a restart resumes it, exactly as + // digest.sweepEnabled is (editor/instrumentation.ts). + sweepEnabled: boolean; + // The sweep's SCOPE, persisted alongside the flag rather than only in the + // launching call. The boot hook re-launches from settings alone, so arming + // without recording the scope resurrects a deliberately-bounded run as a + // corpus-wide one — the bug digestSweep.ts documents at its start/stop pair. + // Empty = every registered lane-tier kind / every channel. + sweepKinds: string[]; + sweepChannels: string[]; + // Re-acquire media for videos whose input is GONE (audio deleted after + // transcription). OFF by default and deliberately so: measured on this corpus, + // 836 videos still have media and ~76,270 would need a re-download — 91x the + // reachable work, against 45 GB free at 97% full. When on, each re-fetched + // file is removed in a `finally` as soon as the backfill has used it, unless + // the video is marked do-not-clean. + allowRedownload: boolean; }; // Configuration for the speaker-diarization capture lane. @@ -876,6 +927,52 @@ function defaults(): SiteSettings { // entirely and the whole section is never read from disk. digest: defaultDigest(), diarization: defaultDiarization(), + backfill: defaultBackfill(), + }; +} + +// The whole default settings object, without touching disk. Exported so a test +// (or any caller that needs a settings-SHAPED value rather than the operator's +// actual configuration) can build one without a settings.json. +export function defaultSiteSettings(): SiteSettings { + return defaults(); +} + +export function defaultBackfill(): BackfillSettings { + return { + enabled: false, + // Idle-only. See BackfillSettings.weight. + weight: 0, + concurrency: 1, + sweepEnabled: false, + sweepKinds: [], + sweepChannels: [], + // See BackfillSettings.allowRedownload — this one holds disk. + allowRedownload: false, + }; +} + +export function sanitizeBackfill(value: unknown): BackfillSettings { + const d = defaultBackfill(); + if (!value || typeof value !== "object") return d; + const r = value as Record<string, unknown>; + const slugs = (v: unknown): string[] => + Array.isArray(v) + ? v.filter((s): s is string => typeof s === "string" && s.trim() !== "") + : []; + return { + enabled: r.enabled === true, + // Clamped rather than rejected: a hand-edited 5 means "as much as possible", + // and reading it as the idle-only 0 would be the opposite of the intent. + weight: + typeof r.weight === "number" && Number.isFinite(r.weight) + ? Math.min(1, Math.max(0, r.weight)) + : d.weight, + concurrency: clampPositiveInt(r.concurrency, d.concurrency, 16), + sweepEnabled: r.sweepEnabled === true, + sweepKinds: slugs(r.sweepKinds), + sweepChannels: slugs(r.sweepChannels), + allowRedownload: r.allowRedownload === true, }; } @@ -888,8 +985,11 @@ export function defaultDiarization(): DiarizationSettings { // OFF — see DiarizationSettings.inlineAfterTranscribe. Diarization is slower // than the transcription it would follow, so inline is the exception. inlineAfterTranscribe: false, - // Measured, not sherpa-onnx's 0.5 — see DiarizationSettings.threshold. - threshold: 0.9, + // Measured, not sherpa-onnx's 0.5 — see DiarizationSettings.threshold. The + // constant lives in lib/diarization.ts because isDiarizationFresh needs it + // to normalize an absent recorded threshold; importing it keeps the default + // and the comparator from drifting apart. + threshold: DEFAULT_DIARIZATION_THRESHOLD, threads: 4, python: "python3", segModel: "", @@ -1137,6 +1237,7 @@ export function getSettings(): SiteSettings { merged.buildPipeline = sanitizeBuildPipeline(merged.buildPipeline); merged.digest = sanitizeDigest(merged.digest); merged.diarization = sanitizeDiarization(merged.diarization); + merged.backfill = sanitizeBackfill(merged.backfill); // Workers. When the file predates the worker model (no `workers` key), // synthesize a default list from the (now-settled) active app + per-app // configs so existing installs behave identically. Otherwise sanitize the @@ -1333,6 +1434,7 @@ export async function writeSettings(next: SiteSettings): Promise<void> { buildPipeline: sanitizeBuildPipeline(next.buildPipeline), digest: sanitizeDigest(next.digest), diarization: sanitizeDiarization(next.diarization), + backfill: sanitizeBackfill(next.backfill), }; const tmp = `${file}.tmp-${process.pid}`; await fs.promises.writeFile(tmp, JSON.stringify(merged, null, 2) + "\n"); diff --git a/common/lib/videoStatus.ts b/common/lib/videoStatus.ts @@ -25,6 +25,16 @@ export type VideoFiles = { // Surfaced separately so the UI can offer a Resume action; audio-check's // own internal snapshots (.part.good, .part.testing) are excluded. partAudioFiles: string[]; + // The raw readdir() listing this whole record was derived from. + // + // Exposed so a caller that needs a question this type does not already answer + // — "is there a source container?", "is there a saved-video pointer?" — can + // ask it WITHOUT a second readdir. That matters in exactly one place and it is + // the expensive one: the channel snapshot reads this per video across 77,000 + // of them, so a per-video re-listing is the difference between an indicator + // that is free and one that doubles the report's I/O. See + // lib/backfillKinds.ts, which classifies missing vs missing-input from it. + entries: string[]; }; export type IndexTranscript = @@ -243,6 +253,7 @@ export async function readVideoFiles( isUntranscribable, audioFiles, partAudioFiles, + entries, }; } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **Backfill is now a first-class thing the system knows about, with a lane of its own and a remainder you can see.** Every derived-data feature that lands runs into the same wall: the corpus that already exists does not have what it needs. Diarization hit it first — the input it wants is audio, which Clean-audio deletes once a video is transcribed — and the only catch-up was a per-channel button written for that one feature, with no way to ask how much of the corpus was missing it and no way to run it alongside new-video work without one starving the other. So a feature now **declares** what it needs and how to tell whether a video has it, and gets the lane, the resource share and the indicator for free. **The number this reports is split in two, and that is the whole design.** On this corpus 836 videos still have their media and about 76,270 do not — 91× more — so a single "remaining" figure would be dominated by work no button can start, and every channel would sit at the top of every list forever. Reachable work and needs-re-acquiring are therefore separate numbers on all four surfaces: the channel's new **Backfill** stage card, a new section on `/actionable` (with the second figure in its own column, exactly as "Est. reclaim" has), a dashboard instrument, and an optional widget strip. **The lane runs on its own queue**, so it is genuinely concurrent with transcription rather than sitting behind it, and how much of the machine it may take is one setting: **0 (the default) means idle-only** — full speed while transcription is quiet, standing aside the instant it isn't — and any value above 0 is a guaranteed share, floored at 1 so a small number is a slow lane and not a stopped one. Catch-up on a corpus that already exists must never slow down new arrivals, and the default enforces that rather than trusting it. There is a corpus-wide **Start Backfill Sweep** on the dashboard that persists its scope along with the flag and resumes itself after a restart. **Staleness now means provenance, not age.** A diarization sidecar records the models and clustering threshold that produced it, and the lane compares that against what the current settings *would* produce — so changing the threshold, which is the single knob most likely to make you want a re-run, finally shows as work instead of leaving the corpus looking finished. A field the record predates compares equal to today's default, so adding one does not invalidate everything on disk. **Re-downloading deleted media is off by default and is the part to read carefully.** With it on, each file is fetched, used, and deleted again immediately in a `finally` — whether the backfill succeeded, failed, or crashed — unless the video is marked "do not clean", and nothing starts at all when free disk is under the configured floor. On a disk at 97% a leak here fills it, which is why those four behaviours have e2e tests of their own. - **Speaker diarization can now be captured while the audio still exists, and Clean audio will no longer delete audio out from under it.** Audio is the one input in this pipeline that goes away: the Clean-audio sweep removes it as soon as a video has a transcript, and nothing downstream can reconstruct it. Everything built on top of speaker turns — attribution, naming the speakers, badges, quote filtering — can be redone from a saved file at any time, but the turns themselves can only be extracted from audio. So this ships the perishable half only: a new **Diarization** section in Settings, off by default, that captures speaker-labelled time ranges to `diarization.json` beside each transcript, and a **Diarize speakers** button in every channel's Cleanup stage that backfills over whatever audio is still on disk. Turning capture on also arms a guard: the sweep now refuses to delete audio for a transcribed video that has no sidecar yet, reports it as `awaiting diarization`, and drops that video out of the channel's "Est. reclaim" so the estimate matches what the sweep will actually take. That hold is the point — it is what keeps the audio alive long enough to be captured — but it does mean enabling this stops reclaiming disk until diarization catches up, and the setting says so. **Inline-after-transcribe is a separate switch, and it is off by default for a measured reason.** Transcription on this box runs at 221 s per audio-hour on the GPU — 16.3× real time, averaged over 3,602 real videos from their own `transcribe-outcome.json` records — while diarization runs on the CPU at roughly 500–680, because no diarization model has been ported to ggml and neither ONNX nor PyTorch has a Vulkan compute path on Linux. Diarization is therefore 2–3× slower than the transcription it would follow, so running it inline makes the whole pipeline 3–4× slower and idles the GPU. The intended shape for a large batch is the opposite: turn capture on so the audio is held, leave inline off so transcription runs at full speed, then backfill. The clustering threshold defaults to **0.9**, not sherpa-onnx's own 0.5, and that is also measured: on a six-minute excerpt of a known two-person interview, 0.5 produced 22 speaker clusters and 0.9 produced 6, with the top two at 40%/40% of talk time — recognizably the two hosts. It still over-splits, which is why every sidecar records the engine, both model names, the library version and the threshold that produced it: over-splitting is the recoverable direction, and a later pass can re-cluster or re-run selectively without ever needing the audio back. The engine sits behind a path (`DIARIZE_BIN` → `scripts/diarize.mjs`, driving sherpa-onnx through `scripts/diarize-sherpa.py`) exactly as the parakeet wrapper does, so swapping in pyannote later is a settings change rather than a code change. A failing diarizer never fails a transcription, and never counts as "done" — the sweep keeps that video's audio. Full spike results, including the disk arithmetic and what is still unmeasured, are in `plans/diarization-spike-results.md`. - **The monitor widget's builder is a floorplan you rearrange, not a list of checkboxes.** The widget rendered its nine strips in an order hardcoded in JSX — controls, last sync, scheduler, disk, cleanable total, needs cleaning, workers, active jobs, needs work — and there was no order field in its config at all, so reference data you glance at once (the two cleanable strips) sat permanently above the thing you were actually watching. The builder couldn't help, because it wasn't spatial: 22 checkboxes in three fieldsets, nothing in the control surface corresponding to where anything landed. `/widget/builder` is now a scale model of the widget. Each section is a draggable card carrying a small CSS-drawn likeness of what that strip really renders — four dots for Workers, a part-filled bar for Active jobs, a meter with the low-disk notch for Disk, stacked rows for the two channel lists — laid out in the columns and order it will appear in, with the switched-off ones in a tray underneath. Drag it with the mouse or use the ↑/↓/◀/▶ buttons on every card; nothing here is mouse-only, and a section is one checkbox that changes place rather than a control that vanishes when you turn it off. **Columns are new** — up to three, each header showing the number that actually decides whether columns are viable (the per-column pixel width at the current preview size, flagged when it falls under 150px), alongside two wider presets, **Wide** 640×260 and **Panel** 720×420, because columns need width to be worth having. They are honored at every size by default, since an arrangement you made on purpose shouldn't quietly undo itself; **Stack columns when narrow** opts a widget you intend to resize into collapsing back to a single stack, on a container query measured against the widget's own box rather than the viewport (which, embedded in an iframe, is the host page's). Per-section display options moved onto the section that owns them — the batch bar, heading progress and ETA on the Active jobs card, worker names on Workers, the Sync flags on Controls — leaving only the genuinely global ones (section titles, hide-when-idle, the gear, absolute times) in a row of their own. The in-widget gear gets the same model in a shape that fits its ~320px overlay: one row per section in render order with move buttons and a column picker, driven by the same registry and the same placement functions as the board rather than a second copy of them. **Every link written before this still works, and still renders in the original order.** The arrangement rides in one new `l=` param (`/widget?l=ctl.disk.wk.jobs-cln.act` — `.` between sections, `-` between columns, both characters a URL leaves literal where a comma would come back as `%2C`), it is omitted entirely whenever the arrangement is the one the visibility flags already imply, and an inherited layout is normalized against those flags on arrival: unknown codes and repeats are dropped, a section since switched off is dropped, and a section switched **on** but unlisted is appended rather than silently disappearing — which is what keeps an already-copied link working after you enable something new from the widget's own gear. Underneath, the nine hardcoded conditionals became a section registry that drives both configuration surfaces, so adding a tenth strip is now a table entry (there is one waiting: `/api/widget/sync` already polls a digest field that nothing renders). The builder's remembered form state resets once, since a saved entry from before this has no layout to carry forward. See `editor/app/widget/lib/{sections,placement,config}.ts`, `editor/app/widget/builder/components/{LayoutBoard,SectionCard,SectionMini}.tsx`, and the normalization tests in `editor/app/widget/lib/placement.test.ts`. - **No video a channel has ever listed can be lost again, and a truncated listing can no longer destroy one.** Until now the only record that a video belonged to a channel was the stored playlist file, which the deep pass overwrites wholesale. A video that appeared in a listing but was never downloaded had no folder on disk and no metadata of its own, so the playlist was the *only* place its URL lived — when it dropped out of the listing it was erased with no trace it had ever existed, and no way to even attempt a direct-link recovery. Each channel now keeps a **roster**: every video id it has ever been seen to contain, with the URL it was seen at, added to and never pruned. It is built the first time a channel syncs, from the playlist and the folders on disk *before* anything is rewritten, so the upgrade itself is the protection. diff --git a/editor/app/actionable/components/InlineActionButton.tsx b/editor/app/actionable/components/InlineActionButton.tsx @@ -14,6 +14,7 @@ import { redownloadIncompleteBucketAction, redownloadShortAudioBucketAction, } from "../../channels/[slug]/incompleteTranscriptActions"; +import { backfillChannelAction } from "../../channels/[slug]/backfillActions"; import { refreshChannelSnapshotAction } from "../../channels/actions"; type Variant = @@ -24,6 +25,7 @@ type Variant = | { kind: "redownloadShortAudio"; slug: string } | { kind: "cleanTranscribedAudio"; slug: string } | { kind: "cleanExtraFormats"; slug: string } + | { kind: "backfillChannel"; slug: string } | { kind: "refreshReport"; slug: string }; type Status = @@ -41,6 +43,7 @@ const LABEL: Record<Variant["kind"], { idle: string; running: string }> = { redownloadShortAudio: { idle: "Re-download (corrected format)", running: "Queuing…" }, cleanTranscribedAudio: { idle: "Clean audio", running: "Queuing…" }, cleanExtraFormats: { idle: "Clean extra formats", running: "Queuing…" }, + backfillChannel: { idle: "Backfill", running: "Queuing…" }, refreshReport: { idle: "Refresh report", running: "Refreshing…" }, }; @@ -80,6 +83,9 @@ async function runAction(variant: Variant): Promise<StreamActionResult> { if (variant.kind === "cleanExtraFormats") { return cleanExtraAudioFormatsAction(variant.slug); } + if (variant.kind === "backfillChannel") { + return backfillChannelAction(variant.slug); + } // refreshReport: returns an ActionResult, not a StreamActionResult — adapt. const result = await refreshChannelSnapshotAction(variant.slug); if (result && "error" in result) { diff --git a/editor/app/actionable/lib/loadActionable.ts b/editor/app/actionable/lib/loadActionable.ts @@ -35,6 +35,7 @@ export type ActionableSummary = { cleanExtraFormats: ActionableRow[]; staleOrMissing: ActionableRow[]; digestWarnings: ActionableRow[]; + backfill: ActionableRow[]; duplicates: DuplicateReport | null; // The human decisions kept alongside the report — a cluster's canonical // choice, "not a duplicate", and the `confirmed` flag that is the only thing @@ -122,6 +123,27 @@ export function actionableDigestWarningsCount(row: ActionableRow): number { return row.snapshot?.buckets.digestWarnings?.length ?? 0; } +// The backfill lane's two numbers, and they are two FUNCTIONS on purpose so no +// caller can accidentally add them. +// +// `reachable` (missing + stale) is what the lane can do today and the only thing +// that decides whether a channel appears in the section at all. `missingInput` +// is the population that needs its media re-acquired first — measured at ~91x +// the reachable count corpus-wide, so filtering on it would put every channel in +// the list forever. That is not a hypothetical: it is the documented reason +// /api/widget/actionable refuses to filter on `noDigest`. +export function actionableBackfillCount(row: ActionableRow): number { + const entries = Object.values(row.snapshot?.backfill ?? {}); + return entries.reduce((n, e) => n + e.missing + e.stale, 0); +} + +export function actionableBackfillMissingInputCount( + row: ActionableRow, +): number { + const entries = Object.values(row.snapshot?.backfill ?? {}); + return entries.reduce((n, e) => n + e.missingInput, 0); +} + // Estimated bytes each cleanup would reclaim (default 0 for snapshots written // before cleanupBytes existed). export function actionableCleanTranscribedBytes(row: ActionableRow): number { @@ -201,6 +223,10 @@ export async function loadActionableSummary( actionableDigestWarningsCount(b) - actionableDigestWarningsCount(a), ); + const backfill = rows + .filter((r) => actionableBackfillCount(r) > 0) + .sort((a, b) => actionableBackfillCount(b) - actionableBackfillCount(a)); + const staleOrMissing = rows .filter(isStaleOrMissing) .sort((a, b) => a.channel.slug.localeCompare(b.channel.slug)); @@ -216,6 +242,7 @@ export async function loadActionableSummary( cleanExtraFormats, staleOrMissing, digestWarnings, + backfill, duplicates, duplicateOverrides, }; diff --git a/editor/app/actionable/page.tsx b/editor/app/actionable/page.tsx @@ -7,6 +7,8 @@ import { actionableCleanExtraFormatsCount, actionableCleanTranscribedBytes, actionableCleanTranscribedCount, + actionableBackfillCount, + actionableBackfillMissingInputCount, actionableDigestWarningsCount, actionableIncompleteTranscriptCount, actionableMissingNeverFetchedCount, @@ -61,7 +63,11 @@ export default async function ActionablePage() { summary.cleanTranscribedAudio.length === 0 && summary.cleanExtraFormats.length === 0 && summary.staleOrMissing.length === 0 && - summary.digestWarnings.length === 0; + summary.digestWarnings.length === 0 && + // HAND-EXHAUSTIVE, and a section missing from here is a section that never + // renders — the whole page collapses to "nothing pending" while its own + // list has rows in it. + summary.backfill.length === 0; const sections: { config: SectionConfig; rows: ActionableRow[] }[] = [ { @@ -176,6 +182,31 @@ export default async function ActionablePage() { }, { config: { + id: "backfill", + title: "Channels missing derived data the corpus predates", + description: + "Videos an enabled backfill feature has nothing on disk for — no record, or one produced by a different engine/model/threshold than the current settings. Run \u201cBackfill\u201d to catch them up. The count is what the lane can do TODAY; the second column is the separate population whose source media has already been deleted, which needs the opt-in re-download to reach at all.", + countLabel: "to backfill", + emptyLabel: "Nothing pending.", + getCount: actionableBackfillCount, + // Same treatment "Est. reclaim" gets, and for a stronger reason: these + // two numbers differ by ~91x on the real corpus, so a single total would + // be dominated by work no button on this page can start. + extraColumn: { + label: "Needs media", + getValue: (r) => + actionableBackfillMissingInputCount(r).toLocaleString(), + }, + primaryAction: (r) => ( + <InlineActionButton + variant={{ kind: "backfillChannel", slug: r.channel.slug }} + /> + ), + }, + rows: summary.backfill, + }, + { + config: { id: "short-audio", title: "Channels with truncated downloads (short audio)", description: diff --git a/editor/app/api/test/resume-backfill-sweep/route.ts b/editor/app/api/test/resume-backfill-sweep/route.ts @@ -0,0 +1,35 @@ +import { NextResponse } from "next/server"; +import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; +import { + getBackfillSweepJobId, + resumeBackfillSweepIfEnabled, +} from "yt-dlp-transcript-common/controller/backfillSweep"; + +export const dynamic = "force-dynamic"; + +// E2E test harness only. Simulates a SERVER RESTART for the backfill sweep, +// which is the one property of an armed sweep that cannot otherwise be tested: +// the sweep is process state, a restart destroys it, and instrumentation.ts's +// boot hook is what brings it back. An e2e suite cannot restart the dev server +// mid-run, so it does the two halves a restart does — drop the in-process +// singleton, then call the hook — and asserts a sweep comes back. +// +// Mounted unconditionally, like the other /api/test routes: the editor is a +// localhost admin tool, not a deployed service. +export async function GET() { + // Half one: lose the live sweep, the way a restart does. Cancelling the job + // clears the singleton through the sweep's own `finally`, so this leaves + // exactly the state a fresh process starts in — the persisted settings flag + // and nothing running. + const before = getBackfillSweepJobId(); + if (before) getRegistry().cancel(before); + + // Half two: the boot hook. Reads settings.backfill.sweepEnabled and its + // persisted scope, and re-launches from those alone. + await resumeBackfillSweepIfEnabled(); + + return NextResponse.json({ + cancelled: before, + resumed: getBackfillSweepJobId(), + }); +} diff --git a/editor/app/api/widget/sync/route.ts b/editor/app/api/widget/sync/route.ts @@ -3,6 +3,7 @@ import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { digestCountOf } from "yt-dlp-transcript-common/controller/channels"; import { getChannelBriefs } from "../../../lib/requestCache"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; +import { laneBackfillKinds } from "yt-dlp-transcript-common/lib/backfillKinds"; import { buildScheduleView } from "yt-dlp-transcript-common/jobs/syncScheduler"; import { readSchedulerState } from "yt-dlp-transcript-common/jobs/syncSchedulerState"; @@ -38,6 +39,26 @@ export type WidgetSyncPayload = { paused: boolean; // settings.digest.digestsPaused sweeping: boolean; // a corpus-wide sweep is armed }; + // Corpus-wide backfill state. Same constraint as `digest` above — scalars only + // — and free for the same reason: every brief already carries its channel's + // snapshot, so this is a sum rather than a corpus walk (which cost 4.4 s on an + // endpoint the widget polls). + // + // `reachable` and `needsMedia` are separate FIELDS, not a total, because on the + // measured corpus they are 835 and ~76,270. A single number here would report + // a backfill as barely begun forever, no matter how much of the reachable work + // was finished. + backfill: { + reachable: number; // missing + stale: what the lane can do now + needsMedia: number; // missing-input: needs an opt-in re-download first + videos: number; // videos in the corpus (the denominator) + enabled: boolean; // settings.backfill.enabled + sweeping: boolean; // a corpus-wide backfill sweep is armed + // Whether any backfill FEATURE is on. Distinct from `enabled`, which is the + // lane: with no feature on there is nothing to report at all, and that must + // not look like "all caught up". + anyKind: boolean; + }; }; // Build the widget sync payload. Exported so the dashboard cockpit can seed its @@ -83,11 +104,19 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> { let digested = 0; let videos = 0; let channelsWithAny = 0; + // Summed off the same snapshots, in the same pass. Kept apart all the way + // through — see the payload type. + let backfillReachable = 0; + let backfillNeedsMedia = 0; for (const c of channels) { videos += c.snapshot?.totals.videos ?? 0; const n = digestCountOf(c.snapshot); digested += n; if (n > 0) channelsWithAny++; + for (const entry of Object.values(c.snapshot?.backfill ?? {})) { + backfillReachable += entry.missing + entry.stale; + backfillNeedsMedia += entry.missingInput; + } } return { @@ -106,6 +135,14 @@ export async function buildWidgetSyncPayload(): Promise<WidgetSyncPayload> { paused: settings.digest.digestsPaused, sweeping: settings.digest.sweepEnabled, }, + backfill: { + reachable: backfillReachable, + needsMedia: backfillNeedsMedia, + videos, + enabled: settings.backfill.enabled, + sweeping: settings.backfill.sweepEnabled, + anyKind: laneBackfillKinds(settings).length > 0, + }, } satisfies WidgetSyncPayload; } diff --git a/editor/app/channels/[slug]/backfillActions.ts b/editor/app/channels/[slug]/backfillActions.ts @@ -0,0 +1,65 @@ +"use server"; + +import { revalidatePath } from "next/cache"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { BACKFILL_QUEUE, resolveQueueKey } from "yt-dlp-transcript-common/lib/queueKeys"; +import { runManagedFunction } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; +import { requestChannelSnapshot } from "yt-dlp-transcript-common/jobs/snapshotScheduler"; +import { + countBackfillWork, + runBackfillBatch, +} from "yt-dlp-transcript-common/controller/backfillBatch"; + +// Run the backfill lane over ONE channel, from its stage card. +// +// Deliberately the same kind, queue and progress metric the sweep's per-channel +// job uses (runBackfillChannelJob), so a hand-clicked run and a swept one are +// the same thing seen from /jobs — that is what makes a sweep inspectable with +// the tools that already exist. It is NOT that function, though, because a +// hand-clicked run must be foreground (the sweep's are `background: true`, +// deliberately behind anything an operator asks for by hand) and takes a queue +// override from the card's QueueControl. +export async function backfillChannelAction( + slug: string, + queueKey?: string, + kindIds?: string[], +): Promise<StreamActionResult> { + const paths = getPaths(); + const kind = "backfill-channel"; + return runManagedFunction({ + kind, + queueKey: resolveQueueKey(BACKFILL_QUEUE, queueKey), + paths, + channelSlug: slug, + spec: { kind, slug, params: { queueKey, kindIds } }, + fn: async (onLog, signal, setProgress, ctx) => { + const before = await countBackfillWork(paths, slug, kindIds); + const batch = await runBackfillBatch({ + channelSlug: slug, + paths, + kindIds, + setProgress, + progressBaseline: 0, + progressTarget: before.reachable, + onLog, + signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), + }); + // BOTH numbers, never their sum — see lib/backfillKinds.ts. + onLog( + `Backfill ${slug}: ${batch.succeeded} done, ${batch.fresh} already current, ` + + `${batch.failed} failed; ${batch.missingInput} still need their media re-acquired` + + (batch.reacquired > 0 + ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` + : "") + + (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + + ".", + ); + requestChannelSnapshot(paths, slug); + revalidatePath(`/channels/${slug}`); + }, + }); +} diff --git a/editor/app/channels/[slug]/components/stages/BackfillStage.tsx b/editor/app/channels/[slug]/components/stages/BackfillStage.tsx @@ -0,0 +1,133 @@ +"use client"; + +// The channel-level view of the backfill lane: what is reachable, what would +// need its media re-acquired, and a button to run it. +// +// Modelled on DigestStage — same bucket-count heading, same StreamActionLog + +// QueueControl shape — so the stage rail reads as one system. +// +// THE TWO NUMBERS ARE NEVER ADDED. The heading counts what this button would +// actually do; the muted sub-line reports what it cannot reach without a +// re-download. Measured on the real corpus those are 835 and ~76,270, and a +// single "77,105 remaining" here would make the card useless the day it shipped +// — the same failure /api/widget/actionable already documents for `noDigest`. + +import { useState } from "react"; +import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; +import { QueueControl } from "../../../../components/QueueControl"; +import { cancelJobAction } from "../../../../jobs/actions"; +import { backfillChannelAction } from "../../backfillActions"; +import { VideoIdList } from "../VideoIdList"; + +export type BackfillKindView = { + id: string; + label: string; + // missing + stale: what the lane can do right now. + reachableIds: string[]; + // Videos whose input is gone. A COUNT only — the id list is corpus-sized and + // deliberately not stored in the snapshot (see BackfillSnapshotEntry). + missingInput: number; + stale: number; +}; + +type Props = { + slug: string; + existingQueues: string[]; + defaultQueueKey: string; + kinds: BackfillKindView[]; + // settings.backfill.allowRedownload — surfaced because it is the difference + // between the sub-line being informational and being actionable. + allowRedownload: boolean; + // No backfill feature is switched on at all. The card still renders (the rail + // is fixed) but says so rather than reporting an empty work list as "done". + anyEnabled: boolean; +}; + +export function BackfillStage({ + slug, + existingQueues, + defaultQueueKey, + kinds, + allowRedownload, + anyEnabled, +}: Props) { + const [queue, setQueue] = useState(defaultQueueKey); + + const reachable = kinds.reduce((n, k) => n + k.reachableIds.length, 0); + const missingInput = kinds.reduce((n, k) => n + k.missingInput, 0); + const allReachableIds = [ + ...new Set(kinds.flatMap((k) => k.reachableIds)), + ].sort(); + + return ( + <div aria-label="backfill section" className="flex flex-col gap-3"> + <div> + <h3 className="text-base font-semibold"> + Backfill derived data ({reachable}) + </h3> + <p className="text-sm text-muted-foreground"> + {anyEnabled + ? "Catch-up for derived data this channel's videos predate. A re-run does only what is still missing or stale, so running it twice costs nothing the second time." + : "No backfill is enabled. Turn one on in Settings and this card will report what the existing corpus is missing."} + </p> + {missingInput > 0 && ( + <p + aria-label="backfill needs re-acquiring" + className="mt-1 text-sm text-muted-foreground" + > + {missingInput.toLocaleString()} more{" "} + {missingInput === 1 ? "video needs" : "videos need"} their media + re-acquired first + {allowRedownload + ? " — re-download is on, so this run will fetch and then delete it, bounded by the free-disk floor." + : " — re-download is off, so this run skips them."} + </p> + )} + </div> + + {kinds.length > 1 && + kinds.map((k) => ( + <p + key={k.id} + aria-label={`backfill kind ${k.id}`} + className="text-xs text-muted-foreground" + > + <span className="font-medium">{k.label}</span>:{" "} + {k.reachableIds.length} reachable + {k.stale > 0 && ` (${k.stale} stale)`} ·{" "} + {k.missingInput.toLocaleString()} needing media + </p> + ))} + + <VideoIdList + slug={slug} + ids={allReachableIds} + ariaLabel="videos needing a backfill list" + emptyAriaLabel="videos needing a backfill empty" + emptyMessage={ + anyEnabled + ? "Nothing reachable to backfill." + : "No backfill is enabled." + } + itemAriaLabel={(id) => `video needing a backfill ${id}`} + /> + + <StreamActionLog + trigger={() => backfillChannelAction(slug, queue)} + cancelAction={cancelJobAction} + buttonLabel="Backfill channel" + runningLabel="Backfilling…" + label="Backfill channel" + extraControls={ + <QueueControl + value={queue} + onChange={setQueue} + defaultQueueKey={defaultQueueKey} + existingQueues={existingQueues} + actionLabel="Backfill channel" + /> + } + /> + </div> + ); +} diff --git a/editor/app/channels/[slug]/lib/stageStatus.ts b/editor/app/channels/[slug]/lib/stageStatus.ts @@ -45,6 +45,7 @@ export type StageId = | "transcode" | "transcribe" | "digest" + | "backfill" | "cleanup" | "diagnostics" | "danger"; @@ -80,6 +81,14 @@ const JOB_KIND_TO_STAGE: Record<string, StageId> = { "digest-channel-local": "digest", "digest-channel-remote": "digest", "digest-share-cluster": "digest", + "backfill-channel": "backfill", + // The sweep is corpus-wide and carries no channelSlug, so it will not normally + // appear in a channel's running jobs — mapped anyway so that if one ever is + // attributed here it lights the right card rather than none. + "backfill-sweep": "backfill", + // The pre-registry per-channel diarization button lands on the channel queue + // but is the same work this card is about, so it lights this card too. + "diarize-channel": "backfill", "clean-audio-transcribed": "cleanup", "purge-superseded-auto-subs": "cleanup", "clean-extra-audio-formats": "cleanup", @@ -364,6 +373,62 @@ export function computeStageStatuses( }), }; + // The backfill lane's work list, summed across every registered kind. + // + // `pending` counts ONLY the reachable half. The needs-re-acquiring population + // is reported in the summary line and never folded in: it is 91x larger on the + // measured corpus, so counting it would hold every channel permanently amber + // for work that cannot be done without an opt-in re-download — precisely the + // trap /api/widget/actionable documents for `noDigest`. + const backfillRunning = runningByStage.has("backfill"); + const backfillEntries = Object.values(snapshot.backfill ?? {}); + const backfillPending = backfillEntries.reduce( + (n, e) => n + e.missing + e.stale, + 0, + ); + const backfillMissingInput = backfillEntries.reduce( + (n, e) => n + e.missingInput, + 0, + ); + const backfillParts: string[] = []; + if (backfillPending > 0) { + backfillParts.push( + pluralize( + backfillPending, + "video needs derived data", + "videos need derived data", + ), + ); + } + if (backfillMissingInput > 0) { + backfillParts.push( + `${backfillMissingInput.toLocaleString()} needing media re-acquired`, + ); + } + const backfill: StageStatus = { + id: "backfill", + title: "Backfill", + pending: backfillPending, + failed: 0, + running: backfillRunning, + defaultOpen: true, + summary: backfillRunning + ? "Running…" + : backfillEntries.length === 0 + ? "No backfill is enabled." + : backfillParts.length > 0 + ? backfillParts.join(" · ") + : "Everything reachable is current.", + tone: pickTone({ + running: backfillRunning, + pending: backfillPending, + failed: 0, + // Neutral rather than "ok" when nothing is enabled: an empty work list + // because a feature is off is not the same as being finished. + fallback: backfillEntries.length === 0 ? "neutral" : "ok", + }), + }; + const cleanupRunning = runningByStage.has("cleanup"); const cleanupParts: string[] = []; if (cleanupPending > 0) { @@ -445,6 +510,7 @@ export function computeStageStatuses( transcode, transcribe, digest, + backfill, cleanup, diagnostics, danger, diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx @@ -51,6 +51,7 @@ import { TRANSCRIPTION_QUEUE, } from "yt-dlp-transcript-common/lib/platform"; import { + BACKFILL_QUEUE, DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE, } from "yt-dlp-transcript-common/lib/queueKeys"; @@ -74,6 +75,8 @@ import { PlaylistStage } from "./components/stages/PlaylistStage"; import { TranscodeStage } from "./components/stages/TranscodeStage"; import { TranscribeStage } from "./components/stages/TranscribeStage"; import { DigestStage } from "./components/stages/DigestStage"; +import { BackfillStage } from "./components/stages/BackfillStage"; +import { laneBackfillKinds } from "yt-dlp-transcript-common/lib/backfillKinds"; import { VideoListPane } from "./components/VideoListPane"; import { VideoListPaneSection } from "./components/VideoListPaneSection"; import { VideoPanel, type VideoFile } from "./videos/[id]/components/VideoPanel"; @@ -311,6 +314,8 @@ export default async function ChannelDetailPage({ }); const transcodeApplies = config.handling === "transcribe" && !!config.audioFormat; + // Enabled lane backfills, for the stage card. A settings read, no I/O. + const backfillKinds = laneBackfillKinds(settings); const stageOrder: StageId[] = [ "configure", @@ -319,6 +324,7 @@ export default async function ChannelDetailPage({ ...(transcodeApplies ? (["transcode"] as const) : []), "transcribe", "digest", + "backfill", "cleanup", "diagnostics", "danger", @@ -403,6 +409,29 @@ export default async function ChannelDetailPage({ remoteEnabled={settings.digest.remoteEnabled} /> ), + backfill: ( + <BackfillStage + slug={slug} + existingQueues={existingQueues} + defaultQueueKey={BACKFILL_QUEUE} + // Labels and hints come from the registry HERE, on the server: it reads + // the filesystem and calls controllers, so a client component must never + // import it. Only ENABLED kinds are passed — a disabled feature reports + // no backfill anywhere. + kinds={backfillKinds.map((kind) => { + const entry = snapshot.backfill?.[kind.id]; + return { + id: kind.id, + label: kind.label, + reachableIds: entry?.ids ?? [], + missingInput: entry?.missingInput ?? 0, + stale: entry?.stale ?? 0, + }; + })} + allowRedownload={settings.backfill.allowRedownload} + anyEnabled={backfillKinds.length > 0} + /> + ), cleanup: ( <CleanupStage slug={slug} diff --git a/editor/app/components/dashboard/PipelineBand.tsx b/editor/app/components/dashboard/PipelineBand.tsx @@ -9,6 +9,7 @@ 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 { BackfillSweepControls } from "../../jobs/components/BackfillSweepControls"; import { syncAllChannelsAction, type SyncAllResult } from "../../channels/actions"; import { fmtTime } from "../../widget/lib/relativeTime"; @@ -45,6 +46,7 @@ export function PipelineBand({ digest && digest.videos > 0 ? (digest.digested / digest.videos) * 100 : null; + const backfill = sync?.backfill ?? null; const lastSyncText = sync == null @@ -120,6 +122,36 @@ export function PipelineBand({ label={<span className="text-destructive">downloads paused</span>} /> )} + {backfill?.anyKind && ( + <Instrument + dotClass={ + backfill.sweeping + ? "bg-success" + : backfill.enabled + ? "bg-success/40" + : "bg-muted-foreground/40" + } + label={ + <> + {"\u21ba backfill "} + <span className="font-medium"> + {backfill.reachable.toLocaleString()} + </span> + {" reachable"} + {backfill.needsMedia > 0 && ( + <span className="text-muted-foreground"> + {" "} + (+{backfill.needsMedia.toLocaleString()} need media) + </span> + )} + {!backfill.enabled && ( + <span className="text-muted-foreground"> · lane off</span> + )} + {backfill.enabled && backfill.sweeping && " · sweeping"} + </> + } + /> + )} {digest && ( <Instrument dotClass={ @@ -164,6 +196,13 @@ export function PipelineBand({ paused={digest?.paused ?? false} onChange={onSynced} /> + {backfill?.anyKind && ( + <BackfillSweepControls + sweeping={backfill.sweeping} + laneEnabled={backfill.enabled} + onChange={onSynced} + /> + )} <SyncAllButton onSynced={onSynced} /> <Link href="/channels/new" diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts @@ -12,6 +12,10 @@ import { startDigestSweep, stopDigestSweep, } from "yt-dlp-transcript-common/controller/digestSweep"; +import { + startBackfillSweep, + stopBackfillSweep, +} from "yt-dlp-transcript-common/controller/backfillSweep"; import { runJobSpec } from "./runJobSpec"; import { buildQueueView } from "./queue/buildQueueView"; @@ -236,6 +240,40 @@ export async function stopDigestSweepAction(): Promise<DigestSweepResult> { } } +// Arm / disarm the corpus-wide BACKFILL sweep. Same pair as the digest sweep +// above and for the same reason it is a pair: the flag persists, so the boot +// hook resumes it, and stopping has to clear the scope as well as the flag. +export type BackfillSweepResult = { + ok: boolean; + jobId?: string; + error?: string; +}; + +export async function startBackfillSweepAction( + kindIds?: string[], + channelSlugs?: string[], +): Promise<BackfillSweepResult> { + try { + const jobId = await startBackfillSweep({ kindIds, channelSlugs }); + 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 stopBackfillSweepAction(): Promise<BackfillSweepResult> { + try { + await stopBackfillSweep(); + 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/components/BackfillSweepControls.tsx b/editor/app/jobs/components/BackfillSweepControls.tsx @@ -0,0 +1,96 @@ +"use client"; + +import { useEffect, useState, useTransition } from "react"; +import { useRouter } from "next/navigation"; +import { + startBackfillSweepAction, + stopBackfillSweepAction, +} from "../actions"; + +// Arm / disarm the corpus-wide backfill sweep. +// +// One button, not two — unlike DigestSweepControls, which pairs a sweep switch +// with a pause. The backfill lane has no separate pause because it does not need +// one: at the default weight it is ALREADY idle-only, standing aside whenever +// transcription works, and the lane switch in Settings is the hold. Adding a +// second control that looks like the digest pause would invite the confusion +// that pair exists to prevent. +// +// Stopping drains rather than cancels: the video in flight finishes instead of +// being thrown away, and the stop reaches the per-channel job the sweep is +// waiting on rather than meaning "after this channel". +export function BackfillSweepControls({ + sweeping, + laneEnabled, + onChange, +}: { + sweeping: boolean; + // settings.backfill.enabled. A sweep can be armed with the lane off — the + // batch then holds at a zero limit rather than doing work — so say so instead + // of letting it look wedged. + laneEnabled: 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 — 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 backfill sweep" : "start backfill sweep"} + onClick={() => + run( + sweeping + ? stopBackfillSweepAction + : () => startBackfillSweepAction(), + ) + } + title={ + sweeping + ? "Stop the corpus-wide backfill. The channel in flight finishes first; nothing already written is lost." + : "Start the corpus-wide backfill: every channel in turn, heaviest first by reachable work. 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 Backfill Sweep" : "Start Backfill Sweep"} + </button> + {sweeping && !laneEnabled && ( + <span + aria-label="backfill lane off" + className="text-xs text-muted-foreground" + > + lane off in Settings — the sweep is holding + </span> + )} + {error && ( + <span role="alert" className="text-xs text-destructive"> + {error} + </span> + )} + </div> + ); +} diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx @@ -227,12 +227,14 @@ const TASK_KIND_VERB: Record<JobTaskKind, string> = { download: "Downloading", transcribe: "Transcribing", digest: "Digesting", + backfill: "Backfilling", }; const METRIC_FILL_BY_TASK: Record<JobTaskKind, string> = { download: "bg-success/60", transcribe: "bg-success", digest: "bg-info", + backfill: "bg-warning", }; function TaskProgressBar({ task }: { task: RunningJobsTask }) { @@ -325,12 +327,14 @@ const METRIC_LABELS: Record<JobProgressMetric, string> = { downloads: "Downloads", transcripts: "Transcripts", digests: "Digests", + backfills: "Backfill", }; const METRIC_FILL: Record<JobProgressMetric, string> = { downloads: "bg-success/60", transcripts: "bg-success", digests: "bg-info", + backfills: "bg-warning", }; function JobProgressBar({ diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts @@ -40,6 +40,7 @@ import { digestChannelAction, type DigestLaneChoice, } from "../channels/[slug]/digestActions"; +import { backfillChannelAction } from "../channels/[slug]/backfillActions"; import { persistKeptAction } from "../channels/[slug]/persistActions"; import { checkPostAvailabilityAction, @@ -271,6 +272,17 @@ export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = { const { queueKey } = params(spec); return diarizeChannelAction(spec.slug, queueKey); }, + // Replays correctly with nothing remembered: the batch re-derives its whole + // work-list from disk, so a replay does what is missing NOW rather than what + // was missing when the record was written. The kind scope is carried through + // so a bookmarked single-kind run stays single-kind. + "backfill-channel": (spec) => { + const { p, queueKey } = params(spec); + const kindIds = Array.isArray(p.kindIds) + ? p.kindIds.filter((k): k is string => typeof k === "string") + : undefined; + return backfillChannelAction(spec.slug, queueKey, kindIds); + }, "check-kept-deleted": (spec) => { const { queueKey } = params(spec); return checkKeptDeletedAction(spec.slug, queueKey); diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts @@ -352,6 +352,24 @@ export async function saveSettingsAction( } : dDiar; + // Backfill lane. Same hidden-marker discipline as the two blocks above, and + // here it protects two things specifically: `allowRedownload`, which an + // unrelated save must never flip ON (it writes media to a 97%-full disk), and + // `sweepEnabled` + its scope, which this form does NOT render at all — those + // are owned by the sweep controls, and reading them from an absent form field + // would disarm a running multi-day sweep on any settings save. + const dBack = getSettings().backfill; + const backfillFormPresent = formData.get("backfillFormPresent") === "1"; + const backfillSettings: SiteSettings["backfill"] = backfillFormPresent + ? { + ...dBack, + enabled: formData.get("backfillEnabled") === "on", + weight: num("backfillWeight", dBack.weight), + concurrency: num("backfillConcurrency", dBack.concurrency), + allowRedownload: formData.get("backfillAllowRedownload") === "on", + } + : dBack; + const dB = defaultBuildPipeline(); const buildModeRaw = String(formData.get("buildMode") ?? "").trim(); const buildPipeline = { @@ -407,6 +425,7 @@ export async function saveSettingsAction( // save must not reset them (and must never silently flip remoteEnabled on). digest: digestSettings, diarization: diarizationSettings, + backfill: backfillSettings, }; try { await writeSettings(next); diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx @@ -798,6 +798,73 @@ export function SettingsForm({ initial, apps, digestApps }: Props) { /> </fieldset> <fieldset className="flex flex-col gap-3 border border-border rounded p-3"> + <legend className="px-1 text-sm font-medium">Backfill lane</legend> + {/* + The marker again, and here it guards two things an unrelated save must + never touch: `allowRedownload` (which writes media to a nearly-full + disk) and the sweep flag + scope, which this form does not render at + all — reading those from absent fields would disarm a running + multi-day sweep on any settings save. + */} + <input type="hidden" name="backfillFormPresent" value="1" readOnly /> + <p className="text-xs text-muted-foreground"> + Catch-up for derived data the existing corpus predates. Each feature + declares what it needs and how to tell whether a video has it; this + section decides how much of the machine the catch-up may use. + </p> + <label className="flex items-start gap-2 text-sm"> + <input + type="checkbox" + name="backfillEnabled" + defaultChecked={initial.backfill.enabled} + className="mt-1" + /> + <span className="flex flex-col gap-1"> + <span className="font-medium">Run the backfill lane</span> + <span className="text-xs text-muted-foreground"> + Off means the lane still <em>reports</em> what is missing — + that is what the indicators are for — but runs nothing. + </span> + </span> + </label> + <Field + label="Resource share" + name="backfillWeight" + defaultValue={String(initial.backfill.weight)} + type="number" + step="0.05" + hint="0 (the default) means idle-only: the lane runs only while transcription is quiet, and stands aside the moment it isn't. Above 0 it takes that fraction of its slots as a guaranteed share, floored at 1 — so a small number is a slow lane, not a stopped one." + /> + <Field + label="Lane concurrency" + name="backfillConcurrency" + defaultValue={String(initial.backfill.concurrency)} + type="number" + hint="Slots the lane may use when it is not standing aside. Default 1 — this is CPU-bound work competing with GPU feeding and the digest sweep for the same threads." + /> + <label className="flex items-start gap-2 text-sm"> + <input + type="checkbox" + name="backfillAllowRedownload" + defaultChecked={initial.backfill.allowRedownload} + className="mt-1" + /> + <span className="flex flex-col gap-1"> + <span className="font-medium"> + Re-acquire media that has already been deleted + </span> + <span className="text-xs text-muted-foreground"> + Off by default, and that default is measured: 836 videos still + have media on disk and ~76,270 would need re-downloading — 91× + the reachable work, against 45 GB free. With this on, each file is + fetched, used, and <strong>deleted again immediately</strong>{" "} + (unless the video is marked &quot;do not clean&quot;), and nothing + starts at all when free space is under the disk floor. + </span> + </span> + </label> + </fieldset> + <fieldset className="flex flex-col gap-3 border border-border rounded p-3"> <legend className="px-1 text-sm font-medium">Social links</legend> <p className="text-xs text-muted-foreground"> Default social links shown in every site&apos;s footer. Each site can diff --git a/editor/app/widget/builder/components/SectionMini.tsx b/editor/app/widget/builder/components/SectionMini.tsx @@ -97,6 +97,18 @@ const MINIS: Record<SectionId, () => React.ReactElement> = { <ListRow chip="bg-info/60" /> </> ), + // A dot plus two rules of different lengths — the strip's shape, and a + // deliberate likeness of its content: two numbers that are NOT the same size + // and are never added together. + backfill: () => ( + <span className="flex items-center gap-[3px]"> + <Dot tone="bg-warning" /> + <span className="flex min-w-0 flex-1 flex-col gap-[2px]"> + <Rule w="45%" tone="bg-warning/60" /> + <Rule w="85%" /> + </span> + </span> + ), }; export function SectionMini({ id }: { id: SectionId }) { diff --git a/editor/app/widget/components/MonitorWidget.tsx b/editor/app/widget/components/MonitorWidget.tsx @@ -225,6 +225,8 @@ export function MonitorWidget({ eta={config.eta} /> ); + case "backfill": + return syncData ? <BackfillStrip data={syncData} /> : null; case "actionable": return ( <ActionableStrip @@ -461,6 +463,41 @@ function SchedulerStrip({ ); } +// The backfill readout: what the lane can reach now, and — as a SEPARATE +// figure, never added to it — how much is only reachable by re-acquiring media +// that has already been deleted. On the measured corpus those are 835 and +// ~76,270, so a single total would report the work as untouched forever. +// +// Renders nothing when no backfill feature is on: an empty work list because a +// feature is switched off must not read as "all caught up". +function BackfillStrip({ data }: { data: WidgetSyncPayload }) { + const b = data.backfill; + if (!b.anyKind) return null; + const done = b.reachable === 0; + return ( + <section + aria-label="Backfill" + className={`flex items-center gap-2 rounded border px-2 py-1 text-xs ${ + done + ? "border-border bg-card text-muted-foreground" + : "border-warning/30 bg-warning-soft text-warning" + }`} + > + <span + className={`inline-block h-2 w-2 shrink-0 rounded-full ${ + b.sweeping ? "bg-success" : done ? "bg-success/40" : "bg-warning" + }`} + /> + <span className="tabular-nums"> + Backfill {b.reachable.toLocaleString()} reachable + {b.needsMedia > 0 && <> · {b.needsMedia.toLocaleString()} need media</>} + {!b.enabled && " · lane off"} + {b.enabled && b.sweeping && " · sweeping"} + </span> + </section> + ); +} + function CleanableStrip({ bytes }: { bytes: number }) { return ( <section @@ -758,12 +795,14 @@ const TASK_KIND_VERB: Record<JobTaskKind, string> = { download: "\u2193", transcribe: "\u270e", digest: "\u00b6", + backfill: "\u21ba", }; const TASK_KIND_FILL: Record<JobTaskKind, string> = { download: "bg-success/60", transcribe: "bg-success", digest: "bg-info", + backfill: "bg-warning", }; // Per-metric glyph for the compact widget line. A Record over JobProgressMetric @@ -773,6 +812,7 @@ const METRIC_PREFIX: Record<JobProgressMetric, string> = { downloads: "\u2193 ", transcripts: "", digests: "\u00b6 ", + backfills: "\u21ba ", }; // One-line textual summary of a job's batch progress, e.g. "↓ 5/10 · ~2m left". diff --git a/editor/app/widget/lib/config.ts b/editor/app/widget/lib/config.ts @@ -64,6 +64,11 @@ export type WidgetConfig = { // Show the auto-sync scheduler status strip (on/off, next due, last run). Off // by default. scheduler: boolean; + // Show the backfill strip: how much derived data the existing corpus is + // missing, split into what the lane can reach now and what needs its media + // re-acquired first. Off by default, like every strip that is only meaningful + // once a feature is turned on. + backfill: boolean; // Readouts show absolute (locale) time instead of relative "5m ago". Off by // default (relative). syncTimeAbsolute: boolean; @@ -103,6 +108,7 @@ const WIDGET_FLAG_DEFAULTS: SectionFlags = { sync: false, lastSync: false, scheduler: false, + backfill: false, syncTimeAbsolute: false, syncConfirm: false, stackNarrow: false, @@ -163,6 +169,7 @@ export function parseWidgetConfig(params: RawParams): WidgetConfig { sync: parseBool(params.sync, WIDGET_DEFAULTS.sync), lastSync: parseBool(params.lastsync, WIDGET_DEFAULTS.lastSync), scheduler: parseBool(params.sched, WIDGET_DEFAULTS.scheduler), + backfill: parseBool(params.backfill, WIDGET_DEFAULTS.backfill), syncTimeAbsolute: parseBool(params.abstime, WIDGET_DEFAULTS.syncTimeAbsolute), syncConfirm: parseBool(params.syncask, WIDGET_DEFAULTS.syncConfirm), stackNarrow: parseBool(params.stack, WIDGET_DEFAULTS.stackNarrow), @@ -213,6 +220,8 @@ export function buildWidgetQuery(config: WidgetConfig): string { sp.set("wnames", config.workerLabels ? "1" : "0"); if (config.controls !== WIDGET_DEFAULTS.controls) sp.set("controls", config.controls ? "1" : "0"); + if (config.backfill !== WIDGET_DEFAULTS.backfill) + sp.set("backfill", config.backfill ? "1" : "0"); if (config.actionable !== WIDGET_DEFAULTS.actionable) sp.set("act", config.actionable ? "1" : "0"); if (config.cleanChannels !== WIDGET_DEFAULTS.cleanChannels) diff --git a/editor/app/widget/lib/sections.ts b/editor/app/widget/lib/sections.ts @@ -21,7 +21,11 @@ export type SectionId = | "cleanChannels" | "workers" | "jobs" - | "actionable"; + | "actionable" + // APPENDED, never inserted: the order of this union feeds DEFAULT_ORDER, and + // DEFAULT_ORDER is what a link with no `l=` renders. Inserting here would + // silently re-arrange every already-copied widget link. + | "backfill"; // The visibility/display flags a section reads, i.e. the whole config minus the // placement it is being placed into. Taking this rather than WidgetConfig lets @@ -157,6 +161,17 @@ export const SECTIONS: SectionDef[] = [ setEnabled: (on) => ({ actionable: on }), options: [], }, + // APPENDED at the end — see the SectionId union. The `code` is part of the URL + // contract too, so it must never be reused or renamed. + { + id: "backfill", + code: "bf", + label: "Backfill", + hint: "Derived data the existing corpus is missing, and how much of it needs its media back.", + enabled: (c) => c.backfill, + setEnabled: (on) => ({ backfill: on }), + options: [], + }, ]; // The order a widget renders in when no layout is set. Unchanged from the diff --git a/editor/e2e/backfill.spec.ts b/editor/e2e/backfill.spec.ts @@ -0,0 +1,586 @@ +import { mkdir, readdir, rm, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import type { APIRequestContext } from "@playwright/test"; +import { baseUrl } from "./baseUrl"; +import { + pathExists, + readJson, + resetData, + resolvePath, + writeSettings, + generateReport, +} from "./helpers"; + +// Active jobs, straight off the JSON API. Reloading /jobs costs hundreds of +// milliseconds per sample, which is wider than the overlap window some of these +// assertions are trying to catch. +async function activeJobs( + request: APIRequestContext, +): Promise<{ kind: string; status: string; queueKey: string }[]> { + const res = await request.get(`${baseUrl}/api/jobs/active`); + const body = await res.json(); + return Array.isArray(body) ? body : (body.jobs ?? []); +} + +// The backfill lane: a declared kind, a lane with a fair share, and a visible +// remainder. +// +// What these specs are for, in order of how much they would cost to get wrong: +// +// 1. THE RE-DOWNLOAD MODE MUST NOT LEAK. It writes media onto a disk with +// ~45 GB free at 97% full, for a population of ~76,000 videos. A file left +// behind after a failed item fills the disk. Three specs pin it: gone after +// success, gone after FAILURE, and kept only when the operator said so. +// 2. THE TWO NUMBERS MUST STAY APART. Reachable work and needs-re-acquiring +// differ by ~91x on the real corpus; a surface that adds them is useless +// the day it ships. +// 3. THE LANE MUST BE CONCURRENT with transcription — that is the entire +// reason it has its own queueKey — and must survive a restart. +// +// The diarization engine is faked via DIARIZE_BIN (e2e/fixtures/bin/ +// fake-diarize.mjs), as in diarization.spec.ts. + +const SLUG = "test-transcribe"; + +function dataRel(videoId: string, file: string): string { + return `test-transcripts/channels/${SLUG}/data/${videoId}/${file}`; +} + +const BASE_SETTINGS = { + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + verifyAvailabilityBeforeClean: false, + syncScheduler: { fullSweepIntervalMinutes: 0 }, +}; + +function backfillSettings(over: { + diarization?: Record<string, unknown>; + backfill?: Record<string, unknown>; + minFreeDiskGB?: number; +} = {}) { + return { + ...BASE_SETTINGS, + ...(over.minFreeDiskGB !== undefined + ? { minFreeDiskGB: over.minFreeDiskGB } + : {}), + diarization: { + enabled: true, + // OFF: the whole point of a backfill lane is that the capture did not + // happen inline, so these specs must not have it quietly done for them. + inlineAfterTranscribe: false, + threshold: 0.5, + threads: 1, + python: "python3", + // Any non-empty pair — an empty one reports "not configured" and the lane + // never runs. + segModel: "/dev/null", + embModel: "/dev/null", + concurrency: 1, + ...over.diarization, + }, + backfill: { + enabled: true, + // 1, not the idle-only 0: these specs assert that work HAPPENS, and at + // weight 0 a transcription running in the same test would legitimately + // park the lane forever. The idle-only default is covered by the pure + // unit test (common/controller/backfillBatch.test.ts), which is the right + // place for it — no pool, no GPU, no timing. + weight: 1, + concurrency: 1, + sweepEnabled: false, + sweepKinds: [], + sweepChannels: [], + allowRedownload: false, + ...over.backfill, + }, + }; +} + +async function seedTranscript(videoId: string): Promise<void> { + await writeFile( + resolvePath(dataRel(videoId, "transcript.json")), + '{"transcription":[{"text":"hi"}]}\n', + ); +} + +// Real audio files in a video dir — what the re-download specs measure. +async function audioFiles(videoId: string): Promise<string[]> { + const dir = resolvePath( + `test-transcripts/channels/${SLUG}/data/${videoId}`, + ); + const entries = await readdir(dir).catch(() => [] as string[]); + return entries.filter( + (e) => e.startsWith("audio.") && !e.endsWith(".info.json"), + ); +} + +// settings.backfill.sweepEnabled, straight off disk — the persisted intent the +// boot hook reads. +async function sweepFlag(): Promise<boolean> { + const s = await readJson<{ backfill?: { sweepEnabled?: boolean } }>( + "test-settings.json", + ).catch(() => ({}) as { backfill?: { sweepEnabled?: boolean } }); + return s.backfill?.sweepEnabled ?? false; +} + +const SLOW = 120_000; + +// (1) The stage card reports the two populations SEPARATELY. +test("the stage card separates reachable work from what needs its media back", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + + // vidA keeps its audio (reachable now); vidB is transcribed with its audio + // already cleaned away (needs re-acquiring). That 1-vs-1 split is the corpus's + // 835-vs-76,270 in miniature, and the card must never show "2". + await seedTranscript("vidA"); + await seedTranscript("vidB"); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + + await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { + force: true, + }); + await generateReport(page, SLUG); + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + const section = page.getByLabel("backfill section"); + await expect(section.getByRole("heading")).toContainText( + "Backfill derived data (1)", + ); + await expect(section.getByLabel("backfill needs re-acquiring")).toContainText( + "1", + ); + // The id list is the reachable half only. + await expect( + section.getByLabel("video needing a backfill vidA"), + ).toBeVisible(); + await expect( + section.getByLabel("video needing a backfill vidB"), + ).toHaveCount(0); +}); + +// (2) The lane runs, and only over what it can reach. +test("running the lane captures the reachable video and skips the one with no media", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + await seedTranscript("vidA"); + await seedTranscript("vidB"); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + await expect + .poll(async () => pathExists(dataRel("vidA", "diarization.json")), { + timeout: 60_000, + }) + .toBe(true); + // Never touched: re-download is off, so a video whose input is gone is + // COUNTED, not attempted. + expect(await pathExists(dataRel("vidB", "diarization.json"))).toBe(false); + await expect(page.getByLabel("Backfill channel output")).toContainText( + "need their media re-acquired", + { timeout: 30_000 }, + ); +}); + +// (3) A stale sidecar is redone. This is the freshness comparator end-to-end: +// before it, diarizeOne short-circuited on mere existence and a threshold change +// left the whole corpus looking done. +test("a sidecar from a different threshold is regenerated", async ({ page }) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + await seedTranscript("vidA"); + // Written at a threshold the current settings no longer use. + await writeFile( + resolvePath(dataRel("vidA", "diarization.json")), + JSON.stringify({ + videoId: "vidA", + generatedAt: "2020-01-01T00:00:00.000Z", + speakers: 1, + turns: [{ start: 0, end: 1, speaker: 0 }], + engine: { + engine: "fake-diarize", + segmentationModel: "null", + embeddingModel: "null", + threshold: 0.1, + }, + }) + "\n", + ); + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + await expect + .poll( + async () => { + const r = await readJson<{ engine: { threshold: number } }>( + dataRel("vidA", "diarization.json"), + ).catch(() => null); + return r?.engine.threshold ?? null; + }, + { timeout: 60_000 }, + ) + .toBe(0.5); +}); + +// (4) THE RE-DOWNLOAD, THE SAFE PATH: fetched, used, and gone again. +test("re-acquired media is deleted after a successful backfill", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings( + backfillSettings({ backfill: { allowRedownload: true } }), + ); + await seedTranscript("vidB"); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + expect(await audioFiles("vidB")).toEqual([]); + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + // The work landed… + await expect + .poll(async () => pathExists(dataRel("vidB", "diarization.json")), { + timeout: 60_000, + }) + .toBe(true); + // …and the media it needed did NOT survive it. This is the property the whole + // mode is built around; without it the corpus-wide run fills the disk. + await expect.poll(async () => audioFiles("vidB"), { timeout: 30_000 }).toEqual( + [], + ); +}); + +// (5) THE RE-DOWNLOAD, THE FAILING PATH — the one a `finally` exists for. vidC +// is the fake diarizer's designated failure, so the backfill throws AFTER the +// media has been fetched. The file must still be gone. +test("re-acquired media is deleted even when the backfill fails", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings( + backfillSettings({ backfill: { allowRedownload: true } }), + ); + await seedTranscript("vidC"); + for (const f of await audioFiles("vidC")) { + await rm(resolvePath(dataRel("vidC", f)), { force: true }); + } + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + await expect(page.getByLabel("Backfill channel output")).toContainText( + /failed|Removed re-acquired/, + { timeout: 60_000 }, + ); + // No sidecar (the engine failed) AND no leftover media. Both halves matter: a + // failure that leaves the file behind is exactly the leak. + expect(await pathExists(dataRel("vidC", "diarization.json"))).toBe(false); + await expect.poll(async () => audioFiles("vidC"), { timeout: 30_000 }).toEqual( + [], + ); +}); + +// (6) …unless the operator said to keep it. do-not-clean is the one exception, +// honoured here for the same reason every cleanup controller honours it. +test("re-acquired media is KEPT when the video is marked do-not-clean", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings( + backfillSettings({ backfill: { allowRedownload: true } }), + ); + await seedTranscript("vidB"); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + await writeFile( + resolvePath(dataRel("vidB", "do-not-clean.json")), + JSON.stringify({ setAt: new Date().toISOString(), note: "e2e" }) + "\n", + ); + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + await expect + .poll(async () => pathExists(dataRel("vidB", "diarization.json")), { + timeout: 60_000, + }) + .toBe(true); + // Kept, and SAID SO in the log — an unexplained file on a full disk is how a + // leak gets discovered the hard way. + await expect(page.getByLabel("Backfill channel output")).toContainText( + "do not clean", + { timeout: 30_000 }, + ); + expect((await audioFiles("vidB")).length).toBeGreaterThan(0); +}); + +// (7) The disk floor refuses to START. Not "stops partway" — a preflight, per +// item, so a long run's twentieth video cannot inherit the first one's headroom. +test("the disk floor refuses to re-acquire anything", async ({ page }) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings( + backfillSettings({ + backfill: { allowRedownload: true }, + // Nothing has this much free. + minFreeDiskGB: 100_000_000, + }), + ); + await seedTranscript("vidB"); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + + await page.goto(`/channels/${SLUG}`); + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + await expect(page.getByLabel("Backfill channel output")).toContainText( + "below the", + { timeout: 60_000 }, + ); + // Nothing fetched, nothing produced. + expect(await audioFiles("vidB")).toEqual([]); + expect(await pathExists(dataRel("vidB", "diarization.json"))).toBe(false); +}); + +// (8) The sweep arms, persists its SCOPE with the flag, and disarms clearing +// both. Persisting the flag without the scope is what turns a deliberately +// bounded run into a corpus-wide one on the next restart; leaving the scope +// behind on stop is what silently narrows the NEXT sweep. Both were real bugs in +// the digest sweep, so both are pinned here. +test("the sweep arms with its scope and disarms clearing it", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + + await page.goto("/"); + const start = page.getByRole("button", { name: "start backfill sweep" }); + await expect(start).toBeEnabled({ timeout: 30_000 }); + await start.click(); + + await expect + .poll( + async () => + ( + await readJson<{ backfill?: { sweepEnabled?: boolean } }>( + "test-settings.json", + ).catch(() => ({}) as { backfill?: { sweepEnabled?: boolean } }) + ).backfill?.sweepEnabled ?? false, + { timeout: 30_000 }, + ) + .toBe(true); + + const stop = page.getByRole("button", { name: "stop backfill sweep" }); + await expect(stop).toBeEnabled({ timeout: 30_000 }); + await stop.click(); + + await expect + .poll( + async () => { + const s = await readJson<{ + backfill?: { + sweepEnabled?: boolean; + sweepChannels?: string[]; + sweepKinds?: string[]; + }; + }>("test-settings.json").catch(() => ({}) as Record<string, never>); + const b = s.backfill ?? {}; + return { + enabled: b.sweepEnabled ?? false, + channels: (b.sweepChannels ?? []).length, + kinds: (b.sweepKinds ?? []).length, + }; + }, + { timeout: 30_000 }, + ) + .toEqual({ enabled: false, channels: 0, kinds: 0 }); +}); + +// (8b) …AND SURVIVES A RESTART. This is the reason the flag is persisted at all: +// a sweep is days of work and will outlive several restarts by construction. +// The suite cannot restart the dev server mid-run, so /api/test/ +// resume-backfill-sweep does the two halves a restart does — drop the live +// sweep, then call the boot hook — and the sweep must come back from the +// persisted intent alone. +test("an armed sweep comes back after a restart", async ({ page, request }) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + + await page.goto("/"); + const start = page.getByRole("button", { name: "start backfill sweep" }); + await expect(start).toBeEnabled({ timeout: 30_000 }); + await start.click(); + await expect + .poll(async () => sweepFlag(), { timeout: 30_000 }) + .toBe(true); + + const res = await request.get(`${baseUrl}/api/test/resume-backfill-sweep`); + const body = (await res.json()) as { + cancelled: string | null; + resumed: string | null; + }; + // A NEW job id: the old one is gone (as a restart would leave it) and the hook + // started a fresh sweep rather than the test merely observing the old one. + expect(body.resumed).toBeTruthy(); + expect(body.resumed).not.toBe(body.cancelled); + + // Clean up, so an armed sweep does not bleed into the next spec. + await page.goto("/"); + const stop = page.getByRole("button", { name: "stop backfill sweep" }); + await expect(stop).toBeEnabled({ timeout: 30_000 }); + await stop.click(); + await expect.poll(async () => sweepFlag(), { timeout: 30_000 }).toBe(false); +}); + +// (9) THE WHOLE POINT OF THE SEPARATE QUEUE KEY: a backfill job and a +// transcription job run AT THE SAME TIME. The registry submits every named queue +// at concurrency 1, so if these shared a key one would sit queued behind the +// other and the "lane" would be a lie. +// +// The transcription is deliberately a SLOWOP video (fake-whisper runs ~7s of +// real wall-time for a dir whose name contains "slowop"), because with the +// instant fixtures both jobs finish before anything can observe the overlap — +// the test would be measuring the poller, not the queue. Polled over the JSON +// API rather than by reloading /jobs for the same reason. +test("a backfill runs concurrently with a transcription", async ({ + page, + request, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + // Backfill work: a transcribed video that still has its audio. + await seedTranscript("vidA"); + // Transcription work that LASTS: a fresh dir with audio and no transcript. + const slowDir = resolvePath( + `test-transcripts/channels/${SLUG}/data/slowop1`, + ); + await mkdir(slowDir, { recursive: true }); + await writeFile(`${slowDir}/audio.mp3`, "fake audio slowop1\n"); + + await page.goto(`/channels/${SLUG}`); + // Transcription first — it is the one that has to still be running when the + // backfill lands. + await page.getByRole("button", { name: "Transcribe stage summary" }).click(); + await page + .getByRole("button", { name: "Transcribe missing", exact: true }) + .click(); + await expect + .poll(async () => activeJobs(request).then((j) => + j.some((x) => x.kind === "whisper-all" && x.status === "running"), + ), { timeout: 30_000 }) + .toBe(true); + + await page.getByRole("button", { name: "Backfill stage summary" }).click(); + await page + .getByRole("button", { name: "Backfill channel", exact: true }) + .click(); + + // Both running at once. This is the assertion the separate queueKey exists for. + const seen = await expect + .poll( + async () => { + const jobs = await activeJobs(request); + const backfill = jobs.find((j) => j.kind === "backfill-channel"); + const transcribe = jobs.find((j) => j.kind === "whisper-all"); + // A backfill that is QUEUED while a transcription runs is the exact + // failure mode a shared key would produce, so record it distinctly. + if (backfill?.status === "running" && transcribe?.status === "running") { + // And the reason it can be: two different keys. + expect(backfill.queueKey).not.toBe(transcribe.queueKey); + return "concurrent"; + } + return backfill?.status ?? "absent"; + }, + { timeout: 30_000, intervals: [100] }, + ) + .toBe("concurrent"); + + // SETTLE BEFORE LEAVING. Ending here with a ~7s slowop whisper still running + // leaks its child into the next spec — the teardown reaper catches it, but + // "the anti-orphan layers did not hold" is a failure of this spec's hygiene, + // not a safety net doing its job. + await expect + .poll( + async () => { + const jobs = await activeJobs(request); + return jobs.filter( + (j) => j.kind === "whisper-all" || j.kind === "backfill-channel", + ).length; + }, + { timeout: 60_000, intervals: [250] }, + ) + .toBe(0); + return seen; +}); + +// (10) /actionable lists the channel, keeps the two numbers in separate columns, +// and — the part that is hand-exhaustive and so easy to miss — does not hide the +// whole page behind "nothing pending". +test("actionable shows the backfill section with both numbers", async ({ + page, +}) => { + test.setTimeout(SLOW); + await resetData("one-transcribe-channel-with-audio"); + await writeSettings(backfillSettings()); + await seedTranscript("vidA"); + await seedTranscript("vidB"); + for (const f of await audioFiles("vidB")) { + await rm(resolvePath(dataRel("vidB", f)), { force: true }); + } + await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { + force: true, + }); + await generateReport(page, SLUG); + + await page.goto("/actionable"); + const section = page.getByLabel("backfill", { exact: true }); + await expect(section).toBeVisible(); + await expect(section).toContainText(SLUG); + // The re-acquire count rides in its own column, exactly as "Est. reclaim" + // does — never added to the count beside it. + await expect(section).toContainText("Needs media"); +}); diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -83,4 +83,18 @@ export async function register() { } catch { /* a sweep that fails to resume must not block server readiness */ } + + // Resume the corpus-wide BACKFILL sweep, if one is armed. Same reasoning as + // the digest sweep above, and safe for the same reason: 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. What a restart destroys is the + // process state — "is a sweep running" — and this is what restores it. + try { + const { resumeBackfillSweepIfEnabled } = await import( + "yt-dlp-transcript-common/controller/backfillSweep" + ); + await resumeBackfillSweepIfEnabled(); + } catch { + /* a sweep that fails to resume must not block server readiness */ + } }