commit 84e00f7a2848c07b598acf659a22efbfed70a529
parent 409313a8d36b4a7a024f5ac94ed52ccd2663a7a3
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 16 May 2026 22:46:04 -0400
smart odysee pipeline
Diffstat:
15 files changed, 1778 insertions(+), 33 deletions(-)
diff --git a/common/lib/channelConfig.ts b/common/lib/channelConfig.ts
@@ -4,6 +4,13 @@ export type ChannelHandling = "youtube" | "transcribe";
export type AudioFormat = "m4a" | "mp3" | "opus";
+export type AudioCheckConfig = {
+ enabled: boolean;
+ intervalSeconds?: number;
+ maxRollbacks?: number;
+ copyTimeoutSeconds?: number;
+};
+
export type ChannelConfig = {
handling: ChannelHandling;
platform?: Platform;
@@ -20,10 +27,36 @@ export type ChannelConfig = {
// omitted, the global SiteSettings value is used. 0 disables the sleep
// for this channel.
sleepBetweenDownloadsSeconds?: number;
+ // Opt-in audio-integrity checking for sources that intermittently serve
+ // corrupt audio mid-download (e.g. Odysee "original" format). When
+ // enabled, the managed downloader periodically validates the in-progress
+ // .part file via ffmpeg and rolls back to the last known-good snapshot
+ // on corruption. transcribe-handling only.
+ audioCheck?: AudioCheckConfig;
};
export const CHANNEL_SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS = 600;
+export const AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS = 60;
+export const AUDIO_CHECK_INTERVAL_MIN_SECONDS = 10;
+export const AUDIO_CHECK_INTERVAL_MAX_SECONDS = 600;
+
+export const AUDIO_CHECK_MAX_ROLLBACKS_DEFAULT = 5;
+export const AUDIO_CHECK_MAX_ROLLBACKS_MIN = 1;
+export const AUDIO_CHECK_MAX_ROLLBACKS_MAX = 20;
+
+export const AUDIO_CHECK_COPY_TIMEOUT_DEFAULT_SECONDS = 30;
+export const AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS = 5;
+export const AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS = 120;
+
+function clampInt(
+ value: number,
+ min: number,
+ max: number,
+): number {
+ return Math.min(Math.max(Math.floor(value), min), max);
+}
+
export const HANDLING_VALUES: ReadonlyArray<ChannelHandling> = [
"youtube",
"transcribe",
@@ -82,5 +115,42 @@ export function parseChannelConfig(raw: unknown): ChannelConfig | null {
CHANNEL_SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS,
);
}
+ if (r.audioCheck && typeof r.audioCheck === "object") {
+ const a = r.audioCheck as Record<string, unknown>;
+ if (typeof a.enabled === "boolean") {
+ const ac: AudioCheckConfig = { enabled: a.enabled };
+ if (
+ typeof a.intervalSeconds === "number" &&
+ Number.isFinite(a.intervalSeconds)
+ ) {
+ ac.intervalSeconds = clampInt(
+ a.intervalSeconds,
+ AUDIO_CHECK_INTERVAL_MIN_SECONDS,
+ AUDIO_CHECK_INTERVAL_MAX_SECONDS,
+ );
+ }
+ if (
+ typeof a.maxRollbacks === "number" &&
+ Number.isFinite(a.maxRollbacks)
+ ) {
+ ac.maxRollbacks = clampInt(
+ a.maxRollbacks,
+ AUDIO_CHECK_MAX_ROLLBACKS_MIN,
+ AUDIO_CHECK_MAX_ROLLBACKS_MAX,
+ );
+ }
+ if (
+ typeof a.copyTimeoutSeconds === "number" &&
+ Number.isFinite(a.copyTimeoutSeconds)
+ ) {
+ ac.copyTimeoutSeconds = clampInt(
+ a.copyTimeoutSeconds,
+ AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS,
+ AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS,
+ );
+ }
+ config.audioCheck = ac;
+ }
+ }
return config;
}
diff --git a/common/lib/downloadOutcome.ts b/common/lib/downloadOutcome.ts
@@ -8,19 +8,33 @@ export type DownloadOutcomeStatus =
| "ok"
| "ok-with-cookies"
| "ok-auto-transcribed"
- | "failed";
+ | "ok-audio-checked"
+ | "failed"
+ | "failed-corrupt-source";
export const DOWNLOAD_OUTCOME_STATUS_VALUES: ReadonlyArray<DownloadOutcomeStatus> = [
"ok",
"ok-with-cookies",
"ok-auto-transcribed",
+ "ok-audio-checked",
"failed",
+ "failed-corrupt-source",
];
export type DownloadAttemptKind =
| "primary"
| "auth-retry"
- | "no-subs-fallback";
+ | "no-subs-fallback"
+ | "audio-checked-primary";
+
+export type AudioCheckProbeVerdict = "clean" | "partial" | "malformed";
+
+export type AudioCheckAttemptStats = {
+ checkpoints: number;
+ rollbacks: number;
+ restarts: number;
+ finalProbeVerdict?: AudioCheckProbeVerdict;
+};
export type DownloadAttempt = {
n: 1 | 2 | 3;
@@ -30,6 +44,7 @@ export type DownloadAttempt = {
ytdlpExitCode: number | null;
availabilityClass?: Availability;
error?: string;
+ audioCheck?: AudioCheckAttemptStats;
};
export type DownloadOutcomeRecord = {
diff --git a/common/ytdlp/audioCheckedDownload.ts b/common/ytdlp/audioCheckedDownload.ts
@@ -0,0 +1,828 @@
+// 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`,
+// sends SIGSTOP, copies (reflink where supported), sends SIGCONT, then
+// truncates the snapshot to the pre-SIGSTOP size to discard any kernel-
+// buffered writes that landed during the snapshot. SIGCONT is always sent
+// in a finally to avoid orphaning a suspended child.
+
+import { constants as fsConstants } from "node:fs";
+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_DEFAULT_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 type { StreamVerdict } from "./ffmpegStreamClassify";
+import { probeAudioStream } from "./ffmpegStreamProbe";
+
+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. 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;
+}
+
+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;
+};
+
+export type AudioCheckOutcomeKind =
+ | "ok"
+ | "failed-corrupt-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 <channelDir>/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;
+ // Optional test-only knob: extra ms to sleep while SIGSTOPped to help
+ // simulate cancellation-during-pause scenarios.
+ debugPauseMs?: number;
+};
+
+// Resolved knobs.
+type Knobs = {
+ intervalMs: number;
+ maxRollbacks: number;
+ copyTimeoutMs: number;
+ debugPauseMs: number;
+ sizeGateBytes: number;
+};
+
+function resolveKnobs(opts: AudioCheckedOpts): Knobs {
+ const cfg = opts.channelConfig.audioCheck;
+ const intervalMsOverride = envIntOverride("AUDIO_CHECK_INTERVAL_MS_OVERRIDE");
+ const sizeGateOverride = envIntOverride("AUDIO_CHECK_SIZE_GATE_OVERRIDE");
+ const debugPauseOverride = envIntOverride("AUDIO_CHECK_DEBUG_PAUSE_MS");
+ return {
+ intervalMs:
+ intervalMsOverride ??
+ (cfg?.intervalSeconds ?? AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS) * 1000,
+ 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,
+ };
+}
+
+// --- 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 };
+ }
+ for (const e of entries) {
+ if (e === "audio.tmp" || e.startsWith("audio.tmp-")) continue;
+ if (e.endsWith(".part")) {
+ if (e.startsWith("audio.")) partFile = path.join(videoDir, e);
+ } else if (e.startsWith("audio.")) {
+ // Skip already-extracted audio.<format> 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<boolean> {
+ try {
+ await stat(p);
+ return true;
+ } catch {
+ return false;
+ }
+}
+
+async function prepareDataTree(
+ channelDir: string,
+ onLog: (s: string) => void,
+): Promise<void> {
+ // 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. Scoped to <channelDir>/data/*.
+ const dataDir = path.join(channelDir, "data");
+ let videoDirs: string[] = [];
+ try {
+ videoDirs = await readdir(dataDir);
+ } catch {
+ return;
+ }
+ 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 (
+ e.endsWith(".part") &&
+ e.startsWith("audio.") &&
+ !e.endsWith(".part.good") &&
+ !e.endsWith(".part.testing")
+ ) {
+ partFile = path.join(dir, e);
+ }
+ }
+ if (!partFile && goodFile) {
+ const restored = goodFile.slice(0, -".good".length);
+ await rename(goodFile, restored);
+ onLog(`Restored prior validated snapshot: ${restored}\n`);
+ }
+ }
+}
+
+async function snapshotPart(
+ partFile: string,
+ testingFile: string,
+ pid: number | undefined,
+ knobs: Knobs,
+ onLog: (s: string) => void,
+ signal: AbortSignal,
+): 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;
+ 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 in case more bytes flushed; we'll truncate the copy
+ // back to preSize after to discard those buffered writes.
+ let postSize = preSize;
+ try {
+ const st = await stat(partFile);
+ postSize = st.size;
+ } catch {
+ // Shouldn't happen — we're stopped — but treat as missing.
+ return { ok: false, reason: "missing-part" };
+ }
+ // Copy with reflink where supported. Race against the copy timeout so we
+ // never hold SIGSTOP indefinitely.
+ const copyOp = (async () => {
+ await copyFile(partFile, testingFile, fsConstants.COPYFILE_FICLONE);
+ })();
+ let timeoutId: ReturnType<typeof setTimeout> | 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 snapshot. Don't await on copyOp
+ // completion; just leave it to settle in the background while we move
+ // on.
+ copyOp.catch(() => {}).then(() => rm(testingFile, { force: true }).catch(() => {}));
+ return { ok: false, reason: "copy-timeout" };
+ }
+ // Truncate snapshot back to preSize.
+ if (postSize > preSize) {
+ const fh = await open(testingFile, "r+");
+ try {
+ await fh.truncate(preSize);
+ } finally {
+ await fh.close();
+ }
+ }
+ // 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<void>((resolve) => {
+ const t = setTimeout(resolve, knobs.debugPauseMs);
+ const onAbort = () => {
+ clearTimeout(t);
+ resolve();
+ };
+ signal.addEventListener("abort", onAbort, { once: true });
+ });
+ }
+ return { ok: true, bytes: preSize };
+ } finally {
+ if (stopped) {
+ 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 discoverPartFile(
+ channelDir: string,
+ signal: AbortSignal,
+): Promise<string | null> {
+ const dataDir = path.join(channelDir, "data");
+ const deadline = Date.now() + PART_DISCOVERY_TIMEOUT_MS;
+ while (Date.now() < deadline) {
+ if (signal.aborted) return null;
+ const subdirs = await readdir(dataDir).catch(() => [] as string[]);
+ for (const sub of subdirs) {
+ const dir = path.join(dataDir, sub);
+ const entries = await readdir(dir).catch(() => [] as string[]);
+ for (const e of entries) {
+ if (
+ e.startsWith("audio.") &&
+ e.endsWith(".part") &&
+ !e.endsWith(".part.testing") &&
+ !e.endsWith(".part.good")
+ ) {
+ return path.join(dir, e);
+ }
+ }
+ }
+ await new Promise<void>((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<AudioCheckOutcome> {
+ const knobs = resolveKnobs(opts);
+ const checkpoints: CheckpointRecord[] = [];
+ let rollbacks = 0;
+ let restarts = 0;
+ let consecutiveRollbacks = 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/<id>/.
+ // Sticky across rollback/restart loops since the id stays the same.
+ let resolvedVideoDir: string | null = null;
+
+ await prepareDataTree(opts.channelDir, opts.onLog);
+
+ // 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") {
+ opts.onLog(
+ `Final probe verdict: malformed. Rolling back to last good snapshot.\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 });
+ }
+ rollbacks++;
+ consecutiveRollbacks++;
+ 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);
+ 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 rm(finalFile, { force: true });
+ }
+ 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<LaunchResult> {
+ // 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;
+ 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<void> {
+ 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;
+ }
+ const testing = testingPath(partFile);
+ const snap = await snapshotPart(
+ partFile,
+ testing,
+ child.pid,
+ knobs,
+ opts.onLog,
+ opts.signal,
+ );
+ 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",
+ });
+ return;
+ }
+ lastCheckedBytes = snap.bytes;
+ // Classify the snapshot.
+ const probe = await probeAudioStream({
+ ffmpegBin: opts.paths.ffmpegBin,
+ file: testing,
+ signal: opts.signal,
+ onLog: opts.onLog,
+ });
+ if (opts.signal.aborted) {
+ await rm(testing, { force: true });
+ return;
+ }
+ if (probe.verdict === "clean" || probe.verdict === "partial") {
+ // Advance: rename testing → good, reset rollback counter.
+ const good = goodPath(partFile);
+ 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 });
+ });
+ consecutiveRollbacks = 0;
+ checkpoints.push({
+ at: new Date().toISOString(),
+ bytes: snap.bytes,
+ verdict: probe.verdict,
+ action: "advance",
+ });
+ opts.onLog(
+ `Checkpoint OK at ${snap.bytes} bytes (verdict=${probe.verdict}).\n`,
+ );
+ 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,
+ });
+ opts.onLog(
+ `Checkpoint MALFORMED at ${snap.bytes} bytes. ${haveGood ? "Rolling back to .good." : "No .good baseline; restarting from 0."}\n`,
+ );
+ // Drop the (failed) snapshot.
+ await rm(testing, { force: true });
+ // Stop the watcher to avoid concurrent decisions.
+ watcherStop = true;
+ // Terminate yt-dlp, await exit, then mutate files.
+ 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 });
+ }
+ }
+
+ // Start the watcher: discover .part, then arm interval.
+ let launchVideoDir: string | null = resolvedVideoDir;
+ const watcherTask = (async () => {
+ partFile = await discoverPartFile(opts.channelDir, opts.signal);
+ if (partFile) {
+ launchVideoDir = path.dirname(partFile);
+ resolvedVideoDir = launchVideoDir;
+ }
+ if (!partFile || watcherStop || childExited) return;
+ while (!watcherStop && !childExited && !opts.signal.aborted) {
+ // Wait the interval (cancellable on abort/exit).
+ await new Promise<void>((resolve) => {
+ const t = setTimeout(resolve, knobs.intervalMs);
+ 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,
+ };
+ }
+
+ // 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.<ext>.
+ const dataDir = path.join(opts.channelDir, "data");
+ const subdirs = await readdir(dataDir).catch(() => [] as string[]);
+ for (const sub of subdirs) {
+ 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,
+ };
+ }
+}
diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts
@@ -10,6 +10,7 @@ import {
type ChannelConfig,
} from "../lib/channelConfig";
import {
+ type AudioCheckAttemptStats,
type DownloadAttempt,
type DownloadOutcomeRecord,
type DownloadOutcomeStatus,
@@ -18,6 +19,7 @@ import { writeDownloadOutcome } from "../lib/downloadOutcome-server";
import type { Paths } from "../lib/paths";
import { transcribeOneVideo } from "../controller/transcribeOne";
import { extractVideoId, isRumbleUrl } from "./runYtdlp";
+import { runAudioCheckedYtdlp } from "./audioCheckedDownload";
const STDERR_TAIL_BYTES = 64 * 1024;
const ARCHIVE_MARKER = "DLOM_ARCHIVE";
@@ -74,6 +76,18 @@ function transcribeHandlingArgs(config: ChannelConfig): string[] {
return args;
}
+// Audio-checked mode: we own the audio extraction, so yt-dlp must NOT run
+// the ExtractAudio postprocessor. We always pass `-c` so resume picks up
+// where the last validated snapshot left off.
+function transcribeHandlingArgsForAudioCheck(_config: ChannelConfig): string[] {
+ return [
+ "--write-info-json",
+ "-f",
+ "bestaudio/worst",
+ "-c",
+ ];
+}
+
function channelConfigArgs(
config: ChannelConfig,
cookies?: string,
@@ -249,26 +263,77 @@ export async function downloadOneManaged(
let lastArchiveLine: string | null = null;
// ---------- Attempt 1: primary ----------
- const primaryArgs = [
- "--ignore-config",
- "--restrict-filenames",
- ...OUTPUT_ARGS,
- ...(opts.channelConfig.handling === "youtube"
- ? youtubeHandlingArgs(opts.channelConfig)
- : transcribeHandlingArgs(opts.channelConfig)),
- "--print",
- `after_video:${ARCHIVE_MARKER} %(extractor)s %(id)s`,
- ...channelConfigArgs(opts.channelConfig),
- "--",
- opts.videoUrl,
- ];
- const primaryRes = await runOneYtdlp(opts, channelDir, primaryArgs);
+ const audioCheckEnabled =
+ opts.channelConfig.handling === "transcribe" &&
+ opts.channelConfig.audioCheck?.enabled === true;
+
+ let primaryRes: AttemptOutcome;
+ let audioCheckStats: AudioCheckAttemptStats | undefined;
+ let audioCheckCorruptSource = false;
+ // Orchestrator resolves data/<id>/ by scanning the on-disk tree. We use
+ // this in preference to extractVideoId(url), which doesn't know yt-dlp's
+ // internal id (e.g. Odysee claim hashes vs URL slugs).
+ let audioCheckVideoDir: string | null = null;
+
+ if (audioCheckEnabled) {
+ const primaryArgs = [
+ "--ignore-config",
+ "--restrict-filenames",
+ ...OUTPUT_ARGS,
+ ...transcribeHandlingArgsForAudioCheck(opts.channelConfig),
+ "--print",
+ `after_video:${ARCHIVE_MARKER} %(extractor)s %(id)s`,
+ ...channelConfigArgs(opts.channelConfig),
+ "--",
+ opts.videoUrl,
+ ];
+ const audioOutcome = await runAudioCheckedYtdlp({
+ paths: opts.paths,
+ channelDir,
+ channelConfig: opts.channelConfig,
+ ytdlpArgs: primaryArgs,
+ onLog: opts.onLog,
+ signal: opts.signal,
+ archiveMarker: ARCHIVE_MARKER,
+ });
+ primaryRes = {
+ exitCode: audioOutcome.ytdlpExitCode,
+ stderrTail: audioOutcome.stderrTail,
+ archiveLine: audioOutcome.archiveLine,
+ };
+ audioCheckStats = {
+ checkpoints: audioOutcome.checkpoints.length,
+ rollbacks: audioOutcome.rollbacks,
+ restarts: audioOutcome.restarts,
+ ...(audioOutcome.finalProbeVerdict
+ ? { finalProbeVerdict: audioOutcome.finalProbeVerdict }
+ : {}),
+ };
+ audioCheckCorruptSource = audioOutcome.kind === "failed-corrupt-source";
+ audioCheckVideoDir = audioOutcome.videoDir;
+ } else {
+ const primaryArgs = [
+ "--ignore-config",
+ "--restrict-filenames",
+ ...OUTPUT_ARGS,
+ ...(opts.channelConfig.handling === "youtube"
+ ? youtubeHandlingArgs(opts.channelConfig)
+ : transcribeHandlingArgs(opts.channelConfig)),
+ "--print",
+ `after_video:${ARCHIVE_MARKER} %(extractor)s %(id)s`,
+ ...channelConfigArgs(opts.channelConfig),
+ "--",
+ opts.videoUrl,
+ ];
+ primaryRes = await runOneYtdlp(opts, channelDir, primaryArgs);
+ }
+
const primaryAvail = attemptSucceeded(primaryRes.exitCode)
? undefined
: parseUnavailableFromStderr(primaryRes.stderrTail);
attempts.push({
n: 1,
- kind: "primary",
+ kind: audioCheckEnabled ? "audio-checked-primary" : "primary",
handling: opts.channelConfig.handling,
usedCookies: false,
ytdlpExitCode: primaryRes.exitCode,
@@ -276,11 +341,17 @@ export async function downloadOneManaged(
error: attemptSucceeded(primaryRes.exitCode)
? undefined
: trimError(primaryRes.stderrTail),
+ ...(audioCheckStats ? { audioCheck: audioCheckStats } : {}),
});
if (primaryRes.archiveLine) lastArchiveLine = primaryRes.archiveLine;
- let lastSucceeded = attemptSucceeded(primaryRes.exitCode);
- if (lastSucceeded) status = "ok";
+ let lastSucceeded =
+ !audioCheckCorruptSource && attemptSucceeded(primaryRes.exitCode);
+ if (lastSucceeded) {
+ status = audioCheckEnabled ? "ok-audio-checked" : "ok";
+ } else if (audioCheckCorruptSource) {
+ status = "failed-corrupt-source";
+ }
// ---------- Attempt 2: auth retry ----------
const shouldAuthRetry =
@@ -329,9 +400,13 @@ export async function downloadOneManaged(
}
// ---------- Attempt 3: no-subs fallback (youtube handling only) ----------
- const videoId =
- (await resolveVideoIdFromUrl(opts.videoUrl, channelDir)) ?? "unknown";
- const videoDir = path.join(channelDir, "data", videoId);
+ const videoDir = audioCheckVideoDir
+ ?? path.join(
+ channelDir,
+ "data",
+ (await resolveVideoIdFromUrl(opts.videoUrl, channelDir)) ?? "unknown",
+ );
+ const videoId = path.basename(videoDir);
if (
lastSucceeded &&
diff --git a/common/ytdlp/ffmpegStreamClassify.ts b/common/ytdlp/ffmpegStreamClassify.ts
@@ -0,0 +1,53 @@
+// Classifies the stderr output of an ffmpeg "decode-only" probe of a partial
+// (or finalized) media file. Used by the audio-checked download flow to
+// detect mid-stream corruption that would render the final transcode useless.
+//
+// Pure-function module: no side effects, no IO. Keeps test specs from having
+// to pull execa/etc. into the test process.
+
+export type StreamVerdict = "clean" | "partial" | "malformed";
+
+// Lines we treat as expected when probing an in-progress download (the
+// trailing bytes are mid-frame, so ffmpeg complains as it hits EOF). These
+// should not cause a rollback by themselves.
+export const PARTIAL_PATTERNS: readonly RegExp[] = [
+ /Invalid data found when processing input/i,
+ /Truncating packet/i,
+ /End of file/i,
+ /unexpected end of file/i,
+ /could not find codec parameters/i,
+ /moov atom not found/i,
+ /Partial frame/i,
+];
+
+// Lines we treat as definitive corruption signals. Drawn from real-world
+// libavcodec output observed against malformed Odysee streams.
+export const MALFORMED_PATTERNS: readonly RegExp[] = [
+ /Sample rate index in program config element does not match/i,
+ /Error while decoding stream/i,
+ /corrupt/i,
+ /invalid NAL/i,
+ /non[- ]?existing PPS/i,
+ /decode_slice_header error/i,
+ /channel element .* is not allocated/i,
+ /Number of bands \(\d+\) exceeds limit/i,
+ /reference picture missing/i,
+];
+
+export function classifyFfmpegStderr(stderr: string): StreamVerdict {
+ const lines = stderr
+ .split(/\r?\n/)
+ .map((l) => l.trim())
+ .filter((l) => l.length > 0);
+ if (lines.length === 0) return "clean";
+ let sawPartial = false;
+ for (const line of lines) {
+ if (MALFORMED_PATTERNS.some((p) => p.test(line))) return "malformed";
+ if (PARTIAL_PATTERNS.some((p) => p.test(line))) {
+ sawPartial = true;
+ continue;
+ }
+ return "malformed";
+ }
+ return sawPartial ? "partial" : "clean";
+}
diff --git a/common/ytdlp/ffmpegStreamProbe.ts b/common/ytdlp/ffmpegStreamProbe.ts
@@ -0,0 +1,55 @@
+import { execa } from "execa";
+import {
+ classifyFfmpegStderr,
+ type StreamVerdict,
+} from "./ffmpegStreamClassify";
+
+export type ProbeAudioStreamOptions = {
+ ffmpegBin: string;
+ file: string;
+ signal: AbortSignal;
+ onLog?: (s: string) => void;
+};
+
+export type ProbeAudioStreamResult = {
+ verdict: StreamVerdict;
+ stderr: string;
+ exitCode: number | null;
+};
+
+export async function probeAudioStream(
+ opts: ProbeAudioStreamOptions,
+): Promise<ProbeAudioStreamResult> {
+ const args = [
+ "-v",
+ "error",
+ "-nostdin",
+ "-xerror",
+ "-i",
+ opts.file,
+ "-f",
+ "null",
+ "-",
+ ];
+ opts.onLog?.(`$ ${opts.ffmpegBin} ${args.join(" ")}\n`);
+ const child = execa(opts.ffmpegBin, args, {
+ cancelSignal: opts.signal,
+ all: false,
+ buffer: false,
+ reject: false,
+ });
+ let stderr = "";
+ child.stderr?.on("data", (c: Buffer) => {
+ const chunk = c.toString("utf8");
+ stderr += chunk;
+ opts.onLog?.(chunk);
+ });
+ // Discard stdout (-f null -); we only care about stderr.
+ child.stdout?.on("data", () => {});
+ const result = await child;
+ return {
+ verdict: classifyFfmpegStderr(stderr),
+ stderr,
+ exitCode: result.exitCode ?? null,
+ };
+}
diff --git a/editor/app/channels/components/ChannelForm.tsx b/editor/app/channels/components/ChannelForm.tsx
@@ -1,4 +1,15 @@
import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig";
+import {
+ AUDIO_CHECK_COPY_TIMEOUT_DEFAULT_SECONDS,
+ AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS,
+ AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS,
+ AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS,
+ AUDIO_CHECK_INTERVAL_MAX_SECONDS,
+ AUDIO_CHECK_INTERVAL_MIN_SECONDS,
+ AUDIO_CHECK_MAX_ROLLBACKS_DEFAULT,
+ AUDIO_CHECK_MAX_ROLLBACKS_MAX,
+ AUDIO_CHECK_MAX_ROLLBACKS_MIN,
+} from "yt-dlp-transcript-common/lib/channelConfig";
type Props = {
action: string | ((formData: FormData) => void | Promise<void>);
@@ -131,6 +142,7 @@ export function ChannelForm({
</span>
</span>
</label>
+ <AudioCheckFields config={c} />
</Section>
<CollapsibleSection title="Advanced">
<label className="flex flex-col gap-1 text-sm">
@@ -230,6 +242,93 @@ function CollapsibleSection({
);
}
+function AudioCheckFields({ config }: { config?: ChannelConfig }) {
+ const ac = config?.audioCheck;
+ const enabled = ac?.enabled === true;
+ return (
+ <div className="flex flex-col gap-3 mt-1">
+ <label className="flex items-start gap-2 text-sm">
+ <input
+ type="checkbox"
+ name="audioCheckEnabled"
+ defaultChecked={enabled}
+ aria-label="audio-check enabled"
+ className="mt-1"
+ />
+ <span className="flex flex-col gap-0.5">
+ <span className="font-medium">
+ Audio-integrity checking{" "}
+ <span className="text-zinc-500 font-normal">
+ (transcribe only, slow)
+ </span>
+ </span>
+ <span className="text-xs text-zinc-500">
+ For sources that intermittently serve corrupt audio (e.g. Odysee
+ "original" format). Periodically probes the in-progress download
+ via ffmpeg and rolls back to the last known-good snapshot on
+ corruption. Disk usage peaks at ~3× file size; final probe decodes
+ the whole file.
+ </span>
+ </span>
+ </label>
+ <label className="flex flex-col gap-1 text-sm pl-6">
+ <span className="font-medium">
+ Check interval (seconds)
+ </span>
+ <input
+ type="number"
+ name="audioCheckIntervalSeconds"
+ defaultValue={ac?.intervalSeconds ?? ""}
+ placeholder={String(AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS)}
+ min={AUDIO_CHECK_INTERVAL_MIN_SECONDS}
+ max={AUDIO_CHECK_INTERVAL_MAX_SECONDS}
+ aria-label="audio-check interval seconds"
+ className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm w-32"
+ />
+ </label>
+ <label className="flex flex-col gap-1 text-sm pl-6">
+ <span className="font-medium">
+ Max consecutive rollbacks
+ </span>
+ <input
+ type="number"
+ name="audioCheckMaxRollbacks"
+ defaultValue={ac?.maxRollbacks ?? ""}
+ placeholder={String(AUDIO_CHECK_MAX_ROLLBACKS_DEFAULT)}
+ min={AUDIO_CHECK_MAX_ROLLBACKS_MIN}
+ max={AUDIO_CHECK_MAX_ROLLBACKS_MAX}
+ aria-label="audio-check max rollbacks"
+ className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm w-32"
+ />
+ <span className="text-xs text-zinc-500">
+ After this many rollbacks against the same checkpoint, the download
+ is marked <code>failed-corrupt-source</code>.
+ </span>
+ </label>
+ <label className="flex flex-col gap-1 text-sm pl-6">
+ <span className="font-medium">
+ Copy timeout (seconds)
+ </span>
+ <input
+ type="number"
+ name="audioCheckCopyTimeoutSeconds"
+ defaultValue={ac?.copyTimeoutSeconds ?? ""}
+ placeholder={String(AUDIO_CHECK_COPY_TIMEOUT_DEFAULT_SECONDS)}
+ min={AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS}
+ max={AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS}
+ aria-label="audio-check copy timeout seconds"
+ className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm w-32"
+ />
+ <span className="text-xs text-zinc-500">
+ Max time SIGSTOP can hold yt-dlp while we snapshot the <code>.part</code>{" "}
+ file. Without reflink (btrfs/xfs), large snapshots may exceed CDN
+ keepalive; the check skips itself rather than holding longer.
+ </span>
+ </label>
+ </div>
+ );
+}
+
function Field({
label,
name,
diff --git a/editor/app/channels/components/parseChannelForm.ts b/editor/app/channels/components/parseChannelForm.ts
@@ -1,5 +1,14 @@
import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig";
import {
+ AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS,
+ AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS,
+ AUDIO_CHECK_INTERVAL_MAX_SECONDS,
+ AUDIO_CHECK_INTERVAL_MIN_SECONDS,
+ AUDIO_CHECK_MAX_ROLLBACKS_MAX,
+ AUDIO_CHECK_MAX_ROLLBACKS_MIN,
+ type AudioCheckConfig,
+} from "yt-dlp-transcript-common/lib/channelConfig";
+import {
PLATFORM_VALUES,
detectPlatform,
type Platform,
@@ -22,6 +31,7 @@ export const CHANNEL_FORM_FIELDS = [
"keepSourceVideo",
"ytdlpExtraArgs",
"sleepBetweenDownloadsSeconds",
+ "audioCheck",
] as const satisfies ReadonlyArray<keyof ChannelConfig>;
// Maps FormData -> ChannelConfig. Throws on invalid input so the server
@@ -74,6 +84,60 @@ export function parseChannelForm(formData: FormData): ParsedChannelForm {
sleepBetweenDownloadsSeconds = n;
}
+ const audioCheckEnabled = formData.get("audioCheckEnabled") != null;
+ let audioCheck: AudioCheckConfig | undefined;
+ if (audioCheckEnabled) {
+ audioCheck = { enabled: true };
+ const intervalRaw = String(
+ formData.get("audioCheckIntervalSeconds") ?? "",
+ ).trim();
+ if (intervalRaw) {
+ const n = Number.parseInt(intervalRaw, 10);
+ if (
+ !Number.isFinite(n) ||
+ n < AUDIO_CHECK_INTERVAL_MIN_SECONDS ||
+ n > AUDIO_CHECK_INTERVAL_MAX_SECONDS
+ ) {
+ throw new Error(
+ `Audio-check interval must be ${AUDIO_CHECK_INTERVAL_MIN_SECONDS}–${AUDIO_CHECK_INTERVAL_MAX_SECONDS} seconds`,
+ );
+ }
+ audioCheck.intervalSeconds = n;
+ }
+ const rollbacksRaw = String(
+ formData.get("audioCheckMaxRollbacks") ?? "",
+ ).trim();
+ if (rollbacksRaw) {
+ const n = Number.parseInt(rollbacksRaw, 10);
+ if (
+ !Number.isFinite(n) ||
+ n < AUDIO_CHECK_MAX_ROLLBACKS_MIN ||
+ n > AUDIO_CHECK_MAX_ROLLBACKS_MAX
+ ) {
+ throw new Error(
+ `Audio-check max rollbacks must be ${AUDIO_CHECK_MAX_ROLLBACKS_MIN}–${AUDIO_CHECK_MAX_ROLLBACKS_MAX}`,
+ );
+ }
+ audioCheck.maxRollbacks = n;
+ }
+ const copyTimeoutRaw = String(
+ formData.get("audioCheckCopyTimeoutSeconds") ?? "",
+ ).trim();
+ if (copyTimeoutRaw) {
+ const n = Number.parseInt(copyTimeoutRaw, 10);
+ if (
+ !Number.isFinite(n) ||
+ n < AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS ||
+ n > AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS
+ ) {
+ throw new Error(
+ `Audio-check copy timeout must be ${AUDIO_CHECK_COPY_TIMEOUT_MIN_SECONDS}–${AUDIO_CHECK_COPY_TIMEOUT_MAX_SECONDS} seconds`,
+ );
+ }
+ audioCheck.copyTimeoutSeconds = n;
+ }
+ }
+
const config: ChannelConfig = {
handling: handlingRaw,
name,
@@ -86,6 +150,7 @@ export function parseChannelForm(formData: FormData): ParsedChannelForm {
if (sleepBetweenDownloadsSeconds != null) {
config.sleepBetweenDownloadsSeconds = sleepBetweenDownloadsSeconds;
}
+ if (audioCheck) config.audioCheck = audioCheck;
return { name, slug, config };
}
diff --git a/editor/e2e/audio-check-classifier.spec.ts b/editor/e2e/audio-check-classifier.spec.ts
@@ -0,0 +1,61 @@
+// Pure-function tests for the ffmpeg stderr classifier. These don't need the
+// dev server, fixtures, or a browser — but the project uses Playwright for
+// everything, so they live here too.
+
+import { test, expect } from "@playwright/test";
+import { classifyFfmpegStderr } from "../../common/ytdlp/ffmpegStreamClassify";
+
+test.describe("classifyFfmpegStderr", () => {
+ test("empty stderr is clean", () => {
+ expect(classifyFfmpegStderr("")).toBe("clean");
+ expect(classifyFfmpegStderr(" \n\n\t\n")).toBe("clean");
+ });
+
+ test("only partial-stream patterns → partial", () => {
+ expect(
+ classifyFfmpegStderr(
+ "[mov,mp4 @ 0x123] Invalid data found when processing input at EOF",
+ ),
+ ).toBe("partial");
+ expect(classifyFfmpegStderr("[aac @ 0x1] Truncating packet of size 137")).toBe(
+ "partial",
+ );
+ expect(
+ classifyFfmpegStderr(
+ "Invalid data found when processing input at EOF\nTruncating packet of size 200",
+ ),
+ ).toBe("partial");
+ });
+
+ test("known malformed patterns → malformed", () => {
+ expect(
+ classifyFfmpegStderr(
+ "[aac @ 0x1] Sample rate index in program config element does not match the sample rate index configured by the container.",
+ ),
+ ).toBe("malformed");
+ expect(classifyFfmpegStderr("Error while decoding stream #0:0")).toBe(
+ "malformed",
+ );
+ expect(classifyFfmpegStderr("[mp4 @ 0x1] corrupt frame")).toBe("malformed");
+ expect(classifyFfmpegStderr("[h264 @ 0x1] invalid NAL unit size")).toBe(
+ "malformed",
+ );
+ });
+
+ test("mixed partial + malformed → malformed wins", () => {
+ expect(
+ classifyFfmpegStderr(
+ "Truncating packet of size 137\n[aac @ 0x1] Error while decoding stream #0:0\nInvalid data found when processing input at EOF",
+ ),
+ ).toBe("malformed");
+ });
+
+ test("unrecognised error → malformed (conservative)", () => {
+ expect(classifyFfmpegStderr("Something completely unexpected")).toBe(
+ "malformed",
+ );
+ expect(
+ classifyFfmpegStderr("FormatIO failed reading from socket"),
+ ).toBe("malformed");
+ });
+});
diff --git a/editor/e2e/audio-check-scenarios.spec.ts b/editor/e2e/audio-check-scenarios.spec.ts
@@ -0,0 +1,201 @@
+import { writeFile } from "node:fs/promises";
+import { test, expect } from "@playwright/test";
+import { pathExists, readJson, resetData, resolvePath } from "./helpers";
+import type {
+ DownloadOutcomeRecord,
+} from "../../common/lib/downloadOutcome";
+
+const CHANNEL = "test-audio-check";
+const CHANNEL_ROOT = `test-transcripts/channels/${CHANNEL}`;
+// urlIdYouTube on "https://odysee.com/abc123" returns "abc123".
+const VIDEO_ID = "abc123";
+
+type FakeSidecar = {
+ mode: "happy" | "corrupt-after" | "corrupt-final" | "corrupt-from-start";
+ recoverOnResume?: boolean;
+ deterministicCorrupt?: boolean;
+ totalChunks?: number;
+ chunkDelayMs?: number;
+ corruptAfterChunk?: number;
+};
+
+async function writeFakeConfig(sidecar: FakeSidecar) {
+ await writeFile(
+ resolvePath(`${CHANNEL_ROOT}/.fake-ytdlp-audio-check.json`),
+ JSON.stringify(sidecar, null, 2),
+ );
+}
+
+async function triggerDownload(page: import("@playwright/test").Page) {
+ await page.goto(`/channels/${CHANNEL}`);
+ await page.getByRole("button", { name: "Download videos" }).click();
+}
+
+async function readOutcome(): Promise<DownloadOutcomeRecord> {
+ return readJson<DownloadOutcomeRecord>(
+ `${CHANNEL_ROOT}/data/${VIDEO_ID}/download-outcome.json`,
+ );
+}
+
+async function waitForOutcome(
+ predicate: (o: DownloadOutcomeRecord) => boolean,
+ message: string,
+): Promise<DownloadOutcomeRecord> {
+ // The orchestrator writes the outcome sidecar after a few additional
+ // steps (final ffmpeg probe + transcode + archive append) past the last
+ // "audio-check download complete" log line. Poll until both the file
+ // exists and matches the predicate.
+ let last: DownloadOutcomeRecord | null = null;
+ await expect.poll(
+ async () => {
+ try {
+ last = await readOutcome();
+ return predicate(last);
+ } catch {
+ return false;
+ }
+ },
+ {
+ timeout: 30_000,
+ message,
+ },
+ ).toBe(true);
+ return last!;
+}
+
+test.describe("audio-checked download scenarios", () => {
+ test("happy path: two clean checkpoints, finalize succeeds", async ({ page }) => {
+ await resetData("audio-check-channel");
+ await writeFakeConfig({ mode: "happy", totalChunks: 8, chunkDelayMs: 150 });
+ await triggerDownload(page);
+
+ const log = page.getByLabel("Download videos output");
+ await expect(log).toContainText("audio-check download complete", {
+ timeout: 30_000,
+ });
+
+ const outcome = await waitForOutcome(
+ (o) => o.status === "ok-audio-checked",
+ "happy path outcome",
+ );
+ expect(outcome.attempts).toHaveLength(1);
+ expect(outcome.attempts[0].kind).toBe("audio-checked-primary");
+ expect(outcome.attempts[0].audioCheck?.rollbacks).toBe(0);
+ expect(outcome.attempts[0].audioCheck?.restarts).toBe(0);
+
+ // Extracted final audio exists; intermediates removed.
+ expect(await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp3`)).toBe(true);
+ expect(
+ await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp4.part`),
+ ).toBe(false);
+ expect(
+ await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp4.part.good`),
+ ).toBe(false);
+ expect(
+ await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp4.part.testing`),
+ ).toBe(false);
+ });
+
+ test("mid-stream corruption rolls back to last good snapshot and recovers", async ({
+ page,
+ }) => {
+ await resetData("audio-check-channel");
+ await writeFakeConfig({
+ mode: "corrupt-after",
+ corruptAfterChunk: 3,
+ recoverOnResume: true,
+ totalChunks: 8,
+ chunkDelayMs: 150,
+ });
+ await triggerDownload(page);
+
+ const log = page.getByLabel("Download videos output");
+ await expect(log).toContainText("Rolling back to .good", {
+ timeout: 30_000,
+ });
+
+ const outcome = await waitForOutcome(
+ (o) => o.status === "ok-audio-checked",
+ "mid-rollback outcome",
+ );
+ expect(outcome.attempts[0].audioCheck?.rollbacks).toBeGreaterThanOrEqual(1);
+ expect(await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp3`)).toBe(true);
+ });
+
+ test("final-extraction probe failure triggers rollback and retry", async ({ page }) => {
+ await resetData("audio-check-channel");
+ await writeFakeConfig({
+ mode: "corrupt-final",
+ recoverOnResume: true,
+ totalChunks: 6,
+ chunkDelayMs: 100,
+ });
+ await triggerDownload(page);
+
+ const log = page.getByLabel("Download videos output");
+ await expect(log).toContainText("Final probe verdict: malformed", {
+ timeout: 30_000,
+ });
+
+ const outcome = await waitForOutcome(
+ (o) => o.status === "ok-audio-checked",
+ "final-malformed outcome",
+ );
+ expect(outcome.attempts[0].audioCheck?.rollbacks).toBeGreaterThanOrEqual(1);
+ expect(await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp3`)).toBe(true);
+ });
+
+ test("deterministic corruption exhausts retry cap and surfaces failed-corrupt-source", async ({
+ page,
+ }) => {
+ await resetData("audio-check-channel");
+ await writeFakeConfig({
+ mode: "corrupt-from-start",
+ deterministicCorrupt: true,
+ totalChunks: 4,
+ chunkDelayMs: 100,
+ });
+ await triggerDownload(page);
+
+ // The fixture sets maxRollbacks=3, so we expect status=failed-corrupt-source
+ // after the cap is exhausted. With these timings most rollbacks come from
+ // the final-stage probe (the .part is gone by the time the watcher fires).
+ const outcome = await waitForOutcome(
+ (o) => o.status === "failed-corrupt-source",
+ "failed-corrupt-source outcome",
+ );
+ expect(outcome.attempts[0].audioCheck?.rollbacks ?? 0)
+ .toBeGreaterThanOrEqual(1);
+ });
+
+ test("first-checkpoint malformed with no baseline restarts from byte 0 and recovers", async ({
+ page,
+ }) => {
+ await resetData("audio-check-channel");
+ await writeFakeConfig({
+ mode: "corrupt-from-start",
+ recoverOnResume: true,
+ totalChunks: 4,
+ chunkDelayMs: 100,
+ });
+ await triggerDownload(page);
+
+ const log = page.getByLabel("Download videos output");
+ // Either the watcher catches mid-stream and logs "restarting from 0",
+ // or the final probe catches it; with deterministicCorrupt=false the
+ // resumed run is clean either way.
+ await expect(log).toContainText(/restarting from 0|Final probe verdict: malformed/, {
+ timeout: 30_000,
+ });
+
+ const outcome = await waitForOutcome(
+ (o) => o.status === "ok-audio-checked",
+ "no-baseline restart outcome",
+ );
+ expect(
+ (outcome.attempts[0].audioCheck?.restarts ?? 0) +
+ (outcome.attempts[0].audioCheck?.rollbacks ?? 0),
+ ).toBeGreaterThanOrEqual(1);
+ expect(await pathExists(`${CHANNEL_ROOT}/data/${VIDEO_ID}/audio.mp3`)).toBe(true);
+ });
+});
diff --git a/editor/e2e/fixtures/bin/fake-ffmpeg.mjs b/editor/e2e/fixtures/bin/fake-ffmpeg.mjs
@@ -1,19 +1,74 @@
#!/usr/bin/env node
-// E2E fake ffmpeg. Mimics the args common/controller/transcode.ts passes:
-// ffmpeg -y -i <src> -vn -c:a <codec> <out>
-// Writes a small placeholder to <out> so the real transcode pipeline (rename
-// from tmp → audio.<fmt>) succeeds.
-import { writeFile } from "node:fs/promises";
+// E2E fake ffmpeg. Recognises two invocation shapes:
+//
+// 1. Transcode: ffmpeg -y -i <src> -vn -c:a <codec> <out>
+// -> writes a small placeholder to <out>. If <src> contains the
+// CORRUPT_MARKER sentinel, exits 1 with a libavcodec-style stderr
+// line so the audio-check classifier flags it as malformed.
+//
+// 2. Probe: ffmpeg -v error -nostdin -xerror -i <file> -f null -
+// -> reads <file>; if it contains CORRUPT_MARKER, writes a
+// libavcodec-style error to stderr and exits 1. Else exits 0
+// silently. Used by the audio-checked download orchestrator.
+import { readFile, writeFile } from "node:fs/promises";
+
+const CORRUPT_MARKER = "__CORRUPT__";
const argv = process.argv.slice(2);
+
+function arg(flag) {
+ const i = argv.indexOf(flag);
+ if (i < 0) return undefined;
+ return argv[i + 1];
+}
+
+async function fileHasCorruptMarker(file) {
+ try {
+ const buf = await readFile(file);
+ return buf.includes(CORRUPT_MARKER);
+ } catch {
+ return false;
+ }
+}
+
+const src = arg("-i") ?? "<unknown>";
+
+// Probe shape: ends with `-f null -`. Stdout suppressed; only stderr matters.
+const last = argv[argv.length - 1];
+const secondLast = argv[argv.length - 2];
+if (secondLast === "-f" && (last === "null" || last === "null,")) {
+ if (await fileHasCorruptMarker(src)) {
+ process.stderr.write(
+ `[aac @ 0x1234] Sample rate index in program config element does not match the sample rate index configured by the container.\n`,
+ );
+ process.exit(1);
+ }
+ process.exit(0);
+}
+// Also accept the older form `-f null -` with `-` as last positional.
+if (last === "-" && argv[argv.length - 3] === "-f" && secondLast === "null") {
+ if (await fileHasCorruptMarker(src)) {
+ process.stderr.write(
+ `[aac @ 0x1234] Sample rate index in program config element does not match the sample rate index configured by the container.\n`,
+ );
+ process.exit(1);
+ }
+ process.exit(0);
+}
+
+// Transcode shape: out is the last positional (not a flag).
const out = argv[argv.length - 1];
if (!out || out.startsWith("-")) {
process.stderr.write(`[fake-ffmpeg] missing output path\n`);
process.exit(2);
}
-const iIdx = argv.indexOf("-i");
-const src = iIdx >= 0 ? argv[iIdx + 1] : "<unknown>";
+if (await fileHasCorruptMarker(src)) {
+ process.stderr.write(
+ `[aac @ 0x1234] Sample rate index in program config element does not match the sample rate index configured by the container.\n`,
+ );
+ process.exit(1);
+}
await writeFile(out, `fake-ffmpeg transcoded from ${src}\n`);
process.stdout.write(`[fake-ffmpeg] wrote ${out}\n`);
diff --git a/editor/e2e/fixtures/bin/fake-ytdlp.mjs b/editor/e2e/fixtures/bin/fake-ytdlp.mjs
@@ -12,10 +12,17 @@
// rest of the editor's controllers stay happy.
//
// Cwd is always the channel root (matches the editor's runYtdlp).
-import { mkdir, writeFile, readFile, appendFile } from "node:fs/promises";
-import { existsSync } from "node:fs";
+import { mkdir, writeFile, readFile, appendFile, stat, rename } from "node:fs/promises";
+import { existsSync, openSync, writeSync, closeSync } from "node:fs";
import path from "node:path";
+// Used in audio-check scenarios. The fake-ffmpeg companion treats files
+// containing this string as malformed when probing.
+const CORRUPT_MARKER = "__CORRUPT__";
+// Tests use small chunks so the size gate (defaults to 1 MiB; tests override
+// via AUDIO_CHECK_SIZE_GATE_OVERRIDE) can be exceeded each chunk. ~16 KiB.
+const CHUNK_BYTES = 16 * 1024;
+
const argv = process.argv.slice(2);
function has(flag) {
@@ -133,6 +140,142 @@ async function modeDownloadOneUrl(url, opts) {
process.stdout.write(`[fake-ytdlp] download complete\n`);
}
+// Audio-check mode: yt-dlp args drop -x but keep -f bestaudio/worst. We
+// simulate a chunked HTTP download by writing audio.mp4.part in pieces with
+// short sleeps between, optionally inserting CORRUPT_MARKER. Supports
+// SIGSTOP/SIGCONT (the kernel handles it for us; we just need to ensure
+// signal handlers don't block) and -c resume by appending to an existing
+// .part. Re-invocations log to fake-ytdlp.invocations so tests can count
+// restarts/rollbacks.
+async function readAudioCheckSidecar() {
+ // Tests configure the fake's per-run behaviour via this file in cwd
+ // (the channel root). Falls back to env vars for ad-hoc use.
+ try {
+ const raw = await readFile(".fake-ytdlp-audio-check.json", "utf8");
+ return JSON.parse(raw);
+ } catch {
+ return {};
+ }
+}
+
+async function modeAudioCheckOneUrl(url) {
+ const sidecar = await readAudioCheckSidecar();
+ const mode =
+ sidecar.mode ?? process.env.FAKE_YTDLP_AUDIO_CHECK_MODE ?? "happy";
+ const recoverOnResume =
+ sidecar.recoverOnResume === true ||
+ process.env.FAKE_YTDLP_RECOVER_ON_RESUME === "1";
+ const deterministicCorrupt =
+ sidecar.deterministicCorrupt === true ||
+ process.env.FAKE_YTDLP_DETERMINISTIC_CORRUPT === "1";
+ const totalChunks = Number.parseInt(
+ String(
+ sidecar.totalChunks ?? process.env.FAKE_YTDLP_TOTAL_CHUNKS ?? "6",
+ ),
+ 10,
+ );
+ const chunkDelayMs = Number.parseInt(
+ String(
+ sidecar.chunkDelayMs ?? process.env.FAKE_YTDLP_CHUNK_DELAY_MS ?? "200",
+ ),
+ 10,
+ );
+
+ const id = urlIdYouTube(url) ?? "ackid";
+ const videoDir = path.join("data", id);
+ await ensureDir(videoDir);
+ // Minimal metadata so downstream code doesn't choke if it inspects it.
+ if (!existsSync(path.join(videoDir, "metadata.info.json"))) {
+ await writeMetadata(videoDir, id);
+ }
+ const partFile = path.join(videoDir, "audio.mp4.part");
+ const finalFile = path.join(videoDir, "audio.mp4");
+ const isResume = existsSync(partFile);
+ // Count prior audio-check invocations so the fake can simulate "transient
+ // CDN hiccups" — first run emits the marker, subsequent runs don't even
+ // if .part was deleted on rollback (no-baseline-restart scenario).
+ let priorAudioCheckRuns = 0;
+ try {
+ const invs = await readFile("fake-ytdlp.invocations", "utf8");
+ priorAudioCheckRuns = invs.split("\n").filter((l) => l.startsWith("audio-check:")).length;
+ } catch {}
+ await appendFile(
+ "fake-ytdlp.invocations",
+ `audio-check:${mode}:resume=${isResume ? 1 : 0}:url=${url}\n`,
+ );
+
+ // Resume semantics: when `recoverOnResume` is set, ANY re-invocation
+ // (whether resuming via .part or restarting from 0) skips the marker.
+ // `deterministicCorrupt` forces the marker every run regardless.
+ const shouldEmitMarker = (() => {
+ if (mode === "happy") return false;
+ if (deterministicCorrupt) return true;
+ if (recoverOnResume && priorAudioCheckRuns > 0) return false;
+ return true;
+ })();
+
+ // Open .part for append (resume) or create.
+ const fd = openSync(partFile, isResume ? "a" : "w");
+
+ // Resume support for SIGCONT logging — install a no-op handler so that
+ // SIGCONT delivery flips a flag that we log on each chunk.
+ process.on("SIGCONT", () => {
+ try {
+ // Log to a sidecar so tests can prove SIGCONT landed.
+ const f = openSync("fake-ytdlp.signals", "a");
+ writeSync(f, `resumed ${new Date().toISOString()}\n`);
+ closeSync(f);
+ } catch {}
+ });
+
+ const corruptAfter =
+ mode === "corrupt-after"
+ ? Number.parseInt(
+ String(
+ sidecar.corruptAfterChunk ??
+ process.env.FAKE_YTDLP_CORRUPT_AFTER_CHUNK ??
+ "2",
+ ),
+ 10,
+ )
+ : -1;
+
+ // Determine starting chunk by current size.
+ let writtenBytes = isResume ? (await stat(partFile)).size : 0;
+ // Each chunk is CHUNK_BYTES of filler. Compute already-written chunks.
+ const alreadyChunks = Math.floor(writtenBytes / CHUNK_BYTES);
+
+ for (let i = alreadyChunks; i < totalChunks; i++) {
+ if (chunkDelayMs > 0) {
+ await sleep(chunkDelayMs);
+ }
+ let chunk = `chunk-${i} `.padEnd(CHUNK_BYTES, "x");
+ // Stream-corrupt: marker landing at chunk corruptAfter (so checkpoints
+ // before that pass clean, the one after fails). Only for non-"final"
+ // modes.
+ if (
+ shouldEmitMarker &&
+ (mode === "corrupt-from-start" ||
+ (mode === "corrupt-after" && i === corruptAfter))
+ ) {
+ chunk = `${CORRUPT_MARKER} ${chunk}`.padEnd(CHUNK_BYTES, "x");
+ }
+ writeSync(fd, chunk);
+ }
+ closeSync(fd);
+
+ // Final-corrupt: appended marker only into the renamed file.
+ if (mode === "corrupt-final" && shouldEmitMarker) {
+ await rename(partFile, finalFile);
+ await appendFile(finalFile, CORRUPT_MARKER);
+ } else {
+ await rename(partFile, finalFile);
+ }
+ // Print archive line in the shape downloadOneManaged expects.
+ process.stdout.write(`DLOM_ARCHIVE youtube ${id}\n`);
+ process.stdout.write(`[fake-ytdlp] audio-check download complete\n`);
+}
+
async function modeSync(url) {
process.stdout.write(`[fake-ytdlp] --lazy-playlist sync against ${url}\n`);
await appendFile("fake-ytdlp.invocations", "sync:1\n");
@@ -240,6 +383,17 @@ async function main() {
return;
}
+ // Audio-check mode: no -x but has -f bestaudio/worst and -c.
+ if (has("-c") && arg("-f") === "bestaudio/worst") {
+ const url = lastNonFlag();
+ if (!url) {
+ process.stderr.write(`[fake-ytdlp] audio-check mode missing URL\n`);
+ process.exit(2);
+ }
+ await modeAudioCheckOneUrl(url);
+ return;
+ }
+
process.stderr.write(
`[fake-ytdlp] unknown invocation: ${argv.join(" ")}\n`,
);
diff --git a/editor/e2e/fixtures/test-transcripts/audio-check-channel/channels/test-audio-check/config.json b/editor/e2e/fixtures/test-transcripts/audio-check-channel/channels/test-audio-check/config.json
@@ -0,0 +1,13 @@
+{
+ "handling": "transcribe",
+ "name": "Audio Check Test Channel",
+ "platform": "odysee",
+ "url": "https://odysee.com/@example",
+ "audioFormat": "mp3",
+ "audioCheck": {
+ "enabled": true,
+ "intervalSeconds": 10,
+ "maxRollbacks": 3,
+ "copyTimeoutSeconds": 5
+ }
+}
diff --git a/editor/e2e/fixtures/test-transcripts/audio-check-channel/channels/test-audio-check/playlist b/editor/e2e/fixtures/test-transcripts/audio-check-channel/channels/test-audio-check/playlist
@@ -0,0 +1 @@
+https://odysee.com/abc123
diff --git a/editor/package.json b/editor/package.json
@@ -5,8 +5,8 @@
"type": "module",
"scripts": {
"dev": "next dev --port 3001",
- "dev:test": "TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs next dev --port 3011",
- "start:test": "TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs next start --port 3011",
+ "dev:test": "TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 next dev --port 3011",
+ "start:test": "TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 next start --port 3011",
"build": "next build",
"start": "next start --port 3001",
"lint": "eslint",