fix(gateway): keep streaming roots until delivery and cleanup finish (#117505)

Co-authored-by: Peter Steinberger <steipete@macos.shared>
This commit is contained in:
Peter Steinberger
2026-08-01 13:57:55 -07:00
committed by GitHub
parent a3f070472d
commit 4b2d78a2ef
4 changed files with 301 additions and 83 deletions

View File

@@ -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<string>();
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<void>((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<boolean>();
let serverAbortSignal: AbortSignal | undefined;
agentCommand.mockClear();
agentCommand.mockImplementationOnce(
(opts: unknown) =>
new Promise<undefined>((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(

View File

@@ -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,

View File

@@ -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<string>();
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<void>((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<boolean>();
let serverAbortSignal: AbortSignal | undefined;
agentCommand.mockClear();
agentCommand.mockImplementationOnce(
(opts: unknown) =>
new Promise<undefined>((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(

View File

@@ -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({