From 2cd236385b3084f7855afc3eb60c8ab349050a03 Mon Sep 17 00:00:00 2001 From: xingzhou Date: Fri, 17 Jul 2026 06:53:23 +0800 Subject: [PATCH] fix(reef): close provider streams after guard errors (#109196) --- extensions/reef/protocol/guard-adapters.ts | 2 + extensions/reef/protocol/guard.test.ts | 44 ++++++++++++++++------ 2 files changed, 35 insertions(+), 11 deletions(-) diff --git a/extensions/reef/protocol/guard-adapters.ts b/extensions/reef/protocol/guard-adapters.ts index d04e4e463fb1..a28773780742 100644 --- a/extensions/reef/protocol/guard-adapters.ts +++ b/extensions/reef/protocol/guard-adapters.ts @@ -58,6 +58,7 @@ export function createOpenAiGuard(options: AdapterOptions): GuardAdapter { }), }); if (!response.ok) { + await response.body?.cancel().catch(() => undefined); throw new Error(`guard HTTP ${response.status}`); } const envelope = await parseJsonResponse(response); @@ -113,6 +114,7 @@ export function createAnthropicGuard(options: AdapterOptions): GuardAdapter { }), }); if (!response.ok) { + await response.body?.cancel().catch(() => undefined); throw new Error(`guard HTTP ${response.status}`); } const envelope = await parseJsonResponse(response); diff --git a/extensions/reef/protocol/guard.test.ts b/extensions/reef/protocol/guard.test.ts index 8d9ce0d267c8..fe2210d37a73 100644 --- a/extensions/reef/protocol/guard.test.ts +++ b/extensions/reef/protocol/guard.test.ts @@ -169,17 +169,39 @@ describe("provider adapters", () => { }); }); - it("fails closed on non-200 provider responses", async () => { - const guard = createOpenAiGuard({ - apiKey: "test", - pinnedModel: model, - fetch: async () => jsonResponse({ error: "no" }, 500), - }); - await expect(guard.classify(request)).resolves.toMatchObject({ - decision: "deny", - category: "guard_failure", - }); - }); + it.each([ + [ + "OpenAI", + (fetch: FetchLike) => createOpenAiGuard({ apiKey: "test", pinnedModel: model, fetch }), + ], + [ + "Anthropic", + (fetch: FetchLike) => createAnthropicGuard({ apiKey: "test", pinnedModel: model, fetch }), + ], + ])( + "cancels %s non-200 provider response bodies before failing closed", + async (_name, createGuard) => { + let cancelled = false; + const response = new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("partial error body")); + }, + cancel() { + cancelled = true; + }, + }), + { status: 503 }, + ); + const guard = createGuard(async () => response); + + await expect(guard.classify(request)).resolves.toMatchObject({ + decision: "deny", + category: "guard_failure", + }); + expect(cancelled).toBe(true); + }, + ); it("cancels oversized provider response streams before buffering them fully", async () => { const maxBytes = 256 * 1024;