From 182ba45c913a8587a8280f027b182ee0804ece6b Mon Sep 17 00:00:00 2001 From: Wibias <37517432+Wibias@users.noreply.github.com> Date: Tue, 28 Jul 2026 07:54:00 +0200 Subject: [PATCH] refactor(responses): single-pass SSE payload rewrite composition Follow-up to #588: extract shared sse-payload-rewrite shell and compose image-gen restore with item-id repair in one parse/stringify pass instead of chaining two JS pull wrappers. --- src/server/responses-image-gen-repair.ts | 95 ++----------------- src/server/responses-item-id-repair.ts | 95 ++----------------- src/server/responses/core.ts | 24 +++-- src/server/sse-payload-rewrite.ts | 116 +++++++++++++++++++++++ structure/04_transports-and-sidecars.md | 4 +- tests/sse-payload-rewrite.test.ts | 101 ++++++++++++++++++++ 6 files changed, 256 insertions(+), 179 deletions(-) create mode 100644 src/server/sse-payload-rewrite.ts create mode 100644 tests/sse-payload-rewrite.test.ts diff --git a/src/server/responses-image-gen-repair.ts b/src/server/responses-image-gen-repair.ts index 5ebee945eb..78468f1a44 100644 --- a/src/server/responses-image-gen-repair.ts +++ b/src/server/responses-image-gen-repair.ts @@ -1,4 +1,5 @@ import { collectResponsesToolGroups } from "../responses/tool-groups"; +import { relaySseWithPayloadRewrite, type SsePayloadRewrite } from "./sse-payload-rewrite"; interface NamespacedTool { namespace: string; @@ -96,45 +97,12 @@ export function restoreImageGenCallsInJson( return restored.changed ? JSON.stringify(restored.value) : text; } -/** Split one complete SSE event block while retaining its original blank-line delimiter. */ -function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null { - const match = buffer.match(/\r?\n\r?\n/); - if (!match || match.index === undefined) return null; - return { - block: buffer.slice(0, match.index), - delimiter: match[0], - rest: buffer.slice(match.index + match[0].length), - }; -} - -/** Join all data lines from one SSE event according to the event-stream field rules. */ -function sseDataPayload(block: string): string | null { - const data: string[] = []; - for (const line of block.split(/\r?\n/)) { - if (!line.startsWith("data:")) continue; - const value = line.slice(5); - data.push(value.startsWith(" ") ? value.slice(1) : value); - } - return data.length > 0 ? data.join("\n") : null; -} - -/** Replace an SSE event's data field while preserving non-data fields and newline style. */ -function replaceSseDataPayload(block: string, payload: string): string { - const newline = block.includes("\r\n") ? "\r\n" : "\n"; - const lines = block.split(/\r?\n/); - const rewritten: string[] = []; - let replaced = false; - for (const line of lines) { - if (!line.startsWith("data:")) { - rewritten.push(line); - continue; - } - if (!replaced) { - rewritten.push(`data: ${payload}`); - replaced = true; - } - } - return replaced ? rewritten.join(newline) : block; +/** Payload rewrite for composition with other client-facing SSE transforms. */ +export function createImageGenCallRestoreRewrite( + aliases: ReadonlyMap, +): SsePayloadRewrite | undefined { + if (aliases.size === 0) return undefined; + return (payload) => restoreImageGenCallsInJson(payload, aliases); } /** @@ -145,51 +113,6 @@ export function relaySseWithImageGenCallRestore( body: ReadableStream, aliases: ReadonlyMap, ): ReadableStream { - if (aliases.size === 0) return body; - const reader = body.getReader(); - const decoder = new TextDecoder(); - const encoder = new TextEncoder(); - let buffer = ""; - - const emitProcessedBlocks = ( - controller: ReadableStreamDefaultController, - flushFinal = false, - ): void => { - let next: { block: string; delimiter: string; rest: string } | null; - while ((next = nextSseBlock(buffer))) { - buffer = next.rest; - const payload = sseDataPayload(next.block); - const restoredPayload = payload ? restoreImageGenCallsInJson(payload, aliases) : undefined; - const block = payload && restoredPayload !== undefined && restoredPayload !== payload - ? replaceSseDataPayload(next.block, restoredPayload) - : next.block; - controller.enqueue(encoder.encode(block + next.delimiter)); - } - if (flushFinal && buffer.length > 0) { - const payload = sseDataPayload(buffer); - const restoredPayload = payload ? restoreImageGenCallsInJson(payload, aliases) : undefined; - const block = payload && restoredPayload !== undefined && restoredPayload !== payload - ? replaceSseDataPayload(buffer, restoredPayload) - : buffer; - controller.enqueue(encoder.encode(block)); - buffer = ""; - } - }; - - return new ReadableStream({ - async pull(controller) { - const { done, value } = await reader.read(); - if (done) { - buffer += decoder.decode(); - emitProcessedBlocks(controller, true); - controller.close(); - return; - } - buffer += decoder.decode(value, { stream: true }); - emitProcessedBlocks(controller); - }, - cancel(reason) { - reader.cancel(reason).catch(() => {}); - }, - }); + const rewrite = createImageGenCallRestoreRewrite(aliases); + return rewrite ? relaySseWithPayloadRewrite(body, rewrite) : body; } diff --git a/src/server/responses-item-id-repair.ts b/src/server/responses-item-id-repair.ts index 2309d141e4..075dd7ce79 100644 --- a/src/server/responses-item-id-repair.ts +++ b/src/server/responses-item-id-repair.ts @@ -1,5 +1,6 @@ import { randomUUID } from "node:crypto"; import type { ResponsesItemIdRepairConfig } from "../types"; +import { relaySseWithPayloadRewrite, type SsePayloadRewrite } from "./sse-payload-rewrite"; type RepairableItemType = "message" | "reasoning"; @@ -35,44 +36,6 @@ function isPlainObject(value: unknown): value is Record { return !!value && typeof value === "object" && !Array.isArray(value); } -function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null { - const match = buffer.match(/\r?\n\r?\n/); - if (!match || match.index === undefined) return null; - return { - block: buffer.slice(0, match.index), - delimiter: match[0], - rest: buffer.slice(match.index + match[0].length), - }; -} - -function sseDataPayload(block: string): string | null { - const data: string[] = []; - for (const line of block.split(/\r?\n/)) { - if (!line.startsWith("data:")) continue; - const value = line.slice(5); - data.push(value.startsWith(" ") ? value.slice(1) : value); - } - return data.length > 0 ? data.join("\n") : null; -} - -function replaceSseDataPayload(block: string, payload: string): string { - const newline = block.includes("\r\n") ? "\r\n" : "\n"; - const lines = block.split(/\r?\n/); - const rewritten: string[] = []; - let replaced = false; - for (const line of lines) { - if (!line.startsWith("data:")) { - rewritten.push(line); - continue; - } - if (!replaced) { - rewritten.push(`data: ${payload}`); - replaced = true; - } - } - return replaced ? rewritten.join(newline) : block; -} - function asOutputIndex(value: unknown): number | null { return typeof value === "number" && Number.isInteger(value) && value >= 0 ? value : null; } @@ -221,57 +184,19 @@ function repairEventPayload( * type에 재사용해도 function_call id/call_id는 보존된다. opt-in 게이트웨이는 sequential streams에서도 * 고유한 canonical id를 얻지만, 보정이 필요한 경우에만 JS stream 재작성 비용을 지불한다. */ +/** Stateful payload rewrite for composition with other client-facing SSE transforms. */ +export function createResponsesItemIdPayloadRewrite( + config: ResponsesItemIdRepairConfig, +): SsePayloadRewrite { + const state = createRepairState(config); + return (payload) => repairEventPayload(payload, state); +} + export function relaySseWithResponsesItemIdRepair( body: ReadableStream, config: ResponsesItemIdRepairConfig, ): ReadableStream { - const reader = body.getReader(); - const decoder = new TextDecoder(); - const encoder = new TextEncoder(); - const state = createRepairState(config); - let buffer = ""; - - const emitProcessedBlocks = ( - controller: ReadableStreamDefaultController, - flushFinal = false, - ): void => { - let next: { block: string; delimiter: string; rest: string } | null; - while ((next = nextSseBlock(buffer))) { - buffer = next.rest; - const payload = sseDataPayload(next.block); - const repairedPayload = payload ? repairEventPayload(payload, state) : undefined; - const block = payload && repairedPayload !== undefined && repairedPayload !== payload - ? replaceSseDataPayload(next.block, repairedPayload) - : next.block; - controller.enqueue(encoder.encode(block + next.delimiter)); - } - if (flushFinal && buffer.length > 0) { - const payload = sseDataPayload(buffer); - const repairedPayload = payload ? repairEventPayload(payload, state) : undefined; - const block = payload && repairedPayload !== undefined && repairedPayload !== payload - ? replaceSseDataPayload(buffer, repairedPayload) - : buffer; - controller.enqueue(encoder.encode(block)); - buffer = ""; - } - }; - - return new ReadableStream({ - async pull(controller) { - const { done, value } = await reader.read(); - if (done) { - buffer += decoder.decode(); - emitProcessedBlocks(controller, true); - controller.close(); - return; - } - buffer += decoder.decode(value, { stream: true }); - emitProcessedBlocks(controller); - }, - cancel(reason) { - reader.cancel(reason).catch(() => {}); - }, - }); + return relaySseWithPayloadRewrite(body, createResponsesItemIdPayloadRewrite(config)); } export function hasResponsesItemIdRepair(config: ResponsesItemIdRepairConfig | undefined): boolean { diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 337eca035e..eb05cac298 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -132,12 +132,16 @@ import { import { relaySseEagerBounded } from "../relay-eager"; import { decideEagerRelay } from "../../lib/bun-stream-caps"; import { cancelBodyOnAbort } from "../../lib/abort"; -import { hasResponsesItemIdRepair, relaySseWithResponsesItemIdRepair } from "../responses-item-id-repair"; import { + createResponsesItemIdPayloadRewrite, + hasResponsesItemIdRepair, +} from "../responses-item-id-repair"; +import { + createImageGenCallRestoreRewrite, imageGenToolCallAliases, - relaySseWithImageGenCallRestore, restoreImageGenCallsInJson, } from "../responses-image-gen-repair"; +import { composeSsePayloadRewrites, relaySseWithPayloadRewrite } from "../sse-payload-rewrite"; import type { EffectiveSubagentRoster, SpawnAgentSurface } from "../../codex/catalog"; import { buildToolBridgeMaps, collabSurface, injectDeveloperMessage, multiAgentGuidanceText } from "./collaboration"; @@ -1663,13 +1667,19 @@ export async function handleResponses( // win32 must keep the pure native relay (Bun#32111 JS-sink segfault); elsewhere a JS pull // relay is established practice (relayWithAbort, relaySseWithHeartbeat) and lets a // mid-stream reset end with a clean response.failed terminal instead of a raw socket error. - const restoredBody = relaySseWithImageGenCallRestore(nativeBody, imageGenCallAliases); - const repairedBody = hasResponsesItemIdRepair(repairConfig) - ? relaySseWithResponsesItemIdRepair(restoredBody, repairConfig!) - : restoredBody; + // Compose opt-in payload rewrites into one parse/stringify pass (image-gen restore first). + const payloadRewrites = [ + createImageGenCallRestoreRewrite(imageGenCallAliases), + hasResponsesItemIdRepair(repairConfig) + ? createResponsesItemIdPayloadRewrite(repairConfig!) + : undefined, + ].filter((rewrite): rewrite is NonNullable => rewrite !== undefined); + const rewrittenBody = payloadRewrites.length > 0 + ? relaySseWithPayloadRewrite(nativeBody, composeSsePayloadRewrites(...payloadRewrites)) + : nativeBody; const clientBody = process.platform === "win32" && !needsClientRewrite ? nativeBody - : relaySseWithFailedTail(repairedBody, upstream); + : relaySseWithFailedTail(rewrittenBody, upstream); return markNativePassthroughSseResponse(new Response(clientBody, { status: upstreamResponse.status, headers, diff --git a/src/server/sse-payload-rewrite.ts b/src/server/sse-payload-rewrite.ts new file mode 100644 index 0000000000..a9c893e87a --- /dev/null +++ b/src/server/sse-payload-rewrite.ts @@ -0,0 +1,116 @@ +/** + * Shared client-facing SSE payload rewrite shell. + * + * Multiple opt-in transforms (image-gen namespace restore, item-id repair, …) compose into one + * parse/stringify pass so a tee'd stream is not re-framed twice per event. + */ + +export type SsePayloadRewrite = (payload: string) => string; + +/** Split one complete SSE event block while retaining its original blank-line delimiter. */ +export function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null { + const match = buffer.match(/\r?\n\r?\n/); + if (!match || match.index === undefined) return null; + return { + block: buffer.slice(0, match.index), + delimiter: match[0], + rest: buffer.slice(match.index + match[0].length), + }; +} + +/** Join all data lines from one SSE event according to the event-stream field rules. */ +export function sseDataPayload(block: string): string | null { + const data: string[] = []; + for (const line of block.split(/\r?\n/)) { + if (!line.startsWith("data:")) continue; + const value = line.slice(5); + data.push(value.startsWith(" ") ? value.slice(1) : value); + } + return data.length > 0 ? data.join("\n") : null; +} + +/** Replace an SSE event's data field while preserving non-data fields and newline style. */ +export function replaceSseDataPayload(block: string, payload: string): string { + const newline = block.includes("\r\n") ? "\r\n" : "\n"; + const lines = block.split(/\r?\n/); + const rewritten: string[] = []; + let replaced = false; + for (const line of lines) { + if (!line.startsWith("data:")) { + rewritten.push(line); + continue; + } + if (!replaced) { + rewritten.push(`data: ${payload}`); + replaced = true; + } + } + return replaced ? rewritten.join(newline) : block; +} + +/** Apply rewrites left-to-right; empty list is identity. */ +export function composeSsePayloadRewrites(...rewrites: SsePayloadRewrite[]): SsePayloadRewrite { + if (rewrites.length === 0) return (payload) => payload; + if (rewrites.length === 1) return rewrites[0]!; + return (payload) => { + let next = payload; + for (const rewrite of rewrites) next = rewrite(next); + return next; + }; +} + +/** + * Relay an SSE body through a single JS pull wrapper, rewriting each event's data payload in place. + * Non-data fields and framing are preserved; invalid JSON payloads are left to the rewrite callback. + */ +export function relaySseWithPayloadRewrite( + body: ReadableStream, + rewrite: SsePayloadRewrite, +): ReadableStream { + const reader = body.getReader(); + const decoder = new TextDecoder(); + const encoder = new TextEncoder(); + let buffer = ""; + + const emitProcessedBlocks = ( + controller: ReadableStreamDefaultController, + flushFinal = false, + ): void => { + let next: { block: string; delimiter: string; rest: string } | null; + while ((next = nextSseBlock(buffer))) { + buffer = next.rest; + const payload = sseDataPayload(next.block); + const rewrittenPayload = payload ? rewrite(payload) : undefined; + const block = payload && rewrittenPayload !== undefined && rewrittenPayload !== payload + ? replaceSseDataPayload(next.block, rewrittenPayload) + : next.block; + controller.enqueue(encoder.encode(block + next.delimiter)); + } + if (flushFinal && buffer.length > 0) { + const payload = sseDataPayload(buffer); + const rewrittenPayload = payload ? rewrite(payload) : undefined; + const block = payload && rewrittenPayload !== undefined && rewrittenPayload !== payload + ? replaceSseDataPayload(buffer, rewrittenPayload) + : buffer; + controller.enqueue(encoder.encode(block)); + buffer = ""; + } + }; + + return new ReadableStream({ + async pull(controller) { + const { done, value } = await reader.read(); + if (done) { + buffer += decoder.decode(); + emitProcessedBlocks(controller, true); + controller.close(); + return; + } + buffer += decoder.decode(value, { stream: true }); + emitProcessedBlocks(controller); + }, + cancel(reason) { + reader.cancel(reason).catch(() => {}); + }, + }); +} diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index d0d4326613..e4d1fbaf12 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -74,7 +74,9 @@ namespaces do not remove the hosted fallback. Discovery and normalization span b Client-facing API-key responses perform the inverse mapping: JSON output and SSE function-call items restore `{ namespace: "image_gen", name: "" }` so Codex can dispatch the local -extension. Inspection and continuation-cache branches keep the raw upstream alias, allowing stored +extension. When item-id repair is also enabled, both transforms compose in one SSE parse/stringify +pass (`src/server/sse-payload-rewrite.ts`) rather than chaining separate JS pull wrappers. +Inspection and continuation-cache branches keep the raw upstream alias, allowing stored replays to return upstream without leaking a client-only namespace shape. Malformed, empty, and unrelated namespaces remain untouched. ChatGPT forward mode preserves the private namespace and hosted tool because that backend understands their native semantics. diff --git a/tests/sse-payload-rewrite.test.ts b/tests/sse-payload-rewrite.test.ts new file mode 100644 index 0000000000..3bdff9325d --- /dev/null +++ b/tests/sse-payload-rewrite.test.ts @@ -0,0 +1,101 @@ +/** + * Single-pass composition of client-facing SSE payload rewrites (#588 follow-up). + */ +import { describe, expect, test } from "bun:test"; +import { createImageGenCallRestoreRewrite } from "../src/server/responses-image-gen-repair"; +import { createResponsesItemIdPayloadRewrite } from "../src/server/responses-item-id-repair"; +import { + composeSsePayloadRewrites, + relaySseWithPayloadRewrite, +} from "../src/server/sse-payload-rewrite"; + +function streamFromText(text: string): ReadableStream { + const chunk = new TextEncoder().encode(text); + let sent = false; + return new ReadableStream({ + pull(controller) { + if (sent) { + controller.close(); + return; + } + sent = true; + controller.enqueue(chunk); + }, + }); +} + +async function readAll(stream: ReadableStream): Promise { + const reader = stream.getReader(); + const decoder = new TextDecoder(); + let text = ""; + while (true) { + const { done, value } = await reader.read(); + if (done) break; + text += decoder.decode(value, { stream: true }); + } + return text; +} + +describe("SSE payload rewrite composition", () => { + test("applies image-gen restore and item-id repair in one relay pass", async () => { + const upstream = [ + 'event: response.output_item.added\ndata: {"type":"response.output_item.added","output_index":0,"item":{"type":"message","id":"msg_0","role":"assistant"}}\n\n', + 'event: response.output_item.added\ndata: {"type":"response.output_item.added","output_index":1,"item":{"type":"function_call","id":"fc_1","call_id":"call_1","name":"image_gen__imagegen","arguments":"{}"}}\n\n', + 'event: response.completed\ndata: {"type":"response.completed","response":{"id":"resp_1","status":"completed","output":[{"type":"message","id":"msg_0","role":"assistant"},{"type":"function_call","id":"fc_1","call_id":"call_1","name":"image_gen__imagegen","arguments":"{}"}]}}\n\n', + ].join(""); + + let imageGenCalls = 0; + let itemIdCalls = 0; + const imageGen = createImageGenCallRestoreRewrite( + new Map([["image_gen__imagegen", { namespace: "image_gen", name: "imagegen" }]]), + )!; + const itemId = createResponsesItemIdPayloadRewrite({ + message: ["msg_0"], + repairMissingTerminalIds: true, + }); + + const composed = composeSsePayloadRewrites( + (payload) => { + imageGenCalls += 1; + return imageGen(payload); + }, + (payload) => { + itemIdCalls += 1; + return itemId(payload); + }, + ); + + const out = await readAll(relaySseWithPayloadRewrite(streamFromText(upstream), composed)); + expect(imageGenCalls).toBe(3); + expect(itemIdCalls).toBe(3); + expect(imageGenCalls).toBe(itemIdCalls); + + const events = out + .trim() + .split(/\r?\n\r?\n/) + .map(block => block.split(/\r?\n/).find(line => line.startsWith("data:"))?.slice(5).trim()) + .filter((payload): payload is string => !!payload) + .map(payload => JSON.parse(payload) as Record); + + const messageAdded = events[0].item as Record; + expect(messageAdded.id).toMatch(/^msg_ocx_[0-9a-f]+_0$/); + + const functionAdded = events[1].item as Record; + expect(functionAdded).toMatchObject({ + name: "imagegen", + namespace: "image_gen", + call_id: "call_1", + }); + + const completed = events[2].response as { output: Record[] }; + expect(completed.output[0].id).toBe(messageAdded.id); + expect(completed.output[1]).toMatchObject({ + name: "imagegen", + namespace: "image_gen", + }); + }); + + test("compose with no rewrites is identity", () => { + expect(composeSsePayloadRewrites()('{"a":1}')).toBe('{"a":1}'); + }); +});