diff --git a/packages/engine/src/index.ts b/packages/engine/src/index.ts index 4fd3df11fa..98507b85bb 100644 --- a/packages/engine/src/index.ts +++ b/packages/engine/src/index.ts @@ -343,6 +343,7 @@ export { export { runFfmpeg, formatFfmpegError, + describeFfmpegFailure, isExternalFfmpegInterruption, type RunFfmpegOptions, type RunFfmpegResult, diff --git a/packages/engine/src/services/chunkEncoder.test.ts b/packages/engine/src/services/chunkEncoder.test.ts index d664ff2758..3ad3dbb566 100644 --- a/packages/engine/src/services/chunkEncoder.test.ts +++ b/packages/engine/src/services/chunkEncoder.test.ts @@ -430,6 +430,80 @@ describe("encodeFramesChunkedConcat ffmpegEncodeTimeout", () => { }); }); +describe("final mux and faststart ffmpegProcessTimeout", () => { + afterEach(() => vi.useRealTimers()); + + it("reports a mux killed at the configured timeout as a timeout", async () => { + vi.useFakeTimers(); + const { spawn, calls } = createSpawnSpy(); + vi.resetModules(); + vi.doMock("child_process", () => ({ spawn })); + + const { muxVideoWithAudio } = await import("./chunkEncoder.js"); + const muxPromise = muxVideoWithAudio( + "/tmp/video-only.mp4", + "/tmp/audio.aac", + "/tmp/output.mp4", + undefined, + { ffmpegProcessTimeout: 1000 }, + ); + + await flushMuxCodecResolution(); + const proc = calls[0]!.proc; + vi.advanceTimersByTime(999); + expect(proc.kill).not.toHaveBeenCalled(); + vi.advanceTimersByTime(1); + expect(proc.kill).toHaveBeenCalledWith("SIGTERM"); + proc.stderr.emit("data", Buffer.from("frame= 9000 time=00:05:00.00\n")); + emitClose(proc, 255); + + const result = await muxPromise; + expect(result.success).toBe(false); + expect(result.error).toMatch(/^FFmpeg timed out after 1000 ms/); + expect(result.error).toContain("FFMPEG_PROCESS_TIMEOUT_MS"); + expect(result.error).toContain("frame= 9000"); + }); + + it("reports a faststart killed at the configured timeout as a timeout", async () => { + vi.useFakeTimers(); + const { spawn, calls } = createSpawnSpy(); + vi.resetModules(); + vi.doMock("child_process", () => ({ spawn })); + + const { applyFaststart } = await import("./chunkEncoder.js"); + const faststartPromise = applyFaststart("/tmp/video-only.mp4", "/tmp/output.mp4", undefined, { + ffmpegProcessTimeout: 2000, + }); + + const proc = calls[0]!.proc; + vi.advanceTimersByTime(1999); + expect(proc.kill).not.toHaveBeenCalled(); + vi.advanceTimersByTime(1); + expect(proc.kill).toHaveBeenCalledWith("SIGTERM"); + emitClose(proc, 255); + + const result = await faststartPromise; + expect(result.success).toBe(false); + expect(result.error).toMatch(/^FFmpeg timed out after 2000 ms/); + }); + + it("leaves an ordinary mux failure's message as it was", async () => { + const { spawn, calls } = createSpawnSpy(); + vi.resetModules(); + vi.doMock("child_process", () => ({ spawn })); + + const { muxVideoWithAudio } = await import("./chunkEncoder.js"); + const muxPromise = muxVideoWithAudio("/tmp/video-only.mp4", "/tmp/audio.aac", "/tmp/out.mp4"); + await flushMuxCodecResolution(); + calls[0]!.proc.stderr.emit("data", Buffer.from("Invalid data found when processing input\n")); + emitClose(calls[0]!.proc, 1); + + const result = await muxPromise; + expect(result.error).toMatch(/^FFmpeg exited with code 1/); + expect(result.error).not.toContain("timed out"); + }); +}); + describe("muxVideoWithAudio audio codec handling", () => { it("preserves an external interruption from mux", async () => { const { spawn, calls } = createSpawnSpy(); diff --git a/packages/engine/src/services/chunkEncoder.ts b/packages/engine/src/services/chunkEncoder.ts index 28243068a3..777f3b2dbc 100644 --- a/packages/engine/src/services/chunkEncoder.ts +++ b/packages/engine/src/services/chunkEncoder.ts @@ -33,6 +33,7 @@ import { type HdrTransfer, getHdrEncoderColorParams } from "../utils/hdr.js"; import { withEvenDimensionPad } from "../utils/evenDimensions.js"; import { SDR_CAPTURE_TO_BT709_FILTER } from "../utils/sdrCaptureColor.js"; import { + describeFfmpegFailure, ffmpegStatsReader, formatFfmpegError, isExternalFfmpegInterruption, @@ -845,7 +846,7 @@ export async function muxVideoWithAudio( success: result.success, outputPath, durationMs: result.durationMs, - error: !result.success ? formatFfmpegError(result.exitCode, result.stderr) : undefined, + error: !result.success ? describeFfmpegFailure(result, processTimeout) : undefined, failureReason: result.failureReason, }; } @@ -993,7 +994,7 @@ export async function packageHls( success: result.success, outputPath: outputDir, durationMs: result.durationMs, - error: !result.success ? formatFfmpegError(result.exitCode, result.stderr) : undefined, + error: !result.success ? describeFfmpegFailure(result, processTimeout) : undefined, failureReason: result.failureReason, }; } @@ -1041,7 +1042,7 @@ export async function applyFaststart( success: result.success, outputPath, durationMs: result.durationMs, - error: !result.success ? formatFfmpegError(result.exitCode, result.stderr) : undefined, + error: !result.success ? describeFfmpegFailure(result, processTimeout) : undefined, failureReason: result.failureReason, }; } diff --git a/packages/engine/src/services/videoFrameExtractor.concurrency.test.ts b/packages/engine/src/services/videoFrameExtractor.concurrency.test.ts new file mode 100644 index 0000000000..502aa81645 --- /dev/null +++ b/packages/engine/src/services/videoFrameExtractor.concurrency.test.ts @@ -0,0 +1,140 @@ +import { chmodSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; +import type { VideoElement } from "../types.js"; +import { FFMPEG_PATH_ENV, getFfmpegBinary } from "../utils/ffmpegBinaries.js"; +import { runFfmpeg } from "../utils/runFfmpeg.js"; + +// Two CPUs: the extractor runs at most one ffmpeg at a time. +vi.mock("os", async (importOriginal) => ({ + ...(await importOriginal()), + cpus: () => [{}, {}], +})); + +const { extractAllVideoFrames, extractVideoFramesRange } = await import("./videoFrameExtractor.js"); +const { extractMediaMetadata } = await import("../utils/ffprobe.js"); + +describe.skipIf(process.platform === "win32")("extractAllVideoFrames ffmpeg concurrency", () => { + const dir = mkdtempSync(join(tmpdir(), "hf-extract-concurrency-")); + const previousFfmpeg = process.env[FFMPEG_PATH_ENV]; + const shim = join(dir, "ffmpeg-shim.sh"); + + async function synth(name: string, lavfi: string, extra: string[] = []): Promise { + const clip = join(dir, `${name}.mp4`); + const made = await runFfmpeg([ + ...["-y", "-v", "error", "-f", "lavfi", "-i", lavfi, ...extra], + ...["-c:v", "libx264", "-preset", "ultrafast", "-pix_fmt", "yuv420p", clip], + ]); + expect(made.success, made.stderr).toBe(true); + return clip; + } + + // Extracts `clips` through a shim ffmpeg and returns how many ffmpegs were alive as each started. + async function aliveCounts(run: string, clips: string[], end: number, fps: number) { + mkdirSync(join(dir, run, "running"), { recursive: true }); + const videos: VideoElement[] = clips.map((src, i) => ({ + ...{ id: `${run}-${i}`, src, start: 0, end, mediaStart: 0 }, + ...{ loop: false, hasAudio: false }, + })); + process.env.HF_TEST_SHIM_DIR = join(dir, run); + process.env[FFMPEG_PATH_ENV] = shim; + try { + const result = await extractAllVideoFrames(videos, dir, { + fps, + outputDir: join(dir, run, "out"), + }); + expect(result.errors).toEqual([]); + } finally { + if (previousFfmpeg === undefined) delete process.env[FFMPEG_PATH_ENV]; + else process.env[FFMPEG_PATH_ENV] = previousFfmpeg; + delete process.env.HF_TEST_SHIM_DIR; + } + return readFileSync(join(dir, run, "alive"), "utf8") + .trim() + .split(/\s+/) + .map(Number); + } + + beforeAll(() => { + writeFileSync( + shim, + [ + "#!/bin/sh", + 'm="$HF_TEST_SHIM_DIR/running/$$"; : > "$m"', + 'ls "$HF_TEST_SHIM_DIR/running" | wc -l >> "$HF_TEST_SHIM_DIR/alive"', + 'case "$*" in *held-*) until [ -f "$HF_TEST_SHIM_DIR/release" ] || [ ! -d "$HF_TEST_SHIM_DIR" ]; do sleep 0.05; done ;; esac', + `"${getFfmpegBinary()}" "$@"; rc=$?`, + 'rm -f "$m"; exit $rc', + ].join("\n"), + ); + chmodSync(shim, 0o755); + }); + + afterAll(() => rmSync(dir, { recursive: true, force: true })); + + it("runs one clip's ffmpeg at a time when only one slot is free", async () => { + const clips = await Promise.all( + [0, 1, 2, 3].map((i) => synth(`short-${i}`, `testsrc=s=32x32:d=1:r=10,hue=h=${i * 60}`)), + ); + const alive = await aliveCounts("short", clips, 1, 10); + expect(alive).toHaveLength(4); + expect(Math.max(...alive)).toBe(1); + }, 60_000); + + it("shares the slot with the segments of long clips", async () => { + const clips = await Promise.all( + [0, 1].map((i) => synth(`long-${i}`, `testsrc=s=32x32:d=125:r=1,hue=h=${i * 90}`)), + ); + const alive = await aliveCounts("long", clips, 125, 1); + expect(alive.length).toBeGreaterThanOrEqual(4); + expect(Math.max(...alive)).toBe(1); + }, 120_000); + + it("gives a variable-frame-rate clip's two-process pipeline one slot", async () => { + const clips = await Promise.all( + [0, 1].map((i) => + synth(`vfr-${i}`, `color=c=0x${i ? "28C83C" : "C83C28"}:s=32x32:d=2:r=60`, [ + ...["-vf", "select='not(between(n\\,30\\,89))'", "-fps_mode", "vfr"], + ]), + ), + ); + const alive = await aliveCounts("vfr", clips, 1, 30); + expect(alive).toHaveLength(4); + expect(Math.max(...alive)).toBeLessThanOrEqual(2); + }, 60_000); + + it("lets a cancelled extraction leave the queue while another holds the slot", async () => { + const [held, queued] = await Promise.all([ + synth("held-0", "testsrc=s=32x32:d=1:r=10"), + synth("queued-0", "testsrc=s=32x32:d=1:r=10,hue=h=90"), + ]); + await extractMediaMetadata(queued); + const run = join(dir, "cancel"); + mkdirSync(join(run, "running"), { recursive: true }); + const extract = (id: string, src: string, signal?: AbortSignal) => + extractVideoFramesRange(src, id, 0, 1, { fps: 10, outputDir: run }, signal); + process.env.HF_TEST_SHIM_DIR = run; + process.env[FFMPEG_PATH_ENV] = shim; + try { + const holder = extract("held", held); + await vi.waitFor(() => expect(readFileSync(join(run, "alive"), "utf8").trim()).toBe("1"), { + timeout: 10_000, + }); + const cancel = new AbortController(); + const cancelled = extract("queued", queued, cancel.signal); + await new Promise(setImmediate); + cancel.abort(); + await expect(cancelled).rejects.toThrow(/cancelled/); + expect(readFileSync(join(run, "alive"), "utf8").trim()).toBe("1"); + + writeFileSync(join(run, "release"), ""); + expect((await holder).totalFrames).toBe(10); + } finally { + writeFileSync(join(run, "release"), ""); + if (previousFfmpeg === undefined) delete process.env[FFMPEG_PATH_ENV]; + else process.env[FFMPEG_PATH_ENV] = previousFfmpeg; + delete process.env.HF_TEST_SHIM_DIR; + } + }, 30_000); +}); diff --git a/packages/engine/src/services/videoFrameExtractor.ts b/packages/engine/src/services/videoFrameExtractor.ts index 1dfee36d99..bdf26ac4e0 100644 --- a/packages/engine/src/services/videoFrameExtractor.ts +++ b/packages/engine/src/services/videoFrameExtractor.ts @@ -10,6 +10,7 @@ import { audioGroupsById, isMemberGroupHidden, isSelfOrAncestorHidden } from "./ import { copyFileSync, existsSync, linkSync, mkdirSync, rmSync } from "fs"; import { join } from "path"; import { cpus } from "os"; +import { createConcurrencyLimit } from "../utils/concurrencyLimit.js"; import { parseHTML } from "linkedom"; import { resolveProjectRelativeSrc } from "@hyperframes/parsers/asset-resolution"; import { @@ -885,6 +886,10 @@ function restoreSourceColourFilter(metadata: VideoMetadata): string[] { */ const EXTRACTION_SEGMENT_SECONDS = 120; +// Caps extraction ffmpeg runs across all clips; uncapped, dozens of clips starved each other into timeouts. +const EXTRACTION_SLOTS = Math.max(1, Math.floor(cpus().length / 2)); +const inExtractionSlot = createConcurrencyLimit(EXTRACTION_SLOTS); + interface SegmentedExtraction { decodeArgs: string[]; filterAndEncodeArgs: string[]; @@ -923,30 +928,34 @@ async function runSegmentedExtraction(job: SegmentedExtraction): Promise + runFfmpeg( + [ + ...decodeArgs, + "-noaccurate_seek", + "-ss", + String(startTime + firstFrame / fps), + "-i", + videoPath, + // Read one frame past the segment and keep exactly its own frames. + "-t", + String((frames + 1) / fps), + "-frames:v", + String(frames), + "-start_number", + String(firstFrame + 1), + ...filterAndEncodeArgs, + ], + { ...job.runOptions, signal }, + ), + signal, ); results[index] = result; if (!result.success) failed.abort(); } }; - const workers = Math.min(segmentCount, Math.max(1, Math.floor(cpus().length / 2))); + const workers = Math.min(segmentCount, EXTRACTION_SLOTS); await Promise.all(Array.from({ length: workers }, worker)); const ran = results.filter(Boolean); // Report the segment that failed, not the ones stopped because of it. @@ -1106,20 +1115,24 @@ export async function extractVideoFramesRange( ) { // ffmpeg <6.1 ignores frame durations in CFR resampling once any -vf is set, // cutting a trailing still short, so the SDR filters run in a second process. - processResult = await runFfmpegPipeline( - [...args, "-an", "-sn", "-dn", "-c:v", "rawvideo", "-f", "nut", "pipe:1"], - [ - "-f", - "nut", - "-i", - "pipe:0", - "-vf", - [...restoreSourceColourFilter(metadata), ...vfFilters].join(","), - "-fps_mode", - "passthrough", - ...encodeArgs, - ], - runOptions, + processResult = await inExtractionSlot( + () => + runFfmpegPipeline( + [...args, "-an", "-sn", "-dn", "-c:v", "rawvideo", "-f", "nut", "pipe:1"], + [ + "-f", + "nut", + "-i", + "pipe:0", + "-vf", + [...restoreSourceColourFilter(metadata), ...vfFilters].join(","), + "-fps_mode", + "passthrough", + ...encodeArgs, + ], + runOptions, + ), + signal, ); } else { const filterArgs = vfFilters.length > 0 ? ["-vf", vfFilters.join(",")] : []; @@ -1139,7 +1152,10 @@ export async function extractVideoFramesRange( ), runOptions, }) - : await runFfmpeg([...args, ...filterArgs, ...encodeArgs], runOptions); + : await inExtractionSlot( + () => runFfmpeg([...args, ...filterArgs, ...encodeArgs], runOptions), + signal, + ); } if (processResult.failureReason === "external_interruption") { throw new VideoSourceExtractionError( diff --git a/packages/engine/src/utils/concurrencyLimit.test.ts b/packages/engine/src/utils/concurrencyLimit.test.ts new file mode 100644 index 0000000000..9efb48301f --- /dev/null +++ b/packages/engine/src/utils/concurrencyLimit.test.ts @@ -0,0 +1,81 @@ +import { describe, expect, it } from "vitest"; +import { createConcurrencyLimit } from "./concurrencyLimit.js"; + +function deferred() { + let resolve!: () => void; + const promise = new Promise((res) => (resolve = res)); + return { promise, resolve }; +} + +const flush = () => new Promise((resolve) => setImmediate(resolve)); + +describe("createConcurrencyLimit", () => { + it("never runs more than the limit at once, and starts waiters in order", async () => { + const limit = createConcurrencyLimit(2); + const gates = Array.from({ length: 5 }, deferred); + const started: number[] = []; + let running = 0; + let peak = 0; + const all = gates.map((gate, i) => + limit(async () => { + started.push(i); + peak = Math.max(peak, ++running); + await gate.promise; + running--; + }), + ); + await flush(); + expect(started).toEqual([0, 1]); + gates[1]!.resolve(); + await flush(); + expect(started).toEqual([0, 1, 2]); + gates[0]!.resolve(); + gates[2]!.resolve(); + await flush(); + expect(started).toEqual([0, 1, 2, 3, 4]); + gates[3]!.resolve(); + gates[4]!.resolve(); + await Promise.all(all); + expect(peak).toBe(2); + }); + + it("frees the slot when a task fails", async () => { + const limit = createConcurrencyLimit(1); + const failing = limit(() => Promise.reject(new Error("ffmpeg failed"))); + const next = limit(async () => "ran"); + await expect(failing).rejects.toThrow("ffmpeg failed"); + await expect(next).resolves.toBe("ran"); + }); + + it("lets a cancelled waiter leave the queue at once, without taking a slot", async () => { + const limit = createConcurrencyLimit(1); + const first = deferred(); + const started: string[] = []; + const running = limit(async () => { + started.push("first"); + await first.promise; + }); + const cancel = new AbortController(); + const cancelled = limit(async () => void started.push("cancelled"), cancel.signal); + const later = limit(async () => void started.push("later")); + await flush(); + expect(started).toEqual(["first"]); + + cancel.abort(); + await cancelled; + expect(started).toEqual(["first", "cancelled"]); + + first.resolve(); + await Promise.all([running, later]); + expect(started).toEqual(["first", "cancelled", "later"]); + }); + + it("runs a task whose signal already aborted without waiting", async () => { + const limit = createConcurrencyLimit(1); + const first = deferred(); + const running = limit(() => first.promise); + await expect(limit(async () => "ran", AbortSignal.abort())).resolves.toBe("ran"); + first.resolve(); + await running; + }); +}); diff --git a/packages/engine/src/utils/concurrencyLimit.ts b/packages/engine/src/utils/concurrencyLimit.ts new file mode 100644 index 0000000000..3e57e2cff1 --- /dev/null +++ b/packages/engine/src/utils/concurrencyLimit.ts @@ -0,0 +1,36 @@ +/** + * Runs at most `limit` tasks at once; later tasks wait, in order, for a free slot. A waiter whose + * `signal` aborts stops waiting and runs at once without a slot, so a cancel never queues behind others. + */ +export function createConcurrencyLimit( + limit: number, +): (task: () => Promise, signal?: AbortSignal) => Promise { + let active = 0; + const waiting: Array<() => void> = []; + const waitForSlot = (signal?: AbortSignal) => + new Promise((resolve) => { + const granted = () => { + signal?.removeEventListener("abort", aborted); + resolve(true); + }; + const aborted = () => { + waiting.splice(waiting.indexOf(granted), 1); + resolve(false); + }; + waiting.push(granted); + signal?.addEventListener("abort", aborted, { once: true }); + }); + return async (task, signal) => { + if (signal?.aborted) return task(); + if (active < limit) active++; + else if (!(await waitForSlot(signal))) return task(); + try { + return await task(); + } finally { + // Hand the slot straight to the next waiter so a newcomer cannot take it first. + const next = waiting.shift(); + if (next) next(); + else active--; + } + }; +} diff --git a/packages/engine/src/utils/runFfmpeg.test.ts b/packages/engine/src/utils/runFfmpeg.test.ts index 2bd153fad8..402685284e 100644 --- a/packages/engine/src/utils/runFfmpeg.test.ts +++ b/packages/engine/src/utils/runFfmpeg.test.ts @@ -320,3 +320,26 @@ describe("ffmpegStatsReader", () => { expect(stats).toEqual([{ frames: undefined, seconds: 68.04 }]); }); }); + +describe("a run whose signal already aborted", () => { + afterEach(() => { + vi.doUnmock("child_process"); + vi.resetModules(); + }); + + it("starts no ffmpeg, alone or as a pipeline", async () => { + const spawn = vi.fn(); + vi.resetModules(); + vi.doMock("child_process", () => ({ spawn })); + const { runFfmpeg, runFfmpegPipeline } = await import("./runFfmpeg.js"); + const signal = AbortSignal.abort(); + + const single = await runFfmpeg(["-version"], { signal }); + const piped = await runFfmpegPipeline(["-i", "in"], ["-i", "pipe:0"], { signal }); + + expect(spawn).not.toHaveBeenCalled(); + for (const result of [single, piped]) { + expect(result).toMatchObject({ success: false, terminationReason: "abort" }); + } + }); +}); diff --git a/packages/engine/src/utils/runFfmpeg.ts b/packages/engine/src/utils/runFfmpeg.ts index 92cceffefe..4d97042410 100644 --- a/packages/engine/src/utils/runFfmpeg.ts +++ b/packages/engine/src/utils/runFfmpeg.ts @@ -138,6 +138,19 @@ export function formatFfmpegError( : `FFmpeg exited with code ${exitCode}`; } +/** formatFfmpegError, led by the timeout when this process was killed at its deadline. */ +export function describeFfmpegFailure( + result: Pick, + timeoutMs: number, +): string { + const error = formatFfmpegError(result.exitCode, result.stderr); + if (result.terminationReason !== "deadline") return error; + return ( + `FFmpeg timed out after ${timeoutMs} ms (ffmpegProcessTimeout; long renders can raise ` + + `FFMPEG_PROCESS_TIMEOUT_MS). ${error}` + ); +} + function spawnFfmpeg(args: string[]): ChildProcess { // windowsHide: ffmpeg/ffprobe are console-subsystem binaries, so without // this Node opens a visible console window per spawn on Windows. A render @@ -168,7 +181,18 @@ function toResult(outcome: ManagedChildProcessOutcome, stderr = outcome.stderr): return result; } +/** A run whose signal had already aborted starts no process. */ +const ABORTED_BEFORE_START: RunFfmpegResult = { + success: false, + exitCode: null, + signal: null, + stderr: "", + durationMs: 0, + terminationReason: "abort", +}; + export async function runFfmpeg(args: string[], opts?: RunFfmpegOptions): Promise { + if (opts?.signal?.aborted) return { ...ABORTED_BEFORE_START }; const timeout = opts?.timeout ?? DEFAULT_TIMEOUT; const managed = new ManagedChildProcess(spawnFfmpeg(args), { signal: opts?.signal, @@ -187,6 +211,7 @@ export async function runFfmpegPipeline( consumerArgs: string[], opts?: RunFfmpegOptions, ): Promise { + if (opts?.signal?.aborted) return { ...ABORTED_BEFORE_START }; const deadlineAtMs = Date.now() + (opts?.timeout ?? DEFAULT_TIMEOUT); const producer = spawnFfmpeg(producerArgs); const consumer = spawnFfmpeg(consumerArgs); diff --git a/packages/producer/src/services/distributed/assemble.test.ts b/packages/producer/src/services/distributed/assemble.test.ts index 12eb6e54e8..72af446b8e 100644 --- a/packages/producer/src/services/distributed/assemble.test.ts +++ b/packages/producer/src/services/distributed/assemble.test.ts @@ -759,6 +759,32 @@ describe("assemble()", () => { TIMEOUT_MS, ); + it( + "stops ffmpeg at FFMPEG_PROCESS_TIMEOUT_MS and says it timed out", + async () => { + if (!hasFfmpeg) return; + const chunks: ChunkSliceJson[] = [ + { index: 0, startFrame: 0, endFrame: 5 }, + { index: 1, startFrame: 5, endFrame: 10 }, + ]; + const planDir = buildPlanDir("mp4", chunks, 10, false); + const chunkPaths = [join(planDir, "chunk-0.mp4"), join(planDir, "chunk-1.mp4")]; + for (const chunkPath of chunkPaths) makeMp4Chunk(chunkPath, 5); + + const previous = process.env.FFMPEG_PROCESS_TIMEOUT_MS; + process.env.FFMPEG_PROCESS_TIMEOUT_MS = "1"; + try { + await expect( + assemble(planDir, chunkPaths, null, join(planDir, "output.mp4")), + ).rejects.toThrow(/FFmpeg timed out after 1 ms/); + } finally { + if (previous === undefined) delete process.env.FFMPEG_PROCESS_TIMEOUT_MS; + else process.env.FFMPEG_PROCESS_TIMEOUT_MS = previous; + } + }, + TIMEOUT_MS, + ); + it("rejects when chunkPaths.length does not match chunks.json length", async () => { const chunks: ChunkSliceJson[] = [ { index: 0, startFrame: 0, endFrame: 5 }, diff --git a/packages/producer/src/services/distributed/assemble.ts b/packages/producer/src/services/distributed/assemble.ts index 2b2f36509e..81f2fb65a9 100644 --- a/packages/producer/src/services/distributed/assemble.ts +++ b/packages/producer/src/services/distributed/assemble.ts @@ -38,8 +38,10 @@ import { dirname, join } from "node:path"; import { appendRenderProvenanceArgs, applyFaststart, + describeFfmpegFailure, MIXED_AUDIO_FILENAME, muxVideoWithAudio, + resolveConfig, runFfmpeg, } from "@hyperframes/engine"; import { fpsToFfmpegArg } from "@hyperframes/core"; @@ -122,6 +124,7 @@ export async function assemble( const log = options?.logger ?? defaultLogger; const abortSignal = options?.abortSignal; const cfr = options?.cfr === true; + const { ffmpegProcessTimeout } = resolveConfig(); // ── 1. Validate planDir manifest matches chunkPaths shape ────────────── const planJsonPath = join(planDir, "plan.json"); @@ -181,10 +184,13 @@ export async function assemble( const remuxArgs = ["-i", chunkPaths[0]!, "-c", "copy", "-r", fpsArg]; appendRenderProvenanceArgs(remuxArgs, concatOutputPath); remuxArgs.push("-y", concatOutputPath); - const remuxResult = await runFfmpeg(remuxArgs, { signal: abortSignal }); + const remuxResult = await runFfmpeg(remuxArgs, { + signal: abortSignal, + timeout: ffmpegProcessTimeout, + }); if (!remuxResult.success) { throw encoderFailureError("[assemble] ffmpeg single-chunk remux failed", { - error: `exit ${remuxResult.exitCode}: ${remuxResult.stderr.slice(-400)}`, + error: describeFfmpegFailure(remuxResult, ffmpegProcessTimeout), failureReason: remuxResult.failureReason, }); } @@ -216,10 +222,13 @@ export async function assemble( ]; appendRenderProvenanceArgs(concatArgs, concatOutputPath); concatArgs.push("-y", concatOutputPath); - const concatResult = await runFfmpeg(concatArgs, { signal: abortSignal }); + const concatResult = await runFfmpeg(concatArgs, { + signal: abortSignal, + timeout: ffmpegProcessTimeout, + }); if (!concatResult.success) { throw encoderFailureError("[assemble] ffmpeg concat-copy failed", { - error: `exit ${concatResult.exitCode}: ${concatResult.stderr.slice(-400)}`, + error: describeFfmpegFailure(concatResult, ffmpegProcessTimeout), failureReason: concatResult.failureReason, }); } @@ -286,10 +295,13 @@ export async function assemble( ]; appendRenderProvenanceArgs(cfrArgs, cfrOutputPath); cfrArgs.push("-y", cfrOutputPath); - const cfrResult = await runFfmpeg(cfrArgs, { signal: abortSignal }); + const cfrResult = await runFfmpeg(cfrArgs, { + signal: abortSignal, + timeout: ffmpegProcessTimeout, + }); if (!cfrResult.success) { throw encoderFailureError("[assemble] ffmpeg cfr re-encode failed", { - error: `exit ${cfrResult.exitCode}: ${cfrResult.stderr.slice(-400)}`, + error: describeFfmpegFailure(cfrResult, ffmpegProcessTimeout), failureReason: cfrResult.failureReason, }); } @@ -310,6 +322,7 @@ export async function assemble( audioPath, outputPath: paddedAudioPath, signal: abortSignal, + timeoutMs: ffmpegProcessTimeout, }); if (!padTrimResult.success) { throw encoderFailureError("[assemble] audio pad/trim failed", padTrimResult); @@ -336,9 +349,7 @@ export async function assemble( normalizedAudioPath, muxOutputPath, abortSignal, - { - audioCodec: "aac", - }, + { audioCodec: "aac", ffmpegProcessTimeout }, { num: plan.dimensions.fpsNum, den: plan.dimensions.fpsDen }, ); if (!muxResult.success) { @@ -352,7 +363,7 @@ export async function assemble( muxOutputPath, outputPath, abortSignal, - undefined, + { ffmpegProcessTimeout }, { num: plan.dimensions.fpsNum, den: plan.dimensions.fpsDen, diff --git a/packages/producer/src/services/render/audioPadTrim.ts b/packages/producer/src/services/render/audioPadTrim.ts index 9b834b922b..fe24b19609 100644 --- a/packages/producer/src/services/render/audioPadTrim.ts +++ b/packages/producer/src/services/render/audioPadTrim.ts @@ -23,8 +23,9 @@ import { spawn } from "node:child_process"; import { mkdtempSync, renameSync, rmSync } from "node:fs"; import { dirname, join } from "node:path"; import { + DEFAULT_CONFIG, + describeFfmpegFailure, extractAudioMetadata, - formatFfmpegError, getFfprobeBinary, ManagedChildProcess, runFfmpeg, @@ -73,6 +74,8 @@ export interface PadTrimAudioInput { /** Path the helper writes the duration-corrected audio to. */ outputPath: string; signal?: AbortSignal; + /** Deadline for each ffmpeg process; defaults to the engine's `ffmpegProcessTimeout`. */ + timeoutMs?: number; /** * Optional injectables for unit tests. Production callers omit them and * get the real `ffprobe`/`ffmpeg`-backed implementations. @@ -291,12 +294,18 @@ export async function padOrTrimAudioToVideoFrameCount( input.probeVideoFrameInfo ?? ((videoPath: string) => defaultProbeVideoFrameInfo(videoPath, input.signal)); const probeAudio = input.probeAudioInfo ?? defaultProbeAudioInfo; - const runner = input.runFfmpeg ?? ((args: string[]) => defaultRunFfmpeg(args, input.signal)); + const timeoutMs = input.timeoutMs ?? DEFAULT_CONFIG.ffmpegProcessTimeout; + const runner = + input.runFfmpeg ?? ((args: string[]) => defaultRunFfmpeg(args, timeoutMs, input.signal)); // Injected runners usually do not materialize media. Tests that exercise // peak correction opt in with a matching probe; production always uses the // real post-codec measurement. const probeTruePeak = - input.probeAudioTruePeakDbfs ?? (input.runFfmpeg ? undefined : defaultProbeAudioTruePeakDbfs); + input.probeAudioTruePeakDbfs ?? + (input.runFfmpeg + ? undefined + : (audioPath: string, signal?: AbortSignal) => + defaultProbeAudioTruePeakDbfs(audioPath, timeoutMs, signal)); // Probe video and audio in parallel — the two ffprobe invocations are // independent and account for most of this function's wall-clock time. @@ -607,14 +616,15 @@ async function defaultProbeAudioInfo( async function defaultProbeAudioTruePeakDbfs( audioPath: string, + timeoutMs: number, signal?: AbortSignal, ): Promise { const result = await runFfmpeg( ["-i", audioPath, "-map", "0:a:0", "-vn", "-af", "ebur128=peak=true", "-f", "null", "-"], - { signal }, + { signal, timeout: timeoutMs }, ); if (!result.success) { - throw new Error(`[audioPadTrim] ${formatFfmpegError(result.exitCode, result.stderr)}`); + throw new Error(`[audioPadTrim] ${describeFfmpegFailure(result, timeoutMs)}`); } const matches = [...result.stderr.matchAll(/^\s*Peak:\s*([+-]?(?:[\d.]+|inf)) dBFS$/gim)]; const token = matches.at(-1)?.[1]?.toLowerCase(); @@ -627,17 +637,18 @@ async function defaultProbeAudioTruePeakDbfs( async function defaultRunFfmpeg( args: string[], + timeoutMs: number, signal?: AbortSignal, ): Promise<{ success: boolean; error?: string; failureReason?: "external_interruption"; }> { - const result = await runFfmpeg(args, { signal }); + const result = await runFfmpeg(args, { signal, timeout: timeoutMs }); if (result.success) return { success: true }; return { success: false, - error: `[audioPadTrim] ${formatFfmpegError(result.exitCode, result.stderr)}`, + error: `[audioPadTrim] ${describeFfmpegFailure(result, timeoutMs)}`, failureReason: result.failureReason, }; } diff --git a/packages/producer/src/services/render/stages/assembleStage.test.ts b/packages/producer/src/services/render/stages/assembleStage.test.ts index da1ef78c77..faa18b14f7 100644 --- a/packages/producer/src/services/render/stages/assembleStage.test.ts +++ b/packages/producer/src/services/render/stages/assembleStage.test.ts @@ -44,6 +44,7 @@ function makeInput(overrides: Partial = {}): AssembleStageIn audioOutputPath: "/tmp/audio.m4a", outputPath: "/tmp/output.mp4", hasAudio: true, + ffmpegProcessTimeout: 1_234_000, abortSignal: undefined, assertNotAborted: () => {}, ...overrides, @@ -78,13 +79,27 @@ describe("runAssembleStage audio duration parity", () => { audioPath: "/tmp/audio.m4a", outputPath: "/tmp/audio.duration-normalized.m4a", signal: undefined, + timeoutMs: 1_234_000, }); expect(muxVideoWithAudioMock).toHaveBeenCalledWith( "/tmp/video-only.mp4", "/tmp/audio.duration-normalized.m4a", "/tmp/output.mp4", undefined, - { audioCodec: "aac" }, + { audioCodec: "aac", ffmpegProcessTimeout: 1_234_000 }, + { num: 30, den: 1 }, + expect.any(Function), + ); + }); + + it("gives a silent mp4's faststart the render's configured ffmpeg timeout", async () => { + await runAssembleStage(makeInput({ hasAudio: false })); + + expect(applyFaststartMock).toHaveBeenCalledWith( + "/tmp/video-only.mp4", + "/tmp/output.mp4", + undefined, + { ffmpegProcessTimeout: 1_234_000 }, { num: 30, den: 1 }, expect.any(Function), ); @@ -98,6 +113,7 @@ describe("runAssembleStage audio duration parity", () => { audioPath: "/tmp/audio.m4a", outputPath: "/tmp/audio.duration-normalized.m4a", signal: undefined, + timeoutMs: 1_234_000, }); }); @@ -204,12 +220,13 @@ describe("runAssembleStage HLS packaging", () => { audioPath: "/tmp/audio.m4a", outputPath: "/tmp/audio.duration-normalized.m4a", signal: undefined, + timeoutMs: 1_234_000, }); expect(packageHlsMock).toHaveBeenCalledWith( "/tmp/video-only.mp4", "/tmp/audio.duration-normalized.m4a", "/tmp/output-hls", - { segmentSeconds: 6, signal: undefined }, + { segmentSeconds: 6, signal: undefined, ffmpegProcessTimeout: 1_234_000 }, ); expect(muxVideoWithAudioMock).not.toHaveBeenCalled(); expect(applyFaststartMock).not.toHaveBeenCalled(); @@ -222,6 +239,7 @@ describe("runAssembleStage HLS packaging", () => { expect(packageHlsMock).toHaveBeenCalledWith("/tmp/video-only.mp4", null, "/tmp/output-hls", { segmentSeconds: 4, signal: undefined, + ffmpegProcessTimeout: 1_234_000, }); expect(applyFaststartMock).not.toHaveBeenCalled(); }); @@ -229,14 +247,22 @@ describe("runAssembleStage HLS packaging", () => { it("defaults to 4 s segments", async () => { await runAssembleStage(hlsInput()); - expect(packageHlsMock.mock.calls[0]?.[3]).toEqual({ segmentSeconds: 4, signal: undefined }); + expect(packageHlsMock.mock.calls[0]?.[3]).toEqual({ + segmentSeconds: 4, + signal: undefined, + ffmpegProcessTimeout: 1_234_000, + }); }); it("forwards the abort signal to the packager", async () => { const abortSignal = new AbortController().signal; await runAssembleStage(hlsInput({ abortSignal })); - expect(packageHlsMock.mock.calls[0]?.[3]).toEqual({ segmentSeconds: 4, signal: abortSignal }); + expect(packageHlsMock.mock.calls[0]?.[3]).toEqual({ + segmentSeconds: 4, + signal: abortSignal, + ffmpegProcessTimeout: 1_234_000, + }); }); it("throws 'HLS packaging failed' with the ffmpeg error appended", async () => { diff --git a/packages/producer/src/services/render/stages/assembleStage.ts b/packages/producer/src/services/render/stages/assembleStage.ts index 46b560f677..2e9ed28cd1 100644 --- a/packages/producer/src/services/render/stages/assembleStage.ts +++ b/packages/producer/src/services/render/stages/assembleStage.ts @@ -46,6 +46,7 @@ export interface AssembleStageInput { format?: RenderOutputFormat; /** Segment length for `format: "hls"`. Defaults to {@link DEFAULT_HLS_SEGMENT_SECONDS}. */ hlsSegmentSeconds?: number; + ffmpegProcessTimeout: number; abortSignal: AbortSignal | undefined; assertNotAborted: () => void; onProgress?: ProgressCallback; @@ -98,6 +99,7 @@ async function assembleOutput( outputPath, hasAudio, format, + ffmpegProcessTimeout, abortSignal, assertNotAborted, } = input; @@ -116,6 +118,7 @@ async function assembleOutput( audioPath: audioOutputPath, outputPath: normalizedAudioPath, signal: abortSignal, + timeoutMs: ffmpegProcessTimeout, }); assertNotAborted(); if (!normalizeResult.success) { @@ -132,6 +135,7 @@ async function assembleOutput( abortSignal, { audioCodec: "aac", + ffmpegProcessTimeout, }, job.config.fps, onSecondsWritten, @@ -148,7 +152,7 @@ async function assembleOutput( videoOnlyPath, outputPath, abortSignal, - undefined, + { ffmpegProcessTimeout }, job.config.fps, onSecondsWritten, ); @@ -177,6 +181,7 @@ async function runHlsPackaging( { segmentSeconds: input.hlsSegmentSeconds ?? DEFAULT_HLS_SEGMENT_SECONDS, signal: input.abortSignal, + ffmpegProcessTimeout: input.ffmpegProcessTimeout, }, ); input.assertNotAborted(); diff --git a/packages/producer/src/services/renderOrchestrator.ts b/packages/producer/src/services/renderOrchestrator.ts index ac6a2cfb67..0e73cc855f 100644 --- a/packages/producer/src/services/renderOrchestrator.ts +++ b/packages/producer/src/services/renderOrchestrator.ts @@ -5146,6 +5146,7 @@ async function executeRenderPipeline(input: { hasAudio, format: outputFormat, hlsSegmentSeconds, + ffmpegProcessTimeout: cfg.ffmpegProcessTimeout, abortSignal: executionSignal, assertNotAborted, onProgress,