From ed297e65049f9117f38752a3095ffc353ad80c93 Mon Sep 17 00:00:00 2001 From: luvs01 Date: Sat, 8 Aug 2026 15:37:10 +0900 Subject: [PATCH] fix(responses): bound routed compaction responses --- src/adapters/openai-responses.ts | 184 ++++++++++++++++----- tests/responses-compaction-routing.test.ts | 32 ++++ 2 files changed, 171 insertions(+), 45 deletions(-) diff --git a/src/adapters/openai-responses.ts b/src/adapters/openai-responses.ts index 5ed43ba381..90340950f2 100644 --- a/src/adapters/openai-responses.ts +++ b/src/adapters/openai-responses.ts @@ -11,6 +11,83 @@ import { OCX_REASONING_PREFIX } from "../responses/reasoning-envelope"; import { modelRecordValue } from "../reasoning-effort"; import type { TranslatorBudget } from "../lib/translator-budget"; +export const ROUTED_COMPACTION_RESPONSE_MAX_BYTES = 32 * 1024 * 1024; +const ROUTED_COMPACTION_RESPONSE_TOO_LARGE = "upstream compaction response exceeded 32 MiB"; + +function declaredResponseBytes(response: Response): number | null { + const value = response.headers.get("content-length"); + if (value === null || value.trim() === "") return null; + const parsed = Number(value); + return Number.isFinite(parsed) && parsed >= 0 ? parsed : null; +} + +async function readRoutedCompactionBody(response: Response): Promise { + if (!response.body) throw new Error("malformed upstream compaction response"); + const reader = response.body.getReader(); + if ((declaredResponseBytes(response) ?? 0) > ROUTED_COMPACTION_RESPONSE_MAX_BYTES) { + await reader.cancel(ROUTED_COMPACTION_RESPONSE_TOO_LARGE).catch(() => undefined); + throw new Error(ROUTED_COMPACTION_RESPONSE_TOO_LARGE); + } + const chunks: Uint8Array[] = []; + let total = 0; + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + total += value.byteLength; + if (total > ROUTED_COMPACTION_RESPONSE_MAX_BYTES) { + await reader.cancel(ROUTED_COMPACTION_RESPONSE_TOO_LARGE).catch(() => undefined); + throw new Error(ROUTED_COMPACTION_RESPONSE_TOO_LARGE); + } + chunks.push(value); + } + } catch (error) { + await reader.cancel(error).catch(() => undefined); + throw error; + } + const body = new Uint8Array(total); + let offset = 0; + for (const chunk of chunks) { + body.set(chunk, offset); + offset += chunk.byteLength; + } + return new TextDecoder().decode(body); +} + +async function boundedRoutedCompactionStream(response: Response): Promise> { + if (!response.body) throw new Error("passthrough adapter received no response body"); + const reader = response.body.getReader(); + if ((declaredResponseBytes(response) ?? 0) > ROUTED_COMPACTION_RESPONSE_MAX_BYTES) { + await reader.cancel(ROUTED_COMPACTION_RESPONSE_TOO_LARGE).catch(() => undefined); + throw new Error(ROUTED_COMPACTION_RESPONSE_TOO_LARGE); + } + let total = 0; + return new ReadableStream({ + async pull(controller) { + try { + const { done, value } = await reader.read(); + if (done) { + controller.close(); + return; + } + total += value.byteLength; + if (total > ROUTED_COMPACTION_RESPONSE_MAX_BYTES) { + await reader.cancel(ROUTED_COMPACTION_RESPONSE_TOO_LARGE).catch(() => undefined); + controller.error(new Error(ROUTED_COMPACTION_RESPONSE_TOO_LARGE)); + return; + } + controller.enqueue(value); + } catch (error) { + await reader.cancel(error).catch(() => undefined); + controller.error(error); + } + }, + async cancel(reason) { + await reader.cancel(reason).catch(() => undefined); + }, + }); +} + // Headers relayed verbatim from the caller in OAuth-passthrough ("forward") mode. // Exported so the web-search sidecar reuses the exact same forwarded-auth set for its ChatGPT call. export const FORWARD_HEADERS = [ @@ -1214,50 +1291,62 @@ export function createResponsesPassthroughAdapter(provider: OcxProviderConfig): let doneText = ""; let snapshot = ""; let usage: OcxUsage | undefined; - for await (const event of decodeServerSentEvents(response.body, { translatorBudget: budget })) { - let payload: unknown; - try { payload = JSON.parse(event.data); } catch { continue; } - if (!isPlainObject(payload)) continue; - switch (payload.type) { - case "response.output_text.delta": - if (typeof payload.delta === "string") { - const next = deltas + payload.delta; - const previousBytes = budgetEncoder.encode(deltas).byteLength; - const reservation = budget.reserveTransient(budgetEncoder.encode(next).byteLength, { kind: "retained_collectors" }); - deltas = next; - reservation.commitRetained(); - budget.releaseRetained(previousBytes, { kind: "retained_collectors" }); - } - break; - case "response.output_text.done": - if (typeof payload.text === "string") { - const next = doneText + payload.text; - const previousBytes = budgetEncoder.encode(doneText).byteLength; - const reservation = budget.reserveTransient(budgetEncoder.encode(next).byteLength, { kind: "retained_collectors" }); - doneText = next; - reservation.commitRetained(); - budget.releaseRetained(previousBytes, { kind: "retained_collectors" }); - } - break; - case "response.failed": - case "error": - yield { type: "error", message: responsesErrorMessage(payload.response ?? payload) }; - return; - case "response.incomplete": - yield { type: "incomplete", reason: responsesErrorMessage(payload.response ?? payload) }; - return; - case "response.completed": - { - const next = responsesPayloadText(payload.response); - const previousBytes = budgetEncoder.encode(snapshot).byteLength; - const reservation = budget.reserveTransient(budgetEncoder.encode(next).byteLength, { kind: "retained_collectors" }); - snapshot = next; - reservation.commitRetained(); - budget.releaseRetained(previousBytes, { kind: "retained_collectors" }); - } - usage = usageFromResponsesPayload(payload.response); - break; + let boundedBody: ReadableStream; + try { + boundedBody = await boundedRoutedCompactionStream(response); + } catch (error) { + yield { type: "error", message: error instanceof Error ? error.message : ROUTED_COMPACTION_RESPONSE_TOO_LARGE }; + return; + } + try { + for await (const event of decodeServerSentEvents(boundedBody, { translatorBudget: budget })) { + let payload: unknown; + try { payload = JSON.parse(event.data); } catch { continue; } + if (!isPlainObject(payload)) continue; + switch (payload.type) { + case "response.output_text.delta": + if (typeof payload.delta === "string") { + const next = deltas + payload.delta; + const previousBytes = budgetEncoder.encode(deltas).byteLength; + const reservation = budget.reserveTransient(budgetEncoder.encode(next).byteLength, { kind: "retained_collectors" }); + deltas = next; + reservation.commitRetained(); + budget.releaseRetained(previousBytes, { kind: "retained_collectors" }); + } + break; + case "response.output_text.done": + if (typeof payload.text === "string") { + const next = doneText + payload.text; + const previousBytes = budgetEncoder.encode(doneText).byteLength; + const reservation = budget.reserveTransient(budgetEncoder.encode(next).byteLength, { kind: "retained_collectors" }); + doneText = next; + reservation.commitRetained(); + budget.releaseRetained(previousBytes, { kind: "retained_collectors" }); + } + break; + case "response.failed": + case "error": + yield { type: "error", message: responsesErrorMessage(payload.response ?? payload) }; + return; + case "response.incomplete": + yield { type: "incomplete", reason: responsesErrorMessage(payload.response ?? payload) }; + return; + case "response.completed": + { + const next = responsesPayloadText(payload.response); + const previousBytes = budgetEncoder.encode(snapshot).byteLength; + const reservation = budget.reserveTransient(budgetEncoder.encode(next).byteLength, { kind: "retained_collectors" }); + snapshot = next; + reservation.commitRetained(); + budget.releaseRetained(previousBytes, { kind: "retained_collectors" }); + } + usage = usageFromResponsesPayload(payload.response); + break; + } } + } catch (error) { + yield { type: "error", message: error instanceof Error ? error.message : ROUTED_COMPACTION_RESPONSE_TOO_LARGE }; + return; } // Gateways differ in which of these they emit; prefer the authoritative // completed snapshot so text is never double-counted. @@ -1269,8 +1358,13 @@ export function createResponsesPassthroughAdapter(provider: OcxProviderConfig): async parseResponse(response: Response, budget: TranslatorBudget): Promise { let payload: unknown; - try { payload = await response.json(); } catch { - return [{ type: "error", message: "malformed upstream compaction response" }]; + try { payload = JSON.parse(await readRoutedCompactionBody(response)); } catch (error) { + return [{ + type: "error", + message: error instanceof Error && error.message === ROUTED_COMPACTION_RESPONSE_TOO_LARGE + ? error.message + : "malformed upstream compaction response", + }]; } budget.chargeRetained(new TextEncoder().encode(JSON.stringify(payload)).byteLength, { kind: "retained_collectors" }); if (!isPlainObject(payload)) { diff --git a/tests/responses-compaction-routing.test.ts b/tests/responses-compaction-routing.test.ts index 4aedbe110c..2ac0e1b9e1 100644 --- a/tests/responses-compaction-routing.test.ts +++ b/tests/responses-compaction-routing.test.ts @@ -482,6 +482,38 @@ describe("routed compaction for key-mode openai-responses (#422)", () => { }); describe("compaction terminal handling (#422)", () => { + for (const stream of [false, true]) { + test(`rejects and cancels an oversized ${stream ? "streaming" : "JSON"} upstream response`, async () => { + let cancelled = false; + const chunk = new Uint8Array(1024 * 1024); + const body = new ReadableStream({ + pull(controller) { + controller.enqueue(chunk); + }, + cancel() { + cancelled = true; + }, + }); + globalThis.fetch = (async () => new Response(body, { + status: 200, + headers: { + "content-type": stream ? "text/event-stream" : "application/json", + }, + })) as typeof fetch; + + const res = await handleResponses( + compactionRequest(baseCompactionBody({ stream })), + keyProviderConfig(), + { model: "", provider: "" }, + ); + const responseText = await res.text(); + + expect(cancelled).toBe(true); + expect(responseText).toContain("upstream compaction response exceeded 32 MiB"); + expect(responseText).not.toContain('"type":"compaction"'); + }); + } + test("an upstream failure does not become an empty compaction", async () => { globalThis.fetch = (async () => jsonResponse({ id: "resp_1",