diff --git a/packages/cli/src/telemetry/events.ts b/packages/cli/src/telemetry/events.ts index c6273b2de..88196e247 100644 --- a/packages/cli/src/telemetry/events.ts +++ b/packages/cli/src/telemetry/events.ts @@ -53,7 +53,8 @@ export interface RenderObservabilityTelemetryPayload { captureDePreRouterWorkers?: number; captureDeSelfVerifyFallback?: boolean; captureDeFallbackReason?: string; - /** Non-DE parallel-streaming router outcome ("screenshot" | "beginframe"). */ + /** Non-DE parallel-streaming router outcome ("screenshot" | "beginframe" — + * routed; "eligible_off" — would route but the kill switch is off). */ captureParallelStream?: string; observabilityExtractVideoCount?: number; observabilityExtractedVideoCount?: number; diff --git a/packages/cli/src/telemetry/renderObservability.test.ts b/packages/cli/src/telemetry/renderObservability.test.ts index 8b5b3c99f..0a64cc866 100644 --- a/packages/cli/src/telemetry/renderObservability.test.ts +++ b/packages/cli/src/telemetry/renderObservability.test.ts @@ -91,6 +91,13 @@ describe("renderObservabilityTelemetryPayload — non-DE parallel-stream router" expect(payload.captureParallelStream).toBe("beginframe"); }); + it("maps the passive eligible_off cohort-sizing signal", () => { + const payload = renderObservabilityTelemetryPayload( + makeSummary({ captureParallelStream: "eligible_off" }), + ); + expect(payload.captureParallelStream).toBe("eligible_off"); + }); + it("stays undefined when the router never fired", () => { const payload = renderObservabilityTelemetryPayload(makeSummary({})); expect(payload.captureParallelStream).toBeUndefined(); diff --git a/packages/producer/src/services/render/observability.ts b/packages/producer/src/services/render/observability.ts index ca542532e..c180f0481 100644 --- a/packages/producer/src/services/render/observability.ts +++ b/packages/producer/src/services/render/observability.ts @@ -70,11 +70,14 @@ export interface RenderCaptureObservability { dePreRouterWorkers?: number; /** * Non-DE parallel-streaming router outcome (HF_CAPTURE_PARALLEL_STREAM): - * set when a multi-worker screenshot/BeginFrame render was routed through - * the interleaved streaming encoder instead of the parallel disk path. - * The value is the capture mode that streamed. Absent = not routed. + * "screenshot" | "beginframe" — the render passed every gate AND the kill + * switch was on, so it was routed through the interleaved streaming encoder + * (the value is the capture mode that streamed); "eligible_off" — the render + * passed every gate EXCEPT the kill switch (passive cohort-sizing signal for + * the default-off soak: how many renders WOULD route if enabled). Absent = + * ineligible regardless of the switch. */ - captureParallelStream?: "screenshot" | "beginframe"; + captureParallelStream?: "screenshot" | "beginframe" | "eligible_off"; protocolTimeoutMs?: number; pageNavigationTimeoutMs?: number; playerReadyTimeoutMs?: number; diff --git a/packages/producer/src/services/renderOrchestrator.ts b/packages/producer/src/services/renderOrchestrator.ts index 8328bfeba..74dcb3d17 100644 --- a/packages/producer/src/services/renderOrchestrator.ts +++ b/packages/producer/src/services/renderOrchestrator.ts @@ -2245,13 +2245,17 @@ export async function executeRenderJob( // Non-DE parallel-streaming router — see shouldStreamParallelCapture. // Mutually exclusive with the DE inversion/router above by construction // (both DE predicates require useDrawElement; this requires its negation). - const captureParallelStreamEligible = shouldStreamParallelCapture({ - routerEnabled: process.env.HF_CAPTURE_PARALLEL_STREAM === "true", + const captureParallelStreamRouterEnabled = process.env.HF_CAPTURE_PARALLEL_STREAM === "true"; + const captureParallelStreamArgs = { workerCount, useDrawElement: cfg.useDrawElement, outputFormat, streamingOk: shouldUseStreamingEncode(cfg, outputFormat, 1, job.duration), layeredOrEffectRoute: hasHdrContent || compiled.hasShaderTransitions, + }; + const captureParallelStreamEligible = shouldStreamParallelCapture({ + routerEnabled: captureParallelStreamRouterEnabled, + ...captureParallelStreamArgs, }); if (captureParallelStreamEligible) { captureParallelStreamForced = true; @@ -2273,6 +2277,13 @@ export async function executeRenderJob( "worker_resolution", `parallel ${captureParallelStream} capture routed to streaming`, ); + } else if (shouldStreamParallelCapture({ routerEnabled: true, ...captureParallelStreamArgs })) { + // The kill switch is the ONLY failed gate: emit a passive cohort-sizing + // signal (capture_parallel_stream = "eligible_off") so the default-off + // soak can measure how many fleet renders WOULD route before anyone + // enables the flag. Observability-only — no behavior change, no log + // noise on the default path. + updateCaptureObservability({ captureParallelStream: "eligible_off" }); } if (workerCount > 1 && probeSession) {