From 4b2d78a2efcfec9009745d5fa2847e3971913d8a Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 1 Aug 2026 13:57:55 -0700 Subject: [PATCH] fix(gateway): keep streaming roots until delivery and cleanup finish (#117505) Co-authored-by: Peter Steinberger --- src/gateway/openai-http.test.ts | 187 ++++++++++++++++++------- src/gateway/openai-http.ts | 19 ++- src/gateway/openresponses-http.test.ts | 159 ++++++++++++++++++--- src/gateway/openresponses-http.ts | 19 ++- 4 files changed, 301 insertions(+), 83 deletions(-) diff --git a/src/gateway/openai-http.test.ts b/src/gateway/openai-http.test.ts index 7f6a3f089ad5..3df486a5e610 100644 --- a/src/gateway/openai-http.test.ts +++ b/src/gateway/openai-http.test.ts @@ -15,9 +15,13 @@ import { FailoverError } from "../agents/failover-error.js"; import { HISTORY_CONTEXT_MARKER } from "../auto-reply/reply/history.js"; import { CURRENT_MESSAGE_MARKER } from "../auto-reply/reply/mentions.js"; import { resetConfigRuntimeState } from "../config/config.js"; -import { emitAgentEvent } from "../infra/agent-events.js"; +import { emitAgentEvent, onAgentEvent } from "../infra/agent-events.js"; import { enqueueCommandInLane } from "../process/command-queue.js"; -import { getActiveGatewayRootWorkCount } from "../process/gateway-work-admission.js"; +import { + getActiveGatewayRootWorkCount, + isGatewaySubordinateWorkAdmissionClosed, + waitForActiveGatewayRootWork, +} from "../process/gateway-work-admission.js"; import { createDeferred } from "../test-utils/deferred.js"; import { withEnvAsync } from "../test-utils/env.js"; import { buildAssistantDeltaResult } from "./test-helpers.agent-results.js"; @@ -2168,40 +2172,99 @@ describe("OpenAI-compatible HTTP API (e2e)", () => { ])( "separates streamed content from the terminal finish for an official SDK $name", async ({ fail, expected }) => { - agentCommand.mockClear(); - if (fail) { - agentCommand.mockRejectedValueOnce(new Error("private upstream failure")); - } else { - agentCommand.mockImplementationOnce((async (opts: unknown) => - buildAssistantDeltaResult({ - opts, - emit: emitAgentEvent, - deltas: ["he", "llo"], - text: expected, - })) as never); - } - - const stream = await createOpenAiChatClient(enabledPort).chat.completions.create({ - model: "openclaw", - messages: [{ role: "user", content: "Return a complete streamed response." }], - stream: true, + const idleRootCount = getActiveGatewayRootWorkCount(); + const terminalAdmission = createDeferred<{ + active: number; + drained: { drained: boolean; active: number }; + }>(); + const wireResponse = createDeferred(); + const continueAgent = createDeferred(); + let activeRunId: string | undefined; + const unsubscribe = onAgentEvent((event) => { + if ( + event.runId !== activeRunId || + event.stream !== "lifecycle" || + event.data?.phase !== "end" + ) { + return; + } + // Agent settlement schedules the terminal in a later microtask. The + // request must remain admitted until Node finishes the actual stream. + const active = getActiveGatewayRootWorkCount(); + void waitForActiveGatewayRootWork(0).then((drained) => { + terminalAdmission.resolve({ active, drained }); + }); }); - const choices: Array<{ - delta: { content?: string | null }; - finish_reason: string | null; - }> = []; - for await (const chunk of stream) { - choices.push(...chunk.choices); + agentCommand.mockClear(); + agentCommand.mockImplementationOnce((async (opts: unknown) => { + activeRunId = (opts as { runId?: string }).runId; + if (!activeRunId) { + throw new Error("expected a streaming chat-completion run ID"); + } + await continueAgent.promise; + if (fail) { + throw new Error("private upstream failure"); + } + return buildAssistantDeltaResult({ + opts, + emit: emitAgentEvent, + deltas: ["he", "llo"], + text: expected, + }); + }) as never); + + try { + const client = new OpenAI({ + apiKey: "test", + baseURL: `http://127.0.0.1:${enabledPort}/v1`, + defaultHeaders: { "x-openclaw-scopes": "operator.write" }, + maxRetries: 0, + fetch: async (input, init) => { + const response = await fetch(input, init); + void response.clone().text().then(wireResponse.resolve, wireResponse.reject); + return response; + }, + }); + const stream = await client.chat.completions.create({ + model: "openclaw", + messages: [{ role: "user", content: "Return a complete streamed response." }], + stream: true, + }); + await vi.waitFor(() => expect(agentCommand).toHaveBeenCalledTimes(1)); + await new Promise((resolve) => { + setImmediate(resolve); + }); + continueAgent.resolve(); + + const choices: Array<{ + delta: { content?: string | null }; + finish_reason: string | null; + }> = []; + for await (const chunk of stream) { + choices.push(...chunk.choices); + } + + const [admission, wire] = await Promise.all([ + terminalAdmission.promise, + wireResponse.promise, + ]); + expect(admission.active).toBe(idleRootCount + 1); + expect(admission.drained).toEqual({ drained: false, active: idleRootCount + 1 }); + expect(parseSseDataLines(wire).at(-1)).toBe("[DONE]"); + + const contentChoices = choices.filter((choice) => typeof choice.delta.content === "string"); + expect(contentChoices.map((choice) => choice.delta.content).join("")).toBe(expected); + expect(contentChoices.every((choice) => choice.finish_reason === null)).toBe(true); + + const terminalChoices = choices.filter((choice) => choice.finish_reason === "stop"); + expect(terminalChoices).toHaveLength(1); + expect(terminalChoices[0]?.delta).toEqual({}); + expect(choices.at(-1)).toEqual(terminalChoices[0]); + await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(idleRootCount)); + } finally { + continueAgent.resolve(); + unsubscribe(); } - - const contentChoices = choices.filter((choice) => typeof choice.delta.content === "string"); - expect(contentChoices.map((choice) => choice.delta.content).join("")).toBe(expected); - expect(contentChoices.every((choice) => choice.finish_reason === null)).toBe(true); - - const terminalChoices = choices.filter((choice) => choice.finish_reason === "stop"); - expect(terminalChoices).toHaveLength(1); - expect(terminalChoices[0]?.delta).toEqual({}); - expect(choices.at(-1)).toEqual(terminalChoices[0]); }, ); @@ -3239,21 +3302,26 @@ describe("OpenAI-compatible HTTP API (e2e)", () => { it("aborts agent command when streaming client disconnects", { timeout: 15_000 }, async () => { const port = enabledPort; + const idleRootCount = getActiveGatewayRootWorkCount(); + const agentAborted = createDeferred(); + const finishAgentCleanup = createDeferred(); + const cleanupAdmissionClosed = createDeferred(); let serverAbortSignal: AbortSignal | undefined; agentCommand.mockClear(); - agentCommand.mockImplementationOnce( - (opts: unknown) => - new Promise((resolve) => { - const signal = (opts as { abortSignal?: AbortSignal } | undefined)?.abortSignal; - serverAbortSignal = signal; - if (signal?.aborted) { - resolve(undefined); - return; - } - signal?.addEventListener("abort", () => resolve(undefined), { once: true }); - }), - ); + agentCommand.mockImplementationOnce(async (opts: unknown) => { + const signal = (opts as { abortSignal?: AbortSignal } | undefined)?.abortSignal; + serverAbortSignal = signal; + if (signal?.aborted) { + agentAborted.resolve(); + } else { + signal?.addEventListener("abort", () => agentAborted.resolve(), { once: true }); + } + await agentAborted.promise; + cleanupAdmissionClosed.resolve(isGatewaySubordinateWorkAdmissionClosed()); + await finishAgentCleanup.promise; + return undefined; + }); const clientReq = http.request({ hostname: "127.0.0.1", @@ -3278,14 +3346,27 @@ describe("OpenAI-compatible HTTP API (e2e)", () => { expect(agentCommand).toHaveBeenCalledTimes(1); }); - clientReq.destroy(); + try { + clientReq.destroy(); - await vi.waitFor( - () => { - expect(serverAbortSignal?.aborted).toBe(true); - }, - { timeout: 5_000, interval: 50 }, - ); + await vi.waitFor(() => expect(serverAbortSignal?.aborted).toBe(true), { + timeout: 5_000, + interval: 50, + }); + expect(await cleanupAdmissionClosed.promise).toBe(false); + expect(getActiveGatewayRootWorkCount()).toBe(idleRootCount + 1); + expect(await waitForActiveGatewayRootWork(0)).toEqual({ + drained: false, + active: idleRootCount + 1, + }); + } finally { + finishAgentCleanup.resolve(); + } + + await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(idleRootCount), { + timeout: 5_000, + interval: 50, + }); }); it( diff --git a/src/gateway/openai-http.ts b/src/gateway/openai-http.ts index b42d3213965e..9699ea33bfe4 100644 --- a/src/gateway/openai-http.ts +++ b/src/gateway/openai-http.ts @@ -1305,18 +1305,27 @@ export async function handleOpenAiHttpRequest( } }); + // Agent cleanup and deferred SSE delivery have independent lifetimes; + // shutdown must wait until both have settled, whichever finishes last. + const releaseAgentRootWork = retainGatewayRootWorkAdmissionContinuation(); + const releaseResponseRootWork = retainGatewayRootWorkAdmissionContinuation(); + const releaseStreamRootWork = () => { + res.off("finish", releaseStreamRootWork); + res.off("close", releaseStreamRootWork); + releaseResponseRootWork?.(); + }; + res.once("finish", releaseStreamRootWork); + res.once("close", releaseStreamRootWork); + stopWatchingDisconnect = watchClientDisconnect(req, res, abortController, () => { closed = true; unsubscribe(); + releaseStreamRootWork(); }); wroteRole = true; writeAssistantRoleChunk(res, { runId, model }); - // The streamed run outlives this handler, whose root-work admission is - // released on return. Without retaining it, subordinate session/lane - // admissions inherit a released lease and fail as GatewayDrainingError. - const releaseRootWork = retainGatewayRootWorkAdmissionContinuation(); void (async () => { try { const result = await agentCommandFromIngress(commandInput, defaultRuntime, deps); @@ -1447,7 +1456,7 @@ export async function handleOpenAiHttpRequest( }); requestFinalize(); } finally { - releaseRootWork?.(); + releaseAgentRootWork?.(); if (!closed) { emitAgentEvent({ runId, diff --git a/src/gateway/openresponses-http.test.ts b/src/gateway/openresponses-http.test.ts index 9d7aad072a13..ce908fc7a1e4 100644 --- a/src/gateway/openresponses-http.test.ts +++ b/src/gateway/openresponses-http.test.ts @@ -12,7 +12,11 @@ import { CURRENT_MESSAGE_MARKER } from "../auto-reply/reply/mentions.js"; import { resetConfigRuntimeState } from "../config/config.js"; import { emitAgentEvent, onAgentEvent } from "../infra/agent-events.js"; import { enqueueCommandInLane } from "../process/command-queue.js"; -import { getActiveGatewayRootWorkCount } from "../process/gateway-work-admission.js"; +import { + getActiveGatewayRootWorkCount, + isGatewaySubordinateWorkAdmissionClosed, + waitForActiveGatewayRootWork, +} from "../process/gateway-work-admission.js"; import { createDeferred } from "../test-utils/deferred.js"; import { withEnvAsync } from "../test-utils/env.js"; import { IMAGE_ONLY_USER_MESSAGE } from "./agent-prompt.js"; @@ -1544,6 +1548,103 @@ describe("OpenResponses HTTP API (e2e)", () => { expect(agentCommand).toHaveBeenCalledTimes(1); }); + it.each([ + { name: "completed", failed: false }, + { name: "failed", failed: true }, + ])( + "keeps the $name response admitted until its deferred SSE terminal is written", + async ({ name, failed }) => { + const idleRootCount = getActiveGatewayRootWorkCount(); + const terminalAdmission = createDeferred<{ + active: number; + drained: { drained: boolean; active: number }; + }>(); + const wireResponse = createDeferred(); + const continueAgent = createDeferred(); + let activeRunId: string | undefined; + const unsubscribe = onAgentEvent((event) => { + if ( + event.runId !== activeRunId || + event.stream !== "lifecycle" || + event.data?.phase !== "end" + ) { + return; + } + // Restart drains run before the next-turn SSE finalizer, so inspect the real + // root owner at that boundary instead of observing eventual client delivery. + queueMicrotask(() => { + const active = getActiveGatewayRootWorkCount(); + void waitForActiveGatewayRootWork(0).then((drained) => { + terminalAdmission.resolve({ active, drained }); + }); + }); + }); + + agentCommand.mockClear(); + agentCommand.mockImplementationOnce((async (opts: unknown) => { + activeRunId = (opts as { runId?: string }).runId; + if (!activeRunId) { + throw new Error("expected a streaming response run ID"); + } + await continueAgent.promise; + emitAgentEvent({ runId: activeRunId, stream: "assistant", data: { delta: "answer" } }); + if (failed) { + emitAgentEvent({ + runId: activeRunId, + stream: "lifecycle", + data: { phase: "error", error: "provider request failed" }, + }); + } + return { + payloads: [{ text: "answer" }], + meta: { agentMeta: { usage: { input: 3, output: 2, total: 5 } } }, + }; + }) as never); + + try { + const client = new OpenAI({ + apiKey: "test", + baseURL: `http://127.0.0.1:${enabledPort}/v1`, + defaultHeaders: { "x-openclaw-scopes": "operator.write" }, + maxRetries: 0, + fetch: async (input, init) => { + const response = await fetch(input, init); + void response.clone().text().then(wireResponse.resolve, wireResponse.reject); + return response; + }, + }); + const stream = client.responses.stream({ + model: "openclaw", + input: "Keep the stream owned.", + }); + const terminalEvents: string[] = []; + stream.on("response.completed", () => terminalEvents.push("response.completed")); + stream.on("response.failed", () => terminalEvents.push("response.failed")); + const finalResponse = stream.finalResponse(); + await vi.waitFor(() => expect(agentCommand).toHaveBeenCalledTimes(1)); + await new Promise((resolve) => { + setImmediate(resolve); + }); + continueAgent.resolve(); + + const [admission, response, wire] = await Promise.all([ + terminalAdmission.promise, + finalResponse, + wireResponse.promise, + ]); + + expect(admission.active).toBe(idleRootCount + 1); + expect(admission.drained).toEqual({ drained: false, active: idleRootCount + 1 }); + expect(response.status).toBe(name); + expect(terminalEvents).toEqual([`response.${name}`]); + expect(parseSseEvents(wire).at(-1)?.data).toBe("[DONE]"); + await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(idleRootCount)); + } finally { + unsubscribe(); + } + }, + ); + it("maps provider format failures to OpenResponses 400 failed responses", async () => { const port = enabledPort; @@ -2807,21 +2908,26 @@ describe("OpenResponses HTTP API (e2e)", () => { it("aborts agent command when streaming client disconnects", { timeout: 15_000 }, async () => { const port = enabledPort; + const idleRootCount = getActiveGatewayRootWorkCount(); + const agentAborted = createDeferred(); + const finishAgentCleanup = createDeferred(); + const cleanupAdmissionClosed = createDeferred(); let serverAbortSignal: AbortSignal | undefined; agentCommand.mockClear(); - agentCommand.mockImplementationOnce( - (opts: unknown) => - new Promise((resolve) => { - const signal = (opts as { abortSignal?: AbortSignal } | undefined)?.abortSignal; - serverAbortSignal = signal; - if (signal?.aborted) { - resolve(undefined); - return; - } - signal?.addEventListener("abort", () => resolve(undefined), { once: true }); - }), - ); + agentCommand.mockImplementationOnce(async (opts: unknown) => { + const signal = (opts as { abortSignal?: AbortSignal } | undefined)?.abortSignal; + serverAbortSignal = signal; + if (signal?.aborted) { + agentAborted.resolve(); + } else { + signal?.addEventListener("abort", () => agentAborted.resolve(), { once: true }); + } + await agentAborted.promise; + cleanupAdmissionClosed.resolve(isGatewaySubordinateWorkAdmissionClosed()); + await finishAgentCleanup.promise; + return undefined; + }); const clientReq = http.request({ hostname: "127.0.0.1", @@ -2849,14 +2955,27 @@ describe("OpenResponses HTTP API (e2e)", () => { { timeout: 5_000, interval: 50 }, ); - clientReq.destroy(); + try { + clientReq.destroy(); - await vi.waitFor( - () => { - expect(serverAbortSignal?.aborted).toBe(true); - }, - { timeout: 5_000, interval: 50 }, - ); + await vi.waitFor(() => expect(serverAbortSignal?.aborted).toBe(true), { + timeout: 5_000, + interval: 50, + }); + expect(await cleanupAdmissionClosed.promise).toBe(false); + expect(getActiveGatewayRootWorkCount()).toBe(idleRootCount + 1); + expect(await waitForActiveGatewayRootWork(0)).toEqual({ + drained: false, + active: idleRootCount + 1, + }); + } finally { + finishAgentCleanup.resolve(); + } + + await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(idleRootCount), { + timeout: 5_000, + interval: 50, + }); }); it( diff --git a/src/gateway/openresponses-http.ts b/src/gateway/openresponses-http.ts index ef7782184cea..e9c27cec9658 100644 --- a/src/gateway/openresponses-http.ts +++ b/src/gateway/openresponses-http.ts @@ -1158,15 +1158,24 @@ export async function handleOpenResponsesHttpRequest( } }); + // Agent cleanup and deferred SSE delivery have independent lifetimes; + // shutdown must wait until both have settled, whichever finishes last. + const releaseAgentRootWork = retainGatewayRootWorkAdmissionContinuation(); + const releaseResponseRootWork = retainGatewayRootWorkAdmissionContinuation(); + const releaseStreamRootWork = () => { + res.off("finish", releaseStreamRootWork); + res.off("close", releaseStreamRootWork); + releaseResponseRootWork?.(); + }; + res.once("finish", releaseStreamRootWork); + res.once("close", releaseStreamRootWork); + stopWatchingDisconnect = watchClientDisconnect(req, res, abortController, () => { closed = true; unsubscribe(); + releaseStreamRootWork(); }); - // The streamed run outlives this handler, whose root-work admission is - // released on return. Without retaining it, subordinate session/lane - // admissions inherit a released lease and fail as GatewayDrainingError. - const releaseRootWork = retainGatewayRootWorkAdmissionContinuation(); void (async () => { try { const result = await runResponsesAgentCommand({ @@ -1419,7 +1428,7 @@ export async function handleOpenResponsesHttpRequest( data: { phase: "error" }, }); } finally { - releaseRootWork?.(); + releaseAgentRootWork?.(); if (!closed) { // Emit lifecycle end to trigger completion emitAgentEvent({