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

it("keeps one code for the producer's browser start-up counts", () => {
expect(normalizeStageCode("Starting browsers (0/6 ready)")).toBe("starting_browsers");
expect(normalizeStageCode("Starting browsers (5/6 ready)")).toBe("starting_browsers");
});

it("slugifies an unrecognized stage string instead of bucketing it as unknown", () => {
expect(normalizeStageCode("Some New Stage!")).toBe("some_new_stage");
});
Expand Down
2 changes: 2 additions & 0 deletions packages/cli/src/commands/render.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1649,6 +1649,8 @@ const KNOWN_STAGE_CODES: Readonly<Record<string, string>> = {
export function normalizeStageCode(stage: string): string {
const known = KNOWN_STAGE_CODES[stage];
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()
.toLowerCase()
Expand Down
42 changes: 40 additions & 2 deletions packages/producer/src/services/render/shared.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { describe, expect, it } from "bun:test";
import { resolveBrowserMediaEnd } from "./shared.js";
import { describe, expect, it, setSystemTime } from "bun:test";
import type { ParallelProgress } from "@hyperframes/engine";
import type { RenderJob } from "../renderOrchestrator.js";
import { reportWorkerStartup, resolveBrowserMediaEnd } from "./shared.js";

describe("resolveBrowserMediaEnd", () => {
it("prefers a runtime duration over a stale compiler-clamped end", () => {
Expand All @@ -15,3 +17,39 @@ describe("resolveBrowserMediaEnd", () => {
expect(resolveBrowserMediaEnd(0, 5.04, 0)).toBe(5.04);
});
});

describe("reportWorkerStartup", () => {
it("counts ready workers and drops ids past a smaller retry's worker count", () => {
const job = { progress: 25 } as RenderJob;
const stages: string[] = [];
const phase = (workerId: number, name: string, activeWorkers: number) => {
setSystemTime(Date.now() + 1_000);
reportWorkerStartup(
job,
{
activeWorkers,
latestWorkerPhase: { workerId, phase: name },
} as unknown as ParallelProgress,
(_job, stage) => {
stages.push(stage);
},
);
};
try {
phase(0, "browser_launch", 3);
phase(2, "frame_capture", 3);
phase(0, "frame_capture", 3);
phase(0, "browser_launch", 2);
phase(1, "frame_capture", 2);
} finally {
setSystemTime();
}
expect(stages).toEqual([
"Starting browsers (0/3 ready)",
"Starting browsers (1/3 ready)",
"Starting browsers (2/3 ready)",
"Starting browsers (0/2 ready)",
"Starting browsers (1/2 ready)",
]);
});
});
60 changes: 60 additions & 0 deletions packages/producer/src/services/render/shared.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import type {
AudioElement,
ExtractedFrames,
ImageElement,
ParallelProgress,
VideoElement,
} from "@hyperframes/engine";
import type { CompiledComposition } from "../htmlCompiler.js";
Expand Down Expand Up @@ -265,6 +266,65 @@ export function updateJobStatus(
if (onProgress) void onProgress(job, stage);
}

const PROGRESS_REPORT_INTERVAL_MS = 250;
const lastFrameReportAt = new WeakMap<RenderJob, number>();
const lastStartupReportAt = new WeakMap<RenderJob, number>();

// The job updates on every call; the callback fires on the first call per job, when forced,
// and at most once per interval in between.
function reportThrottled(
lastReportAt: WeakMap<RenderJob, number>,
job: RenderJob,
stage: string,
progress: number,
onProgress: ProgressCallback | undefined,
force: boolean,
): void {
const now = Date.now();
const last = lastReportAt.get(job);
const due = force || last === undefined || now - last >= PROGRESS_REPORT_INTERVAL_MS;
if (due) lastReportAt.set(job, now);
updateJobStatus(job, "rendering", stage, progress, due ? onProgress : undefined);
}

/** Capture-loop progress: the first frame, the last frame, and at most four reports a second between. */
export function reportFrameProgress(
job: RenderJob,
stage: string,
progress: number,
onProgress: ProgressCallback | undefined,
isLastFrame: boolean,
): void {
reportThrottled(lastFrameReportAt, job, stage, progress, onProgress, isLastFrame);
}

const workerPhasesByJob = new WeakMap<RenderJob, Map<number, string>>();

/** Browser warm-up before the first frame, from the parallel workers' phase events. */
export function reportWorkerStartup(
job: RenderJob,
progress: ParallelProgress,
onProgress: ProgressCallback | undefined,
): void {
const phase = progress.latestWorkerPhase;
if (!phase) return;
const phases = workerPhasesByJob.get(job) ?? new Map<number, string>();
workerPhasesByJob.set(job, phases);
phases.set(phase.workerId, phase.phase);
// ponytail: ids past a smaller retry's worker count are stale, so they are not counted
const ready = [...phases].filter(
([id, p]) => id < progress.activeWorkers && (p === "frame_capture" || p === "frame_encode"),
).length;
reportThrottled(
lastStartupReportAt,
job,
`Starting browsers (${ready}/${progress.activeWorkers} ready)`,
job.progress,
onProgress,
false,
);
}

/**
* Build a `resolver(framePath)` closure that maps an absolute path to
* a frame inside `compiledDir` into a server-relative URL the producer's
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ import {
partitionTransitionFrames,
seekInjectAndQueryStacking,
} from "./captureHdrFrameShared.js";
import { updateJobStatus } from "../shared.js";
import { reportFrameProgress } from "../shared.js";

export interface HybridLoopInput {
job: RenderJob;
Expand Down Expand Up @@ -189,16 +189,13 @@ export async function runHybridLayeredFrameLoop(input: HybridLoopInput): Promise
reorderBuffer.advanceTo(frameIdx + 1);
framesWritten += 1;
job.framesRendered = framesWritten;
if (framesWritten % 10 === 0 || framesWritten === totalFrames) {
const frameProgress = framesWritten / totalFrames;
updateJobStatus(
job,
"rendering",
`Layered composite frame ${framesWritten}/${job.totalFrames}`,
Math.round(25 + frameProgress * 55),
onProgress,
);
}
reportFrameProgress(
job,
`Layered composite frame ${framesWritten}/${job.totalFrames}`,
Math.round(25 + (framesWritten / totalFrames) * 55),
onProgress,
framesWritten === totalFrames,
);
};
const poolRef = shaderPool;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ import {
type LayeredTransitionBuffers,
seekInjectAndQueryStacking,
} from "./captureHdrFrameShared.js";
import { updateJobStatus } from "../shared.js";
import { reportFrameProgress } from "../shared.js";

export interface SequentialLoopInput {
job: RenderJob;
Expand Down Expand Up @@ -219,15 +219,12 @@ export async function runSequentialLayeredFrameLoop(input: SequentialLoopInput):
log,
});
job.framesRendered = i + 1;
if ((i + 1) % 10 === 0 || i + 1 === totalFrames) {
const frameProgress = (i + 1) / totalFrames;
updateJobStatus(
job,
"rendering",
`Layered composite frame ${i + 1}/${job.totalFrames}`,
Math.round(25 + frameProgress * 55),
onProgress,
);
}
reportFrameProgress(
job,
`Layered composite frame ${i + 1}/${job.totalFrames}`,
Math.round(25 + ((i + 1) / totalFrames) * 55),
onProgress,
i + 1 === totalFrames,
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ import type { ProducerLogger } from "../../../logger.js";
import type { ProgressCallback, RenderJob } from "../../renderOrchestrator.js";
import { wrapCaptureStageError } from "../captureStageError.js";
import { ensureFrameWritten } from "./captureHdrFrameShared.js";
import { updateJobStatus } from "../shared.js";
import { reportFrameProgress } from "../shared.js";
import { encoderFailureError } from "../encoderInterruption.js";
import type { SdrSegmentedCapturePlan } from "../capturePlan.js";
import { planSegments, type SegmentSlice } from "../segmentPlan.js";
Expand Down Expand Up @@ -212,14 +212,14 @@ async function captureSegmentFrames(
ctx.job.framesRendered = i + 1;
lastProgressAt = Date.now();

updateJobStatus(
reportFrameProgress(
ctx.job,
"rendering",
`Streaming frame ${i + 1}/${ctx.totalFrames} (segment ${segment.index + 1}/${ctx.segmentCount}` +
(ctx.skipped > 0 ? `, skipped ${ctx.skipped}` : "") +
")",
Math.round(25 + ((i + 1) / ctx.totalFrames) * 55),
ctx.onProgress,
i + 1 === ctx.totalFrames,
);
}
}
Expand Down
37 changes: 16 additions & 21 deletions packages/producer/src/services/render/stages/captureStage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,8 @@
* - `lastBrowserConsole` is set to the buffer of whichever session was
* active last (probe session in the parallel close path; sequential
* session in the sequential path).
* - `job.framesRendered` is updated at the same per-frame / per-progress
* points; the same `Capturing frame N/M` `updateJobStatus` payloads
* fire at 30-frame and completion checkpoints (parallel) or every
* frame (sequential).
* - `job.framesRendered` is updated at every per-frame / per-progress
* point; `Capturing frame N/M` goes through `reportFrameProgress`.
*
* Known follow-up: this stage imports `executeDiskCaptureWithAdaptiveRetry`
* from `renderOrchestrator.ts`, which itself imports the stage — a runtime
Expand Down Expand Up @@ -70,7 +68,7 @@ import {
type RenderJob,
} from "../../renderOrchestrator.js";
import { wrapCaptureStageError } from "../captureStageError.js";
import { updateJobStatus } from "../shared.js";
import { reportFrameProgress, reportWorkerStartup } from "../shared.js";
import type { SdrDiskCapturePlan } from "../capturePlan.js";

export interface CaptureStageInput {
Expand Down Expand Up @@ -366,23 +364,20 @@ export async function runCaptureStage(input: CaptureStageInput): Promise<Capture
frameRangeStart: frameRange?.startFrame,
dedupPerfs,
onProgress: (progress) => {
if (progress.latestWorkerPhase) {
if (progress.capturedFrames === 0) reportWorkerStartup(job, progress, onProgress);
return;
}
job.framesRendered = progress.capturedFrames;
checkDiskProjection(progress.capturedFrames);
const frameProgress = progress.capturedFrames / progress.totalFrames;
const progressPct = 25 + frameProgress * 45;

if (
progress.capturedFrames % 30 === 0 ||
progress.capturedFrames === progress.totalFrames
) {
updateJobStatus(
job,
"rendering",
`Capturing frame ${progress.capturedFrames}/${progress.totalFrames} (${progress.activeWorkers} workers)`,
Math.round(progressPct),
onProgress,
);
}
reportFrameProgress(
job,
`Capturing frame ${progress.capturedFrames}/${progress.totalFrames} (${progress.activeWorkers} workers)`,
Math.round(25 + frameProgress * 45),
onProgress,
progress.capturedFrames === progress.totalFrames,
);
},
cfg: captureCfg,
log,
Expand Down Expand Up @@ -577,12 +572,12 @@ async function captureSessionFrames(
// capture error wrapper below must remain separate from finally so it
// can throw with the browser console before cleanup overwrites flow.
// fallow-ignore-next-line code-duplication
updateJobStatus(
reportFrameProgress(
job,
"rendering",
`Capturing frame ${fileIndex + 1}/${rangeFrames}`,
Math.round(25 + ((fileIndex + 1) / rangeFrames) * 45),
onProgress,
fileIndex + 1 === rangeFrames,
);
};

Expand Down
Loading
Loading