From a23e63326d923d256e6e523f5b985de0ba3f6300 Mon Sep 17 00:00:00 2001 From: Vance Ingalls Date: Wed, 8 Jul 2026 23:06:50 -0700 Subject: [PATCH] fix(producer): restore HF_DE_PARALLEL_STREAM on every render exit path MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Address unanimous review feedback on PR #2095 (Miga, Magi REQUEST_CHANGES, Rames D Jusso): the router's process.env.HF_DE_PARALLEL_STREAM mutation was only cleared on the DE self-verify-retry branch — every other exit (happy path, any non-DE-verify error, abort) left it set, leaking the parallel-streaming opt-in into the next render sharing the same process (a regression/benchmark harness, or any batch host). Capture the prior value before mutating and restore it in the outer `finally` (executeRenderJob's own top-level try/finally, which runs on every exit path by construction), not just the narrow retry branch. Note `deParallelRouter` and the new `deParallelStreamEnvBefore` had to move above the outer `try` — a `let` declared inside `try {}` is not visible in the sibling `finally {}` block in JS, so the original placement alongside the other DE state would not have compiled once referenced from the finally. Also: renamed the shared `preInversionWorkerCount` local to `preRoutingWorkerCount` (Miga + Rames nit — it now serves both the inversion and the router), and pinned worker count 3 explicitly overrides calibration by design (documented per Miga/Rames's question, not a bug). Verified end-to-end (not just unit tests): ran two executeRenderJob calls back-to-back in one process, router-eligible then not — env is restored to undefined after render 1 and stays clean through render 2, the exact leak scenario the reviews described. 117 orchestrator tests pass (3 new, covering the restoreEnv primitive directly). --- .../src/services/renderOrchestrator.test.ts | 32 +++++++- .../src/services/renderOrchestrator.ts | 74 +++++++++++++++---- 2 files changed, 88 insertions(+), 18 deletions(-) diff --git a/packages/producer/src/services/renderOrchestrator.test.ts b/packages/producer/src/services/renderOrchestrator.test.ts index 147258f1b..58b522b3c 100644 --- a/packages/producer/src/services/renderOrchestrator.test.ts +++ b/packages/producer/src/services/renderOrchestrator.test.ts @@ -30,6 +30,7 @@ import { shouldDiscardProbeSessionForPageSideCompositing, resolveInversionRetryPlan, resolveParallelRouterRetryPlan, + restoreEnv, shouldPreferParallelDrawElement, shouldPreferSingleWorkerDrawElement, shouldUseStreamingEncode, @@ -1829,8 +1830,35 @@ describe("resolveParallelRouterRetryPlan (self-verify retry rollback)", () => { deParallelRouter: "reverted", }); } finally { - if (prevEnv === undefined) delete process.env.HF_DE_PARALLEL_STREAM; - else process.env.HF_DE_PARALLEL_STREAM = prevEnv; + restoreEnv("HF_DE_PARALLEL_STREAM", prevEnv); } }); }); + +describe("restoreEnv (DE parallel-router leak-safety primitive)", () => { + const VAR = "HF_DE_PARALLEL_STREAM_TEST_VAR"; + const prev = process.env[VAR]; + + afterEach(() => { + if (prev === undefined) delete process.env[VAR]; + else process.env[VAR] = prev; + }); + + it("deletes the var when the captured prior value was undefined", () => { + process.env[VAR] = "true"; + restoreEnv(VAR, undefined); + expect(process.env[VAR]).toBeUndefined(); + }); + + it("restores the exact prior string value, including a falsy-looking one", () => { + process.env[VAR] = "true"; + restoreEnv(VAR, "false"); + expect(process.env[VAR]).toBe("false"); + }); + + it("is a no-op shape when there was nothing to restore", () => { + delete process.env[VAR]; + restoreEnv(VAR, undefined); + expect(process.env[VAR]).toBeUndefined(); + }); +}); diff --git a/packages/producer/src/services/renderOrchestrator.ts b/packages/producer/src/services/renderOrchestrator.ts index 34f3a0a89..5a8c8299e 100644 --- a/packages/producer/src/services/renderOrchestrator.ts +++ b/packages/producer/src/services/renderOrchestrator.ts @@ -955,6 +955,19 @@ function replaceBodyWithRenderClone(body: HTMLElement, renderClone: Element): vo body.appendChild(renderClone); } +/** + * Restore an env var to a previously-captured value (`undefined` deletes it + * rather than setting the literal string "undefined"). Used to make the DE + * parallel-router's HF_DE_PARALLEL_STREAM mutation safe to leave uncleared on + * any exit path other than the one that explicitly calls this — see the + * outer `finally` in executeRenderJob (review: an unrestored env var leaks + * into the next render sharing this process). + */ +export function restoreEnv(name: string, prev: string | undefined): void { + if (prev === undefined) delete process.env[name]; + else process.env[name] = prev; +} + export function shouldUseStreamingEncode( cfg: Pick, outputFormat: NonNullable, @@ -1311,6 +1324,13 @@ export async function executeRenderJob( // between declaration and the try-block (currently impossible, but // defensible if more setup ever lands here) can't leak the interval. let memSampler: MemorySampler | null = null; + // Same reason these two live out here rather than with the rest of the DE + // state below (which is fine staying try-scoped, nothing outside the try + // reads it): a `let` declared inside `try {}` is NOT visible in the + // sibling `finally {}` block in JS — they're independent block scopes — + // so the leak-safety restore in `finally` needs these declared here. + let deParallelRouter: "routed" | "reverted" | undefined; + let deParallelStreamEnvBefore: string | undefined; try { memSampler = createMemorySampler(); @@ -1435,11 +1455,12 @@ export async function executeRenderJob( // "inverted" = fired and held; "reverted" = fired but the self-verify // retry rolled back to the parallel path; undefined = never fired. let deWorkerInversion: "inverted" | "reverted" | undefined; - // "routed" = the parallel router fired and held; "reverted" = fired but - // the self-verify retry rolled back; undefined = never fired. Mutually - // exclusive with deWorkerInversion — the router takes priority when both - // would be eligible (see shouldPreferParallelDrawElement). - let deParallelRouter: "routed" | "reverted" | undefined; + // deParallelRouter: "routed" = the parallel router fired and held; + // "reverted" = fired but the self-verify retry rolled back; undefined = + // never fired. Mutually exclusive with deWorkerInversion — the router + // takes priority when both would be eligible (see + // shouldPreferParallelDrawElement). Declared above the try (with + // deParallelStreamEnvBefore) so the outer `finally` can read it. let deSelfVerifyFallback = false; let deFallbackReason: string | undefined; let deDrainStats: import("./render/stages/captureStreamingStage.js").DeDrainStats | undefined; @@ -1916,13 +1937,20 @@ export async function executeRenderJob( // INIT-time gate at capture (css-effects / at-risk, ~1.5% of local // renders) render single-worker screenshot streaming — slower than // parallel would have been, accepted for the routing win everywhere else. - // `preInversionWorkerCount` lets the self-verify retry return to the - // parallel path when the drawElement bet loses. - const preInversionWorkerCount = workerCount; + // `preRoutingWorkerCount` lets the self-verify retry return to the + // parallel path when the drawElement bet loses — shared by both the + // inversion and the router below, whichever fires (mutually exclusive). + const preRoutingWorkerCount = workerCount; // Router takes priority over the single-worker inversion when both would // fire — its higher frame threshold means this only ever picks up long- // tail comps the inversion's own benchmark didn't cover (see - // shouldPreferParallelDrawElement). + // shouldPreferParallelDrawElement). Pins to a fixed worker count exactly + // like the inversion pins to 1 — calibration is skipped for both (see + // the capture_calibration gate above), so this deliberately overrides + // whatever a calibrated resolution would have chosen (e.g. 2, on a + // resource-constrained host): the benchmark validated par3 specifically, + // not "whatever calibration picks above 1", and the self-verify retry is + // the safety net if 3 workers tips a given host over. if (deParallelRouterEligible && workerCount > 1) { deParallelRouter = "routed"; // Fixed at 3, not calibration-derived: the benchmark validated exactly @@ -1933,9 +1961,14 @@ export async function executeRenderJob( log.info( "[Render] Fast capture: verified parallel drawElement streaming preferred over " + `single-worker inversion (${totalFrames} frames >= ${deParallelMinFrames}; ` + - "benchmark-validated at 3 workers). Set HF_DE_PARALLEL_ROUTER=false or --workers N to override.", + "benchmark-validated at 3 workers, pinned regardless of calibration). " + + "Set HF_DE_PARALLEL_ROUTER=false or --workers N to override.", ); workerCount = ROUTER_WORKER_COUNT; + // Captured BEFORE the mutation so the outer `finally` can restore the + // exact prior value on every exit path, not just the self-verify + // retry below (review) — see deParallelStreamEnvBefore's declaration. + deParallelStreamEnvBefore = process.env.HF_DE_PARALLEL_STREAM; process.env.HF_DE_PARALLEL_STREAM = "true"; } else if (deInversionEligible && workerCount > 1) { deWorkerInversion = "inverted"; @@ -2254,21 +2287,22 @@ export async function executeRenderJob( deSelfVerifyFallback: true, }); probeSession = null; - // HF_DE_PARALLEL_STREAM must be cleared BEFORE resolveParallelRouterRetryPlan + // HF_DE_PARALLEL_STREAM must be restored BEFORE resolveParallelRouterRetryPlan // recomputes useStreamingEncode, or shouldUseStreamingEncode's own env // check would keep resolving to the parallel-streaming shape on the // retry instead of the well-tested parallel-disk fallback. - if (deParallelRouter === "routed") delete process.env.HF_DE_PARALLEL_STREAM; + if (deParallelRouter === "routed") + restoreEnv("HF_DE_PARALLEL_STREAM", deParallelStreamEnvBefore); const inversionRetryPlan = resolveInversionRetryPlan({ deWorkerInversion, - preInversionWorkerCount, + preInversionWorkerCount: preRoutingWorkerCount, cfg, outputFormat, durationSeconds: job.duration, }); const parallelRouterRetryPlan = resolveParallelRouterRetryPlan({ deParallelRouter, - preRouterWorkerCount: preInversionWorkerCount, + preRouterWorkerCount: preRoutingWorkerCount, cfg, outputFormat, durationSeconds: job.duration, @@ -2519,9 +2553,9 @@ export async function executeRenderJob( compileGate: deCompileGate, clampReason: deClampReason, workerInversion: deWorkerInversion, - preInversionWorkers: deWorkerInversion ? preInversionWorkerCount : undefined, + preInversionWorkers: deWorkerInversion ? preRoutingWorkerCount : undefined, parallelRouter: deParallelRouter, - preRouterWorkers: deParallelRouter ? preInversionWorkerCount : undefined, + preRouterWorkers: deParallelRouter ? preRoutingWorkerCount : undefined, selfVerifyFallback: deSelfVerifyFallback, fallbackReason: deFallbackReason, drainStats: deDrainStats, @@ -2683,5 +2717,13 @@ export async function executeRenderJob( throw error; } finally { memSampler?.stop(); + // Guaranteed restore regardless of exit path (success, any thrown error, + // abort) — the DE self-verify retry path also restores this mid-render + // (so ITS OWN recomputation sees the right value), but that branch is + // only one of several ways this render can end; this is the actual + // leak fix (review). + if (deParallelStreamEnvBefore !== undefined || deParallelRouter) { + restoreEnv("HF_DE_PARALLEL_STREAM", deParallelStreamEnvBefore); + } } }