// Per-attempt orchestrator for audio-checked downloads. // // When `channelConfig.audioCheck.enabled` is true on a `transcribe`-handling // channel, this replaces the simple yt-dlp invocation for the primary // attempt. It launches yt-dlp WITHOUT `-x` (so yt-dlp produces the raw // video container), runs periodic ffmpeg probes over snapshots of the // in-progress `.part`, rolls back to a known-good snapshot on corruption, // and finally runs its own ffmpeg extract to the configured audio format. // // Snapshot discipline: every periodic check captures `fstat(.part).size` and // sends SIGSTOP. In the default paused mode the child stays stopped across the // integrity probe, so the probe reads the `.part` in place — no copy needed — // and only a clean checkpoint copies `.part` → `.good` (reflink where // supported), truncated to the pre-SIGSTOP size to discard any kernel-buffered // writes that landed as SIGSTOP took effect. Legacy resume-during-probe mode // instead copies a `.testing` snapshot up front (the child keeps writing during // the probe) and promotes it to `.good` on success. SIGCONT is always sent in a // finally to avoid orphaning a suspended child. import { removeMediaFile } from "../lib/mediaTier-server"; import { constants as fsConstants } from "node:fs"; import { isPartAudioFile, isRealAudioFile } from "../lib/mediaFiles"; import { copyFile, open, readdir, rename, rm, stat, } from "node:fs/promises"; import path from "node:path"; import { execa } from "execa"; import { AUDIO_CHECK_COPY_TIMEOUT_DEFAULT_SECONDS, AUDIO_CHECK_INTERVAL_BACKOFF_FACTOR_DEFAULT, AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS, AUDIO_CHECK_INTERVAL_MIN_SECONDS, AUDIO_CHECK_INTERVAL_RECOVER_AFTER_CLEAN, AUDIO_CHECK_INTERVAL_RECOVER_STEP_SECONDS, AUDIO_CHECK_MAX_ROLLBACKS_DEFAULT, type AudioFormat, type ChannelConfig, } from "../lib/channelConfig"; import type { Paths } from "../lib/paths"; import { transcodeAudio } from "../controller/transcode"; import { backoffInterval, formatInterval, recoverInterval, type CadenceParams, } from "./audioCheckCadence"; import type { StreamVerdict } from "./ffmpegStreamClassify"; import { probeAudioStream } from "./ffmpegStreamProbe"; import { DLOM_PROBE_MARKER } from "../jobs/progressParsers"; const STDERR_TAIL_BYTES = 64 * 1024; const DEFAULT_SIZE_GATE_BYTES = 1024 * 1024; // 1 MiB minimum growth between checks const PART_DISCOVERY_INTERVAL_MS = 250; const PART_DISCOVERY_TIMEOUT_MS = 10_000; // Test-only overrides, E2E_-prefixed and set by editor/playwright.config.ts on // its test server. Tests need fast intervals and small file sizes to // exercise the loop without burning real seconds and megabytes. The // AUDIO_CHECK_INTERVAL_MIN_SECONDS clamp in the form parser protects against // users accidentally setting this in production; the env knobs bypass the // clamp for test fixtures only. function envIntOverride(name: string): number | null { const raw = process.env[name]; if (!raw) return null; const n = Number.parseInt(raw, 10); return Number.isFinite(n) && n >= 0 ? n : null; } // Float env override for the backoff factor. Only accepts a value in (0, 1] — // a factor >= 1 would never shrink the interval, and <= 0 is nonsensical. function envFloatOverride(name: string): number | null { const raw = process.env[name]; if (!raw) return null; const n = Number.parseFloat(raw); return Number.isFinite(n) && n > 0 && n <= 1 ? n : null; } // Tri-state boolean env override: unset -> null (use config), "1"/"true" -> // true, anything else -> false. Lets comparison runs flip the resume-during- // probe behavior without editing channel config. function envBoolOverride(name: string): boolean | null { const raw = process.env[name]; if (raw == null || raw === "") return null; return raw === "1" || raw.toLowerCase() === "true"; } export type CheckpointAction = | "advance" | "rollback" | "restart" | "skip-size-gate" | "skip-copy-timeout" | "skip-missing-part"; export type CheckpointRecord = { at: string; bytes: number; verdict: StreamVerdict | null; action: CheckpointAction; // Wall-clock duration of the ffmpeg integrity probe, in ms. Present only for // checkpoints where a real probe ran (advance / rollback / restart / // skip-copy-timeout); undefined for size-gate / missing-part skips. durationMs?: number; // The live adaptive-cadence interval (ms) that was in effect when this // checkpoint fired — i.e. the wait before it, and the value the malformed / // clean decision then adjusts. Stamped on advance / rollback / restart so the // AIMD curve is observable (and deterministically assertable) from the // outcome record. Undefined on skip records that don't touch the cadence. intervalMs?: number; }; export type AudioCheckOutcomeKind = | "ok" | "failed-corrupt-source" // A download that COMPLETED (yt-dlp exit 0, all bytes) but whose final // integrity probe is malformed, even after one re-download. Re-downloading // a complete file can't change a deterministic verdict, so we stop, KEEP the // downloaded container on disk for inspection, and flag it as a distinct // terminal state (separate from failed-corrupt-source, which never finished). | "corrupt-full-source" | "aborted" | "ytdlp-error"; export type AudioCheckOutcome = { kind: AudioCheckOutcomeKind; ytdlpExitCode: number | null; stderrTail: string; archiveLine: string | null; checkpoints: CheckpointRecord[]; rollbacks: number; restarts: number; finalProbeVerdict?: StreamVerdict; // Resolved at orchestrator runtime by scanning /data/*/audio.*.part. // May be null if discovery never succeeded (e.g. yt-dlp failed before // creating the data dir). videoDir: string | null; }; export type AudioCheckedOpts = { paths: Paths; channelDir: string; channelConfig: ChannelConfig; ytdlpArgs: string[]; // already built by caller (no -x) onLog: (s: string) => void; signal: AbortSignal; archiveMarker: string; // On-disk subdir name under channelDir/data/ this launch is expected to // write to. For most platforms it's extractVideoId(url). Null when the // caller can't predict the dir name (Rumble: yt-dlp's internal id only // shows up after metadata.info.json is written). The orchestrator scopes // .part discovery and final-file resolution to this dir; when null it // restricts to subdirs that didn't exist before launch. Without this, // channels with many stale .parts from prior interrupted attempts can // probe the wrong file and let a malformed download finalize. expectedVideoIdHint: string | null; // Optional test-only knob: extra ms to sleep while SIGSTOPped to help // simulate cancellation-during-pause scenarios. debugPauseMs?: number; }; // Resolved knobs. type Knobs = { // The configured probe interval — the AIMD ceiling and starting point. The // live cadence (currentIntervalMs in runAudioCheckedYtdlp) shrinks below this // on malformed checkpoints and relaxes back up toward it, never past it. intervalMs: number; maxRollbacks: number; copyTimeoutMs: number; debugPauseMs: number; sizeGateBytes: number; // Adaptive-cadence knobs (AIMD). backoffFactor multiplies the live interval on // each malformed checkpoint (down to backoffFloorMs); after recoverAfterClean // consecutive clean checkpoints the interval steps back up by recoverStepMs. backoffFactor: number; backoffFloorMs: number; recoverStepMs: number; recoverAfterClean: number; // Legacy behavior when true: SIGCONT immediately after the snapshot copy // and probe concurrently. Default false: hold the child stopped across the // probe so no would-be-discarded bytes are downloaded. resumeDuringProbe: boolean; }; function resolveKnobs(opts: AudioCheckedOpts): Knobs { const cfg = opts.channelConfig.audioCheck; const intervalMsOverride = envIntOverride("E2E_AUDIO_CHECK_INTERVAL_MS"); const sizeGateOverride = envIntOverride("E2E_AUDIO_CHECK_SIZE_GATE"); const debugPauseOverride = envIntOverride("E2E_AUDIO_CHECK_DEBUG_PAUSE_MS"); const resumeDuringProbeOverride = envBoolOverride( "AUDIO_CHECK_RESUME_DURING_PROBE", ); const backoffFactorOverride = envFloatOverride("AUDIO_CHECK_BACKOFF_FACTOR"); const floorMsOverride = envIntOverride("E2E_AUDIO_CHECK_INTERVAL_FLOOR_MS"); const recoverStepMsOverride = envIntOverride("E2E_AUDIO_CHECK_RECOVER_STEP_MS"); const recoverAfterOverride = envIntOverride("E2E_AUDIO_CHECK_RECOVER_AFTER"); const intervalMs = intervalMsOverride ?? (cfg?.intervalSeconds ?? AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS) * 1000; // Floor defaults to the configured MIN, but never exceed the starting // interval itself — otherwise a small (test) interval would have no room to // shrink, and the very first backoff would clamp *up*. const backoffFloorMs = Math.min( floorMsOverride ?? AUDIO_CHECK_INTERVAL_MIN_SECONDS * 1000, intervalMs, ); return { intervalMs, maxRollbacks: cfg?.maxRollbacks ?? AUDIO_CHECK_MAX_ROLLBACKS_DEFAULT, copyTimeoutMs: (cfg?.copyTimeoutSeconds ?? AUDIO_CHECK_COPY_TIMEOUT_DEFAULT_SECONDS) * 1000, debugPauseMs: debugPauseOverride ?? opts.debugPauseMs ?? 0, sizeGateBytes: sizeGateOverride ?? DEFAULT_SIZE_GATE_BYTES, backoffFactor: backoffFactorOverride ?? AUDIO_CHECK_INTERVAL_BACKOFF_FACTOR_DEFAULT, backoffFloorMs, recoverStepMs: recoverStepMsOverride ?? AUDIO_CHECK_INTERVAL_RECOVER_STEP_SECONDS * 1000, recoverAfterClean: recoverAfterOverride ?? AUDIO_CHECK_INTERVAL_RECOVER_AFTER_CLEAN, resumeDuringProbe: resumeDuringProbeOverride ?? cfg?.resumeDuringProbe ?? false, }; } // --- File helpers ----------------------------------------------------------- function partPathFor(file: string): string { return `${file}.part`; } function testingPath(partFile: string): string { return `${partFile}.testing`; } function goodPath(partFile: string): string { return `${partFile}.good`; } function deriveFinalFromPart(partFile: string): string { return partFile.endsWith(".part") ? partFile.slice(0, -".part".length) : partFile; } async function findExistingPartOrFinal( videoDir: string, ): Promise<{ partFile: string | null; finalFile: string | null }> { let partFile: string | null = null; let finalFile: string | null = null; let entries: string[] = []; try { entries = await readdir(videoDir); } catch { return { partFile: null, finalFile: null }; } // Shares the app's media-file predicates rather than matching on the audio. // prefix: `audio.live_chat.json.part` satisfies "starts with audio., ends with // .part" and is a live-chat sidecar, not a resumable download. for (const e of entries) { if (isPartAudioFile(e)) { partFile = path.join(videoDir, e); } else if (isRealAudioFile(e)) { // Skip already-extracted audio. outputs that aren't the source // container; ChannelConfig.audioFormat is the target, so if the file // matches that, it's an output (or a stale one). // We can't be 100% sure here without context, but for the // "discover existing partial" use case it's fine to ignore; the // caller will treat (null, null) as "no recovery state". // Choose the first non-part candidate as a tentative finalFile. if (!finalFile) finalFile = path.join(videoDir, e); } } return { partFile, finalFile }; } async function pathExists(p: string): Promise { try { await stat(p); return true; } catch { return false; } } async function prepareDataTree( channelDir: string, expectedVideoIdHint: string | null, onLog: (s: string) => void, precheck?: { paths: Paths; signal: AbortSignal; onCheckpoint: (rec: CheckpointRecord) => void; }, ): Promise { // Drop any leftover .testing snapshots from crashed prior runs, and // promote any orphan .good back to .part so yt-dlp resume picks up where // we left off. When `expectedVideoIdHint` is set, scoped to that single // subdir under /data/ — sibling dirs belong to other videos // that yt-dlp won't touch on this launch, so probing them is wasted work. // When the hint is null (Rumble: yt-dlp's internal id only appears after // metadata.info.json), fall back to scanning every subdir. // // When `precheck` is supplied, also probe each pre-existing .part (i.e. // one we did NOT just restore from a .good — the .good was already // probe-validated at snapshot time). On malformed verdict: roll back to // .good if present, else discard so yt-dlp restarts from 0. This catches // corruption from a prior interrupted attempt that yt-dlp's -c would // otherwise extend rather than detect. const dataDir = path.join(channelDir, "data"); let videoDirs: string[] = []; try { videoDirs = await readdir(dataDir); } catch { return null; } if (expectedVideoIdHint) { videoDirs = videoDirs.filter((d) => d === expectedVideoIdHint); } // The subdir we treated as THIS download's resume target — restored a .good, // or rolled back / discarded a malformed .part. The caller excludes it from // its pre-launch snapshot so post-launch discovery still scopes to it even // when expectedVideoIdHint is null. Clean pre-existing .parts (stale // neighbours from other videos) are NOT reported, preserving the discovery // logic's "ignore unrelated interrupted attempts" guarantee. let resumeDir: string | null = null; for (const sub of videoDirs) { const dir = path.join(dataDir, sub); let entries: string[] = []; try { entries = await readdir(dir); } catch { continue; } let partFile: string | null = null; let goodFile: string | null = null; for (const e of entries) { if (e.endsWith(".part.testing") && e.startsWith("audio.")) { const p = path.join(dir, e); await rm(p, { force: true }); onLog(`Removed stale snapshot: ${p}\n`); } else if (e.endsWith(".part.good") && e.startsWith("audio.")) { goodFile = path.join(dir, e); } else if (isPartAudioFile(e)) { partFile = path.join(dir, e); } } let preExisting = partFile; if (!partFile && goodFile) { const restored = goodFile.slice(0, -".good".length); await rename(goodFile, restored); onLog(`Restored prior validated snapshot: ${restored}\n`); // Don't probe the restored snapshot — .good is clean by construction. preExisting = null; resumeDir = sub; } if (precheck && preExisting && !precheck.signal.aborted) { let bytes = 0; try { const st = await stat(preExisting); bytes = st.size; } catch {} const probe = await probeAudioStream({ ffmpegBin: precheck.paths.ffmpegBin, file: preExisting, signal: precheck.signal, onLog, }); if (precheck.signal.aborted) continue; if (probe.verdict === "clean" || probe.verdict === "partial") { onLog( `Pre-check OK at ${bytes} bytes (verdict=${probe.verdict}): ${preExisting}\n`, ); } else { const good = goodPath(preExisting); const haveGood = await pathExists(good); if (haveGood) { await rm(preExisting, { force: true }); await rename(good, preExisting); onLog( `Pre-check found malformed .part; rolled back to .good: ${preExisting}\n`, ); precheck.onCheckpoint({ at: new Date().toISOString(), bytes, verdict: probe.verdict, action: "rollback", }); resumeDir = sub; } else { await rm(preExisting, { force: true }); onLog( `Pre-check found malformed .part; discarded (no .good baseline): ${preExisting}\n`, ); precheck.onCheckpoint({ at: new Date().toISOString(), bytes, verdict: probe.verdict, action: "restart", }); resumeDir = sub; } } } } return resumeDir; } // SIGCONT a (possibly already-exited) child, swallowing ESRCH. Used by the // caller to resume a child that snapshotPart left stopped via holdStopped. function resumeChild( pid: number | undefined, onLog: (s: string) => void, ): void { if (!pid) return; try { process.kill(pid, "SIGCONT"); } catch (err) { const code = (err as NodeJS.ErrnoException).code; if (code !== "ESRCH") { onLog(`WARN: failed to SIGCONT pid ${pid}: ${(err as Error).message}\n`); } } } // Reflink-copy `src` → `dst`, racing the copy against the copy timeout so we // never hold SIGSTOP indefinitely, then truncate `dst` down to `preSize` to // discard any bytes past the measured snapshot boundary (kernel-buffered writes // that landed as SIGSTOP took effect). Returns "copy-timeout" if the copy // outran the budget; the partial `dst` is then cleaned up best-effort. `dst` is // always a copy of `src`, whose size only grows, so it is >= preSize and the // truncate is always valid. async function copyTruncated( src: string, dst: string, preSize: number, knobs: Knobs, ): Promise<"ok" | "copy-timeout"> { const copyOp = (async () => { await copyFile(src, dst, fsConstants.COPYFILE_FICLONE); })(); let timeoutId: ReturnType | null = null; const timeout = new Promise<"timeout">((resolve) => { timeoutId = setTimeout(() => resolve("timeout"), knobs.copyTimeoutMs); }); const winner = await Promise.race([ copyOp.then(() => "copy" as const), timeout, ]); if (timeoutId) clearTimeout(timeoutId); if (winner === "timeout") { // Best-effort clean up the partial copy. Don't await on copyOp completion; // just leave it to settle in the background while we move on. copyOp.catch(() => {}).then(() => rm(dst, { force: true }).catch(() => {})); return "copy-timeout"; } const fh = await open(dst, "r+"); try { await fh.truncate(preSize); } finally { await fh.close(); } return "ok"; } async function snapshotPart( partFile: string, // When null, the caller probes the held `.part` in place — no copy is taken. // Only valid together with holdStopped (the child must stay SIGSTOPped across // the probe, since the `.part` is read directly). When set, the snapshot is // copied here (legacy resume-during-probe mode). testingFile: string | null, pid: number | undefined, knobs: Knobs, onLog: (s: string) => void, signal: AbortSignal, // When true, leave the child SIGSTOPped on a successful snapshot so the // caller can run the integrity probe before any further bytes download. // The caller then owns resuming (or killing) it. Failure paths always // resume — the never-wedge guarantee is preserved. holdStopped = false, ): Promise<{ ok: true; bytes: number } | { ok: false; reason: "size-gate" | "copy-timeout" | "missing-part" }> { // Stat first; the .part may have been renamed away (yt-dlp finished). let preSize: number; try { const st = await stat(partFile); preSize = st.size; } catch { return { ok: false, reason: "missing-part" }; } if (!pid) { // Child has no pid (already exited). Caller will handle this; we just // bail out of the snapshot without doing anything destructive. return { ok: false, reason: "missing-part" }; } let stopped = false; // Default: resume the child when we leave this function. Flipped to false // only on a successful snapshot when holdStopped is set, handing the resume // responsibility to the caller. let resumeOnExit = true; try { try { process.kill(pid, "SIGSTOP"); stopped = true; } catch (err) { const code = (err as NodeJS.ErrnoException).code; if (code === "ESRCH") return { ok: false, reason: "missing-part" }; throw err; } // Re-stat post-stop purely to confirm the `.part` still exists; the child // is stopped, so it shouldn't have vanished, but treat that as missing. try { await stat(partFile); } catch { return { ok: false, reason: "missing-part" }; } // Legacy resume-during-probe mode: take a reflink copy now, truncated back // to preSize, because the child resumes and keeps writing during the probe. // Default paused mode (testingFile null) skips this — the held `.part` is // probed in place and only a clean checkpoint copies it into `.good`. if (testingFile !== null) { const copied = await copyTruncated(partFile, testingFile, preSize, knobs); if (copied === "copy-timeout") return { ok: false, reason: "copy-timeout" }; } // Honour optional test-only debug pause so cancellation tests can wedge a // cancel into the SIGSTOP window deterministically. if (knobs.debugPauseMs > 0 && !signal.aborted) { await new Promise((resolve) => { const t = setTimeout(resolve, knobs.debugPauseMs); const onAbort = () => { clearTimeout(t); resolve(); }; signal.addEventListener("abort", onAbort, { once: true }); }); } // Successful snapshot. When holding stopped, the caller resumes/kills the // child after probing; leave it suspended. resumeOnExit = !holdStopped; return { ok: true, bytes: preSize }; } finally { if (stopped && resumeOnExit) { try { process.kill(pid, "SIGCONT"); } catch (err) { const code = (err as NodeJS.ErrnoException).code; if (code !== "ESRCH") { onLog( `WARN: failed to SIGCONT pid ${pid}: ${(err as Error).message}\n`, ); } } } } } async function findPartIn(dir: string): Promise { const entries = await readdir(dir).catch(() => [] as string[]); for (const e of entries) { if (isPartAudioFile(e)) return path.join(dir, e); } return null; } async function discoverPartFile( channelDir: string, expectedVideoIdHint: string | null, preLaunchSubdirs: ReadonlySet, signal: AbortSignal, ): Promise { const dataDir = path.join(channelDir, "data"); const deadline = Date.now() + PART_DISCOVERY_TIMEOUT_MS; while (Date.now() < deadline) { if (signal.aborted) return null; // Preferred path: scope to the hinted subdir. extractVideoId(url) gives // the right answer for most platforms; this avoids latching onto stale // .part files in other video subdirs. if (expectedVideoIdHint) { const hit = await findPartIn(path.join(dataDir, expectedVideoIdHint)); if (hit) return hit; } // Fallback (Rumble, or hint mismatch): only consider subdirs that // appeared *after* launch. Pre-existing subdirs are skipped — their // .part files belong to other videos' interrupted attempts. const subdirs = await readdir(dataDir).catch(() => [] as string[]); for (const sub of subdirs) { if (sub === expectedVideoIdHint) continue; // already checked above if (preLaunchSubdirs.has(sub)) continue; const hit = await findPartIn(path.join(dataDir, sub)); if (hit) return hit; } await new Promise((resolve) => { const t = setTimeout(resolve, PART_DISCOVERY_INTERVAL_MS); signal.addEventListener("abort", () => { clearTimeout(t); resolve(); }, { once: true }); }); } return null; } // --- The orchestrator ------------------------------------------------------- export async function runAudioCheckedYtdlp( opts: AudioCheckedOpts, ): Promise { const knobs = resolveKnobs(opts); const checkpoints: CheckpointRecord[] = []; let rollbacks = 0; let restarts = 0; let consecutiveRollbacks = 0; // Adaptive (AIMD) probe cadence, run-scoped so it persists across yt-dlp // relaunches: a malformed checkpoint multiplies it down toward knobs.backoffFloorMs, // and recoverAfterClean consecutive clean checkpoints step it back up toward // the configured knobs.intervalMs. cleanStreak counts the clean run driving recovery. let currentIntervalMs = knobs.intervalMs; let cleanStreak = 0; const cadence: CadenceParams = { ceilingMs: knobs.intervalMs, floorMs: knobs.backoffFloorMs, factor: knobs.backoffFactor, recoverStepMs: knobs.recoverStepMs, recoverAfterClean: knobs.recoverAfterClean, }; // Counts malformed verdicts on the FINAL probe (a complete yt-dlp exit-0 // download), as opposed to mid-download checkpoint rollbacks. Re-downloading // a complete file yields identical bytes and thus the same verdict, so we // allow exactly one re-download before declaring corrupt-full-source. let finalProbeMalformed = 0; let lastExit: number | null = null; let lastStderrTail = ""; let lastArchiveLine: string | null = null; let finalProbeVerdict: StreamVerdict | undefined; // Filled in once the watcher discovers a .part file under data//. // Sticky across rollback/restart loops since the id stays the same. let resolvedVideoDir: string | null = null; const resumeDir = await prepareDataTree( opts.channelDir, opts.expectedVideoIdHint, opts.onLog, { paths: opts.paths, signal: opts.signal, onCheckpoint: (rec) => checkpoints.push(rec), }, ); // Snapshot of subdirs that existed before yt-dlp launched. discoverPartFile // and the success-path fallback scan use this to ignore stale .part files // belonging to other videos' interrupted attempts when the hint dir // doesn't yet exist (Rumble case). Exclude the resume dir the pre-check just // acted on: that subdir IS this download's target (e.g. a malformed .part we // rolled back), so discovery must still scope to it even though it existed // before launch — otherwise the finished file there is never found. const preLaunchSubdirs: ReadonlySet = new Set( ( await readdir(path.join(opts.channelDir, "data")).catch( () => [] as string[], ) ).filter((d) => d !== resumeDir), ); // Main loop: each iteration = one yt-dlp launch. Decides what to do based // on whether the launch ended in success, rollback, restart, or error. while (true) { if (opts.signal.aborted) { return buildOutcome("aborted"); } if (consecutiveRollbacks > knobs.maxRollbacks) { return buildOutcome("failed-corrupt-source"); } const launch = await runOneLaunch(); lastExit = launch.exitCode; if (launch.stderrTail) lastStderrTail = launch.stderrTail; if (launch.archiveLine) lastArchiveLine = launch.archiveLine; if (launch.kind === "aborted") return buildOutcome("aborted"); if (launch.kind === "rollback") { rollbacks++; consecutiveRollbacks++; continue; } if (launch.kind === "restart") { restarts++; consecutiveRollbacks++; continue; } if (launch.kind === "ytdlp-error") { return buildOutcome("ytdlp-error"); } // launch.kind === "success" if (launch.videoDir) resolvedVideoDir = launch.videoDir; const finalFile = launch.finalFile; if (!finalFile) { // yt-dlp claimed success but we couldn't locate the output file. // Treat as a yt-dlp error so the outer attempt machinery can record // failure rather than spinning. return buildOutcome("ytdlp-error"); } // Final probe. const probe = await probeAudioStream({ ffmpegBin: opts.paths.ffmpegBin, file: finalFile, signal: opts.signal, onLog: opts.onLog, }); finalProbeVerdict = probe.verdict; if (probe.verdict === "malformed") { // A cancel that landed during/just after the (potentially long) final // probe must win — don't discard the file or kick off another full // re-download on the user's way out. if (opts.signal.aborted) return buildOutcome("aborted"); finalProbeMalformed++; if (finalProbeMalformed > 1) { // Second malformed final probe on a complete download. Re-downloading // can't change a deterministic verdict, so stop and KEEP the file // (no rename-to-.part, no rm, no transcode) for inspection. Drop only // the audio-check scratch snapshots. opts.onLog( `Final probe verdict: malformed again after re-download. ` + `Keeping the downloaded file and flagging corrupt-full-source: ${finalFile}\n`, ); await rm(goodPath(partPathFor(finalFile)), { force: true }); await rm(testingPath(partPathFor(finalFile)), { force: true }); return buildOutcome("corrupt-full-source"); } opts.onLog( `Final probe verdict: malformed. Re-downloading once before giving up.\n`, ); // Rename back to .part, then restore .good if we have one. const partAgain = partPathFor(finalFile); await rename(finalFile, partAgain).catch(() => {}); const good = goodPath(partAgain); if (await pathExists(good)) { await rm(partAgain, { force: true }); await rename(good, partAgain); } else { await rm(partAgain, { force: true }); } // Counts toward the rollback STAT, but NOT consecutiveRollbacks: the // final-probe retry is governed by finalProbeMalformed (one re-download // then corrupt-full-source), kept independent of the checkpoint-driven // consecutiveRollbacks cap that yields failed-corrupt-source. rollbacks++; continue; } // Clean (or partial — partial at finalize is OK for a complete download // since we never expect strict EOF errors). Extract audio to target // format and finish. const audioFmt: AudioFormat = opts.channelConfig.audioFormat ?? "mp3"; const finalBasename = path.basename(finalFile); const videoDir = path.dirname(finalFile); const targetBasename = `audio.${audioFmt}`; if (finalBasename === targetBasename) { // yt-dlp's bestaudio selector picked a stream already in the target // format (e.g. Odysee serving a native mp3 audio track). The file is // the final output; no transcode and no source-cleanup needed. opts.onLog( `File already in target format (${targetBasename}); skipping transcode.\n`, ); await rm(goodPath(partPathFor(finalFile)), { force: true }); return buildOutcome("ok"); } try { await transcodeAudio({ paths: opts.paths, videoDir, sourceFilename: finalBasename, targetFormat: audioFmt, onLog: opts.onLog, signal: opts.signal, }); } catch (err) { // Treat extract failure like a malformed final probe: rollback if we // can, otherwise surface as ytdlp-error (it isn't really a yt-dlp // error, but it shares the "couldn't finish" outcome shape). opts.onLog( `Final transcode failed: ${(err as Error).message}\n`, ); if (opts.signal.aborted) return buildOutcome("aborted"); const partAgain = partPathFor(finalFile); const exists = await pathExists(finalFile); if (exists) await rename(finalFile, partAgain).catch(() => {}); const good = goodPath(partAgain); if (await pathExists(good)) { await rm(partAgain, { force: true }); await rename(good, partAgain); rollbacks++; consecutiveRollbacks++; continue; } // No baseline; treat as restart. await rm(partAgain, { force: true }); restarts++; consecutiveRollbacks++; continue; } // Cleanup intermediates: drop .part.good, optionally drop the source // container if !keepSourceVideo. await rm(goodPath(partPathFor(finalFile)), { force: true }); if (!opts.channelConfig.keepSourceVideo) { await removeMediaFile(path.dirname(finalFile), path.basename(finalFile)); } return buildOutcome("ok"); } function buildOutcome(kind: AudioCheckOutcomeKind): AudioCheckOutcome { return { kind, ytdlpExitCode: lastExit, stderrTail: lastStderrTail, archiveLine: lastArchiveLine, checkpoints, rollbacks, restarts, videoDir: resolvedVideoDir, ...(finalProbeVerdict ? { finalProbeVerdict } : {}), }; } type LaunchResult = | { kind: "success"; exitCode: number | null; stderrTail: string; archiveLine: string | null; finalFile: string | null; videoDir: string | null } | { kind: "rollback"; exitCode: number | null; stderrTail: string; archiveLine: string | null; videoDir: string | null } | { kind: "restart"; exitCode: number | null; stderrTail: string; archiveLine: string | null; videoDir: string | null } | { kind: "ytdlp-error"; exitCode: number | null; stderrTail: string; archiveLine: string | null; videoDir: string | null } | { kind: "aborted"; exitCode: number | null; stderrTail: string; archiveLine: string | null; videoDir: string | null }; async function runOneLaunch(): Promise { // Derived AbortController: when the caller aborts, we first send // SIGCONT (in case the child is suspended) and then forward to execa. const launchAc = new AbortController(); const onParentAbort = () => { if (child?.pid) { try { process.kill(child.pid, "SIGCONT"); } catch {} } launchAc.abort(); }; if (opts.signal.aborted) { launchAc.abort(); } else { opts.signal.addEventListener("abort", onParentAbort, { once: true }); } opts.onLog(`$ ${opts.paths.ytdlpBin} ${opts.ytdlpArgs.join(" ")}\n`); const child = execa(opts.paths.ytdlpBin, opts.ytdlpArgs, { cwd: opts.channelDir, cancelSignal: launchAc.signal, all: false, buffer: false, reject: false, }); let stderrTail = ""; child.stderr?.on("data", (c: Buffer) => { const chunk = c.toString("utf8"); opts.onLog(chunk); stderrTail = (stderrTail + chunk).slice(-STDERR_TAIL_BYTES); }); let stdoutBuf = ""; let archiveLine: string | null = null; child.stdout?.on("data", (c: Buffer) => { const chunk = c.toString("utf8"); opts.onLog(chunk); stdoutBuf += chunk; let nl: number; while ((nl = stdoutBuf.indexOf("\n")) !== -1) { const line = stdoutBuf.slice(0, nl).trim(); stdoutBuf = stdoutBuf.slice(nl + 1); if (line.startsWith(`${opts.archiveMarker} `)) { archiveLine = line.slice(opts.archiveMarker.length + 1).trim(); } } }); // Promise that resolves when yt-dlp exits. let childExited = false; let childExitCode: number | null = null; const exitPromise = child.then((r) => { childExited = true; childExitCode = r.exitCode ?? null; }); let partFile: string | null = null; let lastCheckedBytes = 0; // Most recent real ffmpeg probe duration (and the one before it), used to // emit the estimated duration of the *next* probe in the DLOM_PROBE start // marker so the UI's scanning bar has a denominator. The parser maintains // its own smoothed trend for ETA inflation; this is just the display est. let lastProbeMs = 0; let prevProbeMs = 0; let watcherStop = false; // Set by checkpoint() as soon as a malformed verdict triggers a rollback // or restart. Persisted on a plain variable rather than a Promise so the // race-vs-yt-dlp-exit ambiguity (the rollback kills yt-dlp, which makes // exitPromise resolve too) doesn't matter — we just check this after the // watcher settles to know if the launch ended in a decision or a natural // exit. let pendingDecision: "rollback" | "restart" | null = null; async function checkpoint(): Promise { if (watcherStop || childExited || opts.signal.aborted) return; if (!partFile) return; let st; try { st = await stat(partFile); } catch { return; // .part vanished — yt-dlp probably renamed; treat as missing. } if (st.size <= lastCheckedBytes + knobs.sizeGateBytes) { checkpoints.push({ at: new Date().toISOString(), bytes: st.size, verdict: null, action: "skip-size-gate", }); return; } // A real probe pause begins here. Announce it so the UI flips the task to // the "probing" phase, carrying the estimated duration (last probe plus // its observed increase). finishProbe/skipProbe close the phase on every // exit path below; durationMs from a real probe feeds the parser's trend. const intervalSec = Math.round(currentIntervalMs / 1000); const probeStartedAt = Date.now(); const estMs = lastProbeMs > 0 ? lastProbeMs + Math.max(0, lastProbeMs - prevProbeMs) : 0; opts.onLog(`${DLOM_PROBE_MARKER} status=start estMs=${Math.round(estMs)}\n`); let probeDone = false; const finishProbe = (reset: boolean): void => { if (probeDone) return; probeDone = true; const durationMs = Date.now() - probeStartedAt; prevProbeMs = lastProbeMs; lastProbeMs = durationMs; opts.onLog( `${DLOM_PROBE_MARKER} status=done durationMs=${durationMs} intervalSec=${intervalSec} reset=${reset ? 1 : 0}\n`, ); }; const skipProbe = (): void => { if (probeDone) return; probeDone = true; opts.onLog(`${DLOM_PROBE_MARKER} status=done skip=1\n`); }; // Default: keep the child SIGSTOPped across the probe below so it never // downloads bytes a malformed verdict would discard (429 minimisation). // resumeDuringProbe restores the legacy resume-immediately behavior. const holdStopped = !knobs.resumeDuringProbe; // Paused mode probes the held `.part` in place (testing === null); legacy // mode copies a `.testing` snapshot since the child keeps writing. const testing = holdStopped ? null : testingPath(partFile); const snap = await snapshotPart( partFile, testing, child.pid, knobs, opts.onLog, opts.signal, holdStopped, ); if (!snap.ok) { checkpoints.push({ at: new Date().toISOString(), bytes: st.size, verdict: null, action: snap.reason === "copy-timeout" ? "skip-copy-timeout" : "skip-missing-part", }); skipProbe(); return; } lastCheckedBytes = snap.bytes; // When holdStopped, the child is suspended for the whole probe below. // Resume it on every exit path EXCEPT the malformed branch, which // intentionally kills it (SIGCONT + SIGTERM). `killed` gates that. let killed = false; try { // Classify the snapshot — or the held `.part` directly in paused mode. const probe = await probeAudioStream({ ffmpegBin: opts.paths.ffmpegBin, file: testing ?? partFile, signal: opts.signal, onLog: opts.onLog, }); if (opts.signal.aborted) { if (testing) await rm(testing, { force: true }); skipProbe(); return; } if (probe.verdict === "clean" || probe.verdict === "partial") { // Advance: record a known-good baseline, reset rollback counter. The // finally resumes the held child, so yt-dlp continues from where it // paused. const good = goodPath(partFile); if (testing) { // Legacy: promote the snapshot we already copied. await rename(testing, good).catch(async () => { // If rename failed (e.g. cross-device on weird setups), try // copyFile+rm fallback. await copyFile(testing, good).catch(() => {}); await rm(testing, { force: true }); }); } else { // Paused mode: the child is still SIGSTOPped, so the `.part` we just // probed clean is stable — copy it into `.good` now (truncated to // the probed size). This is the only copy in the paused path. const copied = await copyTruncated( partFile, good, snap.bytes, knobs, ); if (copied === "copy-timeout") { // Couldn't persist the baseline in time. The bytes are clean but // we have no new snapshot, so don't advance — leave the rollback // counter intact and try again next interval. checkpoints.push({ at: new Date().toISOString(), bytes: snap.bytes, verdict: probe.verdict, action: "skip-copy-timeout", durationMs: Date.now() - probeStartedAt, }); opts.onLog( `Checkpoint clean at ${snap.bytes} bytes but timed out saving .good baseline; not advancing.\n`, ); finishProbe(false); return; } } consecutiveRollbacks = 0; checkpoints.push({ at: new Date().toISOString(), bytes: snap.bytes, verdict: probe.verdict, action: "advance", durationMs: Date.now() - probeStartedAt, intervalMs: currentIntervalMs, }); opts.onLog( `Checkpoint OK at ${snap.bytes} bytes (verdict=${probe.verdict}).\n`, ); // AIMD additive increase: after a run of clean checkpoints, relax the // cadence one step back toward the configured interval (never past it). { const rec = recoverInterval(currentIntervalMs, cleanStreak, cadence); if (rec.intervalMs > currentIntervalMs) { opts.onLog( `Audio-check interval recovering: ${formatInterval(currentIntervalMs)} -> ${formatInterval(rec.intervalMs)} after clean checkpoints\n`, ); } currentIntervalMs = rec.intervalMs; cleanStreak = rec.cleanStreak; } finishProbe(false); return; } // Malformed. Decide rollback vs restart depending on whether we have // a .good baseline. const good = goodPath(partFile); const haveGood = await pathExists(good); pendingDecision = haveGood ? "rollback" : "restart"; checkpoints.push({ at: new Date().toISOString(), bytes: snap.bytes, verdict: probe.verdict, action: pendingDecision, durationMs: Date.now() - probeStartedAt, intervalMs: currentIntervalMs, }); opts.onLog( `Checkpoint MALFORMED at ${snap.bytes} bytes. ${haveGood ? "Rolling back to .good." : "No .good baseline; restarting from 0."}\n`, ); // AIMD multiplicative decrease: a malformed source gets probed more // aggressively for the rest of the run (persists across relaunches), // down to the floor. Resets the clean streak driving recovery. cleanStreak = 0; { const next = backoffInterval(currentIntervalMs, cadence); if (next < currentIntervalMs) { opts.onLog( `Audio-check interval backoff: ${formatInterval(currentIntervalMs)} -> ${formatInterval(next)} after malformed checkpoint\n`, ); currentIntervalMs = next; } } // reset=true: the rollback/restart shrinks the .part, so the next // probe is cheaper — the parser must drop its duration trend. finishProbe(true); // Drop the (failed) snapshot, if legacy mode took one. if (testing) await rm(testing, { force: true }); // Stop the watcher to avoid concurrent decisions. watcherStop = true; // Terminate yt-dlp, await exit, then mutate files. SIGCONT first so a // held-stopped child can actually process the SIGTERM. killed = true; if (!childExited && child.pid) { try { process.kill(child.pid, "SIGCONT"); } catch {} try { child.kill("SIGTERM"); } catch {} } try { await exitPromise; } catch {} // Now safe to manipulate .part. if (haveGood) { await rm(partFile, { force: true }); await rename(good, partFile); } else { await rm(partFile, { force: true }); } } finally { // Safety net: if the probe threw before any explicit finish/skip, still // close the probe phase so the UI bar doesn't stay stuck "probing". skipProbe(); // Resume a child snapshotPart left stopped, unless we just killed it. // No-op in resumeDuringProbe mode (snapshotPart already resumed). if (holdStopped && !killed) resumeChild(child.pid, opts.onLog); } } // Start the watcher: discover .part, then arm interval. let launchVideoDir: string | null = resolvedVideoDir; const watcherTask = (async () => { partFile = await discoverPartFile( opts.channelDir, opts.expectedVideoIdHint, preLaunchSubdirs, opts.signal, ); if (partFile) { launchVideoDir = path.dirname(partFile); resolvedVideoDir = launchVideoDir; } if (!partFile || watcherStop || childExited) return; while (!watcherStop && !childExited && !opts.signal.aborted) { // Wait the (adaptive) interval (cancellable on abort/exit). await new Promise((resolve) => { const t = setTimeout(resolve, currentIntervalMs); const onAbort = () => { clearTimeout(t); resolve(); }; opts.signal.addEventListener("abort", onAbort, { once: true }); }); if (watcherStop || childExited || opts.signal.aborted) break; await checkpoint(); } })().catch((err) => { opts.onLog(`Watcher error: ${(err as Error).message}\n`); }); // Wait for either yt-dlp to exit naturally or an abort. The watcher // can kill yt-dlp on a rollback decision (which makes exitPromise // resolve), so we always wait for the watcher to fully settle after to // see whether a decision was made. const result = await Promise.race([ exitPromise.then(() => ({ tag: "exit" as const })), new Promise<{ tag: "abort" }>((resolve) => { if (opts.signal.aborted) resolve({ tag: "abort" }); else opts.signal.addEventListener( "abort", () => resolve({ tag: "abort" }), { once: true }, ); }), ]); watcherStop = true; // Ensure the watcher task settles before we return so timers/listeners // don't fire across the next loop iteration AND so pendingDecision / // file mutations are committed before we inspect them. await watcherTask; opts.signal.removeEventListener?.("abort", onParentAbort); if (result.tag === "abort") { // Resume in case suspended, then terminate. if (child.pid) { try { process.kill(child.pid, "SIGCONT"); } catch {} } try { child.kill("SIGTERM"); } catch {} try { await exitPromise; } catch {} return { kind: "aborted", exitCode: childExitCode, stderrTail, archiveLine, videoDir: launchVideoDir, }; } // A cancel that raced a rollback decision must win: the user asked to // stop, so don't report a rollback/restart that would loop the caller into // another launch. (The race above can resolve "exit" because the watcher // SIGTERMs yt-dlp on a malformed verdict, so check the signal directly.) if (opts.signal.aborted) { return { kind: "aborted", exitCode: childExitCode, stderrTail, archiveLine, videoDir: launchVideoDir, }; } // yt-dlp has exited. If the watcher decided a rollback/restart, that // wins — the exit was triggered by our SIGTERM, not by yt-dlp finishing. if (pendingDecision) { return { kind: pendingDecision, exitCode: childExitCode, stderrTail, archiveLine, videoDir: launchVideoDir, }; } // Natural yt-dlp exit. Decide success or error from its exit code. const exited = childExitCode; if (exited === 0 || exited === 101) { // Locate the final renamed file. let resolvedFinal: string | null = null; if (partFile) { const candidate = deriveFinalFromPart(partFile); if (await pathExists(candidate)) resolvedFinal = candidate; } if (!resolvedFinal) { // Scan the data tree to find any newly written audio.. Scope // the candidate set the same way discoverPartFile does: prefer the // hinted subdir, then any subdir that didn't exist at launch time. // Pre-existing subdirs from other videos are skipped — picking one // would mis-probe an unrelated file (e.g. another video's already- // finalized audio.mp3) and skip the real download's final probe. const dataDir = path.join(opts.channelDir, "data"); const candidates: string[] = []; if (opts.expectedVideoIdHint) candidates.push(opts.expectedVideoIdHint); const subdirs = await readdir(dataDir).catch(() => [] as string[]); for (const sub of subdirs) { if (sub === opts.expectedVideoIdHint) continue; if (preLaunchSubdirs.has(sub)) continue; candidates.push(sub); } for (const sub of candidates) { const dir = path.join(dataDir, sub); const { partFile: leftover, finalFile } = await findExistingPartOrFinal(dir); if (finalFile && !leftover) { resolvedFinal = finalFile; launchVideoDir = dir; resolvedVideoDir = dir; break; } } } return { kind: "success", exitCode: exited, stderrTail, archiveLine, finalFile: resolvedFinal, videoDir: launchVideoDir, }; } return { kind: "ytdlp-error", exitCode: exited, stderrTail, archiveLine, videoDir: launchVideoDir, }; } }