From d5e0a9a2069ac5655dc208c7261981cd8c59112d Mon Sep 17 00:00:00 2001 From: Danh Thanh Date: Mon, 7 Sep 2026 05:09:13 +0700 Subject: [PATCH] fix(grok): filter Codex metadata frames --- 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 () => {