Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
184 changes: 139 additions & 45 deletions src/adapters/openai-responses.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string> {
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<ReadableStream<Uint8Array>> {
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<Uint8Array>({
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 = [
Expand Down Expand Up @@ -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<Uint8Array>;
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.
Expand All @@ -1269,8 +1358,13 @@ export function createResponsesPassthroughAdapter(provider: OcxProviderConfig):

async parseResponse(response: Response, budget: TranslatorBudget): Promise<AdapterEvent[]> {
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)) {
Expand Down
32 changes: 32 additions & 0 deletions tests/responses-compaction-routing.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -482,6 +482,38 @@
});

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<Uint8Array>({
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",
Expand All @@ -508,7 +540,7 @@
output: [{ type: "message", role: "assistant", content: [{ type: "output_text", text: "partial" }] }],
})) as typeof fetch;

const res = await handleResponses(

Check failure on line 543 in tests/responses-compaction-routing.test.ts

View workflow job for this annotation

GitHub Actions / macos

error: expect(received).toContain(expected)

Expected to contain: "upstream compaction response exceeded 32 MiB" Received: "event: response.created\ndata: {\"type\":\"response.created\",\"sequence_number\":0,\"response\":{\"id\":\"resp_89c3ca7b68e348958ab275f7884c59f7\",\"object\":\"response\",\"created_at\":1786171383,\"status\":\"in_progress\",\"model\":\"some-model\",\"output\":[],\"usage\":null}}\n\nevent: response.failed\ndata: {\"type\":\"response.failed\",\"sequence_number\":1,\"response\":{\"id\":\"resp_89c3ca7b68e348958ab275f7884c59f7\",\"object\":\"response\",\"created_at\":1786171383,\"status\":\"failed\",\"model\":\"some-model\",\"output\":[],\"usage\":null,\"error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"},\"last_error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"}}}\n\ndata: [DONE]\n\n" at <anonymous> (/Users/runner/work/opencodex/opencodex/tests/responses-compaction-routing.test.ts:543:28)

Check failure on line 543 in tests/responses-compaction-routing.test.ts

View workflow job for this annotation

GitHub Actions / test 2/4

error: expect(received).toContain(expected)

Expected to contain: "upstream compaction response exceeded 32 MiB" Received: "event: response.created\ndata: {\"type\":\"response.created\",\"sequence_number\":0,\"response\":{\"id\":\"resp_b7c47e9d3b0a4695a53e32c066a4bb1c\",\"object\":\"response\",\"created_at\":1786171122,\"status\":\"in_progress\",\"model\":\"some-model\",\"output\":[],\"usage\":null}}\n\nevent: response.failed\ndata: {\"type\":\"response.failed\",\"sequence_number\":1,\"response\":{\"id\":\"resp_b7c47e9d3b0a4695a53e32c066a4bb1c\",\"object\":\"response\",\"created_at\":1786171122,\"status\":\"failed\",\"model\":\"some-model\",\"output\":[],\"usage\":null,\"error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"},\"last_error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"}}}\n\ndata: [DONE]\n\n" at <anonymous> (/home/runner/work/opencodex/opencodex/tests/responses-compaction-routing.test.ts:543:28)
compactionRequest(baseCompactionBody()),
keyProviderConfig(),
{ model: "", provider: "" },
Expand Down Expand Up @@ -550,7 +582,7 @@
{ model: "", provider: "" },
);

expect(await res.text()).toContain("\"type\":\"compaction\"");

Check failure on line 585 in tests/responses-compaction-routing.test.ts

View workflow job for this annotation

GitHub Actions / test 3/4

error: expect(received).toContain(expected)

Expected to contain: "upstream compaction response exceeded 32 MiB" Received: "event: response.created\ndata: {\"type\":\"response.created\",\"sequence_number\":0,\"response\":{\"id\":\"resp_9f248c79e8cb4ac0a4735d953a294494\",\"object\":\"response\",\"created_at\":1787130225,\"status\":\"in_progress\",\"model\":\"some-model\",\"output\":[],\"usage\":null}}\n\nevent: response.failed\ndata: {\"type\":\"response.failed\",\"sequence_number\":1,\"response\":{\"id\":\"resp_9f248c79e8cb4ac0a4735d953a294494\",\"object\":\"response\",\"created_at\":1787130225,\"status\":\"failed\",\"model\":\"some-model\",\"output\":[],\"usage\":null,\"error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"},\"last_error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"}}}\n\ndata: [DONE]\n\n" at <anonymous> (/home/runner/work/opencodex/opencodex/tests/responses-compaction-routing.test.ts:585:28)

Check failure on line 585 in tests/responses-compaction-routing.test.ts

View workflow job for this annotation

GitHub Actions / macos

error: expect(received).toContain(expected)

Expected to contain: "upstream compaction response exceeded 32 MiB" Received: "event: response.created\ndata: {\"type\":\"response.created\",\"sequence_number\":0,\"response\":{\"id\":\"resp_e3cef365869a4225b8a6d890d84e3a81\",\"object\":\"response\",\"created_at\":1787140847,\"status\":\"in_progress\",\"model\":\"some-model\",\"output\":[],\"usage\":null}}\n\nevent: response.failed\ndata: {\"type\":\"response.failed\",\"sequence_number\":1,\"response\":{\"id\":\"resp_e3cef365869a4225b8a6d890d84e3a81\",\"object\":\"response\",\"created_at\":1787140847,\"status\":\"failed\",\"model\":\"some-model\",\"output\":[],\"usage\":null,\"error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"},\"last_error\":{\"message\":\"translator live_transient buffer exceeded 33554432 bytes\",\"type\":\"server_error\",\"code\":\"upstream_server_error\"}}}\n\ndata: [DONE]\n\n" at <anonymous> (/Users/runner/work/opencodex/opencodex/tests/responses-compaction-routing.test.ts:585:28)
});
});

Expand Down
Loading