Archilyzer · Source

archilyzer

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

commit dcf1e830a220fa1d05d272fc015803abe7256b60
parent 8955bb94d6a91f0328ce2574d31dc4dc3e6a2187
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 28 Aug 2026 00:14:43 -0400

backfill: re-acquired audio is handed to auto-transcribe, never deleted under it

The previous slice made the backfill re-acquire land an audio.mp3 on an ASR-only
`handling: "youtube"` video, diarize, and delete it again in the item's
`finally`. That file is exactly what `autoQueue.transcription` looks for.

The collision, mechanically. The snapshot bucket `downloadedAutoSubsOnly` is
"ASR VTT AND audio present" (channelSnapshot.ts:997), and the snapshot
regenerates ~1 s after any unit on the channel finishes (autoRunner.ts:871).
There is no per-video lock anywhere — the registry serializes on `queueKey`
only (registry.ts:78-80) — so the runner can start whisper on the very file the
cleanup is about to unlink. Parakeet re-opens the audio once per 480 s window
(scripts/parakeet-stitch.mjs:350), so the unlink surfaces as a mid-run slice
failure, no transcript.json is written, and transcribeOneFromQueue.ts:164-166
appends the id to failed-transcriptions — which the manual per-channel batch
honours permanently (whisperBatch.ts:101-107).

The operator wants the outcome the race was fighting over: re-acquired audio on
an ASR-only video is an opportunity for a real transcript. So this is a
deterministic hand-off, not a lock.

Guard (4) — "a re-fetched file does not survive the item that fetched it" — now
has TWO exceptions, both reported by cleanup() so the batch can name them:

  a. do-not-clean, unchanged.
  b. THE HAND-OFF: `autoQueue.transcription` would draw this video from
     `downloadedAutoSubsOnly`, or a transcription is already running on it.

`decideKeep` is the whole decision as a pure function of eight inputs, and its
order is stated once in a comment rather than spread over control flow:
do-not-clean, in-flight, no audio, not ASR-only, runner off, snoozed, no leaf,
disk. buildCleanup supplies thin I/O around it and re-reads the dir listing HERE
rather than trusting the batch's Candidate — the same "re-checked here"
discipline reacquireMediaFor already applies to resolveDiarizableMedia.

Three things the design had to respect:

- `replaceAutoSubs: false` alone is NOT a refusal. A leaf whose match names
  `downloadedAutoSubsOnly` draws it regardless. So the rule is "no leaf covering
  this channel draws that bucket", where a bucket-less leaf draws
  defaultBucketsForPolicy — the only place the flag enters. That is
  `policyDrawsBucket`, new in autoQueuePolicy.ts, pure, with `matchesChannel`
  exported and widened to Pick<ChannelWork, "slug" | "platform">.
- The in-flight veto reads registry TASKS (digestYield.ts, beside the existing
  lane-level yield), because tasks cover both the auto runner and the manual
  whisper-all batch; getAutoRunnerStatus() would miss the latter. Fails open to
  false — the policy decision is the primary mechanism. Inline transcribe inside
  a download registers no task and stays invisible; accepted.
- The disk bar is the RESUME mark (floor + margin), not the floor:
  evaluateDiskGate({ ...disk, latched: true }).ok on the PURE core. A hand-off is
  a download the backfill was about to give back, and the download runner itself
  resumes only at resumeBytes (autoRunner.ts:663-668) — so the backfill must
  never keep audio the runner would have refused to fetch. diskGate() in enforce
  mode mutates a shared module latch and is never called from a keep decision.

Nothing is lost by keeping the file: nothing automatic deletes audio after a
transcription (cleanAudioFromTranscribed is a manual channel button,
whisperActions.ts:415), so handed-off audio lands in `transcribedWithAudio` and
waits for the operator's Clean-audio sweep like any other download.

The transcription side no longer trusts the engine's error for a file that went
away. `throwIfAudioVanished` re-checks the resolved audio and converts its
absence into TranscribeError("no-audio"), which the queue ALREADY treats as a
skip — no new failure class. It fires at the local rethrow, before the "produced
no output" throw, AND around the remote branch, because remoteTranscribe wraps a
local ENOENT as "transport" (remoteTranscribe.ts:111, :126-133) and a check at
the local rethrow alone would still let the remote path blacklist the video. The
parakeet resume cache is deliberately untouched — .audio.mp3.parakeet/ matches
none of buildCleanup's predicates, and its windows are re-validated against
duration and segmentation when the same audio is fetched again.

Fixed in passing: the `"failed"` outcome has always returned a real partial-file
cleanup that the batch's `finally` never called — it only called cleanup() for
`"fetched"`. It runs now. That is the leak the builder was written for.

The batch counts hand-offs separately (`reacquireHandedOff`, hand-off and
in-flight both) so the reconciliation line stops calling them do-not-clean, and
operationJobs' one-line summary says how many were handed over.

22 new unit tests (835 -> 857): decideKeep one per reason plus both disk edges
and the explicit-leaf case, policyDrawsBucket over all/channel/platform/
operation leaves, evaluateTranscribingVideo's three negatives, and
transcribeOne driven against a shell-script engine that removes its own input.

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

Diffstat:
Mcommon/controller/backfillBatch.ts | 43+++++++++++++++++++++++++++++++++----------
Mcommon/controller/backfillReacquire.test.ts | 223++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/controller/backfillReacquire.ts | 235+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mcommon/controller/digestYield.test.ts | 67+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/digestYield.ts | 45++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/controller/operationJobs.ts | 5++++-
Acommon/controller/transcribeOne.test.ts | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/transcribeOne.ts | 57+++++++++++++++++++++++++++++++++++++++++++++++----------
Mcommon/jobs/autoQueuePolicy.test.ts | 180+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/autoQueuePolicy.ts | 38+++++++++++++++++++++++++++++++++++++-
10 files changed, 912 insertions(+), 54 deletions(-)

diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts @@ -137,6 +137,13 @@ export type BackfillBatchResult = { // gap is a leak and is logged as one. reacquired: number; reacquireCleaned: number; + // Re-acquired media deliberately LEFT on disk for the transcription lane: + // either autoQueue.transcription would draw this video from + // downloadedAutoSubsOnly (the hand-off), or a transcription was already + // running on it. Both mean "left to the transcription lane"; the per-video log + // line distinguishes them. Media kept for do-not-clean is NOT counted here — + // that is still the unexplained-gap case the reconciliation names. + reacquireHandedOff: number; reacquireFailed: number; // True when the run stopped taking new work because free disk fell under the // configured floor. @@ -260,6 +267,7 @@ export async function runBackfillBatch( skipped: 0, reacquired: 0, reacquireCleaned: 0, + reacquireHandedOff: 0, reacquireFailed: 0, diskFloorHit: false, deferred: 0, @@ -684,11 +692,25 @@ export async function runBackfillBatch( } finally { // THE FINALLY THAT KEEPS THE DISK ALIVE. A re-fetched file is removed // whether the backfill succeeded, failed, or threw — on a 97%-full disk a - // leak here fills it. The one exception is a video marked do-not-clean, - // which cleanup() honours and reports. - if (reacquired?.status === "fetched") { - const cleaned = await reacquired.cleanup(); - if (cleaned) result.reacquireCleaned++; + // leak here fills it. Two exceptions, both reported by cleanup(): a video + // marked do-not-clean, and the hand-off to auto-transcribe (see + // backfillReacquire's header). A FAILED re-acquire is cleaned too — a + // download that threw can still have left a partial file, which is exactly + // the leak this exists to stop, and its cleanup was built for that and + // never called until now. + if ( + reacquired?.status === "fetched" || + reacquired?.status === "failed" + ) { + const out = await reacquired.cleanup(); + // Only a "fetched" contributes to the reconciliation below: the other + // statuses were never counted in result.reacquired. + if (reacquired.status === "fetched") { + if (out.status === "removed") result.reacquireCleaned++; + else if (out.status === "kept" && out.reason !== "do-not-clean") { + result.reacquireHandedOff++; + } + } } task?.end(); } @@ -823,13 +845,14 @@ export async function runBackfillBatch( idlePollMs: 3000, }); - // A gap between fetched and cleaned that is not explained by do-not-clean is a - // LEAK, and it is stated in the job log rather than left to be discovered by a - // full disk. + // A gap between fetched and cleaned that is not explained by the hand-off or by + // do-not-clean is a LEAK, and it is stated in the job log rather than left to + // be discovered by a full disk. if (result.reacquired > result.reacquireCleaned) { log( - `Re-acquired ${result.reacquired} file(s), removed ${result.reacquireCleaned}. ` + - `The difference is media kept because its video is marked "do not clean" — ` + + `Re-acquired ${result.reacquired} file(s), removed ${result.reacquireCleaned}, ` + + `handed ${result.reacquireHandedOff} to auto-transcribe. ` + + `Any other difference is media kept because its video is marked "do not clean" — ` + `if that is not what you expect, check the disk.`, ); } diff --git a/common/controller/backfillReacquire.test.ts b/common/controller/backfillReacquire.test.ts @@ -3,7 +3,12 @@ import assert from "node:assert/strict"; import path from "node:path"; import { mkdtemp, rm, writeFile, utimes } from "node:fs/promises"; import { tmpdir } from "node:os"; -import { reacquireConfigFor, refreshCues } from "./backfillReacquire"; +import { + decideKeep, + reacquireConfigFor, + refreshCues, + type KeepInput, +} from "./backfillReacquire"; import { isCuesJsonFresh } from "./normalizeTranscript"; import { CUES_JSON_FILENAME, @@ -136,3 +141,219 @@ test("refreshCues never throws — it runs on the download-failure path too", as await rm(dir, { recursive: true, force: true }); } }); + +// --------------------------------------------------------------------------- +// decideKeep — keep the re-acquired audio for the transcription lane, or remove +// it? Pure, so every branch is asserted here and the call site is left holding +// nothing but thin I/O. +// +// The stake: on the "remove" side this file is the only thing standing between a +// 76,000-video sweep and a full disk; on the "keep" side, deleting an audio file +// that autoQueue.transcription has already started whisper on costs that video's +// whole run AND appends it to failed-transcriptions, which the manual +// per-channel batch honours permanently. + +const GB = 1024 * 1024 * 1024; + +const ALL_LEAF = { + id: "root", + mode: "strict" as const, + children: [{ id: "leaf-all", match: { type: "all" as const } }], +}; + +// The live configuration this was written for: enabled, replaceAutoSubs on, one +// catch-all leaf. +function keepInput(over: Partial<KeepInput> = {}): KeepInput { + return { + doNotClean: false, + transcribing: false, + hasAudio: true, + autoSubsOnly: true, + policy: { + enabled: true, + replaceAutoSubs: true, + snoozeUntil: null, + root: ALL_LEAF, + }, + channel: { slug: "the-channel", platform: "youtube" }, + disk: { freeBytes: 100 * GB, minFreeDiskGB: 20, resumeMarginGB: 5 }, + ...over, + }; +} + +test("decideKeep: the operator's do-not-clean outranks every policy question", () => { + // Everything else says "remove" — no audio, not ASR-only, runner off. + const d = decideKeep( + keepInput({ + doNotClean: true, + hasAudio: false, + autoSubsOnly: false, + policy: { + enabled: false, + replaceAutoSubs: false, + snoozeUntil: null, + root: ALL_LEAF, + }, + }), + ); + assert.deepEqual(d, { keep: true, reason: "do-not-clean" }); +}); + +test("decideKeep: a running transcription keeps the file whatever the policy says", () => { + // The in-flight veto is defense in depth: it covers the manual whisper-all + // batch and any lane the policy check cannot predict. + const d = decideKeep( + keepInput({ + transcribing: true, + autoSubsOnly: false, + policy: { + enabled: false, + replaceAutoSubs: false, + snoozeUntil: null, + root: ALL_LEAF, + }, + }), + ); + assert.deepEqual(d, { keep: true, reason: "in-flight" }); +}); + +test("decideKeep: ASR-only + an enabled policy that covers the channel is the hand-off", () => { + assert.deepEqual(decideKeep(keepInput()), { + keep: true, + reason: "hand-off", + }); +}); + +test("decideKeep: nothing to hand over without audio, or without ASR-only captions", () => { + assert.deepEqual(decideKeep(keepInput({ hasAudio: false })), { + keep: false, + reason: "no-audio", + }); + // A manually-captioned or already-whispered video is not in the bucket, so + // keeping its audio would leak it. + assert.deepEqual(decideKeep(keepInput({ autoSubsOnly: false })), { + keep: false, + reason: "not-auto-subs-only", + }); +}); + +test("decideKeep: audio is not kept for a runner that is off or snoozed", () => { + assert.deepEqual( + decideKeep( + keepInput({ + policy: { + enabled: false, + replaceAutoSubs: true, + snoozeUntil: null, + root: ALL_LEAF, + }, + }), + ), + { keep: false, reason: "policy-off" }, + ); + // A LAPSED snooze is already normalized to null by sanitizePolicy, so a + // non-null value here means the runner really is still idling. + assert.deepEqual( + decideKeep( + keepInput({ + policy: { + enabled: true, + replaceAutoSubs: true, + snoozeUntil: Date.now() + 60_000, + root: ALL_LEAF, + }, + }), + ), + { keep: false, reason: "policy-snoozed" }, + ); +}); + +test("decideKeep: no leaf covering this channel draws the bucket", () => { + // replaceAutoSubs is ON, but the only leaf is for another channel — the video + // would never be claimed, so its audio would sit forever. + const d = decideKeep( + keepInput({ + policy: { + enabled: true, + replaceAutoSubs: true, + snoozeUntil: null, + root: { + id: "root", + mode: "strict", + children: [ + { id: "l", match: { type: "channel", value: "someone-else" } }, + ], + }, + }, + }), + ); + assert.deepEqual(d, { keep: false, reason: "no-leaf" }); +}); + +test("decideKeep: an explicit bucket leaf hands off with replaceAutoSubs OFF", () => { + // The exploration's finding, pinned: `replaceAutoSubs: false` alone is NOT a + // refusal. A leaf naming downloadedAutoSubsOnly draws it regardless, so an + // operator who turned the flag off but kept the leaf still gets the hand-off + // — and would otherwise get the race this whole change removes. + const d = decideKeep( + keepInput({ + policy: { + enabled: true, + replaceAutoSubs: false, + snoozeUntil: null, + root: { + id: "root", + mode: "strict", + children: [ + { + id: "l", + match: { type: "all", bucket: "downloadedAutoSubsOnly" }, + }, + ], + }, + }, + }), + ); + assert.deepEqual(d, { keep: true, reason: "hand-off" }); +}); + +test("decideKeep: the disk bar is the RESUME mark, not the floor", () => { + // 22 GB free with a 20 GB floor and a 5 GB margin: above the floor, below the + // 25 GB the download runner itself would demand before resuming. The backfill + // must not keep audio the runner would have refused to fetch. + assert.deepEqual( + decideKeep( + keepInput({ + disk: { freeBytes: 22 * GB, minFreeDiskGB: 20, resumeMarginGB: 5 }, + }), + ), + { keep: false, reason: "disk-low" }, + ); + // At the mark exactly, it is a hand-off. + assert.deepEqual( + decideKeep( + keepInput({ + disk: { freeBytes: 25 * GB, minFreeDiskGB: 20, resumeMarginGB: 5 }, + }), + ), + { keep: true, reason: "hand-off" }, + ); +}); + +test("decideKeep: a disabled disk gate never refuses the hand-off", () => { + assert.deepEqual( + decideKeep( + keepInput({ + disk: { freeBytes: 1, minFreeDiskGB: 0, resumeMarginGB: 5 }, + }), + ), + { keep: true, reason: "hand-off" }, + ); +}); + +test("decideKeep never mutates its input", () => { + const input = keepInput(); + const before = JSON.stringify(input); + decideKeep(input); + assert.equal(JSON.stringify(input), before); +}); diff --git a/common/controller/backfillReacquire.ts b/common/controller/backfillReacquire.ts @@ -35,43 +35,102 @@ // 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. // -// 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. +// 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 path from "node:path"; import { readdir, rm } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; -import { diskGate } from "../lib/diskSpace"; +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 { isRealAudioFile, isSourceMediaFile } from "../lib/videoStatus"; +import { detectPlatform, type Platform } from "../lib/platform"; +import { isAutoSubsOnly } from "../lib/subtitleProvenance"; +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" + // 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<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> } + | { status: "present"; cleanup: () => Promise<CleanupOutcome> } + // 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<CleanupOutcome> } // Free disk is under the configured floor. A refusal to start, not a failure. - | { status: "disk-floor"; cleanup: () => Promise<boolean> } + | { status: "disk-floor"; cleanup: () => Promise<CleanupOutcome> } // 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> }; + | { status: "gone"; cleanup: () => Promise<CleanupOutcome> } + | { status: "failed"; cleanup: () => Promise<CleanupOutcome> }; -const NOTHING_TO_CLEAN = async (): Promise<boolean> => false; +const NOTHING_TO_CLEAN = async (): Promise<CleanupOutcome> => ({ + 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 @@ -195,7 +254,11 @@ export async function reacquireMediaFor(opts: { // 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), + cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log, { + paths: opts.paths, + channelSlug: opts.channelSlug, + platform: detectPlatform(config.url), + }), }; } @@ -204,12 +267,20 @@ export async function reacquireMediaFor(opts: { log(`Re-acquire ${opts.videoId}: nothing usable landed.`); return { status: "failed", - cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log), + 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), + cleanup: buildCleanup(opts.videoDir, before, opts.videoId, log, { + paths: opts.paths, + channelSlug: opts.channelSlug, + platform: detectPlatform(config.url), + }), }; } @@ -246,20 +317,90 @@ export async function refreshCues( } } -// 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. +// 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" + >; + 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) +// !policyDrawsBucket no leaf covering this channel draws downloadedAutoSubsOnly +// disk the runner itself would refuse to fetch this much +// +// 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 ( + !policyDrawsBucket( + "transcription", + input.policy, + 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 media on purpose. The parakeet resume cache +// (<videoDir>/.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<string>, videoId: string, log: (msg: string) => void, -): () => Promise<boolean> { + ctx: { paths: Paths; channelSlug: string; platform: Platform | null }, +): () => Promise<CleanupOutcome> { return async () => { const after = await readdir(videoDir).catch(() => [] as string[]); const added = after.filter( @@ -270,13 +411,47 @@ function buildCleanup( 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; + 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, + 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 { + 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) { await rm(path.join(videoDir, name), { force: true }).catch((err) => { // Reported, never thrown: this runs in a `finally`, and throwing here @@ -285,6 +460,6 @@ function buildCleanup( }); } log(`Removed re-acquired media for ${videoId} (${added.join(", ")}).`); - return true; + return { status: "removed" }; }; } diff --git a/common/controller/digestYield.test.ts b/common/controller/digestYield.test.ts @@ -1,10 +1,12 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import { + evaluateTranscribingVideo, evaluateTranscriptionActivity, workerContendsForGpu, } from "./digestYield"; import { defaultDigest, sanitizeDigest } from "../lib/settings"; +import type { JobTask } from "../jobs/registry"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/digestYield.test.ts @@ -162,3 +164,68 @@ test("yieldToCpuWorkers defaults off; yieldToTranscription defaults on", () => { ); } }); + +// --------------------------------------------------------------------------- +// evaluateTranscribingVideo — the narrower question: is a transcription running +// on THIS video? The backfill's cleanup asks it before unlinking a re-acquired +// audio file, so a false positive leaks a file and a false negative kills a run. + +// A registry task as the registry actually stores it (label/startedAt are set by +// addTask); only `kind` and `id` are consulted. +function task(kind: JobTask["kind"], id: string): JobTask { + return { id, label: id, kind, startedAt: 0 }; +} +const TRANSCRIBE_TASK = task("transcribe", "vid1"); +const DOWNLOAD_TASK = task("download", "vid1"); + +test("a running job with a transcribe task on this video counts", () => { + assert.equal( + evaluateTranscribingVideo( + [{ status: "running", tasks: [TRANSCRIBE_TASK] }], + "vid1", + ), + true, + ); +}); + +test("a DOWNLOAD task on the same video does not count", () => { + // Same id, different lane: a download holds no engine open on the audio, and + // treating it as one would keep media the backfill should have removed. + assert.equal( + evaluateTranscribingVideo( + [{ status: "running", tasks: [DOWNLOAD_TASK] }], + "vid1", + ), + false, + ); +}); + +test("a finished job does not count, whatever its tasks say", () => { + // Tasks are cleared when a job ends, but a record read mid-transition must not + // be trusted on `tasks` alone. + assert.equal( + evaluateTranscribingVideo( + [{ status: "done", tasks: [TRANSCRIBE_TASK] }], + "vid1", + ), + false, + ); +}); + +test("another video's transcription does not count", () => { + assert.equal( + evaluateTranscribingVideo( + [{ status: "running", tasks: [task("transcribe", "vid2")] }], + "vid1", + ), + false, + ); +}); + +test("no jobs, or a running job with no tasks, is false", () => { + assert.equal(evaluateTranscribingVideo([], "vid1"), false); + assert.equal( + evaluateTranscribingVideo([{ status: "running", tasks: undefined }], "vid1"), + false, + ); +}); diff --git a/common/controller/digestYield.ts b/common/controller/digestYield.ts @@ -28,7 +28,7 @@ // next dispatch sees the busy lane and holds. import { getWorkerPool, type WorkerSummary } from "../jobs/workerPool"; -import { getRegistry } from "../jobs/registry"; +import { getRegistry, type JobRecord } from "../jobs/registry"; import { TRANSCRIPTION_QUEUE } from "../lib/queueKeys"; import { getSettings } from "../lib/settings"; @@ -142,3 +142,46 @@ export function transcriptionActivity(): TranscriptionActivity { yieldToCpuWorkers, }); } + +// --- Is a transcription running on THIS video right now? -------------------- +// +// A second consumer of the same registry, asking a narrower question than the +// lane-level one above: not "is the transcription lane busy" but "is some job +// transcribing this exact video". The backfill's cleanup uses it as defense in +// depth before unlinking a re-acquired audio file — deleting audio out from +// under a running engine costs that video's whole run and appends it to +// failed-transcriptions. +// +// Registry TASKS are the right signal because they cover BOTH producers: the +// auto runner and the manual whisper-all batch both go through transcribeOne's +// tracker (transcribeOne.ts -> jobs/taskHooks.ts -> registry.addTask), whereas +// getAutoRunnerStatus() sees only the runner and would miss the manual batch. +// +// Known blind spot, accepted: an inline transcribe inside a download +// (downloadOneManaged) passes no tracker and registers no task, so it is +// invisible here. The primary mechanism is the policy hand-off below it, not +// this check. +export function evaluateTranscribingVideo( + jobs: ReadonlyArray<Pick<JobRecord, "status" | "tasks">>, + videoId: string, +): boolean { + return jobs.some( + (job) => + job.status === "running" && + (job.tasks ?? []).some( + (t) => t.kind === "transcribe" && t.id === videoId, + ), + ); +} + +// Thin I/O wrapper, same shape as transcriptionActivity() above. Fails OPEN to +// false: "not transcribing" simply falls through to the caller's own policy +// decision, which is the mechanism that actually protects the file. An +// unreadable registry must not change what the backfill does. +export function isTranscribingVideo(videoId: string): boolean { + try { + return evaluateTranscribingVideo(getRegistry().list(), videoId); + } catch { + return false; + } +} diff --git a/common/controller/operationJobs.ts b/common/controller/operationJobs.ts @@ -130,7 +130,10 @@ export async function runBackfillChannelJob( ? `; ${batch.blocked} waiting on a prerequisite backfill` : "") + (batch.reacquired > 0 - ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` + ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` + + (batch.reacquireHandedOff > 0 + ? `, ${batch.reacquireHandedOff} handed to auto-transcribe` + : "") : "") + (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + ".", diff --git a/common/controller/transcribeOne.test.ts b/common/controller/transcribeOne.test.ts @@ -0,0 +1,73 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import path from "node:path"; +import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { transcribeOneVideo } from "./transcribeOne"; +import { TranscribeError } from "./transcribeError"; +import type { Paths } from "../lib/paths"; +import type { Worker } from "../lib/workers"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/transcribeOne.test.ts +// +// Audio that VANISHES MID-RUN is a skip, not a failure. Another lane removing +// the file — today the backfill's re-acquire cleanup, which runs in a `finally` +// and holds no per-video lock — must not append the video to +// failed-transcriptions, which the manual per-channel batch honours forever. +// +// Driven with a shell script as the worker binary (the repo's fake-bin +// convention; this codebase does not module-mock), so the engine's exit is real. + +async function runWith(script: string): Promise<unknown> { + const dir = await mkdtemp(path.join(tmpdir(), "transcribe-one-")); + try { + const bin = path.join(dir, "fake-engine.sh"); + await writeFile(bin, `#!/bin/sh\n${script}\n`, { mode: 0o755 }); + await writeFile(path.join(dir, "audio.mp3"), "not really audio"); + const worker: Worker = { + id: "w1", + name: "fake", + kind: "local", + enabled: true, + priority: 0, + appId: "whisper-cpp", + // `model` set so app.build never reaches getPaths(). + config: { bin, model: "model.bin" }, + }; + return await transcribeOneVideo({ + paths: {} as Paths, + videoDir: dir, + videoId: "vid1", + audioFilename: "audio.mp3", + worker, + onLog: () => {}, + }).then( + (v) => v, + (e) => e, + ); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +test("an engine that fails AFTER its audio disappeared is a no-audio skip", async () => { + const err = await runWith("rm -f audio.mp3; exit 1"); + assert.ok(err instanceof TranscribeError, `not a TranscribeError: ${err}`); + assert.equal(err.failureClass, "no-audio"); + assert.match(String(err.message), /vanished mid-run/); +}); + +test("an engine that just fails is still a real failure", async () => { + // The guard must not swallow genuine engine failures into the skip class — + // that would silently stop the queue ever recording a bad video. + const err = await runWith("exit 1"); + assert.ok(err instanceof Error, `not an Error: ${err}`); + assert.notEqual((err as TranscribeError).failureClass, "no-audio"); +}); + +test("an engine that exits 0 with no output, audio gone, is a skip too", async () => { + const err = await runWith("rm -f audio.mp3; exit 0"); + assert.ok(err instanceof TranscribeError, `not a TranscribeError: ${err}`); + assert.equal(err.failureClass, "no-audio"); +}); diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -44,6 +44,30 @@ async function resolveAudioFile( return null; } +// The audio we resolved at the top of the run is gone NOW. +// +// Another lane removed it mid-transcription — today that means the backfill's +// re-acquire cleanup, which runs in a `finally` and holds no per-video lock. The +// engine's own error for this is whatever its next read failed with (parakeet +// re-opens the file once per 480 s window, so it surfaces mid-run), and letting +// that reach transcribeOneFromQueue appends the video to failed-transcriptions — +// a permanent blacklist for the manual per-channel batch — for a file problem +// that will be gone on the next attempt. +// +// So re-check the file and convert it into the "no-audio" failure class, which +// the queue ALREADY treats as a skip. No new class, no new call site behaviour. +async function throwIfAudioVanished( + videoDir: string, + audio: string, + videoId: string, +): Promise<void> { + if (await pathExists(path.join(videoDir, audio))) return; + throw new TranscribeError( + `audio ${audio} for ${videoId} vanished mid-run (removed by another lane) — skipped, not failed`, + "no-audio", + ); +} + export type TranscribeOneOptions = { paths: Paths; videoDir: string; @@ -125,16 +149,25 @@ export async function transcribeOneVideo( log( `Transcribe ${opts.videoId} start (${resolvedAudio}) via remote [${worker.id}] ${worker.remote?.baseUrl ?? ""}`, ); - const bytes = await transcribeViaRemote({ - worker, - audioPath: path.join(opts.videoDir, resolvedAudio), - audioName: resolvedAudio, - videoId: opts.videoId, - channelSlug: deriveChannelSlug(opts.paths, opts.videoDir) ?? undefined, - onLog: (line: string) => log(line), - onProgress: opts.onProgress, - signal: opts.signal, - }); + // The remote client wraps a local ENOENT on the upload as a "transport" + // failure, so the vanished-audio check has to happen HERE too — the local + // rethrow below would never see it, and the video would be blacklisted. + let bytes: Buffer | null; + try { + bytes = await transcribeViaRemote({ + worker, + audioPath: path.join(opts.videoDir, resolvedAudio), + audioName: resolvedAudio, + videoId: opts.videoId, + channelSlug: deriveChannelSlug(opts.paths, opts.videoDir) ?? undefined, + onLog: (line: string) => log(line), + onProgress: opts.onProgress, + signal: opts.signal, + }); + } catch (err) { + await throwIfAudioVanished(opts.videoDir, resolvedAudio, opts.videoId); + throw err; + } // Upload path: write the pulled bytes. Shared-fs path (bytes === null): the // remote already wrote transcript.json onto the shared mount at this path. if (bytes) await writeFile(transcriptPath, bytes); @@ -193,6 +226,7 @@ export async function transcribeOneVideo( ); return "paused"; } + await throwIfAudioVanished(opts.videoDir, resolvedAudio, opts.videoId); throw err; } const producedPath = path.join(opts.videoDir, build.outputFile); @@ -203,6 +237,9 @@ export async function transcribeOneVideo( ); return "paused"; } + // An engine that exited 0 but wrote nothing may have been reading a file + // that disappeared under it — same check, same reason as the catch above. + await throwIfAudioVanished(opts.videoDir, resolvedAudio, opts.videoId); throw new TranscribeError( `transcription with ${app.id} produced no ${build.outputFile} in ${opts.videoDir}`, "transcription", diff --git a/common/jobs/autoQueuePolicy.test.ts b/common/jobs/autoQueuePolicy.test.ts @@ -11,6 +11,7 @@ import { selectableBucketsForKind, emptyAutoQueueRuntime, flattenLeaves, + policyDrawsBucket, sanitizeAutoQueue, selectNextWork, } from "./autoQueuePolicy"; @@ -824,3 +825,182 @@ test("sanitize: operation wins and DROPS bucket, so the tree cannot be ambiguous // A blank operation is not an operation. assert.deepEqual(leaves[2].match, { type: "all" }); }); + +// --------------------------------------------------------------------------- +// policyDrawsBucket — "would the runner actually pick this up?", asked from +// outside the runner. The backfill hand-off (controller/backfillReacquire.ts) +// keeps a re-acquired audio file only when the answer is yes, so a wrong answer +// here either leaks a file onto a 97%-full disk or deletes one out from under a +// running engine. + +const LEAF_ALL: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [{ id: "leaf-all", match: { type: "all" } }], +}; + +test("policyDrawsBucket: an opt-in bucket needs replaceAutoSubs or an explicit leaf", () => { + const off = { replaceAutoSubs: false, root: LEAF_ALL }; + const on = { replaceAutoSubs: true, root: LEAF_ALL }; + const ch = { slug: "chan", platform: "youtube" as const }; + + // A bucket-less catch-all draws the DEFAULT union, which excludes the opt-in + // bucket until replaceAutoSubs appends it. + assert.equal( + policyDrawsBucket("transcription", off, ch, "downloadedAutoSubsOnly"), + false, + ); + assert.equal( + policyDrawsBucket("transcription", on, ch, "downloadedAutoSubsOnly"), + true, + ); + // The always-drawn buckets are drawn either way. + assert.equal( + policyDrawsBucket("transcription", off, ch, "downloadedNoTranscript"), + true, + ); + + // The second way in: a leaf naming the bucket, with the flag still off. This + // is the case that makes `replaceAutoSubs: false` NOT a refusal. + const explicit = { + replaceAutoSubs: false, + root: { + id: "root", + mode: "strict" as const, + children: [ + { + id: "leaf-bucket", + match: { + type: "all" as const, + bucket: "downloadedAutoSubsOnly", + }, + }, + ], + }, + }; + assert.equal( + policyDrawsBucket("transcription", explicit, ch, "downloadedAutoSubsOnly"), + true, + ); + // …and it draws ONLY that bucket. + assert.equal( + policyDrawsBucket("transcription", explicit, ch, "downloadedNoTranscript"), + false, + ); +}); + +test("policyDrawsBucket: channel and platform leaves cover only what they match", () => { + const byChannel = { + replaceAutoSubs: true, + root: { + id: "root", + mode: "strict" as const, + children: [ + { id: "l", match: { type: "channel" as const, value: "mine" } }, + ], + }, + }; + assert.equal( + policyDrawsBucket( + "transcription", + byChannel, + { slug: "mine", platform: "youtube" }, + "downloadedAutoSubsOnly", + ), + true, + ); + assert.equal( + policyDrawsBucket( + "transcription", + byChannel, + { slug: "other", platform: "youtube" }, + "downloadedAutoSubsOnly", + ), + false, + ); + + const byPlatform = { + replaceAutoSubs: true, + root: { + id: "root", + mode: "strict" as const, + children: [ + { id: "l", match: { type: "platform" as const, value: "rumble" } }, + ], + }, + }; + assert.equal( + policyDrawsBucket( + "transcription", + byPlatform, + { slug: "mine", platform: "rumble" }, + "downloadedAutoSubsOnly", + ), + true, + ); + assert.equal( + policyDrawsBucket( + "transcription", + byPlatform, + { slug: "mine", platform: "youtube" }, + "downloadedAutoSubsOnly", + ), + false, + ); + // An unknown platform is not a wildcard. + assert.equal( + policyDrawsBucket( + "transcription", + byPlatform, + { slug: "mine", platform: null }, + "downloadedAutoSubsOnly", + ), + false, + ); +}); + +test("policyDrawsBucket: an operation leaf draws no bucket at all", () => { + const ops = { + replaceAutoSubs: true, + root: { + id: "root", + mode: "strict" as const, + children: [ + { + id: "l", + match: { type: "all" as const, operation: "diarization" }, + }, + ], + }, + }; + assert.equal( + policyDrawsBucket( + "transcription", + ops, + { slug: "mine", platform: "youtube" }, + "downloadedAutoSubsOnly", + ), + false, + ); + assert.equal( + policyDrawsBucket( + "transcription", + ops, + { slug: "mine", platform: "youtube" }, + "downloadedNoTranscript", + ), + false, + ); +}); + +test("policyDrawsBucket: an empty tree draws nothing", () => { + assert.equal( + policyDrawsBucket( + "transcription", + defaultAutoQueue().transcription, + { slug: "mine", platform: "youtube" }, + "downloadedNoTranscript", + ), + false, + ); +}); diff --git a/common/jobs/autoQueuePolicy.ts b/common/jobs/autoQueuePolicy.ts @@ -222,6 +222,35 @@ export function defaultBucketsForPolicy( : bucketsForKind(kind); } +// Would this runner kind, under this policy, draw `bucket` for this channel? +// +// The question a lane outside the runner has to ask before it hands work over: +// "if I leave this video in that bucket, will the runner actually pick it up?" +// It is NOT `policy.replaceAutoSubs` — that flag is only one of the two ways in. +// A leaf naming an operation draws no bucket at all; a leaf naming a bucket +// draws only that one (per-channel opt-in, no flag needed); a bucket-less leaf +// draws defaultBucketsForPolicy — which is where replaceAutoSubs enters, and the +// only place it does. +// +// Pure, like everything else here: the caller supplies the channel's slug and +// platform, never a snapshot. +export function policyDrawsBucket( + kind: "transcription" | "download", + policy: Pick<AutoQueuePolicy, "replaceAutoSubs" | "root">, + channel: Pick<ChannelWork, "slug" | "platform">, + bucket: string, +): boolean { + const defaults = defaultBucketsForPolicy(kind, policy); + return flattenLeaves(policy.root).some( + (l) => + !l.match.operation && + matchesChannel(l.match, channel) && + (l.match.bucket + ? l.match.bucket === bucket + : defaults.includes(bucket)), + ); +} + export const AUTO_QUEUE_MAX_WORKERS_MAX = 64; // --- Runtime fairness state (persisted best-effort by autoQueueState.ts) ---- @@ -304,7 +333,14 @@ export function flattenLeaves(node: AutoQueueNode): AutoQueueLeaf[] { return out; } -function matchesChannel(match: AutoQueueMatch, ch: ChannelWork): boolean { +// Exported for policyDrawsBucket's callers outside this module (the backfill +// hand-off asks it about one channel it holds only a slug and a platform for), +// which is why the parameter is the narrowest shape that answers the question +// rather than a whole ChannelWork. +export function matchesChannel( + match: AutoQueueMatch, + ch: Pick<ChannelWork, "slug" | "platform">, +): boolean { if (match.type === "all") return true; if (match.type === "channel") return !!match.value && ch.slug === match.value; if (match.type === "platform") {