Skip to content

Commit 5b310a7

Browse files
authored
fix(agents): release abandoned provider streams
Fix streamed provider cleanup so abandoned managed fetch bodies no longer keep undici sockets open, and cancel Anthropic/Gemini SSE readers deterministically when parsing exits early. Keep the FinalizationRegistry abort path as a last-resort GC safety net for unmanaged/abandoned responses, while parser-owned paths cancel readers explicitly on thrown errors or malformed events. Also records the browser-only Control UI redactor alias in the optional deadcode allowlist and keeps mocked exec supervisor tests off shell snapshot wrapping after the branch was rebased onto default shell snapshots. Fixes #67461 Verification: - node scripts/run-vitest.mjs src/agents/provider-transport-fetch.test.ts src/agents/anthropic-transport-stream.test.ts extensions/google/transport-stream.test.ts src/agents/bash-tools.test.ts src/agents/bash-tools.exec.path.test.ts test/scripts/test-live-shard.test.ts - pnpm check:test-types - node scripts/run-oxlint-shards.mjs --threads=8 - .agents/skills/autoreview/scripts/autoreview --mode branch --base origin/main --parallel-tests "node scripts/run-vitest.mjs src/agents/provider-transport-fetch.test.ts src/agents/anthropic-transport-stream.test.ts extensions/google/transport-stream.test.ts src/agents/bash-tools.test.ts src/agents/bash-tools.exec.path.test.ts test/scripts/test-live-shard.test.ts" - git diff --check origin/main...HEAD - PR CI on a1db789 Co-authored-by: samzong <[email protected]> Signed-off-by: samzong <[email protected]>
1 parent 31c83c6 commit 5b310a7

11 files changed

Lines changed: 239 additions & 29 deletions

extensions/google/transport-stream.test.ts

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -110,6 +110,22 @@ function buildRawSseResponse(sse: string): Response {
110110
});
111111
}
112112

113+
function buildOpenRawSseResponse(params: { sse: string; onCancel: () => void }): Response {
114+
const encoder = new TextEncoder();
115+
const body = new ReadableStream<Uint8Array>({
116+
start(controller) {
117+
controller.enqueue(encoder.encode(params.sse));
118+
},
119+
cancel() {
120+
params.onCancel();
121+
},
122+
});
123+
return new Response(body, {
124+
status: 200,
125+
headers: { "content-type": "text/event-stream" },
126+
});
127+
}
128+
113129
function buildDelayedSecondSseResponse(params: {
114130
first: unknown;
115131
second: unknown;
@@ -583,6 +599,37 @@ describe("google transport stream", () => {
583599
expect(result.errorMessage).toBe("Google SSE stream returned malformed JSON");
584600
});
585601

602+
it("cancels open Gemini SSE bodies when parsing fails", async () => {
603+
let cancelCalled = false;
604+
guardedFetchMock.mockResolvedValueOnce(
605+
buildOpenRawSseResponse({
606+
sse: "data: {not json\n\n",
607+
onCancel: () => {
608+
cancelCalled = true;
609+
},
610+
}),
611+
);
612+
613+
const streamFn = createGoogleGenerativeAiTransportStreamFn();
614+
const stream = await Promise.resolve(
615+
streamFn(
616+
buildGeminiModel(),
617+
{
618+
messages: [{ role: "user", content: "hello", timestamp: 0 }],
619+
} as unknown as Parameters<typeof streamFn>[1],
620+
{
621+
apiKey: "gemini-api-key",
622+
} as Parameters<typeof streamFn>[2],
623+
),
624+
);
625+
626+
const result = await stream.result();
627+
628+
expect(result.stopReason).toBe("error");
629+
expect(result.errorMessage).toBe("Google SSE stream returned malformed JSON");
630+
expect(cancelCalled).toBe(true);
631+
});
632+
586633
it("retries Gemini 3 requests with lean thinking when the first attempt has no first response", async () => {
587634
vi.stubEnv("OPENCLAW_GOOGLE_GEMINI_FIRST_RESPONSE_RETRY_MS", "10");
588635
guardedFetchMock

extensions/google/transport-stream.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1098,6 +1098,7 @@ async function* parseGoogleSseChunks(
10981098
const reader = response.body.getReader();
10991099
const decoder = new TextDecoder();
11001100
let buffer = "";
1101+
let completed = false;
11011102
const abortHandler = () => {
11021103
void reader.cancel().catch(() => undefined);
11031104
};
@@ -1109,6 +1110,7 @@ async function* parseGoogleSseChunks(
11091110
}
11101111
const { done, value } = await reader.read();
11111112
if (done) {
1113+
completed = true;
11121114
break;
11131115
}
11141116
buffer += decoder.decode(value, { stream: true }).replace(/\r/g, "");
@@ -1134,6 +1136,10 @@ async function* parseGoogleSseChunks(
11341136
}
11351137
} finally {
11361138
signal?.removeEventListener("abort", abortHandler);
1139+
if (!completed) {
1140+
await reader.cancel(signal?.reason).catch(() => undefined);
1141+
}
1142+
reader.releaseLock();
11371143
}
11381144
}
11391145

scripts/deadcode-unused-files.allowlist.mjs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ export const KNIP_OPTIONAL_UNUSED_FILE_ALLOWLIST = [
4848
"src/plugins/contracts/tts-contract-suites.ts",
4949
"src/plugins/runtime-sidecar-paths-baseline.ts",
5050
"src/tasks/task-registry-control.runtime.ts",
51+
"ui/src/ui/browser-redact.ts",
5152
"extensions/qa-lab/src/auth-profile.fixture.ts",
5253
"extensions/qa-lab/src/codex-plugin.fixture.ts",
5354
];

src/agents/anthropic-transport-stream.test.ts

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,25 @@ function createRawSseResponse(body: string): Response {
5555
});
5656
}
5757

58+
function createOpenRawSseResponse(params: {
59+
body: string;
60+
onCancel: (reason: unknown) => void;
61+
}): Response {
62+
const encoder = new TextEncoder();
63+
const stream = new ReadableStream<Uint8Array>({
64+
start(controller) {
65+
controller.enqueue(encoder.encode(params.body));
66+
},
67+
cancel(reason) {
68+
params.onCancel(reason);
69+
},
70+
});
71+
return new Response(stream, {
72+
status: 200,
73+
headers: { "content-type": "text/event-stream" },
74+
});
75+
}
76+
5877
function delay<T>(ms: number, value: T): Promise<T> {
5978
return new Promise((resolve) => {
6079
setTimeout(() => resolve(value), ms);
@@ -1861,6 +1880,28 @@ describe("anthropic transport stream", () => {
18611880
expect(cancelReason).toBe(abortReason);
18621881
});
18631882

1883+
it("cancels open SSE bodies when Anthropic stream consumers throw", async () => {
1884+
let cancelCalled = false;
1885+
guardedFetchMock.mockResolvedValueOnce(
1886+
createOpenRawSseResponse({
1887+
body: 'data: {"type":"error","error":{"message":"stream exploded"}}\n\n',
1888+
onCancel: () => {
1889+
cancelCalled = true;
1890+
},
1891+
}),
1892+
);
1893+
1894+
const result = await runTransportStream(
1895+
makeAnthropicTransportModel(),
1896+
{ messages: [{ role: "user", content: "hello" }] } as AnthropicStreamContext,
1897+
{ apiKey: "sk-ant-api" } as AnthropicStreamOptions,
1898+
);
1899+
1900+
expect(result.stopReason).toBe("error");
1901+
expect(result.errorMessage).toBe("stream exploded");
1902+
expect(cancelCalled).toBe(true);
1903+
});
1904+
18641905
it("maps adaptive thinking effort for Claude 4.6 transport runs", async () => {
18651906
const model = makeAnthropicTransportModel({
18661907
id: "claude-opus-4-6",

src/agents/anthropic-transport-stream.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -617,10 +617,12 @@ async function* parseAnthropicSseBody(
617617
const reader = body.getReader();
618618
const decoder = new TextDecoder();
619619
let buffer = "";
620+
let completed = false;
620621
try {
621622
while (true) {
622623
const { done, value } = await readAnthropicSseChunk(reader, signal);
623624
if (done) {
625+
completed = true;
624626
break;
625627
}
626628
buffer = `${buffer}${decoder.decode(value, { stream: true })}`.replaceAll("\r\n", "\n");
@@ -651,6 +653,9 @@ async function* parseAnthropicSseBody(
651653
}
652654
}
653655
} finally {
656+
if (!completed) {
657+
await reader.cancel(signal?.reason).catch(() => undefined);
658+
}
654659
reader.releaseLock();
655660
}
656661
}

src/agents/bash-tools.exec.path.test.ts

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -144,7 +144,8 @@ describe("exec PATH login shell merge", () => {
144144
});
145145

146146
beforeEach(() => {
147-
envSnapshot = captureEnv(["PATH", "SHELL"]);
147+
envSnapshot = captureEnv(["OPENCLAW_EXEC_SHELL_SNAPSHOT", "PATH", "SHELL"]);
148+
process.env.OPENCLAW_EXEC_SHELL_SNAPSHOT = "0";
148149
shellEnvMocks.getShellPathFromLoginShell.mockReset();
149150
shellEnvMocks.getShellPathFromLoginShell.mockReturnValue("/custom/bin:/opt/bin");
150151
shellEnvMocks.resolveShellEnvFallbackTimeoutMs.mockReset();
@@ -289,6 +290,17 @@ describe("exec PATH login shell merge", () => {
289290
});
290291

291292
describe("exec host env validation", () => {
293+
let envSnapshot: ReturnType<typeof captureEnv>;
294+
295+
beforeEach(() => {
296+
envSnapshot = captureEnv(["OPENCLAW_EXEC_SHELL_SNAPSHOT"]);
297+
process.env.OPENCLAW_EXEC_SHELL_SNAPSHOT = "0";
298+
});
299+
300+
afterEach(() => {
301+
envSnapshot.restore();
302+
});
303+
292304
it("blocks LD_/DYLD_ env vars on host execution", async () => {
293305
const tool = createExecTool({ host: "gateway", security: "full", ask: "off" });
294306

src/agents/bash-tools.test.ts

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -217,8 +217,8 @@ const NOTIFY_POLL_OPTIONS = {
217217
timeout: NOTIFY_EVENT_TIMEOUT_MS,
218218
interval: POLL_INTERVAL_MS,
219219
};
220-
const SHELL_ENV_KEYS = ["SHELL"] as const;
221-
const PATH_SHELL_ENV_KEYS = ["PATH", "SHELL"] as const;
220+
const SHELL_ENV_KEYS = ["OPENCLAW_EXEC_SHELL_SNAPSHOT", "SHELL"] as const;
221+
const PATH_SHELL_ENV_KEYS = ["OPENCLAW_EXEC_SHELL_SNAPSHOT", "PATH", "SHELL"] as const;
222222
const PROCESS_STATUS_RUNNING = "running";
223223
const PROCESS_STATUS_COMPLETED = "completed";
224224
const PROCESS_STATUS_FAILED = "failed";
@@ -348,6 +348,7 @@ async function pollProcessSession(params: {
348348
};
349349
}
350350
function applyDefaultShellEnv() {
351+
process.env.OPENCLAW_EXEC_SHELL_SNAPSHOT = "0";
351352
if (!isWin && defaultShell) {
352353
process.env.SHELL = defaultShell;
353354
}
@@ -771,6 +772,8 @@ describe("exec exit codes", () => {
771772
});
772773

773774
describe("exec notifyOnExit", () => {
775+
useCapturedEnv([...SHELL_ENV_KEYS], applyDefaultShellEnv);
776+
774777
beforeEach(() => {
775778
resetHeartbeatWakeStateForTests();
776779
});

src/agents/provider-transport-fetch.test.ts

Lines changed: 96 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -20,27 +20,58 @@ const {
2020
resolveProviderRequestPolicyConfigMock,
2121
shouldUseEnvHttpProxyForUrlMock,
2222
withTrustedEnvProxyGuardedFetchModeMock,
23-
} = vi.hoisted(() => ({
24-
buildProviderRequestDispatcherPolicyMock: vi.fn<
25-
(_request?: unknown) => { mode: "direct" } | undefined
26-
>(() => undefined),
27-
fetchWithSsrFGuardMock: vi.fn(),
28-
ensureModelProviderLocalServiceMock: vi.fn(),
29-
mergeModelProviderRequestOverridesMock: vi.fn((current, overrides) => ({
30-
...current,
31-
...overrides,
32-
})),
33-
resolveProviderRequestPolicyConfigMock: vi.fn<() => ProviderRequestPolicyConfigMockResult>(
34-
() => ({
35-
allowPrivateNetwork: false,
36-
}),
37-
),
38-
shouldUseEnvHttpProxyForUrlMock: vi.fn(() => false),
39-
withTrustedEnvProxyGuardedFetchModeMock: vi.fn((params: Record<string, unknown>) => ({
40-
...params,
41-
mode: "trusted_env_proxy",
42-
})),
43-
}));
23+
managedStreamCleanupRegistrations,
24+
} = vi.hoisted(() => {
25+
const managedStreamCleanupRegistrations: Array<{
26+
callback: (held: { finalize: () => Promise<void> }) => void;
27+
held: { finalize: () => Promise<void> };
28+
token: object;
29+
}> = [];
30+
31+
class MockFinalizationRegistry {
32+
constructor(private callback: (held: { finalize: () => Promise<void> }) => void) {}
33+
34+
register(_target: object, held: { finalize: () => Promise<void> }, token?: object) {
35+
managedStreamCleanupRegistrations.push({
36+
callback: this.callback,
37+
held,
38+
token: token ?? {},
39+
});
40+
}
41+
42+
unregister(token: object) {
43+
const index = managedStreamCleanupRegistrations.findIndex((entry) => entry.token === token);
44+
if (index >= 0) {
45+
managedStreamCleanupRegistrations.splice(index, 1);
46+
}
47+
}
48+
}
49+
50+
vi.stubGlobal("FinalizationRegistry", MockFinalizationRegistry);
51+
52+
return {
53+
buildProviderRequestDispatcherPolicyMock: vi.fn<
54+
(_request?: unknown) => { mode: "direct" } | undefined
55+
>(() => undefined),
56+
fetchWithSsrFGuardMock: vi.fn(),
57+
ensureModelProviderLocalServiceMock: vi.fn(),
58+
mergeModelProviderRequestOverridesMock: vi.fn((current, overrides) => ({
59+
...current,
60+
...overrides,
61+
})),
62+
resolveProviderRequestPolicyConfigMock: vi.fn<() => ProviderRequestPolicyConfigMockResult>(
63+
() => ({
64+
allowPrivateNetwork: false,
65+
}),
66+
),
67+
shouldUseEnvHttpProxyForUrlMock: vi.fn(() => false),
68+
withTrustedEnvProxyGuardedFetchModeMock: vi.fn((params: Record<string, unknown>) => ({
69+
...params,
70+
mode: "trusted_env_proxy",
71+
})),
72+
managedStreamCleanupRegistrations,
73+
};
74+
});
4475

4576
vi.mock("../infra/net/fetch-guard.js", () => ({
4677
fetchWithSsrFGuard: fetchWithSsrFGuardMock,
@@ -82,6 +113,7 @@ function latestTrustedEnvProxyParams(): Record<string, unknown> {
82113

83114
describe("buildGuardedModelFetch", () => {
84115
beforeEach(() => {
116+
managedStreamCleanupRegistrations.length = 0;
85117
fetchWithSsrFGuardMock.mockReset().mockResolvedValue({
86118
response: new Response("ok", { status: 200 }),
87119
finalUrl: "https://api.openai.com/v1/responses",
@@ -151,6 +183,49 @@ describe("buildGuardedModelFetch", () => {
151183
await vi.waitFor(() => expect(release).toHaveBeenCalledTimes(1));
152184
});
153185

186+
it("releases guarded fetch slots when streamed bodies are abandoned", async () => {
187+
const release = vi.fn(async () => undefined);
188+
const encoder = new TextEncoder();
189+
fetchWithSsrFGuardMock.mockResolvedValue({
190+
response: new Response(
191+
new ReadableStream<Uint8Array>({
192+
start(controller) {
193+
controller.enqueue(encoder.encode("chunk-1"));
194+
controller.enqueue(encoder.encode("chunk-2"));
195+
},
196+
}),
197+
{ status: 200 },
198+
),
199+
finalUrl: "https://api.anthropic.com/v1/messages",
200+
release,
201+
});
202+
const model = {
203+
id: "claude-sonnet-4-6",
204+
provider: "anthropic",
205+
api: "anthropic-messages",
206+
baseUrl: "https://api.anthropic.com",
207+
} as unknown as Model<"anthropic-messages">;
208+
209+
const fetcher = buildGuardedModelFetch(model, undefined, { sanitizeSse: false });
210+
let response = await fetcher("https://api.anthropic.com/v1/messages", {
211+
method: "POST",
212+
headers: { "content-type": "application/json" },
213+
body: '{"stream":true}',
214+
});
215+
const reader = response.body?.getReader();
216+
expect(reader).toBeDefined();
217+
const firstChunk = await reader?.read();
218+
expect(firstChunk?.done).toBe(false);
219+
220+
response = undefined as unknown as Response;
221+
const registration = managedStreamCleanupRegistrations.at(-1);
222+
expect(registration).toBeDefined();
223+
await registration?.held.finalize();
224+
225+
expect(release).toHaveBeenCalledTimes(1);
226+
expect(managedStreamCleanupRegistrations).toHaveLength(0);
227+
});
228+
154229
it("passes model request headers to local service health probes", async () => {
155230
const model = {
156231
id: "deepseek-v4-flash",

0 commit comments

Comments
 (0)