Skip to content

Commit 85dfb59

Browse files
committed
feat(producer): emit eligible_off cohort-sizing signal when parallel-stream router is off
1 parent 8c39822 commit 85dfb59

4 files changed

Lines changed: 29 additions & 7 deletions

File tree

packages/cli/src/telemetry/events.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,8 @@ export interface RenderObservabilityTelemetryPayload {
4040
capturePlayerReadyTimeoutMs?: number;
4141
captureTransientRetries?: number;
4242
captureMemoryExhaustionDetected?: boolean;
43-
/** Non-DE parallel-streaming router outcome ("screenshot" | "beginframe"). */
43+
/** Non-DE parallel-streaming router outcome ("screenshot" | "beginframe" —
44+
* routed; "eligible_off" — would route but the kill switch is off). */
4445
captureParallelStream?: string;
4546
observabilityExtractVideoCount?: number;
4647
observabilityExtractedVideoCount?: number;

packages/cli/src/telemetry/renderObservability.test.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,13 @@ describe("renderObservabilityTelemetryPayload — non-DE parallel-stream router"
4747
expect(payload.captureParallelStream).toBe("beginframe");
4848
});
4949

50+
it("maps the passive eligible_off cohort-sizing signal", () => {
51+
const payload = renderObservabilityTelemetryPayload(
52+
makeSummary({ captureParallelStream: "eligible_off" }),
53+
);
54+
expect(payload.captureParallelStream).toBe("eligible_off");
55+
});
56+
5057
it("stays undefined when the router never fired", () => {
5158
const payload = renderObservabilityTelemetryPayload(makeSummary({}));
5259
expect(payload.captureParallelStream).toBeUndefined();

packages/producer/src/services/render/observability.ts

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -48,11 +48,14 @@ export interface RenderCaptureObservability {
4848
deParallelRouter?: "routed" | "reverted";
4949
/**
5050
* Non-DE parallel-streaming router outcome (HF_CAPTURE_PARALLEL_STREAM):
51-
* set when a multi-worker screenshot/BeginFrame render was routed through
52-
* the interleaved streaming encoder instead of the parallel disk path.
53-
* The value is the capture mode that streamed. Absent = not routed.
51+
* "screenshot" | "beginframe" — the render passed every gate AND the kill
52+
* switch was on, so it was routed through the interleaved streaming encoder
53+
* (the value is the capture mode that streamed); "eligible_off" — the render
54+
* passed every gate EXCEPT the kill switch (passive cohort-sizing signal for
55+
* the default-off soak: how many renders WOULD route if enabled). Absent =
56+
* ineligible regardless of the switch.
5457
*/
55-
captureParallelStream?: "screenshot" | "beginframe";
58+
captureParallelStream?: "screenshot" | "beginframe" | "eligible_off";
5659
protocolTimeoutMs?: number;
5760
pageNavigationTimeoutMs?: number;
5861
playerReadyTimeoutMs?: number;

packages/producer/src/services/renderOrchestrator.ts

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2032,13 +2032,17 @@ export async function executeRenderJob(
20322032
// Non-DE parallel-streaming router — see shouldStreamParallelCapture.
20332033
// Mutually exclusive with the DE inversion/router above by construction
20342034
// (both DE predicates require useDrawElement; this requires its negation).
2035-
const captureParallelStreamEligible = shouldStreamParallelCapture({
2036-
routerEnabled: process.env.HF_CAPTURE_PARALLEL_STREAM === "true",
2035+
const captureParallelStreamRouterEnabled = process.env.HF_CAPTURE_PARALLEL_STREAM === "true";
2036+
const captureParallelStreamArgs = {
20372037
workerCount,
20382038
useDrawElement: cfg.useDrawElement,
20392039
outputFormat,
20402040
streamingOk: shouldUseStreamingEncode(cfg, outputFormat, 1, job.duration),
20412041
layeredOrEffectRoute: hasHdrContent || compiled.hasShaderTransitions,
2042+
};
2043+
const captureParallelStreamEligible = shouldStreamParallelCapture({
2044+
routerEnabled: captureParallelStreamRouterEnabled,
2045+
...captureParallelStreamArgs,
20422046
});
20432047
if (captureParallelStreamEligible) {
20442048
captureParallelStreamForced = true;
@@ -2060,6 +2064,13 @@ export async function executeRenderJob(
20602064
"worker_resolution",
20612065
`parallel ${captureParallelStream} capture routed to streaming`,
20622066
);
2067+
} else if (shouldStreamParallelCapture({ routerEnabled: true, ...captureParallelStreamArgs })) {
2068+
// The kill switch is the ONLY failed gate: emit a passive cohort-sizing
2069+
// signal (capture_parallel_stream = "eligible_off") so the default-off
2070+
// soak can measure how many fleet renders WOULD route before anyone
2071+
// enables the flag. Observability-only — no behavior change, no log
2072+
// noise on the default path.
2073+
updateCaptureObservability({ captureParallelStream: "eligible_off" });
20632074
}
20642075

20652076
if (workerCount > 1 && probeSession) {

0 commit comments

Comments
 (0)