From 4d334592c17dc450ebe07f5cc5ca4e2496656c49 Mon Sep 17 00:00:00 2001 From: metaphorics <152830360+metaphorics@users.noreply.github.com> Date: Wed, 29 Jul 2026 18:30:58 +0900 Subject: [PATCH] perf(ai): compress Codex SSE request bodies with zstd (cherry picked from commit 0b830d33eb9dcf6290db0f45cd1cd0995dbd7fe0) --- packages/ai/CHANGELOG.md | 4 + .../src/providers/openai-codex-responses.ts | 25 ++- .../test/openai-codex-responses-lite.test.ts | 12 +- packages/ai/test/openai-codex-stream.test.ts | 35 ++-- packages/ai/test/openai-codex-zstd.test.ts | 157 ++++++++++++++++++ 5 files changed, 216 insertions(+), 17 deletions(-) create mode 100644 packages/ai/test/openai-codex-zstd.test.ts diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 9c5eee224..f25ae4647 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Changed + +- Codex SSE request bodies are now zstd-compressed (level 3, `content-encoding: zstd`), matching the official Codex client; disable with `PI_CODEX_ZSTD=0`. Compression failure falls back to plain JSON. + ## [17.1.8] - 2026-07-28 ### Fixed diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 26684c1dd..6b2825af4 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -3906,6 +3906,24 @@ async function getOrCreateCodexWebSocketConnection( return state.connection; } +/** + * Compress an SSE request body with zstd. Returns `undefined` when + * compression is disabled or fails, in which case the caller sends the + * plain JSON string without a `content-encoding` header. + */ +function compressCodexRequestBody(bodyJson: string): Uint8Array | undefined { + if (!$flag("PI_CODEX_ZSTD", true)) return undefined; + try { + return Bun.zstdCompressSync(bodyJson, { level: 3 }); + } catch (error) { + CODEX_DEBUG && + logger.debug("[codex] codex request body compression failed", { + error: error instanceof Error ? error.message : String(error), + }); + return undefined; + } +} + async function openCodexSseEventStream( url: string, requestHeaders: Record | undefined, @@ -3956,12 +3974,17 @@ async function openCodexSseEventStream( clearPreResponseTimeout = undefined; } }; + const bodyJson = JSON.stringify(body); + const compressedBody = compressCodexRequestBody(bodyJson); + if (compressedBody !== undefined) { + headers.set("content-encoding", "zstd"); + } let response: Response; try { response = await fetchWithRetry(url, { method: "POST", headers, - body: JSON.stringify(body), + body: compressedBody ?? bodyJson, signal, prepareInit: () => { const watchdog = armPreResponseTimeout(signal, firstEventTimeoutMs); diff --git a/packages/ai/test/openai-codex-responses-lite.test.ts b/packages/ai/test/openai-codex-responses-lite.test.ts index 0b31a2aae..3f4ff7a74 100644 --- a/packages/ai/test/openai-codex-responses-lite.test.ts +++ b/packages/ai/test/openai-codex-responses-lite.test.ts @@ -91,6 +91,16 @@ function requireRecord(value: unknown, label: string): Record { return value; } +/** + * Decode a captured Codex SSE request body. The provider zstd-compresses the + * body by default, so a binary payload is decompressed before JSON parsing. + */ +function decodeCodexRequestBody(body: RequestInit["body"]): string { + if (typeof body === "string") return body; + if (body instanceof Uint8Array) return new TextDecoder().decode(Bun.zstdDecompressSync(body)); + throw new Error("expected a string or binary Codex request body"); +} + function parseTurnMetadata(clientMetadata: Record): Record { const encoded = clientMetadata["x-codex-turn-metadata"]; if (typeof encoded !== "string") throw new Error("expected x-codex-turn-metadata"); @@ -110,7 +120,7 @@ function createCodexFetchMock(sse: string, onRequest: (captured: CapturedCodexRe if (url.endsWith("/responses")) { onRequest({ headers: init?.headers instanceof Headers ? init.headers : new Headers(init?.headers), - body: typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : {}, + body: JSON.parse(decodeCodexRequestBody(init?.body)) as Record, }); return new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" } }); } diff --git a/packages/ai/test/openai-codex-stream.test.ts b/packages/ai/test/openai-codex-stream.test.ts index a61eeb3e8..7d7980044 100644 --- a/packages/ai/test/openai-codex-stream.test.ts +++ b/packages/ai/test/openai-codex-stream.test.ts @@ -101,6 +101,16 @@ function requireRecord(value: unknown, label: string): Record { return value; } +/** + * Decode a captured Codex SSE request body. The provider zstd-compresses the + * body by default, so a binary payload is decompressed before JSON parsing. + */ +function decodeCodexRequestBody(body: RequestInit["body"]): string { + if (typeof body === "string") return body; + if (body instanceof Uint8Array) return new TextDecoder().decode(Bun.zstdDecompressSync(body)); + throw new Error("expected a string or binary Codex request body"); +} + function parseTurnMetadata(clientMetadata: Record): Record { const encoded = clientMetadata["x-codex-turn-metadata"]; if (typeof encoded !== "string") throw new Error("expected x-codex-turn-metadata"); @@ -441,7 +451,7 @@ describe("openai-codex streaming", () => { const model = { ...createCodexTestModel("https://chatgpt.com/backend-api"), preferWebsockets: false }; let capturedBody: Record | undefined; const fetchMock = vi.fn(async (_input: string | URL, init?: RequestInit) => { - capturedBody = typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : undefined; + capturedBody = JSON.parse(decodeCodexRequestBody(init?.body)) as Record; return new Response(createCompletedCodexSse("Hello"), { status: 200, headers: { "content-type": "text/event-stream" }, @@ -471,10 +481,8 @@ describe("openai-codex streaming", () => { const model = { ...createCodexTestModel("https://chatgpt.com/backend-api"), preferWebsockets: false }; let capturedText: unknown; const fetchMock: FetchImpl = async (_input, init) => { - if (typeof init?.body === "string") { - const parsed: { text?: unknown } = JSON.parse(init.body); - capturedText = parsed.text; - } + const parsed: { text?: unknown } = JSON.parse(decodeCodexRequestBody(init?.body)); + capturedText = parsed.text; return new Response(createCompletedCodexSse("Hello"), { status: 200, headers: { "content-type": "text/event-stream" }, @@ -1800,7 +1808,7 @@ describe("openai-codex streaming", () => { `data: ${JSON.stringify({ type: "response.completed", response: { status: "completed", service_tier: "default", usage: { input_tokens: 5, output_tokens: 3, total_tokens: 8, input_tokens_details: { cached_tokens: 0 } } } })}`, ].join("\n\n")}\n\n`; const fetchMock = vi.fn(async (_input: string | URL, init?: RequestInit) => { - capturedBody = JSON.parse(String(init?.body)) as Record; + capturedBody = JSON.parse(decodeCodexRequestBody(init?.body)) as Record; return new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" }, @@ -2267,7 +2275,7 @@ describe("openai-codex streaming", () => { expect(headers?.get("x-client-request-id")).toBe(sessionId); // Verify sessionId is set in request body as prompt_cache_key - const body = typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : null; + const body = JSON.parse(decodeCodexRequestBody(init?.body)) as Record; expect(body?.prompt_cache_key).toBe(sessionId); return new Response(stream, { @@ -2324,8 +2332,7 @@ describe("openai-codex streaming", () => { } if (url === "https://chatgpt.com/backend-api/codex/responses") { capturedHeaders = init?.headers instanceof Headers ? init.headers : new Headers(init?.headers); - capturedBody = - typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : undefined; + capturedBody = JSON.parse(decodeCodexRequestBody(init?.body)) as Record; return new Response(createCompletedCodexSse("Hello"), { status: 200, headers: { "content-type": "text/event-stream" }, @@ -2362,8 +2369,7 @@ describe("openai-codex streaming", () => { const fetchMock = vi.fn(async (input: string | URL, init?: RequestInit) => { const url = typeof input === "string" ? input : input.toString(); if (url === "https://chatgpt.com/backend-api/codex/responses") { - capturedBody = - typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : undefined; + capturedBody = JSON.parse(decodeCodexRequestBody(init?.body)) as Record; return new Response(createCompletedCodexSse("Hello"), { status: 200, headers: { "content-type": "text/event-stream" }, @@ -2455,7 +2461,7 @@ describe("openai-codex streaming", () => { return new Response("PROMPT", { status: 200, headers: { etag: '"etag"' } }); } if (url === "https://chatgpt.com/backend-api/codex/responses") { - const body = typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : null; + const body = JSON.parse(decodeCodexRequestBody(init?.body)) as Record; expect(body?.reasoning).toEqual({ effort: "low", summary: "auto" }); return new Response(stream, { @@ -2766,8 +2772,7 @@ describe("openai-codex streaming", () => { continuationHeaders = init?.headers instanceof Headers ? init.headers : new Headers(init?.headers); expect(continuationHeaders.get("x-codex-turn-state")).toBe("ws-turn-state-1"); expect(continuationHeaders.get("x-models-etag")).toBe("models-etag-1"); - if (typeof init?.body !== "string") throw new Error("expected an SSE request body"); - const body: unknown = JSON.parse(init.body); + const body: unknown = JSON.parse(decodeCodexRequestBody(init?.body)); continuationRequest = requireRecord(body, "SSE continuation request"); return new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" } }); }); @@ -5053,7 +5058,7 @@ describe("openai-codex SSE statelessness", () => { function createCapturingFetch(sentRequests: Array>): FetchImpl { return vi.fn(async (_input: string | URL | Request, init?: RequestInit) => { - sentRequests.push(JSON.parse(String(init?.body)) as Record); + sentRequests.push(JSON.parse(decodeCodexRequestBody(init?.body)) as Record); return new Response(createStatefulCodexSse(`Answer ${sentRequests.length}`, `resp_${sentRequests.length}`), { status: 200, headers: { "content-type": "text/event-stream" }, diff --git a/packages/ai/test/openai-codex-zstd.test.ts b/packages/ai/test/openai-codex-zstd.test.ts new file mode 100644 index 000000000..372462b7e --- /dev/null +++ b/packages/ai/test/openai-codex-zstd.test.ts @@ -0,0 +1,157 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; +import { streamOpenAICodexResponses } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; +import type { Context, FetchImpl, Model } from "@oh-my-pi/pi-ai/types"; +import { __resetProxyCache } from "@oh-my-pi/pi-ai/utils/proxy"; +import { buildModel } from "@oh-my-pi/pi-catalog/build"; +import * as piUtils from "@oh-my-pi/pi-utils"; + +const { getAgentDir, setAgentDir, TempDir } = piUtils; + +const originalAgentDir = getAgentDir(); +const originalCodexZstd = Bun.env.PI_CODEX_ZSTD; +const originalProxyEnv: Record = { + PI_PROXY: Bun.env.PI_PROXY, + HTTPS_PROXY: Bun.env.HTTPS_PROXY, + https_proxy: Bun.env.https_proxy, + ALL_PROXY: Bun.env.ALL_PROXY, + all_proxy: Bun.env.all_proxy, + NO_PROXY: Bun.env.NO_PROXY, + no_proxy: Bun.env.no_proxy, +}; +const TEST_INSTALLATION_ID = "00000000-0000-4000-8000-000000000001"; + +function restoreEnv(name: string, value: string | undefined): void { + if (value === undefined) { + delete Bun.env[name]; + return; + } + Bun.env[name] = value; +} + +beforeEach(() => { + for (const key in originalProxyEnv) delete Bun.env[key]; + __resetProxyCache(); + vi.spyOn(piUtils, "getInstallId").mockReturnValue(TEST_INSTALLATION_ID); +}); + +afterEach(() => { + setAgentDir(originalAgentDir); + restoreEnv("PI_CODEX_ZSTD", originalCodexZstd); + for (const key in originalProxyEnv) restoreEnv(key, originalProxyEnv[key]); + __resetProxyCache(); + vi.restoreAllMocks(); +}); + +function createCodexTestToken(accountId = "acc_test"): string { + const payload = Buffer.from( + JSON.stringify({ "https://api.openai.com/auth": { chatgpt_account_id: accountId } }), + "utf8", + ).toBase64(); + return `aaa.${payload}.bbb`; +} + +function createCodexTestModel(): Model<"openai-codex-responses"> { + return buildModel({ + id: "gpt-5.3-codex-spark", + name: "GPT-5.3 Codex Spark", + api: "openai-codex-responses", + provider: "openai-codex", + baseUrl: "https://chatgpt.com/backend-api", + reasoning: true, + preferWebsockets: false, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 128000, + maxTokens: 128000, + }); +} + +function createCodexTestContext(): Context { + return { + systemPrompt: ["You are a helpful assistant."], + messages: [{ role: "user", content: "Say hello", timestamp: Date.now() }], + }; +} + +function createCompletedCodexSse(text: string): string { + return `${[ + `data: ${JSON.stringify({ type: "response.content_part.added", part: { type: "output_text", text: "" } })}`, + `data: ${JSON.stringify({ type: "response.output_text.delta", delta: text })}`, + `data: ${JSON.stringify({ type: "response.output_item.done", item: { type: "message", id: "msg_1", role: "assistant", status: "completed", content: [{ type: "output_text", text }] } })}`, + `data: ${JSON.stringify({ type: "response.completed", response: { status: "completed", usage: { input_tokens: 5, output_tokens: 3, total_tokens: 8, input_tokens_details: { cached_tokens: 0 } } } })}`, + ].join("\n\n")}\n\n`; +} + +// A fixed replacement payload pins the outgoing wire body so the serialized +// JSON is byte-deterministic across the compress/decompress round-trip. +const PINNED_PAYLOAD: Record = { + model: "gpt-5.3-codex-spark", + input: [{ role: "user", content: [{ type: "input_text", text: "Say hello" }] }], + stream: true, + prompt_cache_key: "zstd-test-cache-key", +}; + +interface CapturedRequest { + body: RequestInit["body"]; + headers: Headers; +} + +async function runAndCaptureRequest(): Promise { + const tempDir = TempDir.createSync("@pi-codex-zstd-"); + setAgentDir(tempDir.path()); + const token = createCodexTestToken(); + const model = createCodexTestModel(); + + let captured: CapturedRequest | undefined; + const fetchMock = vi.fn(async (_input: string | URL | Request, init?: RequestInit) => { + captured = { + body: init?.body, + headers: init?.headers instanceof Headers ? init.headers : new Headers(init?.headers), + }; + return new Response(createCompletedCodexSse("Hello"), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + }); + + const result = await streamOpenAICodexResponses(model, createCodexTestContext(), { + apiKey: token, + fetch: fetchMock as FetchImpl, + onPayload: async () => PINNED_PAYLOAD, + }).result(); + + expect(result.stopReason).toBe("stop"); + if (captured === undefined) throw new Error("expected the SSE request to reach fetch"); + return captured; +} + +describe("codex SSE request body zstd compression", () => { + it("compresses the request body with zstd and sets content-encoding by default", async () => { + delete Bun.env.PI_CODEX_ZSTD; + + const { body, headers } = await runAndCaptureRequest(); + + expect(headers.get("content-encoding")).toBe("zstd"); + expect(headers.get("content-type")).toContain("application/json"); + if (!(body instanceof Uint8Array)) throw new Error("expected a compressed binary body"); + // A zstd frame begins with the magic number 0xFD2FB528 (little-endian). + expect(body[0]).toBe(0x28); + expect(body[1]).toBe(0xb5); + expect(body[2]).toBe(0x2f); + expect(body[3]).toBe(0xfd); + + const decompressed = new TextDecoder().decode(Bun.zstdDecompressSync(body)); + expect(decompressed).toBe(JSON.stringify(PINNED_PAYLOAD)); + }); + + it("sends the plain JSON string without content-encoding when PI_CODEX_ZSTD=0", async () => { + Bun.env.PI_CODEX_ZSTD = "0"; + + const { body, headers } = await runAndCaptureRequest(); + + expect(headers.has("content-encoding")).toBe(false); + expect(headers.get("content-type")).toContain("application/json"); + expect(typeof body).toBe("string"); + expect(body).toBe(JSON.stringify(PINNED_PAYLOAD)); + }); +});