import path from "node:path"; import fs from "fs-extra"; import { execa } from "execa"; import type { Paths } from "../lib/paths"; import type { Worker } from "../lib/workers"; import { getTranscriptionApp } from "../lib/transcriptionApps"; import { detectTranscriptFormat } from "../lib/whisper"; import { WORKER_DEGRADE_THRESHOLD, getWorkerPool, type WorkerFilter, } from "../jobs/workerPool"; import type { TaskTracker } from "../jobs/taskHooks"; import type { SchedulerTier } from "../jobs/jobKinds"; import { normalizeTranscript } from "./normalizeTranscript"; import { diarizeOneVideo } from "./diarizeOne"; import { getSettings } from "../lib/settings"; import { pingRemoteHealth, transcribeViaRemote } from "./remoteTranscribe"; import { TranscribeError } from "./transcribeError"; import { findSourceMedia, isRealAudioFile } from "../lib/videoStatus"; import { pickPreferredAudio } from "../lib/mediaFiles"; import { resolveSavedVideo } from "../lib/savedVideo-server"; import { writeTranscribeOutcome } from "../lib/transcribeOutcome-server"; import { writeFileAtomic } from "../lib/jsonFile-server"; const { pathExists, readdir, rename } = fs; async function resolveAudioFile( videoDir: string, requested: string, strict: boolean, ): Promise { if (await pathExists(path.join(videoDir, requested))) return requested; if (strict) return null; const entries = await readdir(videoDir); const preferred = pickPreferredAudio(entries.filter(isRealAudioFile)); if (preferred) return preferred; // No extracted audio on disk — fall back to a persisted source video // container. Transcribers that ffmpeg-slice their input (parakeet) read a // container directly; this is what makes a kept-but-cleaned video, or a // redownload-to-archive that only fetched the container, still transcribable. const container = findSourceMedia(entries); if (container) return container; // The container may have been moved into the saved-video store (Phase 3), // leaving a saved-video.json pointer. Resolve it and return a path relative to // videoDir so the engine (cwd = videoDir) and the remote uploader // (path.join(videoDir, …)) both read the stored file correctly. const saved = await resolveSavedVideo(videoDir); if (saved) return path.relative(videoDir, saved); 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 { 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; videoId: string; audioFilename: string; strictAudio?: boolean; // The worker that runs this transcription. Callers obtain one from the global // worker pool (getWorkerPool().acquire()). A local worker runs its engine here // via execa; a remote worker delegates to another app instance. worker: Worker; onLog?: (msg: string) => void; // Parsed progress callback, used by the remote branch to forward the remote's // {fraction, detail} to the per-task bar (the local branch reports progress via // onLog → the engine's parser instead). onProgress?: (patch: { fraction?: number; detail?: string }) => void; signal?: AbortSignal; // Graceful "pause & keep progress": when aborted and the engine supports it // (parakeet), the child is SIGTERM'd so it finishes the current window and // exits without producing a transcript — its completed windows are cached for // resume. transcribeOneVideo returns "paused" so the batch leaves the video // untranscribed (next run resumes). Ignored by engines without partial support. partialSignal?: AbortSignal; // True skips the inline diarization pass even when settings turn it on. A // one-off file transcription (controller/transcribeFile.ts) runs in a scratch // dir that is deleted afterwards, so a diarization there is work thrown away. skipInlineDiarization?: boolean; // True asks the engine to keep word timestamps in transcript.json (see // TranscribeBuildInput.words). Only the one-off file transcription asks. words?: boolean; }; export type TranscribeOneOutcome = "transcribed" | "already-exists" | "paused"; // TranscribeError lives in its own module so the remote client can throw it // without an import cycle; re-exported here for existing import sites. export { TranscribeError } from "./transcribeError"; export async function transcribeOneVideo( opts: TranscribeOneOptions, ): Promise { const log = opts.onLog ?? ((m: string) => console.log(m)); const transcriptPath = path.join(opts.videoDir, "transcript.json"); if (await pathExists(transcriptPath)) { log(`Transcription for ${opts.videoId} already exists`); return "already-exists"; } const resolvedAudio = await resolveAudioFile( opts.videoDir, opts.audioFilename, opts.strictAudio === true, ); if (!resolvedAudio) { // A missing audio file is a download/source problem, not a transcription // failure. Throw "no-audio" so the batch skips it (it transcribes once a // good audio file exists) instead of permanently recording a failure. if (opts.strictAudio) { throw new TranscribeError( `no ${opts.audioFilename} found in ${opts.videoDir} (strict audio mode)`, "no-audio", ); } throw new TranscribeError( `no audio file found in ${opts.videoDir}`, "no-audio", ); } if (resolvedAudio !== opts.audioFilename) { log(`Using ${resolvedAudio} instead of ${opts.audioFilename}`); } const worker = opts.worker; const start = Date.now(); // Each app writes "" relative to cwd (the ext differs: whisper // appends ".json", chough writes the exact -o path). Use a tmp base and rename // the app-declared output file on success so a SIGTERM mid-write can't leave a // half-baked transcript.json. const tmpBase = `transcript.tmp-${process.pid}`; let outputFormat: "whisper-json" | "chough-json" | "vtt"; if (worker.kind === "remote") { // Delegate to another instance of this app over HTTP: upload the audio, poll // progress, pull back the transcript.json bytes, and write them locally. The // remote ran its OWN worker pool to produce them; we normalize here. log( `Transcribe ${opts.videoId} start (${resolvedAudio}) via remote [${worker.id}] ${worker.remote?.baseUrl ?? ""}`, ); // 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 — atomically, so a crash mid-write // never leaves a truncated transcript.json for the indexer to read. Shared-fs // path (bytes === null): the remote already wrote transcript.json onto the // shared mount at this path. if (bytes) await writeFileAtomic(transcriptPath, bytes); // The remote didn't tell us the engine format; sniff it from the file (the // same content sniff the indexer uses for mixed corpora). const raw = await fs.readFile(transcriptPath, "utf8").catch(() => ""); outputFormat = detectTranscriptFormat(raw) ?? "whisper-json"; } else { const app = getTranscriptionApp(worker.appId); const appConfig = worker.config ?? {}; const bin = appConfig.bin?.trim() || app.defaultBin(); log( `Transcribe ${opts.videoId} start (${resolvedAudio}) via ${app.id} [${worker.id}]`, ); const build = app.build({ audioFile: resolvedAudio, outputBase: tmpBase, config: appConfig, ...(opts.words ? { words: true } : {}), }); const child = execa(bin, build.argv, { cwd: opts.videoDir, cancelSignal: opts.signal, all: true, buffer: false, env: build.env ? { ...process.env, ...build.env } : undefined, }); child.all?.on("data", (c: Buffer) => log(c.toString("utf8"))); // "Pause & keep progress": for an engine that handles SIGTERM by stopping // gracefully and caching completed work (parakeet), send it on partialSignal. // execa's cancelSignal is the separate hard-cancel `opts.signal`, so this // SIGTERM lets the child exit 0 and `await child` resolves normally. The // engine writes NO transcript on pause (it stitches only when fully done), // so a missing output + a requested pause is a resumable pause, not a failure. let pauseRequested = false; if (app.supportsPartialStop && opts.partialSignal) { const stop = () => { pauseRequested = true; log(`Pausing ${worker.id} on ${opts.videoId} (progress kept for resume)…`); child.kill("SIGTERM"); }; if (opts.partialSignal.aborted) stop(); else opts.partialSignal.addEventListener("abort", stop, { once: true }); } try { await child; } catch (err) { // A requested pause (SIGTERM via partialSignal) that doesn't exit within // execa's forceKillAfterDelay gets force-SIGKILLed, so `await child` // REJECTS instead of resolving. That's still a pause — no transcript was // produced and the engine cached its completed windows for resume — NOT a // failure. Return "paused" so the batch skips it (retries/resumes next // run) rather than blacklisting the video in failed-transcriptions. if (pauseRequested) { log( `Transcribe ${opts.videoId} paused (stopped before exit); completed windows cached — re-run to resume.`, ); return "paused"; } await throwIfAudioVanished(opts.videoDir, resolvedAudio, opts.videoId); throw err; } const producedPath = path.join(opts.videoDir, build.outputFile); if (!(await pathExists(producedPath))) { if (pauseRequested) { log( `Transcribe ${opts.videoId} paused; completed windows cached — re-run to resume.`, ); 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", ); } await rename(producedPath, transcriptPath); outputFormat = build.outputFormat; } log( `Transcribe ${opts.videoId} done in ${((Date.now() - start) / 1000).toFixed(2)}s`, ); // Persist when this transcript was acquired so the stats build can chart // "content added over time" by transcription date. Best-effort: a failed // sidecar write must not fail an otherwise-successful transcription (the build // falls back to the transcript file mtime when the sidecar is absent). try { await writeTranscribeOutcome(opts.videoDir, { videoId: opts.videoId, transcribedAt: new Date().toISOString(), worker: { id: worker.id, appId: worker.appId, kind: worker.kind }, durationMs: Date.now() - start, }); } catch (err) { log( `Warning: failed to write transcribe-outcome.json for ${opts.videoId}: ${(err as Error).message}`, ); } const channelSlug = deriveChannelSlug(opts.paths, opts.videoDir); if (channelSlug) { try { await normalizeTranscript({ videoDir: opts.videoDir, channelSlug, formatHint: outputFormat, log, }); } catch (err) { log( `Warning: failed to write transcript.cues.json for ${opts.videoId}: ${(err as Error).message}`, ); } } // Speaker diarization, while the audio is still on disk. THIS is the only // moment it can happen for free: cleanAudioFromTranscribed deletes the audio // once a video is transcribed, and nothing downstream can reconstruct it. // // Off by default and deliberately so — diarization is slower than the // transcription it follows, so the batch workflow leaves this off and uses the // backfill pass instead (see DiarizationSettings.inlineAfterTranscribe). When // it IS on, it stays non-fatal: the established idiom here is a try/catch that // logs and continues, and a diarization problem must never fail a // transcription that already succeeded. // // The "paused" early returns above correctly bypass this — there is no // transcript yet, so there is nothing to diarize alongside. const diarization = getSettings().diarization; if ( !opts.skipInlineDiarization && diarization.enabled && diarization.inlineAfterTranscribe ) { try { const outcome = await diarizeOneVideo({ paths: opts.paths, videoDir: opts.videoDir, videoId: opts.videoId, settings: diarization, onLog: log, signal: opts.signal, }); if (outcome === "failed") { log( `Warning: diarization failed for ${opts.videoId}; audio is retained for a later pass.`, ); } } catch (err) { log( `Warning: diarization threw for ${opts.videoId}: ${(err as Error).message}`, ); } } return "transcribed"; } function deriveChannelSlug(paths: Paths, videoDir: string): string | null { const rel = path.relative(paths.channelsDir, videoDir); if (!rel || rel.startsWith("..")) return null; const [slug] = rel.split(path.sep); return slug || null; } // How many workers a single video will try before giving up, when each attempt // fails with a transport-class error (remote unreachable, etc.). A // transcription-class failure stops immediately — retrying the same audio // elsewhere is pointless. Also bounds the loop if a worker keeps flapping. export const MAX_WORKER_ATTEMPTS = 3; export type TranscribeWithWorkerOptions = { paths: Paths; videoDir: string; videoId: string; audioFilename: string; strictAudio?: boolean; // Optional per-task progress tracker; when provided the acquired worker drives // which progress parser is used and which worker the Workers page shows. tracker?: TaskTracker; taskId?: string; taskLabel?: string; onLog?: (msg: string) => void; // Hard cancel. Also unblocks a parked acquire. signal?: AbortSignal; // Soft cancel (drain): when aborted, a parked acquire unblocks and this throws // an AbortError so the caller skips the video without starting it. drainSignal?: AbortSignal; // Auto-runner units pass true so they park BEHIND any manual (foreground) // acquire in the worker pool — a manual transcribe preempts queued auto work. background?: boolean; // The pool tier outright, over `background`: a short one-off file // transcription asks "urgent", ahead of every parked batch (transcribeFile). tier?: SchedulerTier; // Narrows WHICH workers may take this video (the pool's `only`): a one-off // file transcription keeps to local workers, or to the one it was told to use. only?: WorkerFilter; // Called with the worker each attempt runs on, before it starts — how a // caller learns which engine and model produced the transcript. onWorker?: (worker: Worker) => void; // See TranscribeOneOptions.skipInlineDiarization. skipInlineDiarization?: boolean; // See TranscribeOneOptions.words. words?: boolean; }; // Acquire a worker from the global pool and transcribe one video through it, // retrying on a different worker when a transport-class failure occurs (a remote // went away). This is the single entry point batches and single-video actions // share so pool accounting (lease release, success/failure marks) is consistent. // Throws TranscribeError on a transcription-class failure or exhausted retries, // and an AbortError when cancelled/drained while parked. export async function transcribeWithWorker( opts: TranscribeWithWorkerOptions, ): Promise { const log = opts.onLog ?? ((m: string) => console.log(m)); const pool = getWorkerPool(); // Unblock a parked acquire on either a hard cancel or a drain. const acquireSignal = opts.signal && opts.drainSignal ? AbortSignal.any([opts.signal, opts.drainSignal]) : (opts.signal ?? opts.drainSignal); let lastErr: unknown; for (let attempt = 0; attempt < MAX_WORKER_ATTEMPTS; attempt++) { if (opts.signal?.aborted || opts.drainSignal?.aborted) throw abortError(); let lease; try { lease = await pool.acquire(acquireSignal, { background: opts.background, ...(opts.tier ? { tier: opts.tier } : {}), ...(opts.only ? { only: opts.only } : {}), }); } catch (err) { // An `only` no configured worker passes is refused at once — that is // not a cancel, and must not read as one. if (opts.only && !acquireSignal?.aborted) throw err; throw abortError(); // cancelled or drained while parked } const worker = lease.worker; opts.onWorker?.(worker); const task = opts.tracker?.start({ id: opts.taskId ?? opts.videoId, label: opts.taskLabel ?? opts.videoId, kind: "transcribe", appId: worker.appId, workerId: worker.id, }); // For a partial-capable local engine, register a graceful "stop & keep // partial" handle the Workers page can trigger for this worker. const partialCapable = worker.kind === "local" && getTranscriptionApp(worker.appId).supportsPartialStop === true; const partialController = partialCapable ? new AbortController() : undefined; if (partialController) { pool.setActiveStop(worker.id, () => partialController.abort()); } try { const outcome = await transcribeOneVideo({ paths: opts.paths, videoDir: opts.videoDir, videoId: opts.videoId, audioFilename: opts.audioFilename, strictAudio: opts.strictAudio, worker, onLog: task ? task.onLog : log, onProgress: task ? task.update : undefined, signal: opts.signal, partialSignal: partialController?.signal, skipInlineDiarization: opts.skipInlineDiarization, words: opts.words, }); pool.markSuccess(worker.id); return outcome; } catch (err) { if (opts.signal?.aborted) throw err; const failureClass = err instanceof TranscribeError ? err.failureClass : "transcription"; if (failureClass === "transport") { lastErr = err; // A remote that fails a health check is genuinely down: degrade it now so // this video's retry AND later videos skip it, rather than burning all // attempts on it. A reachable remote (transient 5xx) or a local worker // uses the gentler consecutive-failure threshold. if (worker.kind === "remote" && !(await pingRemoteHealth(worker))) { if (pool.markDegraded(worker.id)) { log( `Worker ${worker.id} is unreachable — disabled for this run; re-enable it from the Workers page once it's back.`, ); } } else if (pool.markFailure(worker.id)) { // The degrade rule is the pool's WORKER_DEGRADE_THRESHOLD, not this // loop's per-video attempt cap — the two happen to be equal today, // and the message used to name the wrong one. log( `Worker ${worker.id} auto-disabled after ${WORKER_DEGRADE_THRESHOLD} consecutive failures; re-enable it from the Workers page.`, ); } log( `Worker ${worker.id} transport failure on ${opts.videoId}: ${String(err)} — retrying on another worker`, ); continue; } throw err; } finally { if (partialController) pool.clearActiveStop(worker.id); task?.end(); lease.release(); } } throw ( lastErr ?? new TranscribeError( `no worker could transcribe ${opts.videoId}`, "transport", ) ); } function abortError(): Error { const err = new Error("transcription aborted"); err.name = "AbortError"; return err; }