fix(telemetry): expose stalled render stages (#2220)

* fix(telemetry): expose stalled render stages

* fix(telemetry): preserve capture data on terminal stage events

* fix(telemetry): fix calibration TDZ crash, tag encode/assemble, extend heartbeat cadence

capture_calibration referenced captureStageObservationData before its
declaration (later in the same scope), which would throw a ReferenceError
for any render hitting the calibration path. Hoist the closure and split
workerCount's declaration from its resolution so calibration can safely
read it as undefined before capture strategy resolves worker count.

Also address the two non-blocking review items: wire encode/assemble
stages through captureStageObservationData for consistent tagging, and
extend the heartbeat schedule to repeat every 120s after the initial
30/60/120s ramp instead of going dark on stalls beyond two minutes.
This commit is contained in:
Miguel Ángel
2026-07-10 23:59:08 -04:00
committed by GitHub
parent d51fa7eba2
commit 87618eef4c
8 changed files with 383 additions and 23 deletions
+58
View File
@@ -41,6 +41,7 @@ const trackingState = vi.hoisted(() => ({
// `.ts` source under vitest — mocked here so the CLI-trial tests can
// control it directly instead of inheriting that environment quirk.
shouldTrack: true,
renderObservations: [] as Array<Record<string, unknown>>,
}));
const preflightState = vi.hoisted(() => ({
@@ -110,6 +111,9 @@ vi.mock("../telemetry/client.js", () => ({
vi.mock("../telemetry/events.js", () => ({
trackRenderComplete: vi.fn(),
trackRenderError: vi.fn(),
trackRenderObservation: vi.fn((props: Record<string, unknown>) => {
trackingState.renderObservations.push(props);
}),
}));
vi.mock("../browser/ffmpeg.js", () => ({
@@ -165,6 +169,7 @@ describe("renderLocal browser GPU config", () => {
configState.failWrites = 0;
configState.writeConfigCalls = [];
trackingState.shouldTrack = true;
trackingState.renderObservations = [];
resetTrialState();
savedEnv.clear();
savedEnv.set("HYPERFRAMES_FFMPEG_PATH", process.env.HYPERFRAMES_FFMPEG_PATH);
@@ -210,6 +215,59 @@ describe("renderLocal browser GPU config", () => {
});
}, 15_000);
it("forwards render stage start and end lifecycle events to telemetry", async () => {
producerState.executeImpl = async (job) => {
const logger = (job.config as { logger: { info: (message: string, meta: object) => void } })
.logger;
logger.info("[Render:trace]", {
renderJobId: "render-lifecycle",
phase: "capture_streaming",
status: "start",
elapsedMs: 100,
workerCount: 1,
captureMode: "screenshot",
captureOperation: "captureScreenshot",
framesCompleted: 12,
totalFrames: 900,
});
logger.info("[Render:trace]", {
renderJobId: "render-lifecycle",
phase: "capture_streaming",
status: "end",
elapsedMs: 250,
durationMs: 150,
});
};
await renderLocal("/tmp/project", "/tmp/out.mp4", {
fps: { num: 30, den: 1 },
quality: "standard",
format: "mp4",
gpu: false,
browserGpuMode: "software",
hdrMode: "auto",
quiet: true,
skipFeedback: true,
});
expect(trackingState.renderObservations).toEqual([
expect.objectContaining({
renderJobId: "render-lifecycle",
phase: "capture_streaming",
status: "start",
captureOperation: "captureScreenshot",
framesCompleted: 12,
totalFrames: 900,
}),
expect.objectContaining({
renderJobId: "render-lifecycle",
phase: "capture_streaming",
status: "end",
durationMs: 150,
}),
]);
});
it("forwards browserGpuMode='auto' into producer config (probe-then-choose)", async () => {
await renderLocal("/tmp/project", "/tmp/out.mp4", {
fps: { num: 30, den: 1 },
+8 -1
View File
@@ -1578,7 +1578,9 @@ function metaBoolean(meta: Record<string, unknown> | undefined, key: string): bo
function trackRenderTraceFromLog(message: string, meta: Record<string, unknown> | undefined): void {
if (message !== "[Render:trace]") return;
const status = metaString(meta, "status");
if (status !== "checkpoint" && status !== "error") return;
if (status !== "start" && status !== "end" && status !== "checkpoint" && status !== "error") {
return;
}
trackRenderObservation({
source: "cli",
renderJobId: metaString(meta, "renderJobId"),
@@ -1595,6 +1597,11 @@ function trackRenderTraceFromLog(message: string, meta: Record<string, unknown>
usePageSideCompositing: metaBoolean(meta, "usePageSideCompositing"),
hasHdrContent: metaBoolean(meta, "hasHdrContent"),
captureMode: metaString(meta, "captureMode"),
captureOperation: metaString(meta, "captureOperation"),
framesCompleted: metaNumber(meta, "framesCompleted"),
totalFrames: metaNumber(meta, "totalFrames"),
heartbeatIndex: metaNumber(meta, "heartbeatIndex"),
stageElapsedMs: metaNumber(meta, "stageElapsedMs"),
videoCount: metaNumber(meta, "videoCount"),
extractedVideoCount: metaNumber(meta, "extractedVideoCount"),
totalFramesExtracted: metaNumber(meta, "totalFramesExtracted"),
+12
View File
@@ -316,6 +316,12 @@ describe("render telemetry events", () => {
phase: "capture_hdr_layered",
status: "error",
compositionHash: "abc123",
captureMode: "screenshot",
captureOperation: "captureScreenshot",
framesCompleted: 12,
totalFrames: 900,
heartbeatIndex: 1,
stageElapsedMs: 30_000,
message: "Navigation failed for C:\\Users\\Alice\\project\\video.mov?not-a-query",
});
@@ -324,6 +330,12 @@ describe("render telemetry events", () => {
expect.objectContaining({
render_job_id: "render-123",
composition_hash: "abc123",
capture_mode: "screenshot",
capture_operation: "captureScreenshot",
frames_completed: 12,
total_frames: 900,
heartbeat_index: 1,
stage_elapsed_ms: 30_000,
message: "Navigation failed for [path]",
}),
);
+10
View File
@@ -347,6 +347,11 @@ export function trackRenderObservation(props: {
usePageSideCompositing?: boolean;
hasHdrContent?: boolean;
captureMode?: string;
captureOperation?: string;
framesCompleted?: number;
totalFrames?: number;
heartbeatIndex?: number;
stageElapsedMs?: number;
videoCount?: number;
extractedVideoCount?: number;
totalFramesExtracted?: number;
@@ -373,6 +378,11 @@ export function trackRenderObservation(props: {
use_page_side_compositing: props.usePageSideCompositing,
has_hdr_content: props.hasHdrContent,
capture_mode: props.captureMode,
capture_operation: props.captureOperation,
frames_completed: props.framesCompleted,
total_frames: props.totalFrames,
heartbeat_index: props.heartbeatIndex,
stage_elapsed_ms: props.stageElapsedMs,
video_count: props.videoCount,
extracted_video_count: props.extractedVideoCount,
total_frames_extracted: props.totalFramesExtracted,
@@ -2,6 +2,7 @@ import { describe, expect, it, vi } from "vitest";
import { CaptureStageError, getCaptureStageBrowserConsole } from "./captureStageError.js";
import {
computeCompositionObservabilityHash,
observeRenderStage,
RenderObservabilityRecorder,
sanitizeObservationMessage,
summarizeBrowserDiagnostics,
@@ -95,6 +96,169 @@ describe("CaptureStageError", () => {
});
describe("RenderObservabilityRecorder", () => {
it("emits capped heartbeats while a stage is still running and stops after settlement", async () => {
vi.useFakeTimers();
const log = makeLog();
const recorder = new RenderObservabilityRecorder({
pipelineStartMs: Date.now(),
log,
renderJobId: "render-hang",
});
let resolveStage: (() => void) | undefined;
let framesCompleted = 0;
const stage = observeRenderStage(
recorder,
"capture_streaming",
{
workerCount: 1,
captureMode: "screenshot",
captureOperation: "captureScreenshot",
totalFrames: 900,
get framesCompleted() {
return framesCompleted;
},
},
() =>
new Promise<void>((resolve) => {
resolveStage = resolve;
}),
);
framesCompleted = 12;
await vi.advanceTimersByTimeAsync(30_000);
expect(log.info).toHaveBeenCalledWith(
"[Render:trace]",
expect.objectContaining({
phase: "capture_streaming",
status: "checkpoint",
message: "stage still running",
heartbeatIndex: 1,
stageElapsedMs: 30_000,
captureMode: "screenshot",
captureOperation: "captureScreenshot",
framesCompleted: 12,
totalFrames: 900,
}),
);
await vi.advanceTimersByTimeAsync(90_000);
const heartbeatCalls = log.info.mock.calls.filter(
([message, meta]) => message === "[Render:trace]" && meta?.message === "stage still running",
);
expect(heartbeatCalls).toHaveLength(3);
resolveStage?.();
await stage;
const endCall = log.info.mock.calls.find(
// fallow-ignore-next-line complexity
([message, meta]) =>
message === "[Render:trace]" &&
meta?.phase === "capture_streaming" &&
meta?.status === "end",
);
expect(endCall?.[1]).toEqual(
expect.objectContaining({
framesCompleted: 12,
totalFrames: 900,
captureMode: "screenshot",
captureOperation: "captureScreenshot",
workerCount: 1,
}),
);
await vi.advanceTimersByTimeAsync(240_000);
expect(
log.info.mock.calls.filter(
([message, meta]) =>
message === "[Render:trace]" && meta?.message === "stage still running",
),
).toHaveLength(3);
vi.useRealTimers();
});
it("keeps emitting heartbeats every 120s once the initial ramp is exhausted", async () => {
vi.useFakeTimers();
const log = makeLog();
const recorder = new RenderObservabilityRecorder({
pipelineStartMs: Date.now(),
log,
renderJobId: "render-long-hang",
});
let resolveStage: (() => void) | undefined;
const stage = observeRenderStage(
recorder,
"capture_streaming",
{ captureMode: "screenshot" },
() =>
new Promise<void>((resolve) => {
resolveStage = resolve;
}),
);
// Ramp: 30s, 60s, 120s → 3 heartbeats. Then steady 120s cadence: 240s, 360s.
await vi.advanceTimersByTimeAsync(360_000);
const heartbeatCalls = log.info.mock.calls.filter(
([message, meta]) => message === "[Render:trace]" && meta?.message === "stage still running",
);
expect(heartbeatCalls.map(([, meta]) => meta?.heartbeatIndex)).toEqual([1, 2, 3, 4, 5]);
expect(heartbeatCalls.map(([, meta]) => meta?.stageElapsedMs)).toEqual([
30_000, 60_000, 120_000, 240_000, 360_000,
]);
resolveStage?.();
await stage;
vi.useRealTimers();
});
it("clears pending heartbeats when a stage rejects", async () => {
vi.useFakeTimers();
const log = makeLog();
const recorder = new RenderObservabilityRecorder({
pipelineStartMs: Date.now(),
log,
renderJobId: "render-error",
});
await expect(
observeRenderStage(
recorder,
"capture_disk",
{
workerCount: 2,
totalFrames: 42,
framesCompleted: 7,
captureMode: "screenshot",
captureOperation: "captureScreenshot",
},
async () => {
throw new Error("capture failed");
},
),
).rejects.toThrow("capture failed");
const errorCall = log.info.mock.calls.find(
// fallow-ignore-next-line complexity
([message, meta]) =>
message === "[Render:trace]" && meta?.phase === "capture_disk" && meta?.status === "error",
);
expect(errorCall?.[1]).toEqual(
expect.objectContaining({
framesCompleted: 7,
totalFrames: 42,
captureMode: "screenshot",
captureOperation: "captureScreenshot",
workerCount: 2,
}),
);
await vi.advanceTimersByTimeAsync(240_000);
expect(
log.info.mock.calls.some(
([message, meta]) =>
message === "[Render:trace]" && meta?.message === "stage still running",
),
).toBe(false);
vi.useRealTimers();
});
// fallow-ignore-next-line complexity
it("records bounded phase events and summarizes browser diagnostics", () => {
const log = makeLog();
@@ -118,6 +118,7 @@ const MAX_EVENTS = 160;
const ALLOWED_STRING_DATA_KEYS = new Set([
"browserGpuMode",
"captureMode",
"captureOperation",
"compositionHash",
"effectiveHdr",
"format",
@@ -363,6 +364,20 @@ export class RenderObservabilityRecorder {
}
}
/** Heartbeat ramp before falling back to a steady repeat cadence. */
const HEARTBEAT_RAMP_MS = [30_000, 60_000, 120_000];
const HEARTBEAT_REPEAT_MS = 120_000;
const HEARTBEAT_RAMP_END_MS =
HEARTBEAT_RAMP_MS[HEARTBEAT_RAMP_MS.length - 1] ?? HEARTBEAT_REPEAT_MS;
/** Target elapsed-ms for the Nth heartbeat (0-indexed): ramp, then steady repeat so long stalls keep emitting breadcrumbs instead of going dark after the ramp. */
function heartbeatTargetMs(index: number): number {
const rampTarget = HEARTBEAT_RAMP_MS[index];
if (rampTarget !== undefined) return rampTarget;
const overflow = index - HEARTBEAT_RAMP_MS.length + 1;
return HEARTBEAT_RAMP_END_MS + overflow * HEARTBEAT_REPEAT_MS;
}
export async function observeRenderStage<T>(
recorder: RenderObservabilityRecorder,
phase: string,
@@ -370,12 +385,35 @@ export async function observeRenderStage<T>(
fn: () => Promise<T>,
): Promise<T> {
const startedAt = recorder.stageStart(phase, data);
let heartbeatCount = 0;
let lastFiredAtMs = 0;
let heartbeatTimer: ReturnType<typeof setTimeout> | undefined;
const scheduleNextHeartbeat = () => {
const targetMs = heartbeatTargetMs(heartbeatCount);
heartbeatTimer = setTimeout(() => {
lastFiredAtMs = targetMs;
heartbeatCount += 1;
recorder.checkpoint(phase, "stage still running", {
...data,
heartbeatIndex: heartbeatCount,
stageElapsedMs: Date.now() - startedAt,
});
scheduleNextHeartbeat();
}, targetMs - lastFiredAtMs);
heartbeatTimer.unref?.();
};
scheduleNextHeartbeat();
const clearHeartbeats = () => {
clearTimeout(heartbeatTimer);
};
try {
const result = await fn();
recorder.stageEnd(phase, startedAt);
clearHeartbeats();
recorder.stageEnd(phase, startedAt, data);
return result;
} catch (error) {
recorder.stageError(phase, startedAt, error);
clearHeartbeats();
recorder.stageError(phase, startedAt, error, data);
throw error;
}
}
@@ -30,6 +30,7 @@ import {
shouldDiscardProbeSessionForPageSideCompositing,
resolveInversionRetryPlan,
resolveParallelRouterRetryPlan,
resetCaptureAttemptProgress,
shouldRetryViaPinnedFallback,
shouldPreferParallelDrawElement,
shouldPreferSingleWorkerDrawElement,
@@ -115,6 +116,14 @@ describe("extractStandaloneEntryFromIndex", () => {
});
describe("captureAttemptMadeProgress", () => {
it("resets completed frames before a fallback attempt starts", () => {
const job = { framesRendered: 900 };
resetCaptureAttemptProgress(job);
expect(job.framesRendered).toBe(0);
});
it("retries when the attempt captured at least one frame toward its target", () => {
// targeted 100 frames, 40 still missing -> 60 captured -> worth retrying the rest
expect(captureAttemptMadeProgress(100, 40)).toBe(true);
@@ -110,6 +110,7 @@ import {
observeRenderStage,
type RenderCaptureObservability,
type RenderExtractionObservability,
type RenderObservationData,
type RenderObservabilitySummary,
} from "./render/observability.js";
import { type HdrPerfCollector, type HdrPerfSummary } from "./render/hdrPerf.js";
@@ -660,6 +661,10 @@ export function captureAttemptMadeProgress(
return remainingFrameCount < attemptTargetFrameCount;
}
export function resetCaptureAttemptProgress(job: { framesRendered?: number }): void {
job.framesRendered = 0;
}
export function isRecoverableParallelCaptureError(error: unknown): boolean {
const message = normalizeErrorMessage(error);
return (
@@ -1934,6 +1939,51 @@ export async function executeRenderJob(
process.env.HF_DE_PARALLEL_STREAM === "true",
routerEnabled: deParallelRouterEnabled,
});
// Declared ahead of resolution (assigned below, after calibration) so
// captureStageObservationData can close over it for the calibration
// stage itself — reads as undefined until resolveRenderWorkerCount runs.
let workerCount: number;
const captureStageObservationData = (
extra: RenderObservationData = {},
): RenderObservationData => ({
...extra,
get workerCount() {
return workerCount;
},
get forceScreenshot() {
return captureForceScreenshot;
},
get totalFrames() {
return totalFrames;
},
get framesCompleted() {
return job.framesRendered ?? 0;
},
get captureMode() {
return (
probeSession?.captureMode ??
(captureForceScreenshot
? "screenshot"
: cfg.useDrawElement
? "drawelement"
: "beginframe")
);
},
get captureOperation() {
if ((job.framesRendered ?? 0) >= totalFrames) return "encode";
const mode =
probeSession?.captureMode ??
(captureForceScreenshot
? "screenshot"
: cfg.useDrawElement
? "drawelement"
: "beginframe");
if (mode === "screenshot") return "captureScreenshot";
if (mode === "drawelement") return "drawElement";
return "beginFrame";
},
});
if (
job.config.workers === undefined &&
totalFrames >= 60 &&
@@ -1945,7 +1995,7 @@ export async function executeRenderJob(
const outcome = await observeRenderStage(
observability,
"capture_calibration",
{ forceScreenshot: captureForceScreenshot },
captureStageObservationData({ forceScreenshot: captureForceScreenshot }),
() =>
runCaptureCalibration({
cfg,
@@ -1986,7 +2036,7 @@ export async function executeRenderJob(
// Low-memory safe-mode's single-worker pin lives inside
// resolveRenderWorkerCount so its "why workers=N" logging stays coherent.
let workerCount = resolveRenderWorkerCount(
workerCount = resolveRenderWorkerCount(
totalFrames,
job.config.workers,
cfg,
@@ -2229,7 +2279,7 @@ export async function executeRenderJob(
const effectiveQuality = job.config.crf ?? preset.quality;
const effectiveBitrate = job.config.crf != null ? undefined : job.config.videoBitrate;
job.framesRendered = 0;
resetCaptureAttemptProgress(job);
// ── Z-ordered multi-layer compositing ─────────────────────────────────
// Per frame: query all elements' z-order, group into layers (DOM or HDR),
@@ -2248,7 +2298,7 @@ export async function executeRenderJob(
const hdrRes = await observeRenderStage(
observability,
"capture_hdr_layered",
{ workerCount, forceScreenshot: captureForceScreenshot, hasHdrContent },
captureStageObservationData({ hasHdrContent }),
() =>
runCaptureHdrStage({
job,
@@ -2297,11 +2347,12 @@ export async function executeRenderJob(
let streamingHandled = false;
if (useStreamingEncode) {
const captureFrameStart = Date.now();
const invokeStreaming = () =>
observeRenderStage(
const invokeStreaming = () => {
resetCaptureAttemptProgress(job);
return observeRenderStage(
observability,
"capture_streaming",
{ workerCount, forceScreenshot: captureForceScreenshot },
captureStageObservationData(),
() =>
runCaptureStreamingStage({
fileServer: activeFileServer,
@@ -2339,6 +2390,7 @@ export async function executeRenderJob(
dedupPerfs,
}),
);
};
let streamingRes;
try {
streamingRes = await invokeStreaming();
@@ -2510,11 +2562,12 @@ export async function executeRenderJob(
if (!streamingHandled) {
// ── Disk-based capture (original flow) ────────────────────────────
resetCaptureAttemptProgress(job);
const captureFrameStart = Date.now();
const captureRes = await observeRenderStage(
observability,
"capture_disk",
{ workerCount, forceScreenshot: captureForceScreenshot, needsAlpha },
captureStageObservationData({ needsAlpha }),
() =>
runCaptureStage({
fileServer: activeFileServer,
@@ -2555,7 +2608,12 @@ export async function executeRenderJob(
const encodeRes = await observeRenderStage(
observability,
"encode",
{ hasAudio, isPngSequence, isGif, chunkedEncode: enableChunkedEncode },
captureStageObservationData({
hasAudio,
isPngSequence,
isGif,
chunkedEncode: enableChunkedEncode,
}),
() =>
runEncodeStage({
job,
@@ -2603,17 +2661,21 @@ export async function executeRenderJob(
// directory deliverable, and gif is written directly to outputPath by the
// two-pass palette encoder.
if (!isPngSequence && !isGif) {
const assembleRes = await observeRenderStage(observability, "assemble", { hasAudio }, () =>
runAssembleStage({
job,
videoOnlyPath,
audioOutputPath,
outputPath,
hasAudio,
abortSignal,
assertNotAborted,
onProgress,
}),
const assembleRes = await observeRenderStage(
observability,
"assemble",
captureStageObservationData({ hasAudio }),
() =>
runAssembleStage({
job,
videoOnlyPath,
audioOutputPath,
outputPath,
hasAudio,
abortSignal,
assertNotAborted,
onProgress,
}),
);
perfStages.assembleMs = assembleRes.assembleMs;
} else {