From b2703f87017fd36789bac6689b1ec86ca00d9950 Mon Sep 17 00:00:00 2001 From: JUN <243035832+lidge-jun@users.noreply.github.com> Date: Mon, 7 Sep 2026 07:37:03 +0900 Subject: [PATCH 1/2] fix(grok): filter Codex control frames for strict Responses clients Carry PR #3816 at d5e0a9a2069ac5655dc208c7261981cd8c59112d. Keep original frames on the proxy inspection branch and scope projection to the existing Grok HTTP/SSE client marker. Co-authored-by: Danh Thanh --- src/server/grok-responses-control-frame.ts | 38 +++++++++++ src/server/responses/core.ts | 10 ++- .../responses-snapshot-repair-server.test.ts | 64 ++++++++++++++++++- 3 files changed, 106 insertions(+), 6 deletions(-) create mode 100644 src/server/grok-responses-control-frame.ts diff --git a/src/server/grok-responses-control-frame.ts b/src/server/grok-responses-control-frame.ts new file mode 100644 index 0000000000..cb572acee9 --- /dev/null +++ b/src/server/grok-responses-control-frame.ts @@ -0,0 +1,38 @@ +import { sseDataPayload, type SseBlockRewrite } from "./sse-payload-rewrite"; + +const GROK_CONTROL_FRAME_TYPES: Record = { + "codex.rate_limits": true, + "codex.response.metadata": true, +}; + +/** + * Hide Codex-only control frames from Grok's strict Responses decoder. + * + * The inspection branch still sees these frames before this client-facing + * rewrite, so quota accounting and response metadata remain available to the + * proxy while Grok receives only its declared Responses event variants. + */ +export function createGrokResponsesControlFrameBlockRewrite(): SseBlockRewrite { + return (block) => { + const eventName = block + .split(/\r?\n/) + .find(line => line.startsWith("event:")) + ?.slice("event:".length) + .trim(); + if (GROK_CONTROL_FRAME_TYPES[eventName ?? ""] === true) return []; + + const payload = sseDataPayload(block); + if (payload === null || payload === "[DONE]") return [block]; + + let event: unknown; + try { + event = JSON.parse(payload); + } catch { + return [block]; + } + if (!event || typeof event !== "object" || Array.isArray(event) || !("type" in event)) return [block]; + return typeof event.type === "string" && GROK_CONTROL_FRAME_TYPES[event.type] === true + ? [] + : [block]; + }; +} diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 3c539c6d8e..e87136b67a 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -374,6 +374,7 @@ import { type UpstreamHostAdmissionLease, } from "../../codex/upstream-host-health"; import { createGrokResponsesSparseTerminalBlockRewrite } from "../grok-responses-snapshot-repair"; +import { createGrokResponsesControlFrameBlockRewrite } from "../grok-responses-control-frame"; import { createResponsesSnapshotBlockRewrite, hasResponsesSnapshotRepair, @@ -5494,9 +5495,9 @@ async function handleResponsesInner( // Grok Build renders deltas live but reconstructs its durable assistant // turn from the completed response snapshot. Native Responses streams // may instead carry the complete items in output_item.done, so the - // explicit Grok compatibility marker enables strict terminal-only repair. + // explicit Grok compatibility marker enables strict client compatibility rewrites. // The provider's broader snapshot/lifecycle repair remains opt-in. - const grokClientSnapshotRepairEnabled = logCtx.surface === "grok"; + const grokClientCompatibilityEnabled = logCtx.surface === "grok"; const snapshotRepairEnabled = hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair); const githubCopilotRepairEnabled = route.providerName === "github-copilot"; const responseModelRewrite = parsed._responseModelId !== undefined @@ -5545,7 +5546,10 @@ async function handleResponsesInner( githubCopilotRepairEnabled ? createGithubCopilotResponsesBlockRewrite(translatorBudget) : undefined, - grokClientSnapshotRepairEnabled + grokClientCompatibilityEnabled + ? createGrokResponsesControlFrameBlockRewrite() + : undefined, + grokClientCompatibilityEnabled ? createGrokResponsesSparseTerminalBlockRewrite(translatorBudget) : undefined, snapshotRepairEnabled diff --git a/tests/responses/responses-snapshot-repair-server.test.ts b/tests/responses/responses-snapshot-repair-server.test.ts index 214806b141..a7f3ec167c 100644 --- a/tests/responses/responses-snapshot-repair-server.test.ts +++ b/tests/responses/responses-snapshot-repair-server.test.ts @@ -58,12 +58,26 @@ const CODEX_SPARSE_TERMINAL_EVENTS = [ }, ]; -function sparseSseBody(events: readonly Record[] = SPARSE_EVENTS): ReadableStream { +const GROK_CONTROL_FRAME_EVENTS = [ + { + type: "codex.rate_limits", + rate_limits: { primary: { used_percent: 12, window_minutes: 60, reset_at: 123 } }, + }, + { type: "codex.response.metadata", headers: { "x-models-etag": "fixture" } }, + { type: "response.created", response: { id: "resp_control" } }, + { type: "response.completed", response: { id: "resp_control", status: "completed", output: [] } }, +]; + +function sparseSseBody( + events: readonly Record[] = SPARSE_EVENTS, + includeEventNames = false, +): ReadableStream { return new ReadableStream({ start(controller) { const encoder = new TextEncoder(); for (const event of events) { - controller.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`)); + const eventLine = includeEventNames ? `event: ${event.type}\n` : ""; + controller.enqueue(encoder.encode(`${eventLine}data: ${JSON.stringify(event)}\n\n`)); } controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.close(); @@ -74,6 +88,7 @@ function sparseSseBody(events: readonly Record[] = SPARSE_EVENT function stubSparseGateway( origin: string, events: readonly Record[] = SPARSE_EVENTS, + includeEventNames = false, ): void { globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => { const requestUrl = typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url; @@ -82,7 +97,7 @@ function stubSparseGateway( return Response.json({ data: [] }); } if (url.origin === origin && url.pathname.endsWith("/responses")) { - return new Response(sparseSseBody(events), { + return new Response(sparseSseBody(events, includeEventNames), { status: 200, headers: { "content-type": "text/event-stream" }, }); @@ -326,6 +341,49 @@ describe("responsesSnapshotRepair through /v1/responses", () => { await server.stop(true); } }); + test("the Grok marker filters Codex control frames at the client boundary", async () => { + const gateway = "https://grok-control-frame.example.test"; + stubSparseGateway(gateway, GROK_CONTROL_FRAME_EVENTS, true); + saveConfig({ + port: 0, + defaultProvider: "sparse", + providers: { + sparse: { + adapter: "openai-responses", + baseUrl: `${gateway}/v1`, + authMode: "key", + apiKey: "test-key", + }, + }, + } as OcxConfig); + + const server = startServer(0); + try { + const request = (grokMarker: boolean) => originalFetch(new URL("/v1/responses", server.url), { + method: "POST", + headers: { + "content-type": "application/json", + ...(grokMarker ? { "x-opencodex-grok": "1" } : {}), + }, + body: JSON.stringify({ model: "sparse-model", input: "hi", stream: true }), + }); + + const grokResponse = await request(true); + expect(grokResponse.status).toBe(200); + const grokText = await grokResponse.text(); + expect(grokText).not.toContain("codex.rate_limits"); + expect(grokText).not.toContain("codex.response.metadata"); + expect(grokText).toContain('"type":"response.completed"'); + + const ordinaryResponse = await request(false); + expect(ordinaryResponse.status).toBe(200); + const ordinaryText = await ordinaryResponse.text(); + expect(ordinaryText).toContain("codex.rate_limits"); + expect(ordinaryText).toContain("codex.response.metadata"); + } finally { + await server.stop(true); + } + }); }); test("sparse JSON completion inference precedes function repair in client output and replay", async () => { From 336c621a35b63e4dec5bff235f897f92adeae2c7 Mon Sep 17 00:00:00 2001 From: JUN <243035832+lidge-jun@users.noreply.github.com> Date: Mon, 7 Sep 2026 07:38:20 +0900 Subject: [PATCH 2/2] fix(grok): honor SSE event order and empty resets Use the last event field and preserve significant whitespace while retaining independent JSON control-type filtering. Add discriminator, order, reset, and preservation regressions to the existing Responses test file. Tests were authored but not run; validation is delegated to final combined remote CI. Co-authored-by: Danh Thanh --- src/server/grok-responses-control-frame.ts | 17 ++++-- .../responses-snapshot-repair-server.test.ts | 60 ++++++++++++++++++- 2 files changed, 69 insertions(+), 8 deletions(-) diff --git a/src/server/grok-responses-control-frame.ts b/src/server/grok-responses-control-frame.ts index cb572acee9..e910daf99d 100644 --- a/src/server/grok-responses-control-frame.ts +++ b/src/server/grok-responses-control-frame.ts @@ -14,12 +14,17 @@ const GROK_CONTROL_FRAME_TYPES: Record = { */ export function createGrokResponsesControlFrameBlockRewrite(): SseBlockRewrite { return (block) => { - const eventName = block - .split(/\r?\n/) - .find(line => line.startsWith("event:")) - ?.slice("event:".length) - .trim(); - if (GROK_CONTROL_FRAME_TYPES[eventName ?? ""] === true) return []; + let eventName = ""; + // SSE overwrites the event type on every event field, including empty resets. + // Like sseDataPayload, remove only one optional ASCII space after the colon. + for (const line of block.split(/\r?\n/)) { + if (line === "event") eventName = ""; + else if (line.startsWith("event:")) { + const value = line.slice("event:".length); + eventName = value.startsWith(" ") ? value.slice(1) : value; + } + } + if (GROK_CONTROL_FRAME_TYPES[eventName] === true) return []; const payload = sseDataPayload(block); if (payload === null || payload === "[DONE]") return [block]; diff --git a/tests/responses/responses-snapshot-repair-server.test.ts b/tests/responses/responses-snapshot-repair-server.test.ts index a7f3ec167c..f6e4ac0e92 100644 --- a/tests/responses/responses-snapshot-repair-server.test.ts +++ b/tests/responses/responses-snapshot-repair-server.test.ts @@ -6,6 +6,7 @@ import { saveConfig } from "../../src/config"; import { startServer } from "../../src/server"; import { handleResponses } from "../../src/server/responses"; import { isEagerRelaySseResponse } from "../../src/server/relay"; +import { createGrokResponsesControlFrameBlockRewrite } from "../../src/server/grok-responses-control-frame"; import type { OcxConfig } from "../../src/types"; import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; import { removeTreeWithRetry } from "../helpers/remove-tree"; @@ -118,6 +119,61 @@ afterEach(async () => { removeTreeWithRetry(TEST_DIR); }); +for (const controlType of ["codex.rate_limits", "codex.response.metadata"]) { + describe(`Grok control frame ${controlType}`, () => { + test.each(["{}", "not-json"])("filters an event-only discriminator with payload %s", payload => { + const rewrite = createGrokResponsesControlFrameBlockRewrite(); + expect(rewrite(`event: ${controlType}\ndata: ${payload}`)).toEqual([]); + }); + + test("filters a data-only discriminator without an event field", () => { + const rewrite = createGrokResponsesControlFrameBlockRewrite(); + expect(rewrite(`data: {"type":"${controlType}"}`)).toEqual([]); + }); + + test.each(["{}", "not-json"])("filters the last event field with payload %s", payload => { + const rewrite = createGrokResponsesControlFrameBlockRewrite(); + expect(rewrite(`event: message\nevent: ${controlType}\ndata: ${payload}`)).toEqual([]); + }); + + test("preserves completion when the last event field overrides a control type", () => { + const block = `event: ${controlType}\nevent: response.completed\ndata: {"type":"response.completed","response":{"id":"r1","status":"completed","output":[]}}`; + expect(createGrokResponsesControlFrameBlockRewrite()(block)).toEqual([block]); + }); + + test.each(["event:", "event: ", "event"])("honors the empty reset %s", reset => { + const block = `event: ${controlType}\n${reset}\ndata: {}`; + expect(createGrokResponsesControlFrameBlockRewrite()(block)).toEqual([block]); + }); + + test("still filters the JSON type after an empty event reset", () => { + const block = `event: ${controlType}\nevent:\ndata: {"type":"${controlType}"}`; + expect(createGrokResponsesControlFrameBlockRewrite()(block)).toEqual([]); + }); + + test.each([`event: ${controlType}`, `event:\t${controlType}`, `event: ${controlType} `])( + "preserves significant event-value whitespace in %s", + eventLine => { + const block = `${eventLine}\ndata: {}`; + expect(createGrokResponsesControlFrameBlockRewrite()(block)).toEqual([block]); + }, + ); + + test("recognizes a CRLF event field without an optional space", () => { + expect(createGrokResponsesControlFrameBlockRewrite()(`event:message\r\nevent:${controlType}\r\ndata: {}`)).toEqual([]); + }); + + test("does not retain the event type across blocks or consume ordinary content", () => { + const rewrite = createGrokResponsesControlFrameBlockRewrite(); + expect(rewrite(`event: ${controlType}\ndata: {}`)).toEqual([]); + for (const block of ["data: {}", "data: not-json", ": heartbeat", "data: [DONE]", + `data: {"type":"response.output_text.delta","delta":"${controlType}"}`]) { + expect(rewrite(block)).toEqual([block]); + } + }); + }); +} + describe("responsesSnapshotRepair through /v1/responses", () => { test.skipIf(process.platform !== "darwin")( "Darwin eager-relay applies snapshot repair inline before bytes reach the client", @@ -341,9 +397,9 @@ describe("responsesSnapshotRepair through /v1/responses", () => { await server.stop(true); } }); - test("the Grok marker filters Codex control frames at the client boundary", async () => { + test.each([true, false])("the Grok marker filters Codex control frames at the client boundary (event names: %s)", async includeEventNames => { const gateway = "https://grok-control-frame.example.test"; - stubSparseGateway(gateway, GROK_CONTROL_FRAME_EVENTS, true); + stubSparseGateway(gateway, GROK_CONTROL_FRAME_EVENTS, includeEventNames); saveConfig({ port: 0, defaultProvider: "sparse",