import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; import type { AudioFormat } from "../lib/channelConfig"; import { isVideoDownloaded, isVideoTranscribed, readVideoFiles, } from "../lib/videoStatus"; import { isAutoSubsOnly } from "../lib/subtitleProvenance"; import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; import { runPool } from "../jobs/concurrentRunner"; import { pruneFailedTranscriptions } from "./failedTranscriptions"; import { resolveShardItems } from "./shard"; import { countNotYetTranscribed } from "./channels"; import { getWorkerPool } from "../jobs/workerPool"; import type { TaskTracker } from "../jobs/taskHooks"; import type { JobProgress } from "../jobs/registry"; const { readdir, readFile, ensureFile } = fs; export type WhisperBatchOptions = { channelSlug: string; paths: Paths; audioFilename?: string; audioFormat?: AudioFormat; strictAudioFormat?: boolean; reverse?: boolean; shardTotal?: number; shardIndex?: number; // When set, resolve + persist the shard slice and return immediately without // transcribing anything. Backs the "Save shard" button so a slice can be // pinned (and confirmed saved) before dividing work across synced machines. saveShardOnly?: boolean; // When sharding is active, refine the job's progress total to this slice // (progressBaseline + shard items still needing a transcript) so the Active // Jobs bar measures the shard, not the whole channel. setProgress?: (snap: JobProgress) => void; progressBaseline?: number; // When set, only consider these video IDs (intersected with what's on // disk). Used by bucket-scoped actions like "Transcribe downloaded audio" // where the caller has already determined the exact set. ids?: string[]; onLog?: (msg: string) => void; signal?: AbortSignal; // Soft-cancel: when aborted, stop starting NEW transcriptions but let // in-flight ones finish. Separate from the hard `signal`. drainSignal?: AbortSignal; // When provided, each transcription is tracked as a per-operation task with // its own parsed progress bar on the Active Jobs screen. tracker?: TaskTracker; // Replace-auto-captions lane (always paired with an explicit `ids` set from // the downloadedAutoSubsOnly bucket): treat only transcript.json as "already // transcribed", so videos whose sole transcript is a YouTube ASR VTT are // transcribed rather than skipped. Per-video provenance is re-checked here, so // a stale bucket entry (or a hand-passed id) can't clobber a manual caption // track. replaceAutoSubs?: boolean; }; export type WhisperBatchResult = { attempted: number; succeeded: number; failed: number; skipped: number; }; export async function runWhisperBatch({ channelSlug, paths, audioFilename, audioFormat, strictAudioFormat = false, reverse = false, shardTotal, shardIndex, saveShardOnly = false, setProgress, progressBaseline, ids, onLog, signal, drainSignal, tracker, replaceAutoSubs = false, }: WhisperBatchOptions): Promise { const log = onLog ?? ((m: string) => console.log(m)); const channelDir = path.join(paths.channelsDir, channelSlug); const dataDir = path.join(channelDir, "data"); const failureListFile = path.join(channelDir, "failed-transcriptions"); // Refresh the global worker pool from current settings so this batch sees the // latest worker config. Reconcile never drops an in-use worker, so this is safe // even if another batch is somehow in flight. The pool's total free slots — // not a pLimit — is what throttles this batch's parallelism. const pool = getWorkerPool(); pool.reconfigure(); const resolvedAudioFilename = audioFilename ?? (audioFormat ? `audio.${audioFormat}` : "audio.mp3"); const strict = strictAudioFormat && audioFormat !== undefined; await ensureFile(failureListFile); const originalFailures = (await readFile(failureListFile, "utf-8")) .split("\n") .filter(Boolean); const failedSet = new Set(originalFailures); const allDirs = await readdir(dataDir); // Auto-prune mislabeled failures: any id on the failed list whose dir has no // real audio file was never download-ready — a corrupt-source / failed-extract // download (not a transcription failure) that got recorded here. Drop those so // they retry once re-downloaded and so the count reflects real failures. const onDiskDirs = new Set(allDirs); const bogusFailures: string[] = []; for (const id of failedSet) { if (!onDiskDirs.has(id)) continue; if (!isVideoDownloaded(await readVideoFiles(path.join(dataDir, id)))) { bogusFailures.push(id); } } if (bogusFailures.length > 0) { const { pruned } = await pruneFailedTranscriptions( failureListFile, new Set(bogusFailures), ); for (const id of bogusFailures) failedSet.delete(id); log( `Pruned ${pruned} undownloaded video(s) from failed-transcriptions (no audio on disk — will retry once re-downloaded).`, ); } // When the caller supplies an explicit id set (bucket-scoped run), filter // to that intersection so we never try to transcribe ids whose dirs aren't // on disk. Otherwise consider every video dir. const candidateIds = ids ? (() => { const onDisk = new Set(allDirs); return ids.filter((id) => onDisk.has(id)); })() : allDirs; // Shard only the *remaining* work — videos this run would actually attempt: // no transcript yet (neither whisper's transcript.json nor a yt-dlp VTT) and // not on the failed list. This is what makes a shard slice meaningful (e.g. // 6 untranscribed videos split 2 ways → 3 each, not 29 of the whole channel). const fullItems: string[] = []; for (const id of candidateIds) { if (failedSet.has(id)) continue; const vp = path.join(dataDir, id); const files = await readVideoFiles(vp, { checkUntranscribable: true }); if (replaceAutoSubs) { // Replace-auto-captions lane: only OUR transcript counts as done, and the // VTT never counts as a downloaded artifact (real audio is required). // Re-verify provenance per video so a stale bucket entry can't schedule a // human-authored caption track for replacement. if (files.hasWhisper) continue; if (files.audioFiles.length === 0) continue; if (!(await isAutoSubsOnly(vp, files))) { log(`Skipping ${id}: transcript is not YouTube auto-captions.`); continue; } fullItems.push(id); continue; } // Already transcribed if whisper ran (transcript.json) OR yt-dlp wrote an // English VTT (transcript.en.vtt or a regional/auto fallback like en-US). if (isVideoTranscribed(files)) continue; // Not downloaded (no real audio file): a download candidate, not a // transcription candidate. Attempting it would throw "no audio file found" // and mislabel a download/source problem as a failed transcription. if (!isVideoDownloaded(files)) continue; fullItems.push(id); } const shardResult = await resolveShardItems({ paths, slug: channelSlug, op: "transcribe-missing", fullItems, totalShards: shardTotal, shardIndex: shardIndex, onLog: log, }); const shardItems = shardResult.items; if (saveShardOnly) { log( `Shard transcribe-missing: save-only — slice ${(shardResult.config?.shardIndex ?? shardIndex ?? 0) + 1}/${shardResult.config?.totalShards ?? shardTotal} (${shardItems.length} item(s)) persisted; not transcribing.`, ); return { attempted: 0, succeeded: 0, failed: 0, skipped: 0 }; } // With a shard active, refine the job's progress total to this slice: // baseline (channel-wide transcriptCount at run start) + the shard's items // still needing a transcript. `current` is read channel-wide on the Active // Jobs screen, so a baseline-relative, subset-sized range still reaches 100% // when the shard finishes. if ( shardResult.source !== "full" && setProgress && progressBaseline !== undefined ) { const remaining = await countNotYetTranscribed(paths, channelSlug, shardItems); setProgress({ metric: "transcripts", initial: progressBaseline, target: progressBaseline + remaining, }); log( `Shard transcribe-missing: progress scoped to ${remaining} of ${shardItems.length} shard item(s) needing a transcript.`, ); } const videoDirs = [...shardItems]; if (reverse) videoDirs.reverse(); let attempted = 0; let succeededCount = 0; let failed = 0; let skipped = 0; // The same runPool primitive the cross-channel auto-runner uses, instead of // Promise.all over every video at once. It keeps up to `limit()` units in // flight (the pool's eligible slot count), pulling the next as each finishes — // so the batch fills the pool without dumping a backlog of parked acquires, and // a drain stops pulling while letting in-flight units finish (no event-loop // spin). When all workers are disabled the limit is 0 and the batch PAUSES (it // does not fail), resuming when a worker is re-enabled. Per-video gating + // transcribe is shared with the auto-runner via transcribeOneFromQueue. const NEVER = new AbortController().signal; let idx = 0; const runOne = async (videoDir: string, runSignal: AbortSignal): Promise => { const res = await transcribeOneFromQueue({ paths, channelSlug, videoId: videoDir, audioFilename: resolvedAudioFilename, strict, failedSet, tracker, onLog: log, signal: runSignal, drainSignal, replaceAutoSubs, }); if (res.attempted) attempted++; if (res.outcome === "transcribed") succeededCount++; else if (res.outcome === "failed") failed++; else skipped++; }; await runPool({ next: async () => (idx < videoDirs.length ? videoDirs[idx++] : null), run: runOne, // Eligible (enabled, non-degraded) worker slots — the real throughput // ceiling. 0 ⇒ paused until a worker is re-enabled (runPool idle-waits). limit: () => pool.summary().filter((s) => s.state === "enabled" && !s.degraded).length, signal: signal ?? NEVER, drainSignal: drainSignal ?? NEVER, finite: true, idlePollMs: 3000, }); return { attempted, succeeded: succeededCount, failed, skipped }; }