From b4908bc45c5fb247082afa22c5469539c629921e Mon Sep 17 00:00:00 2001 From: Miguel Angel Simon Sierra Date: Thu, 24 Sep 2026 06:26:14 -0400 Subject: [PATCH 1/2] fix(producer): report render progress from the first frame, at most four times a second Capture loops reported every 30th (or 10th) frame, or every single frame, so a render sat silent through browser warm-up and the first 30 frames, then flooded on the sequential paths. One helper now owns the cadence for every capture stage: first frame, last frame, and every 250 ms in between. Browser GPU check, capture calibration and worker start-up each report their own step. --- .../src/services/render/shared.test.ts | 42 +++++++++++- .../producer/src/services/render/shared.ts | 48 +++++++++++++ .../render/stages/captureHdrHybridLoop.ts | 19 +++--- .../render/stages/captureHdrSequentialLoop.ts | 19 +++--- .../render/stages/captureSegmentedStage.ts | 6 +- .../services/render/stages/captureStage.ts | 37 +++++----- .../stages/captureStreamingStage.test.ts | 68 ++++++++++++++++++- .../render/stages/captureStreamingStage.ts | 39 +++++------ .../src/services/renderOrchestrator.ts | 2 + 9 files changed, 208 insertions(+), 72 deletions(-) diff --git a/packages/producer/src/services/render/shared.test.ts b/packages/producer/src/services/render/shared.test.ts index 04e47bc90d..95df6f98e7 100644 --- a/packages/producer/src/services/render/shared.test.ts +++ b/packages/producer/src/services/render/shared.test.ts @@ -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", () => { @@ -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)", + ]); + }); +}); diff --git a/packages/producer/src/services/render/shared.ts b/packages/producer/src/services/render/shared.ts index 15eaf8226f..3459fa7720 100644 --- a/packages/producer/src/services/render/shared.ts +++ b/packages/producer/src/services/render/shared.ts @@ -29,6 +29,7 @@ import type { AudioElement, ExtractedFrames, ImageElement, + ParallelProgress, VideoElement, } from "@hyperframes/engine"; import type { CompiledComposition } from "../htmlCompiler.js"; @@ -265,6 +266,53 @@ export function updateJobStatus( if (onProgress) void onProgress(job, stage); } +const FRAME_PROGRESS_INTERVAL_MS = 250; +const lastFrameReportAt = new WeakMap(); + +/** + * Capture-loop progress: the job updates on every call, the callback fires on + * the first call, the last frame, and at most once per 250 ms in between. + */ +export function reportFrameProgress( + job: RenderJob, + stage: string, + progress: number, + onProgress: ProgressCallback | undefined, + isLastFrame: boolean, +): void { + const now = Date.now(); + const last = lastFrameReportAt.get(job); + const due = isLastFrame || last === undefined || now - last >= FRAME_PROGRESS_INTERVAL_MS; + if (due) lastFrameReportAt.set(job, now); + updateJobStatus(job, "rendering", stage, progress, due ? onProgress : undefined); +} + +const workerPhasesByJob = new WeakMap>(); + +/** 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(); + 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; + reportFrameProgress( + 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 diff --git a/packages/producer/src/services/render/stages/captureHdrHybridLoop.ts b/packages/producer/src/services/render/stages/captureHdrHybridLoop.ts index 0c232da66a..099b3bf1be 100644 --- a/packages/producer/src/services/render/stages/captureHdrHybridLoop.ts +++ b/packages/producer/src/services/render/stages/captureHdrHybridLoop.ts @@ -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; @@ -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; diff --git a/packages/producer/src/services/render/stages/captureHdrSequentialLoop.ts b/packages/producer/src/services/render/stages/captureHdrSequentialLoop.ts index 19644eb8ae..37539f922e 100644 --- a/packages/producer/src/services/render/stages/captureHdrSequentialLoop.ts +++ b/packages/producer/src/services/render/stages/captureHdrSequentialLoop.ts @@ -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; @@ -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, + ); } } diff --git a/packages/producer/src/services/render/stages/captureSegmentedStage.ts b/packages/producer/src/services/render/stages/captureSegmentedStage.ts index 9b9700e10e..301de38df2 100644 --- a/packages/producer/src/services/render/stages/captureSegmentedStage.ts +++ b/packages/producer/src/services/render/stages/captureSegmentedStage.ts @@ -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"; @@ -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, ); } } diff --git a/packages/producer/src/services/render/stages/captureStage.ts b/packages/producer/src/services/render/stages/captureStage.ts index 1d988b69f4..ae436f8b35 100644 --- a/packages/producer/src/services/render/stages/captureStage.ts +++ b/packages/producer/src/services/render/stages/captureStage.ts @@ -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 @@ -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 { @@ -366,23 +364,20 @@ export async function runCaptureStage(input: CaptureStageInput): Promise { + 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, @@ -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, ); }; diff --git a/packages/producer/src/services/render/stages/captureStreamingStage.test.ts b/packages/producer/src/services/render/stages/captureStreamingStage.test.ts index edb1399f81..588ee04691 100644 --- a/packages/producer/src/services/render/stages/captureStreamingStage.test.ts +++ b/packages/producer/src/services/render/stages/captureStreamingStage.test.ts @@ -1,5 +1,5 @@ // fallow-ignore-file code-duplication -import { afterAll, describe, expect, it, mock } from "bun:test"; +import { afterAll, describe, expect, it, mock, setSystemTime } from "bun:test"; import { mkdtempSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -31,6 +31,9 @@ let hangParallelUntilAbort = false; // a capture call (the watchdog's job); session_init = still booting (not). let hangParallelPhase: "session_init" | "frame_capture" = "frame_capture"; let failWorkerTransient = false; +// Parallel capture that reports this many frames, advancing the clock by parallelFrameMs each. +let parallelFrames = 0; +let parallelFrameMs = 0; let injectedWorkerFailure: Error | null = null; const reorderAbortCalls: unknown[] = []; let hangSequentialUntilStall = false; @@ -139,6 +142,32 @@ mock.module("@hyperframes/engine", () => ({ signal?.addEventListener("abort", fail, { once: true }); }); } + if (parallelFrames > 0) { + const report = (capturedFrames: number, phase?: string) => + onProgress?.({ + totalFrames: parallelFrames, + capturedFrames, + activeWorkers: 2, + workerProgress: new Map(), + latestWorkerPhase: phase && { + workerId: 0, + phase, + browserExecutable: "chrome", + browserVersion: "Chrome/152.0.7977.30", + canvasDrawElement: true, + gpuBackend: "swiftshader", + }, + }); + report(0, "browser_launch"); + report(0, "frame_capture"); + let now = Date.now(); + for (let frame = 1; frame <= parallelFrames; frame++) { + now += parallelFrameMs; + setSystemTime(now); + report(frame); + } + setSystemTime(); + } return []; }, getCapturePerfSummary: () => ({}), @@ -380,6 +409,43 @@ describe("runCaptureStreamingStage", () => { expect((caught as Error).message).not.toContain("stalled"); }); + async function streamParallelFrames(frames: number, frameMs: number): Promise { + parallelFrames = frames; + parallelFrameMs = frameMs; + const stages: string[] = []; + const { runCaptureStreamingStage } = await import("./captureStreamingStage.js"); + const baseInput = createInput({ forceScreenshot: false, ffmpegStreamingTimeout: 3_600_000 }); + try { + await runCaptureStreamingStage({ + ...baseInput, + totalFrames: frames, + plan: { ...baseInput.plan, workerCount: 2, forceParallelStream: true }, + onProgress: (_job: unknown, stage: string) => { + stages.push(stage); + }, + }); + } finally { + parallelFrames = 0; + } + return stages; + } + + it("reports browser warm-up and the first frame of a slow capture, then steadily", async () => { + const stages = await streamParallelFrames(40, 400); + expect(stages.slice(0, 3)).toEqual([ + "Starting browsers (0/2 ready)", + "Streaming frame 1/40 (2 workers)", + "Streaming frame 2/40 (2 workers)", + ]); + expect(stages).toHaveLength(41); + }); + + it("does not flood the callback when frames arrive faster than the report interval", async () => { + const stages = await streamParallelFrames(300, 1); + expect(stages.filter((stage) => stage.startsWith("Streaming frame")).length).toBeLessThan(5); + expect(stages.at(-1)).toBe("Streaming frame 300/300 (2 workers)"); + }); + it("releases the writer with the dead worker's own error, not a stall", async () => { // A transient Chrome death in one interleaved worker. Before, peers parked // in the ordered writer until the watchdog relabelled it a stall a minute diff --git a/packages/producer/src/services/render/stages/captureStreamingStage.ts b/packages/producer/src/services/render/stages/captureStreamingStage.ts index 3199074b0b..fd07fcc450 100644 --- a/packages/producer/src/services/render/stages/captureStreamingStage.ts +++ b/packages/producer/src/services/render/stages/captureStreamingStage.ts @@ -26,8 +26,7 @@ * - `lastBrowserConsole` is set to the buffer of whichever session * was active last (probe close path, or sequential session finally). * - `job.framesRendered` is updated per-frame; `Streaming frame N/M` - * `updateJobStatus` payloads fire at the same 30-frame and - * completion checkpoints (parallel) or every frame (sequential). + * goes through `reportFrameProgress`. * - Encoder close + result inspection happens inside the stage; a * `Streaming encode failed: ...` error throws on `success: false`. * - Defensive cleanup of `streamingEncoder` happens in the stage's @@ -76,7 +75,7 @@ import type { ProgressCallback, RenderJob } from "../../renderOrchestrator.js"; import { wrapCaptureStageError } from "../captureStageError.js"; import { pushWorkerDedupPerfs } from "../perfSummary.js"; import { ensureFrameWritten } from "./captureHdrFrameShared.js"; -import { updateJobStatus } from "../shared.js"; +import { reportFrameProgress, reportWorkerStartup } from "../shared.js"; import { encoderFailureError } from "../encoderInterruption.js"; import type { SdrStreamingCapturePlan } from "../capturePlan.js"; @@ -490,12 +489,12 @@ async function runWorkerEncodePipelineLoop( reorderBuffer.advanceTo(prev.idx + 1); job.framesRendered = prev.idx + 1; lastProgressAt = Date.now(); - updateJobStatus( + reportFrameProgress( job, - "rendering", `Streaming frame ${prev.idx + 1}/${totalFrames}`, Math.round(25 + ((prev.idx + 1) / totalFrames) * 55), onProgress, + prev.idx + 1 === totalFrames, ); }; @@ -517,12 +516,12 @@ async function runWorkerEncodePipelineLoop( reorderBuffer.advanceTo(item.idx + 1); job.framesRendered = item.idx + 1; lastProgressAt = Date.now(); - updateJobStatus( + reportFrameProgress( job, - "rendering", `Streaming frame ${item.idx + 1}/${totalFrames}`, Math.round(25 + ((item.idx + 1) / totalFrames) * 55), onProgress, + item.idx + 1 === totalFrames, ); } }; @@ -813,6 +812,7 @@ export async function runCaptureStreamingStage( canvasDrawElement: phase.canvasDrawElement, gpuBackend: phase.gpuBackend, }); + if (progress.capturedFrames === 0) reportWorkerStartup(job, progress, onProgress); return; } if (progress.capturedFrames > lastCapturedFrames) { @@ -822,20 +822,13 @@ export async function runCaptureStreamingStage( } job.framesRendered = progress.capturedFrames; const frameProgress = progress.capturedFrames / progress.totalFrames; - const progressPct = 25 + frameProgress * 55; - - if ( - progress.capturedFrames % 30 === 0 || - progress.capturedFrames === progress.totalFrames - ) { - updateJobStatus( - job, - "rendering", - `Streaming frame ${progress.capturedFrames}/${progress.totalFrames} (${workerCount} workers)`, - Math.round(progressPct), - onProgress, - ); - } + reportFrameProgress( + job, + `Streaming frame ${progress.capturedFrames}/${progress.totalFrames} (${workerCount} workers)`, + Math.round(25 + frameProgress * 55), + onProgress, + progress.capturedFrames === progress.totalFrames, + ); }, onFrameBuffer, // Interleaved DE workers each need their own browser PROCESS: @@ -973,12 +966,12 @@ export async function runCaptureStreamingStage( // capture error wrapper below must remain separate from finally so it // can throw with the browser console before encoder cleanup runs. // fallow-ignore-next-line code-duplication - updateJobStatus( + reportFrameProgress( job, - "rendering", `Streaming frame ${i + 1}/${totalFrames}`, Math.round(progress), onProgress, + i + 1 === totalFrames, ); } } diff --git a/packages/producer/src/services/renderOrchestrator.ts b/packages/producer/src/services/renderOrchestrator.ts index 0295bf2a1c..f9c4694375 100644 --- a/packages/producer/src/services/renderOrchestrator.ts +++ b/packages/producer/src/services/renderOrchestrator.ts @@ -3494,6 +3494,7 @@ async function executeRenderPipeline(input: { const framesDir = join(workDir, "captured-frames"); if (!existsSync(framesDir)) mkdirSync(framesDir, { recursive: true }); + updateJobStatus(job, "rendering", "Checking browser GPU", 25, onProgress); const resolvedBrowserGpuMode = await resolveBrowserGpuMode(cfg.browserGpuMode, { chromePath: resolveHeadlessShellPath(cfg), browserTimeout: cfg.browserTimeout, @@ -3893,6 +3894,7 @@ async function executeRenderPipeline(input: { !deInversionEligible && !deParallelRouterEligible ) { + updateJobStatus(job, "rendering", "Measuring capture speed", 25, onProgress); const outcome = await observeRenderStage( observability, "capture_calibration", From 74548b5c72ee11d8589e642f165834b48d01b8dd Mon Sep 17 00:00:00 2001 From: Miguel Angel Simon Sierra Date: Thu, 24 Sep 2026 07:34:29 -0400 Subject: [PATCH 2/2] fix(producer): always report the first captured frame and keep one telemetry code for browser start-up --- packages/cli/src/commands/render.test.ts | 5 +++ packages/cli/src/commands/render.ts | 2 ++ .../producer/src/services/render/shared.ts | 34 +++++++++++++------ .../stages/captureStreamingStage.test.ts | 7 ++-- 4 files changed, 35 insertions(+), 13 deletions(-) diff --git a/packages/cli/src/commands/render.test.ts b/packages/cli/src/commands/render.test.ts index 5eec74b758..039a270fa4 100644 --- a/packages/cli/src/commands/render.test.ts +++ b/packages/cli/src/commands/render.test.ts @@ -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"); }); diff --git a/packages/cli/src/commands/render.ts b/packages/cli/src/commands/render.ts index 71499db68d..7553b06073 100644 --- a/packages/cli/src/commands/render.ts +++ b/packages/cli/src/commands/render.ts @@ -1649,6 +1649,8 @@ const KNOWN_STAGE_CODES: Readonly> = { 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() diff --git a/packages/producer/src/services/render/shared.ts b/packages/producer/src/services/render/shared.ts index 3459fa7720..e1a52b92d3 100644 --- a/packages/producer/src/services/render/shared.ts +++ b/packages/producer/src/services/render/shared.ts @@ -266,27 +266,38 @@ export function updateJobStatus( if (onProgress) void onProgress(job, stage); } -const FRAME_PROGRESS_INTERVAL_MS = 250; +const PROGRESS_REPORT_INTERVAL_MS = 250; const lastFrameReportAt = new WeakMap(); +const lastStartupReportAt = new WeakMap(); -/** - * Capture-loop progress: the job updates on every call, the callback fires on - * the first call, the last frame, and at most once per 250 ms in between. - */ -export function reportFrameProgress( +// 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, job: RenderJob, stage: string, progress: number, onProgress: ProgressCallback | undefined, - isLastFrame: boolean, + force: boolean, ): void { const now = Date.now(); - const last = lastFrameReportAt.get(job); - const due = isLastFrame || last === undefined || now - last >= FRAME_PROGRESS_INTERVAL_MS; - if (due) lastFrameReportAt.set(job, 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>(); /** Browser warm-up before the first frame, from the parallel workers' phase events. */ @@ -304,7 +315,8 @@ export function reportWorkerStartup( const ready = [...phases].filter( ([id, p]) => id < progress.activeWorkers && (p === "frame_capture" || p === "frame_encode"), ).length; - reportFrameProgress( + reportThrottled( + lastStartupReportAt, job, `Starting browsers (${ready}/${progress.activeWorkers} ready)`, job.progress, diff --git a/packages/producer/src/services/render/stages/captureStreamingStage.test.ts b/packages/producer/src/services/render/stages/captureStreamingStage.test.ts index 588ee04691..b2e334899a 100644 --- a/packages/producer/src/services/render/stages/captureStreamingStage.test.ts +++ b/packages/producer/src/services/render/stages/captureStreamingStage.test.ts @@ -442,8 +442,11 @@ describe("runCaptureStreamingStage", () => { it("does not flood the callback when frames arrive faster than the report interval", async () => { const stages = await streamParallelFrames(300, 1); - expect(stages.filter((stage) => stage.startsWith("Streaming frame")).length).toBeLessThan(5); - expect(stages.at(-1)).toBe("Streaming frame 300/300 (2 workers)"); + const frames = stages.filter((stage) => stage.startsWith("Streaming frame")); + expect(frames.length).toBeLessThan(5); + // Frame 1 lands 1 ms after the start-up report and is still reported. + expect(frames[0]).toBe("Streaming frame 1/300 (2 workers)"); + expect(frames.at(-1)).toBe("Streaming frame 300/300 (2 workers)"); }); it("releases the writer with the dead worker's own error, not a stall", async () => {