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
10 changes: 10 additions & 0 deletions packages/cli/src/commands/render.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2047,6 +2047,16 @@ describe("normalizeStageCode", () => {
expect(normalizeStageCode("Starting browsers (5/6 ready)")).toBe("starting_browsers");
});

it("keeps one code per stage whatever its live frame counts", () => {
expect(normalizeStageCode("Encoding frame 600/600")).toBe("encoding_video");
expect(normalizeStageCode("Encoding frame 12/90")).toBe("encoding_video");
expect(normalizeStageCode("Capturing frame 120/600 (6 workers)")).toBe("capturing_frame");
expect(normalizeStageCode("Streaming frame 3/40 (segment 1/2, 2 workers)")).toBe(
"streaming_frame",
);
expect(normalizeStageCode("Assembling final video")).toBe("assembling_final_video");
});

it("slugifies an unrecognized stage string instead of bucketing it as unknown", () => {
expect(normalizeStageCode("Some New Stage!")).toBe("some_new_stage");
});
Expand Down
20 changes: 13 additions & 7 deletions packages/cli/src/commands/render.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ import {
errorBox,
} from "../ui/format.js";
import { warnIfWebmAlphaDropped } from "../utils/webmAlphaCheck.js";
import { renderProgress } from "../ui/progress.js";
import { renderProgress, renderMachineProgress } from "../ui/progress.js";
import {
trackRenderComplete,
trackRenderError,
Expand Down Expand Up @@ -801,6 +801,7 @@ async function renderDocker(
outputDir: resolve(outputDir),
outputFilename,
platform,
hostStdoutIsTty: process.stdout.isTTY === true,
options: {
fps: options.fps,
quality: options.quality,
Expand Down Expand Up @@ -1095,8 +1096,9 @@ async function executeLocalRender(

const onProgress = options.quiet
? undefined
: (progressJob: { progress: number }, message: string) => {
: (progressJob: Pick<RenderJob, "progress" | "stageProgress">, message: string) => {
renderProgress(progressJob.progress, message);
renderMachineProgress(progressJob.progress, progressJob.stageProgress);
};

try {
Expand Down Expand Up @@ -1654,15 +1656,19 @@ const KNOWN_STAGE_CODES: Readonly<Record<string, string>> = {
"Render complete": "render_complete",
"Render cancelled": "render_cancelled",
pipeline: "pipeline",
"Starting browsers": "starting_browsers",
// Was "Encoding video" before encode reported frames; keeps the same bucket.
"Encoding frame": "encoding_video",
};

// Live counts in a progress label ("Capturing frame 120/600 (6 workers)") would make a code per render.
const STAGE_COUNTS = /\([^)]*\)|\d+\/\d+/g;

export function normalizeStageCode(stage: string): string {
const known = KNOWN_STAGE_CODES[stage];
const base = stage.replace(STAGE_COUNTS, "").trim();
const known = KNOWN_STAGE_CODES[base];
if (known) return known;
// The producer's "Starting browsers (k/n ready)" carries live counts; keep one code for it.
if (stage.startsWith("Starting browsers")) return "starting_browsers";
const slug = stage
.trim()
const slug = base
.toLowerCase()
.replace(/[^a-z0-9]+/g, "_")
.replace(/^_+|_+$/g, "");
Expand Down
34 changes: 33 additions & 1 deletion packages/cli/src/ui/progress.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { afterEach, describe, expect, it } from "vitest";

import { renderProgress } from "./progress.js";
import { renderMachineProgress, renderProgress } from "./progress.js";

const originalWrite = process.stdout.write.bind(process.stdout);
const originalIsTTY = process.stdout.isTTY;
Expand Down Expand Up @@ -31,3 +31,35 @@ describe("renderProgress", () => {
expect(output).not.toContain("\r");
});
});

describe("renderMachineProgress", () => {
const capture = (isTTY: boolean) => {
let output = "";
Object.defineProperty(process.stdout, "isTTY", { value: isTTY, configurable: true });
process.stdout.write = ((chunk: string | Uint8Array) => {
output += String(chunk);
return true;
}) as typeof process.stdout.write;
renderMachineProgress(83.4, { code: "encode", done: 812, total: 1800 });
return output;
};

it("writes one @hf-progress line for a parent reading piped stdout", () => {
expect(capture(false)).toBe(
'@hf-progress {"code":"encode","done":812,"total":1800,"pct":83}\n',
);
});

it("stays silent in a terminal", () => {
expect(capture(true)).toBe("");
});

it("stays silent in a container whose host prints to a terminal", () => {
process.env.HYPERFRAMES_STDOUT_IS_TTY = "1";
try {
expect(capture(false)).toBe("");
} finally {
delete process.env.HYPERFRAMES_STDOUT_IS_TTY;
}
});
});
6 changes: 6 additions & 0 deletions packages/cli/src/ui/progress.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { c } from "./colors.js";
import type { RenderJob } from "@hyperframes/producer";

const { stdout } = process;

Expand All @@ -18,3 +19,8 @@ export function renderProgress(percent: number, stage: string, row?: number): vo
stdout.write(`\r\x1b[2K${line}`);
}
}

export function renderMachineProgress(percent: number, stage: RenderJob["stageProgress"]): void {
if (stdout.isTTY || process.env.HYPERFRAMES_STDOUT_IS_TTY === "1" || !stage) return;
stdout.write(`@hf-progress ${JSON.stringify({ ...stage, pct: Math.round(percent) })}\n`);
}
8 changes: 8 additions & 0 deletions packages/cli/src/utils/dockerRunArgs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,14 @@ describe("buildDockerRunArgs", () => {
`);
});

it("tells the container its output reaches a terminal only when the host's stdout is one", () => {
const flags = (hostStdoutIsTty?: boolean) =>
buildDockerRunArgs({ ...FIXED_INPUT, hostStdoutIsTty, options: BASE }).slice(0, 4);
expect(flags(true)).toEqual(["run", "--rm", "-e", "HYPERFRAMES_STDOUT_IS_TTY=1"]);
expect(flags(false)).toEqual(["run", "--rm", "--platform", "linux/amd64"]);
expect(flags(undefined)).toEqual(["run", "--rm", "--platform", "linux/amd64"]);
});

it("omits --workers when auto sizing should happen inside the container", () => {
const args = buildDockerRunArgs({ ...FIXED_INPUT, options: BASE });
expect(args).not.toContain("--workers");
Expand Down
2 changes: 2 additions & 0 deletions packages/cli/src/utils/dockerRunArgs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ export interface DockerRunArgsInput {
outputDir: string;
/** Filename within `outputDir` (joined to /output inside the container). */
outputFilename: string;
hostStdoutIsTty?: boolean;
/**
* Docker `--platform` value (`linux/amd64` or `linux/arm64`). When omitted,
* resolves to the host architecture via `resolveDockerPlatform()`. Pinning
Expand Down Expand Up @@ -109,6 +110,7 @@ export function buildDockerRunArgs(input: DockerRunArgsInput): string[] {
return [
"run",
"--rm",
...(input.hostStdoutIsTty ? ["-e", "HYPERFRAMES_STDOUT_IS_TTY=1"] : []),
"--platform",
platform,
"--shm-size=2g",
Expand Down
42 changes: 37 additions & 5 deletions packages/engine/src/services/chunkEncoder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,12 @@ import {
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 { formatFfmpegError, isExternalFfmpegInterruption, runFfmpeg } from "../utils/runFfmpeg.js";
import {
ffmpegStatsReader,
formatFfmpegError,
isExternalFfmpegInterruption,
runFfmpeg,
} from "../utils/runFfmpeg.js";
import { extractAudioMetadata } from "../utils/ffprobe.js";
import { type Fps, fpsToFfmpegArg, fpsToNumber } from "@hyperframes/core";
import type { EncoderOptions, EncodeResult, MuxResult } from "./chunkEncoder.types.js";
Expand Down Expand Up @@ -491,13 +496,19 @@ export function buildEncoderArgs(
return args;
}

const framesReader = (onFrames?: (frames: number) => void) =>
onFrames && ffmpegStatsReader(({ frames }) => frames !== undefined && onFrames(frames));
const secondsReader = (onSeconds?: (seconds: number) => void) =>
onSeconds && ffmpegStatsReader(({ seconds }) => seconds !== undefined && onSeconds(seconds));

export async function encodeFramesFromDir(
framesDir: string,
framePattern: string,
outputPath: string,
options: EncoderOptions,
signal?: AbortSignal,
config?: Partial<Pick<EngineConfig, "ffmpegEncodeTimeout">>,
onFramesEncoded?: (frames: number) => void,
): Promise<EncodeResult> {
const startTime = Date.now();

Expand Down Expand Up @@ -527,7 +538,11 @@ export async function encodeFramesFromDir(
const inputArgs = ["-framerate", fpsToFfmpegArg(options.fps), "-i", inputPath];
const args = buildEncoderArgs(options, inputArgs, outputPath, gpuEncoder);
const encodeTimeout = config?.ffmpegEncodeTimeout ?? DEFAULT_CONFIG.ffmpegEncodeTimeout;
const result = await runFfmpeg(args, { signal, timeout: encodeTimeout });
const result = await runFfmpeg(args, {
signal,
timeout: encodeTimeout,
onStderr: framesReader(onFramesEncoded),
});
if (result.terminationReason === "abort") {
return {
success: false,
Expand Down Expand Up @@ -636,6 +651,7 @@ export async function encodeFramesChunkedConcat(
chunkSizeFrames: number,
signal?: AbortSignal,
config?: Partial<Pick<EngineConfig, "ffmpegEncodeTimeout">>,
onFramesEncoded?: (frames: number) => void,
): Promise<EncodeResult> {
const start = Date.now();
const files = readdirSync(framesDir)
Expand Down Expand Up @@ -693,7 +709,13 @@ export async function encodeFramesChunkedConcat(
if (options.useGpu) gpuEncoder = await getCachedGpuEncoder();
const args = buildEncoderArgs(options, inputArgs, chunkPath, gpuEncoder);
const encodeTimeout = config?.ffmpegEncodeTimeout ?? DEFAULT_CONFIG.ffmpegEncodeTimeout;
const processResult = await runFfmpeg(args, { signal, timeout: encodeTimeout });
const processResult = await runFfmpeg(args, {
signal,
timeout: encodeTimeout,
onStderr: framesReader(
onFramesEncoded && ((frames) => onFramesEncoded(startNumber + frames)),
),
});
const chunkResult = {
success: processResult.success,
error: processResult.success
Expand Down Expand Up @@ -750,6 +772,7 @@ export async function muxVideoWithAudio(
signal?: AbortSignal,
config?: MuxVideoWithAudioOptions,
fps?: Fps,
onSecondsWritten?: (seconds: number) => void,
): Promise<MuxResult> {
const outputDir = dirname(outputPath);
if (!existsSync(outputDir)) mkdirSync(outputDir, { recursive: true });
Expand Down Expand Up @@ -804,7 +827,11 @@ export async function muxVideoWithAudio(
args.push("-y", outputPath);

const processTimeout = config?.ffmpegProcessTimeout ?? DEFAULT_CONFIG.ffmpegProcessTimeout;
const result = await runFfmpeg(args, { signal, timeout: processTimeout });
const result = await runFfmpeg(args, {
signal,
timeout: processTimeout,
onStderr: secondsReader(onSecondsWritten),
});

if (signal?.aborted) {
return {
Expand Down Expand Up @@ -977,6 +1004,7 @@ export async function applyFaststart(
signal?: AbortSignal,
config?: Partial<Pick<EngineConfig, "ffmpegProcessTimeout">>,
fps?: Fps,
onSecondsWritten?: (seconds: number) => void,
): Promise<MuxResult> {
// faststart is MP4-only (moves moov atom to file start for streaming).
// WebM and MOV don't need it — skip the re-mux.
Expand All @@ -995,7 +1023,11 @@ export async function applyFaststart(
args.push("-y", outputPath);

const processTimeout = config?.ffmpegProcessTimeout ?? DEFAULT_CONFIG.ffmpegProcessTimeout;
const result = await runFfmpeg(args, { signal, timeout: processTimeout });
const result = await runFfmpeg(args, {
signal,
timeout: processTimeout,
onStderr: secondsReader(onSecondsWritten),
});

if (signal?.aborted) {
return {
Expand Down
21 changes: 21 additions & 0 deletions packages/engine/src/services/streamingEncoder.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -604,6 +604,27 @@ describe("spawnStreamingEncoder lifecycle and cleanup", () => {
expect(result.fileSize).toBe(0); // No real ffmpeg, no file written
});

it("reports ffmpeg's encoded-frame count while close() waits, and stops at exit", async () => {
const { spawn, calls } = createSpawnSpy();
vi.resetModules();
vi.doMock("child_process", () => ({ spawn }));

const { spawnStreamingEncoder } = await import("./streamingEncoder.js");
const dir = mkdtempSync(join(tmpdir(), "se-encoded-"));
const encoder = await spawnStreamingEncoder(join(dir, "out.mp4"), baseOptions);
const proc = calls[0]!.proc;
proc.stderr.emit("data", Buffer.from("frame= 10 fps=20 q=28.0 size=1kB time=00:00:00.33\r"));

const heard: number[] = [];
const closePromise = encoder.close((frames) => heard.push(frames));
proc.stderr.emit("data", Buffer.from("frame= 25 fps=20 q=28.0 size=2kB time=00:00:00.83\r"));
process.nextTick(() => proc.emit("close", 0));
await closePromise;
proc.stderr.emit("data", Buffer.from("frame= 30 fps=20 q=28.0 size=2kB time=00:00:01.00\r"));

expect(heard).toEqual([10, 25]);
});

it("returns a failure result (does NOT throw) when ffmpeg exits non-zero before close()", async () => {
const { spawn, calls } = createSpawnSpy();
vi.resetModules();
Expand Down
22 changes: 19 additions & 3 deletions packages/engine/src/services/streamingEncoder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,11 @@ import {
buildVideoToolboxRateControlArgs,
mapPresetForGpuEncoder,
} from "../utils/gpuEncoder.js";
import { formatFfmpegError, isExternalFfmpegInterruption } from "../utils/runFfmpeg.js";
import {
ffmpegStatsReader,
formatFfmpegError,
isExternalFfmpegInterruption,
} from "../utils/runFfmpeg.js";
import { getFfmpegBinary } from "../utils/ffmpegBinaries.js";
import { getHdrEncoderColorParams } from "../utils/hdr.js";
import { withEvenDimensionPad } from "../utils/evenDimensions.js";
Expand Down Expand Up @@ -187,7 +191,8 @@ export interface StreamingEncoder {
* calls would interleave frame bytes on the pipe and race the drain wait.
*/
writeFrame: (buffer: Buffer) => Promise<boolean>;
close: () => Promise<StreamingEncoderResult>;
/** `onFramesEncoded` hears ffmpeg's encoded-frame count while close waits for it to finish. */
close: (onFramesEncoded?: (frames: number) => void) => Promise<StreamingEncoderResult>;
getExitStatus: () => "running" | "success" | "error";
/**
* The FFmpeg failure reason (exit code + tail of stderr), or `undefined`
Expand Down Expand Up @@ -497,9 +502,16 @@ export async function spawnStreamingEncoder(
// libx264 printed its summary and exited 255, observable as
// "Streaming encode failed: FFmpeg exited with code 255" with audio:0kB).
const streamingTimeout = config?.ffmpegStreamingTimeout ?? DEFAULT_CONFIG.ffmpegStreamingTimeout;
let framesEncoded = 0;
let onFramesEncoded: ((frames: number) => void) | undefined;
const managed = new ManagedChildProcess(ffmpeg, {
signal,
inactivityTimeoutMs: streamingTimeout,
onStderr: ffmpegStatsReader(({ frames }) => {
if (frames === undefined) return;
framesEncoded = frames;
onFramesEncoded?.(frames);
}),
});
const exitPromise = managed.wait().then((outcome) => {
exitCode = outcome.exitCode;
Expand Down Expand Up @@ -604,7 +616,7 @@ export async function spawnStreamingEncoder(
return true;
},

close: async (): Promise<StreamingEncoderResult> => {
close: async (onEncoded?: (frames: number) => void): Promise<StreamingEncoderResult> => {
// INVARIANT: close() is idempotent. The renderOrchestrator HDR cleanup
// path tracks an `encoderClosed` flag and may still re-call close() in
// the outer finally if the inner cleanup raised before the flag flipped.
Expand All @@ -616,6 +628,10 @@ export async function spawnStreamingEncoder(
// repeated calls. If you change this method, preserve idempotency or
// a regression here will silently double-close ffmpeg and produce
// harder-to-trace errors at the orchestrator layer.
if (onEncoded && exitStatus === "running") {
onFramesEncoded = onEncoded;
onEncoded(framesEncoded);
}
const stdin = ffmpeg.stdin;
if (stdin && !stdin.destroyed) {
await new Promise<void>((resolve) => {
Expand Down
Loading
Loading