Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions packages/engine/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -343,6 +343,7 @@ export {
export {
runFfmpeg,
formatFfmpegError,
describeFfmpegFailure,
isExternalFfmpegInterruption,
type RunFfmpegOptions,
type RunFfmpegResult,
Expand Down
74 changes: 74 additions & 0 deletions packages/engine/src/services/chunkEncoder.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
7 changes: 4 additions & 3 deletions packages/engine/src/services/chunkEncoder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
};
}
Expand Down Expand Up @@ -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,
};
}
Expand Down Expand Up @@ -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,
};
}
140 changes: 140 additions & 0 deletions packages/engine/src/services/videoFrameExtractor.concurrency.test.ts
Original file line number Diff line number Diff line change
@@ -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<typeof import("os")>()),
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<string> {
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);
});
84 changes: 50 additions & 34 deletions packages/engine/src/services/videoFrameExtractor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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[];
Expand Down Expand Up @@ -923,30 +928,34 @@ async function runSegmentedExtraction(job: SegmentedExtraction): Promise<RunFfmp
while (next < segmentCount && !signal.aborted) {
const index = next++;
const { firstFrame, frames } = segments[index]!;
const result = await 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 },
const result = await inExtractionSlot(
() =>
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.
Expand Down Expand Up @@ -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(",")] : [];
Expand All @@ -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(
Expand Down
Loading
Loading