// 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. // // Five 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. The shared diskGate (floor + hysteresis) 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. // 5. AUDIO, NOT THE CHANNEL'S DEFAULT. A `handling: "youtube"` channel's own // download is --skip-download --write-subs --write-auto-subs: run with the // channel's stored config, a re-acquire re-fetches the captions the video // already has and lands nothing a diarizer can read. That is measured, not // hypothetical — between 2026-08-22 and 08-26 it spent ~16,000 fetches // across eight channels for zero diarizations. reacquireConfigFor forces // transcribe-handling for the one video; the channel's stored config is // untouched. // // And one thing that is not a guard: EVERY fetch rewrites metadata.info.json // (the prefetch runs before any attempt, and it runs even when the attempt then // fails), which leaves transcript.cues.json stale by mtime. So the download is // followed by a re-normalize on both paths — see refreshCues at its call site. // // There are TWO exceptions to (4), and cleanup() reports which one fired so the // batch can say so in its log rather than leaving it to be discovered as a // mystery on a full disk. // // a. 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. // b. THE HAND-OFF. On a subtitle channel the file we are about to delete is an // `audio.mp3` next to YouTube auto-captions — which is precisely what // `autoQueue.transcription` looks for: the snapshot bucket // `downloadedAutoSubsOnly` is "ASR VTT **and** audio present", and the // snapshot regenerates ~1 s after any unit on the channel finishes. There // is no per-video lock anywhere, so the runner can start whisper on the // very file this cleanup is about to unlink. Rather than race it, hand it // over: keep the audio when the transcription policy WOULD draw this video // from that bucket (decideKeep, below — a fresh disk read at the resume // mark included), or when a transcription is already running on it. // // Nothing is dropped by keeping it. The runner's own transcription leaves // the audio in `transcribedWithAudio`, exactly like any other download, and // it waits there for the operator's Clean-audio sweep — nothing automatic // deletes audio after a transcription. The outcome is the good one: a video // whose only transcript was YouTube ASR gets our own. import { removeMediaFile } from "../lib/mediaTier-server"; import path from "node:path"; import { readdir } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; import { diskGate, evaluateDiskGate, getFreeBytes } from "../lib/diskSpace"; import { isPermanentlyGone } from "../lib/availability"; import { resolveEffectiveAvailability } from "../lib/availability-server"; import { isDoNotClean } from "../lib/doNotClean-server"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; import type { ChannelConfig } from "../lib/channelConfig"; import { detectPlatform, type Platform } from "../lib/platform"; import { isAutoSubsOnly } from "../lib/subtitleProvenance"; import { compileLaneRoot, isChannelPaused, isDefaultChannelPriority, type ChannelPriority, } from "../lib/channelPriority"; import { policyDrawsBucket, type AutoQueuePolicy, } from "../jobs/autoQueuePolicy"; import { isTranscribingVideo } from "./digestYield"; import { isRealAudioFile, isSourceMediaFile, readVideoFiles, } from "../lib/videoStatus"; import { readChannelConfig } from "./channels"; import { findVideoSourceUrl } from "./undownloadedVideos"; import { downloadOneManaged } from "../ytdlp/downloadOneManaged"; import { normalizeTranscript } from "./normalizeTranscript"; import { resolveDiarizableMedia } from "./diarizeOne"; // Why cleanup() left the media on disk. Both are deliberate; neither is a leak. export type KeepReason = // The operator marked this video "do not clean". | "do-not-clean" // A transcription task is running on this video right now. | "in-flight" // The channel is PAUSED for transcription. A hold is never a stop: the audio // is kept for the pause to end, not deleted because it started. | "paused" // autoQueue.transcription would draw this video from downloadedAutoSubsOnly. | "hand-off"; // Why cleanup() went ahead and removed it. Only ever used for the log/test // vocabulary of decideKeep — the batch sees `{ status: "removed" }`. export type RefuseReason = | "no-audio" | "not-auto-subs-only" | "policy-off" | "policy-snoozed" | "no-leaf" | "disk-low"; export type KeepDecision = | { keep: true; reason: KeepReason } | { keep: false; reason: RefuseReason }; export type CleanupOutcome = | { status: "removed" } | { status: "nothing-added" } | { status: "kept"; reason: KeepReason }; export type ReacquireOutcome = // Media was already there — nothing fetched, nothing to clean up. | { status: "present"; cleanup: () => Promise } // Fetched. `cleanup()` removes it again, or reports which of the two // exceptions kept it (do-not-clean, or the hand-off to auto-transcribe). | { status: "fetched"; cleanup: () => Promise } // Free disk is under the configured floor. A refusal to start, not a failure. | { status: "disk-floor"; cleanup: () => Promise } // Deleted / members-only / private, or no resolvable source URL. Re-trying // this video will not help. | { status: "gone"; cleanup: () => Promise } | { status: "failed"; cleanup: () => Promise }; const NOTHING_TO_CLEAN = async (): Promise => ({ status: "nothing-added", }); // A backfill wants AUDIO. A `handling: "youtube"` channel's own download is // --skip-download --write-subs --write-auto-subs: it would re-fetch the captions // the video already has and land nothing a diarizer can read — which is exactly // what happened to ~16,000 videos on eight channels, 2026-08-22 -> 08-26. Same // override, same reason, as autoRunner's replaceAutoSubs unit and // downloadOneManaged's own no-subs fallback: transcribe-handling for this one // video, the channel's stored config untouched. export function reacquireConfigFor(config: ChannelConfig): ChannelConfig { return config.handling === "transcribe" ? config : { ...config, handling: "transcribe" }; } export async function reacquireMediaFor(opts: { paths: Paths; channelSlug: string; videoId: string; videoDir: string; onLog?: (msg: string) => void; signal?: AbortSignal; }): Promise { 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(); // The channel's own volume: a relocated channel re-acquires onto the platter. const gate = await diskGate(opts.paths, settings, { dir: path.join(opts.paths.channelsDir, opts.channelSlug, "data"), }); if (!gate.ok) { log(`Not re-acquiring ${opts.videoId}: ${gate.message}.`); 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 }; } // Guard (5): the config this re-acquire actually downloads with. Passed to // findVideoSourceUrl as well as to downloadOneManaged, the way autoRunner // passes its override to both. resolveCookiePolicy reads cookie fields only, // so it is unaffected either way — same object, for consistency. const downloadConfig = reacquireConfigFor(config); if (downloadConfig !== config) { log( `Re-acquiring media for ${opts.videoId} — subtitle channel, downloading audio (handling override: transcribe).`, ); } const url = await findVideoSourceUrl( opts.paths, opts.channelSlug, opts.videoId, downloadConfig, ); 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})…`); let failed = false; try { await downloadOneManaged({ channelSlug: opts.channelSlug, channelConfig: downloadConfig, paths: opts.paths, videoUrl: url, onLog: (s) => log(s.trimEnd()), signal: opts.signal ?? new AbortController().signal, cookiePolicy: resolveCookiePolicy(settings, downloadConfig), 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. On a // subtitle channel it is guard (5)'s transcribe override that makes there // BE audio to extract in the first place. extractImmediately: true, keepSourceVideoOverride: false, }); } catch (err) { failed = true; log(`Re-acquire ${opts.videoId} failed: ${(err as Error).message}`); } // The fetch rewrote metadata.info.json, which the cues embed and which // isCuesJsonFresh compares against. Without this the very next lane in the // same sweep (attribution) skips the video as no-transcript and the digest // lane defers it, until an operator runs Normalize. Same step transcribeOne // takes after a transcription, for the same reason. Cheap: `fresh` when // nothing moved. Digest freshness is provenance-keyed, so this invalidates // nothing. // // On the FAILURE path too: the metadata prefetch runs before any attempt, so // a download that threw has already moved the mtime the cues are compared // against. That is the 16,081. await refreshCues(opts.videoDir, opts.channelSlug, config, opts.videoId, log); if (failed) { // 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, { paths: opts.paths, channelSlug: opts.channelSlug, platform: detectPlatform(config.url), }), }; } // 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, { paths: opts.paths, channelSlug: opts.channelSlug, platform: detectPlatform(config.url), }), }; } return { status: "fetched", cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log, { paths: opts.paths, channelSlug: opts.channelSlug, platform: detectPlatform(config.url), }), }; } // Re-parse the raw transcript into transcript.cues.json after a fetch has moved // metadata.info.json under it. // // NEVER RETHROWS: this runs on the failure path as well, and a normalize problem // must not replace the real outcome of the re-acquire. It is also deliberately // NOT part of cleanup(): cues.json is a record, not media, and buildCleanup // never touches non-media files. // // Exported for its unit only — reacquireMediaFor's own path needs a network. export async function refreshCues( videoDir: string, channelSlug: string, config: ChannelConfig, videoId: string, log: (msg: string) => void, ): Promise { try { const outcome = await normalizeTranscript({ videoDir, channelSlug, configName: config.name, log, }); if (outcome.status === "wrote") { log(`Re-normalized transcript.cues.json for ${videoId}.`); } } catch (err) { log( `Warning: could not re-normalize transcript.cues.json for ${videoId}: ${(err as Error).message}`, ); } } // Everything the keep decision needs, as data. Pure so every branch can be // asserted without a disk, a registry or a settings file — the I/O that fills it // in lives in buildCleanup below. export type KeepInput = { doNotClean: boolean; // isTranscribingVideo(videoId) — a transcribe task on a running job. transcribing: boolean; // files.audioFiles.length > 0, RE-READ after the fetch. hasAudio: boolean; // isAutoSubsOnly(videoDir, files) — the same rule the snapshot bucket uses. autoSubsOnly: boolean; policy: Pick< AutoQueuePolicy, "enabled" | "replaceAutoSubs" | "snoozeUntil" | "root" >; // The channel-priority document. Two questions are asked of it, and both are // about the tree the runner would ACTUALLY dispatch from — see decideKeep. priority: ChannelPriority; channel: { slug: string; platform: Platform | null }; disk: { freeBytes: number; minFreeDiskGB: number; resumeMarginGB: number }; }; // Keep this re-acquired media, or remove it? // // The order matters and is stated once, here, rather than being spread over the // call site's control flow: // // do-not-clean the operator's marker outranks every policy question // transcribing something is reading the file right now; do not unlink it // !hasAudio nothing a transcriber could use landed — nothing to hand // !autoSubsOnly not the bucket's shape; a hand-off would never be drawn // !policy.enabled the runner is off; audio kept for it would sit forever // snoozeUntil idem, temporarily (a lapsed snooze is already normalized // to null by sanitizePolicy, so non-null means still on) // paused the CHANNEL is held for transcription — KEEP (below) // !policyDrawsBucket no leaf covering this channel draws downloadedAutoSubsOnly // disk the runner itself would refuse to fetch this much // // CHANNEL PRIORITY ENTERS TWICE, AND THE FIRST ONE IS WHY. Once the priority // model says anything, `autoQueue.transcription.root` is COMPILED from it // (lib/channelPriority.ts) and a channel paused for transcription has NO LEAF // in that tree at all. Asked naively, `policyDrawsBucket` would then answer // "no leaf" for exactly the channels the operator just put on hold, and this // function would UNLINK their audio — a pause turning into data loss, silently, // one re-acquired file at a time. So the pause is asked FIRST and it KEEPS: // a hold is never a stop (controller/operationBatch.ts:22-25), and audio kept // for a paused lane is audio waiting for the pause to end. // // The second is the tree itself: while a model exists the runner does not // dispatch from the STORED root, so neither may this decision. The compiled // root for this one channel is the same answer the real compiled tree gives — // every non-paused channel gets one bare leaf plus the trailing catch-all, and // WHICH tier group holds it changes the order, never whether a leaf draws the // bucket — so the compile is done with this slug alone rather than by listing // 68 configs inside a cleanup. // // The disk bar is the RESUME mark (floor + margin), not the floor: a hand-off is // a download the backfill was about to give back, and the download runner resumes // only at resumeBytes — so the backfill must never keep audio the runner would // have refused to fetch. `latched: true` reuses evaluateDiskGate's hysteresis // math on the PURE core; diskGate() in enforce mode mutates a shared module latch // and must never be called from a keep decision. A disabled gate // (minFreeDiskGB: 0) short-circuits inside evaluateDiskGate and never refuses. export function decideKeep(input: KeepInput): KeepDecision { if (input.doNotClean) return { keep: true, reason: "do-not-clean" }; if (input.transcribing) return { keep: true, reason: "in-flight" }; if (!input.hasAudio) return { keep: false, reason: "no-audio" }; if (!input.autoSubsOnly) return { keep: false, reason: "not-auto-subs-only" }; if (!input.policy.enabled) return { keep: false, reason: "policy-off" }; if (input.policy.snoozeUntil != null) { return { keep: false, reason: "policy-snoozed" }; } if (isChannelPaused(input.priority, input.channel.slug, "transcription")) { return { keep: true, reason: "paused" }; } const root = isDefaultChannelPriority(input.priority) ? input.policy.root : compileLaneRoot( "transcription", input.priority, [input.channel.slug], [], ); if ( !policyDrawsBucket( "transcription", { ...input.policy, root }, input.channel, "downloadedAutoSubsOnly", ) ) { return { keep: false, reason: "no-leaf" }; } if (!evaluateDiskGate({ ...input.disk, latched: true }).ok) { return { keep: false, reason: "disk-low" }; } return { keep: true, reason: "hand-off" }; } // Remove exactly the media this fetch added — or report which of the two // exceptions kept it. // // 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. The parakeet resume cache // (/.audio.mp3.parakeet/) matches none of these predicates either, so // it survives both branches — and its cached windows are re-validated against // the audio's duration and segmentation when the same file is fetched again, so // a stale cache cannot corrupt a later transcription. function buildCleanup( videoDir: string, before: Set, videoId: string, log: (msg: string) => void, ctx: { paths: Paths; channelSlug: string; platform: Platform | null }, ): () => Promise { 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 { status: "nothing-added" }; // Thin I/O, then one pure decision. The dir listing is RE-READ here rather // than carried from the caller for the same reason reacquireMediaFor // re-checks resolveDiarizableMedia: the classification that sent us here // predates the fetch, and this is the moment the answer has to be true. const settings = getSettings(); const files = await readVideoFiles(videoDir, { checkUntranscribable: true, }); const decision = decideKeep({ doNotClean: await isDoNotClean(videoDir), transcribing: isTranscribingVideo(videoId), hasAudio: files.audioFiles.length > 0, autoSubsOnly: await isAutoSubsOnly(videoDir, files), policy: settings.autoQueue.transcription, priority: settings.channelPriority, channel: { slug: ctx.channelSlug, platform: ctx.platform }, disk: { freeBytes: await getFreeBytes(ctx.paths.transcriptsDir), minFreeDiskGB: settings.minFreeDiskGB, resumeMarginGB: settings.resumeMarginGB, }, }); if (decision.keep) { const list = added.join(", "); if (decision.reason === "do-not-clean") { log( `Keeping re-acquired media for ${videoId}: marked "do not clean" (${list}).`, ); } else if (decision.reason === "in-flight") { log( `Keeping re-acquired media for ${videoId}: a transcription is running on it (${list}).`, ); } else if (decision.reason === "paused") { log( `Keeping re-acquired media for ${videoId}: ${ctx.channelSlug} is paused for transcription, and a hold is not a stop (${list}).`, ); } else { log( `Keeping re-acquired media for ${videoId}: handed to auto-transcribe, which will replace the auto-captions (${list}).`, ); } return { status: "kept", reason: decision.reason }; } for (const name of added) { // Through its link when the hook tiered it (release 17). await removeMediaFile(videoDir, name).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 { status: "removed" }; }; }