From 490ab265ba899dbb1da476bea6c8bf78a8284ad2 Mon Sep 17 00:00:00 2001 From: zengLingbiao Date: Sun, 19 Jul 2026 05:55:06 +0800 Subject: [PATCH] fix(agents): preserve sanitized stream cancellation (#110427) * fix(agents): add .catch() to reader.cancel() to prevent unhandled rejection * fix(agents): preserve sanitized stream cancellation Co-authored-by: zenglingbiao --------- Co-authored-by: Peter Steinberger --- src/agents/provider-transport-fetch.test.ts | 52 +++++++++++++++++++++ src/agents/provider-transport-fetch.ts | 25 ++++++---- 2 files changed, 69 insertions(+), 8 deletions(-) diff --git a/src/agents/provider-transport-fetch.test.ts b/src/agents/provider-transport-fetch.test.ts index 1b09bddda86b..11ac8a60b15f 100644 --- a/src/agents/provider-transport-fetch.test.ts +++ b/src/agents/provider-transport-fetch.test.ts @@ -1323,6 +1323,58 @@ describe("buildGuardedModelFetch", () => { expect(items).toEqual([{ ok: true }]); }); + it.each([ + { + name: "JSON-to-SSE synthesis", + contentType: "application/json", + body: '{"ok": true}', + }, + { + name: "SSE sanitization", + contentType: "text/event-stream", + body: 'data: {"ok": true}\n\n', + }, + ])("ignores source cancellation failures during $name", async ({ contentType, body }) => { + const cancel = vi.fn(async () => { + throw new Error("upstream cancellation failed"); + }); + const release = vi.fn(async () => undefined); + const encoder = new TextEncoder(); + fetchWithSsrFGuardMock.mockResolvedValue({ + response: new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode(body)); + }, + cancel, + }), + { headers: { "content-type": contentType } }, + ), + finalUrl: "https://openrouter.ai/api/v1/chat/completions", + release, + }); + const model = { + id: "gpt-5.4", + provider: "openrouter", + api: "openai-completions", + baseUrl: "https://openrouter.ai/api/v1", + } as unknown as Model<"openai-completions">; + + const response = await buildGuardedModelFetch(model)( + "https://openrouter.ai/api/v1/chat/completions", + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "gpt-5.4", stream: true }), + }, + ); + + expect(response.body).not.toBeNull(); + await expect(response.body!.cancel("consumer stopped")).resolves.toBeUndefined(); + expect(cancel).toHaveBeenCalledTimes(1); + expect(release).toHaveBeenCalledTimes(1); + }); + it("does not re-prefix SSE bodies mislabeled as JSON by streaming gateways", async () => { const source = openResponseStreamText( 'data: {"id":"a","choices":[{"index":0,"delta":{"content":"Hi","role":"assistant"}}]}\n\n' + diff --git a/src/agents/provider-transport-fetch.ts b/src/agents/provider-transport-fetch.ts index f458448f6566..91858922ee4b 100644 --- a/src/agents/provider-transport-fetch.ts +++ b/src/agents/provider-transport-fetch.ts @@ -96,6 +96,15 @@ function findSseEventBoundary(buffer: string): { index: number; length: number } return best; } +async function cancelReaderBestEffort( + reader: ReadableStreamDefaultReader | undefined, + reason?: unknown, +): Promise { + // Reader cancellation is cleanup. An upstream cancel failure must not replace + // the wrapper's authoritative stream error or downstream cancellation. + await reader?.cancel(reason).catch(() => undefined); +} + function capNonOkResponseBodyLazily(response: Response, maxBytes: number): Response { const source = response.body; if (!source) { @@ -123,18 +132,18 @@ function capNonOkResponseBodyLazily(response: Response, maxBytes: number): Respo } total = maxBytes; controller.close(); - void reader?.cancel().catch(() => undefined); + void cancelReaderBestEffort(reader); return; } total += chunk.value.byteLength; controller.enqueue(chunk.value); } catch (error) { controller.error(error); - void reader?.cancel(error).catch(() => undefined); + void cancelReaderBestEffort(reader, error); } }, async cancel(reason) { - await reader?.cancel(reason).catch(() => undefined); + await cancelReaderBestEffort(reader, reason); }, }); return new Response(capped, response); @@ -189,12 +198,12 @@ function sanitizeOpenAISdkSseResponse( buffer += decoder.decode(chunk.value, { stream: true }); } } catch (error) { - await reader?.cancel(error).catch(() => {}); + await cancelReaderBestEffort(reader, error); controller.error(error); } }, async cancel(reason) { - await reader?.cancel(reason); + await cancelReaderBestEffort(reader, reason); }, }); const headers = new Headers(response.headers); @@ -277,12 +286,12 @@ function sanitizeOpenAISdkSseResponse( } } } catch (error) { - await reader?.cancel(error).catch(() => {}); + await cancelReaderBestEffort(reader, error); controller.error(error); } }, async cancel(reason) { - await reader?.cancel(reason); + await cancelReaderBestEffort(reader, reason); }, }); @@ -361,7 +370,7 @@ async function classifyOpenAISdkStreamBody(response: Response): Promise undefined); + void cancelReaderBestEffort(reader); } }