From 8ab2881b0acaae1de5aecf0b2dc7fcb76c5aec2b Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 01:24:17 +0900 Subject: [PATCH] fix(server): retain workflow slots for streaming turns --- src/server/index.ts | 4 +--- src/server/lifecycle.ts | 8 +++++++ structure/transports/responses.md | 2 ++ .../active-registry-admission.test.ts | 24 ++++++++++++++++++- 4 files changed, 34 insertions(+), 4 deletions(-) diff --git a/src/server/index.ts b/src/server/index.ts index 61fbcec7d8a..b7c9136ac20 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -515,16 +515,14 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe // still unreadable to a browser dashboard -- which made exposing it pointless. return withCors(workflowDecisionRefusalResponse(workflow, undefined, refusalLog), req, policy); } - const releaseWorkflow = (): void => { if (workflow?.admitted) workflow.lease.release(); }; + if (workflow?.admitted) lease.attach(workflow.lease); let response: Response; try { response = await runAdmittedBodyWork(req, policy, config.maxInboundBodyBytes, () => work(lease), refusalLog); } catch (error) { - releaseWorkflow(); lease.release(); throw error; } - releaseWorkflow(); if (!lease.isTransferred()) { lease.release(); } diff --git a/src/server/lifecycle.ts b/src/server/lifecycle.ts index fc757cc0670..e36bc460682 100644 --- a/src/server/lifecycle.ts +++ b/src/server/lifecycle.ts @@ -35,6 +35,7 @@ export const MAX_ACTIVE_SESSION_LANES = 64; export const SESSION_LANE_ID_BYTES = 32; const turnGate = createAdmissionGate("active_turns", MAX_ACTIVE_TURNS); export interface ActiveTurnLease extends AdmissionLease { + attach(lease: AdmissionLease): void; bindAbortController(ac: AbortController): void; beginCodexAccountSelection(): CodexAccountSelectionAdmission; isTransferred(): boolean; @@ -190,10 +191,15 @@ export function tryAdmitTurn(sessionLaneId?: string): ActiveTurnLease | null { } } const controllers = new Set(); + const attachedLeases = new Set(); let active = true; let transferred = false; let nativeMainClaimed = false; const lease: ActiveTurnLease = { + attach(attachedLease) { + if (!active) attachedLease.release(); + else attachedLeases.add(attachedLease); + }, bindAbortController(ac) { knownTurnControllers.add(ac); if (!active) { @@ -238,6 +244,8 @@ export function tryAdmitTurn(sessionLaneId?: string): ActiveTurnLease | null { if (activeTurns.get(controller) === lease) activeTurns.delete(controller); } controllers.clear(); + for (const attachedLease of attachedLeases) attachedLease.release(); + attachedLeases.clear(); nativeMainTurns.delete(lease); if (opaqueSessionLaneId) { const currentRefCount = activeSessionLaneRefCounts.get(opaqueSessionLaneId); diff --git a/structure/transports/responses.md b/structure/transports/responses.md index a82e8f8cb08..17205cc8e22 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -69,6 +69,8 @@ repository state. Consequently, the active-turn and session-lane gates are concu limits, the translator budget is a live retained-byte limit, the response-state caps are cache retention limits, and the stall watchdog is a silence limit. None is a cumulative continuation or semantic no-progress budget. +Active-turn admission owns workflow admission, so both remain held until a streaming body finishes +or is cancelled. > Decision record: [ADR-0031](../decisions/ADR-0031-responses-http-sse.md) diff --git a/tests/codex-integration/active-registry-admission.test.ts b/tests/codex-integration/active-registry-admission.test.ts index 88bbeb5790c..16b375b9872 100644 --- a/tests/codex-integration/active-registry-admission.test.ts +++ b/tests/codex-integration/active-registry-admission.test.ts @@ -3,6 +3,7 @@ import { mkdtempSync} from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { MAX_ACTIVE_TURNS, abortAndReleaseAllTurns, activeRegistryMetrics, trackStreamLifetime, tryAdmitTurn, unregisterTurn } from "../../src/server/lifecycle"; +import { workflowBudgetSnapshot } from "../../src/lib/workflow-budget"; import { MAX_TRACKED_CODEX_WEBSOCKETS, getTrackedCodexWebSocketCountForAccount, @@ -125,14 +126,20 @@ describe("active registry admission", () => { try { const response = await fetch(new URL("/v1/responses", server.url), { method: "POST", - headers: { "content-type": "application/json" }, + headers: { + "content-type": "application/json", + "x-codex-parent-thread-id": "stream-root", + "thread-id": "stream-child", + }, body: JSON.stringify({ model: "fixture/model", input: "hello", stream: true }), }); expect(response.status).toBe(200); expect(activeRegistryMetrics().activeTurns.active).toBe(before + 1); + expect(workflowBudgetSnapshot("stream-root")?.active).toBe(1); settle(); expect(await response.text()).toBe("chunk"); expect(activeRegistryMetrics().activeTurns.active).toBe(before); + expect(workflowBudgetSnapshot("stream-root")?.active).toBe(0); } finally { settle?.(); await server.stop(true); @@ -177,6 +184,21 @@ describe("active registry admission", () => { expect(activeRegistryMetrics().activeTurns.active).toBe(before); }); + test("a transferred turn retains attached admission until its stream settles", async () => { + const lease = tryAdmitTurn()!; + let attachedReleases = 0; + lease.attach({ release() { attachedReleases += 1; } }); + const source = new ReadableStream({ pull() {} }); + const tracked = trackStreamLifetime(source, new AbortController(), undefined, lease); + + expect(lease.isTransferred()).toBe(true); + expect(attachedReleases).toBe(0); + await tracked.cancel(); + expect(attachedReleases).toBe(1); + lease.release(); + expect(attachedReleases).toBe(1); + }); + test("the 256-turn gate bounds concurrency, not sequential continuation count", () => { const before = activeRegistryMetrics().activeTurns.active; for (let index = 0; index <= MAX_ACTIVE_TURNS; index += 1) {