Skip to content

Commit 406a522

Browse files
authored
refactor(responses): single-pass SSE payload rewrite composition (#602)
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.
1 parent c380ef7 commit 406a522

6 files changed

Lines changed: 256 additions & 179 deletions

File tree

src/server/responses-image-gen-repair.ts

Lines changed: 9 additions & 86 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { collectResponsesToolGroups } from "../responses/tool-groups";
2+
import { relaySseWithPayloadRewrite, type SsePayloadRewrite } from "./sse-payload-rewrite";
23

34
interface NamespacedTool {
45
namespace: string;
@@ -96,45 +97,12 @@ export function restoreImageGenCallsInJson(
9697
return restored.changed ? JSON.stringify(restored.value) : text;
9798
}
9899

99-
/** Split one complete SSE event block while retaining its original blank-line delimiter. */
100-
function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null {
101-
const match = buffer.match(/\r?\n\r?\n/);
102-
if (!match || match.index === undefined) return null;
103-
return {
104-
block: buffer.slice(0, match.index),
105-
delimiter: match[0],
106-
rest: buffer.slice(match.index + match[0].length),
107-
};
108-
}
109-
110-
/** Join all data lines from one SSE event according to the event-stream field rules. */
111-
function sseDataPayload(block: string): string | null {
112-
const data: string[] = [];
113-
for (const line of block.split(/\r?\n/)) {
114-
if (!line.startsWith("data:")) continue;
115-
const value = line.slice(5);
116-
data.push(value.startsWith(" ") ? value.slice(1) : value);
117-
}
118-
return data.length > 0 ? data.join("\n") : null;
119-
}
120-
121-
/** Replace an SSE event's data field while preserving non-data fields and newline style. */
122-
function replaceSseDataPayload(block: string, payload: string): string {
123-
const newline = block.includes("\r\n") ? "\r\n" : "\n";
124-
const lines = block.split(/\r?\n/);
125-
const rewritten: string[] = [];
126-
let replaced = false;
127-
for (const line of lines) {
128-
if (!line.startsWith("data:")) {
129-
rewritten.push(line);
130-
continue;
131-
}
132-
if (!replaced) {
133-
rewritten.push(`data: ${payload}`);
134-
replaced = true;
135-
}
136-
}
137-
return replaced ? rewritten.join(newline) : block;
100+
/** Payload rewrite for composition with other client-facing SSE transforms. */
101+
export function createImageGenCallRestoreRewrite(
102+
aliases: ReadonlyMap<string, NamespacedTool>,
103+
): SsePayloadRewrite | undefined {
104+
if (aliases.size === 0) return undefined;
105+
return (payload) => restoreImageGenCallsInJson(payload, aliases);
138106
}
139107

140108
/**
@@ -145,51 +113,6 @@ export function relaySseWithImageGenCallRestore(
145113
body: ReadableStream<Uint8Array>,
146114
aliases: ReadonlyMap<string, NamespacedTool>,
147115
): ReadableStream<Uint8Array> {
148-
if (aliases.size === 0) return body;
149-
const reader = body.getReader();
150-
const decoder = new TextDecoder();
151-
const encoder = new TextEncoder();
152-
let buffer = "";
153-
154-
const emitProcessedBlocks = (
155-
controller: ReadableStreamDefaultController<Uint8Array>,
156-
flushFinal = false,
157-
): void => {
158-
let next: { block: string; delimiter: string; rest: string } | null;
159-
while ((next = nextSseBlock(buffer))) {
160-
buffer = next.rest;
161-
const payload = sseDataPayload(next.block);
162-
const restoredPayload = payload ? restoreImageGenCallsInJson(payload, aliases) : undefined;
163-
const block = payload && restoredPayload !== undefined && restoredPayload !== payload
164-
? replaceSseDataPayload(next.block, restoredPayload)
165-
: next.block;
166-
controller.enqueue(encoder.encode(block + next.delimiter));
167-
}
168-
if (flushFinal && buffer.length > 0) {
169-
const payload = sseDataPayload(buffer);
170-
const restoredPayload = payload ? restoreImageGenCallsInJson(payload, aliases) : undefined;
171-
const block = payload && restoredPayload !== undefined && restoredPayload !== payload
172-
? replaceSseDataPayload(buffer, restoredPayload)
173-
: buffer;
174-
controller.enqueue(encoder.encode(block));
175-
buffer = "";
176-
}
177-
};
178-
179-
return new ReadableStream<Uint8Array>({
180-
async pull(controller) {
181-
const { done, value } = await reader.read();
182-
if (done) {
183-
buffer += decoder.decode();
184-
emitProcessedBlocks(controller, true);
185-
controller.close();
186-
return;
187-
}
188-
buffer += decoder.decode(value, { stream: true });
189-
emitProcessedBlocks(controller);
190-
},
191-
cancel(reason) {
192-
reader.cancel(reason).catch(() => {});
193-
},
194-
});
116+
const rewrite = createImageGenCallRestoreRewrite(aliases);
117+
return rewrite ? relaySseWithPayloadRewrite(body, rewrite) : body;
195118
}

src/server/responses-item-id-repair.ts

Lines changed: 10 additions & 85 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { randomUUID } from "node:crypto";
22
import type { ResponsesItemIdRepairConfig } from "../types";
3+
import { relaySseWithPayloadRewrite, type SsePayloadRewrite } from "./sse-payload-rewrite";
34

45
type RepairableItemType = "message" | "reasoning";
56

@@ -35,44 +36,6 @@ function isPlainObject(value: unknown): value is Record<string, unknown> {
3536
return !!value && typeof value === "object" && !Array.isArray(value);
3637
}
3738

38-
function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null {
39-
const match = buffer.match(/\r?\n\r?\n/);
40-
if (!match || match.index === undefined) return null;
41-
return {
42-
block: buffer.slice(0, match.index),
43-
delimiter: match[0],
44-
rest: buffer.slice(match.index + match[0].length),
45-
};
46-
}
47-
48-
function sseDataPayload(block: string): string | null {
49-
const data: string[] = [];
50-
for (const line of block.split(/\r?\n/)) {
51-
if (!line.startsWith("data:")) continue;
52-
const value = line.slice(5);
53-
data.push(value.startsWith(" ") ? value.slice(1) : value);
54-
}
55-
return data.length > 0 ? data.join("\n") : null;
56-
}
57-
58-
function replaceSseDataPayload(block: string, payload: string): string {
59-
const newline = block.includes("\r\n") ? "\r\n" : "\n";
60-
const lines = block.split(/\r?\n/);
61-
const rewritten: string[] = [];
62-
let replaced = false;
63-
for (const line of lines) {
64-
if (!line.startsWith("data:")) {
65-
rewritten.push(line);
66-
continue;
67-
}
68-
if (!replaced) {
69-
rewritten.push(`data: ${payload}`);
70-
replaced = true;
71-
}
72-
}
73-
return replaced ? rewritten.join(newline) : block;
74-
}
75-
7639
function asOutputIndex(value: unknown): number | null {
7740
return typeof value === "number" && Number.isInteger(value) && value >= 0 ? value : null;
7841
}
@@ -221,57 +184,19 @@ function repairEventPayload(
221184
* type에 재사용해도 function_call id/call_id는 보존된다. opt-in 게이트웨이는 sequential streams에서도
222185
* 고유한 canonical id를 얻지만, 보정이 필요한 경우에만 JS stream 재작성 비용을 지불한다.
223186
*/
187+
/** Stateful payload rewrite for composition with other client-facing SSE transforms. */
188+
export function createResponsesItemIdPayloadRewrite(
189+
config: ResponsesItemIdRepairConfig,
190+
): SsePayloadRewrite {
191+
const state = createRepairState(config);
192+
return (payload) => repairEventPayload(payload, state);
193+
}
194+
224195
export function relaySseWithResponsesItemIdRepair(
225196
body: ReadableStream<Uint8Array>,
226197
config: ResponsesItemIdRepairConfig,
227198
): ReadableStream<Uint8Array> {
228-
const reader = body.getReader();
229-
const decoder = new TextDecoder();
230-
const encoder = new TextEncoder();
231-
const state = createRepairState(config);
232-
let buffer = "";
233-
234-
const emitProcessedBlocks = (
235-
controller: ReadableStreamDefaultController<Uint8Array>,
236-
flushFinal = false,
237-
): void => {
238-
let next: { block: string; delimiter: string; rest: string } | null;
239-
while ((next = nextSseBlock(buffer))) {
240-
buffer = next.rest;
241-
const payload = sseDataPayload(next.block);
242-
const repairedPayload = payload ? repairEventPayload(payload, state) : undefined;
243-
const block = payload && repairedPayload !== undefined && repairedPayload !== payload
244-
? replaceSseDataPayload(next.block, repairedPayload)
245-
: next.block;
246-
controller.enqueue(encoder.encode(block + next.delimiter));
247-
}
248-
if (flushFinal && buffer.length > 0) {
249-
const payload = sseDataPayload(buffer);
250-
const repairedPayload = payload ? repairEventPayload(payload, state) : undefined;
251-
const block = payload && repairedPayload !== undefined && repairedPayload !== payload
252-
? replaceSseDataPayload(buffer, repairedPayload)
253-
: buffer;
254-
controller.enqueue(encoder.encode(block));
255-
buffer = "";
256-
}
257-
};
258-
259-
return new ReadableStream<Uint8Array>({
260-
async pull(controller) {
261-
const { done, value } = await reader.read();
262-
if (done) {
263-
buffer += decoder.decode();
264-
emitProcessedBlocks(controller, true);
265-
controller.close();
266-
return;
267-
}
268-
buffer += decoder.decode(value, { stream: true });
269-
emitProcessedBlocks(controller);
270-
},
271-
cancel(reason) {
272-
reader.cancel(reason).catch(() => {});
273-
},
274-
});
199+
return relaySseWithPayloadRewrite(body, createResponsesItemIdPayloadRewrite(config));
275200
}
276201

277202
export function hasResponsesItemIdRepair(config: ResponsesItemIdRepairConfig | undefined): boolean {

src/server/responses/core.ts

Lines changed: 17 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -132,12 +132,16 @@ import {
132132
import { relaySseEagerBounded } from "../relay-eager";
133133
import { decideEagerRelay } from "../../lib/bun-stream-caps";
134134
import { cancelBodyOnAbort } from "../../lib/abort";
135-
import { hasResponsesItemIdRepair, relaySseWithResponsesItemIdRepair } from "../responses-item-id-repair";
136135
import {
136+
createResponsesItemIdPayloadRewrite,
137+
hasResponsesItemIdRepair,
138+
} from "../responses-item-id-repair";
139+
import {
140+
createImageGenCallRestoreRewrite,
137141
imageGenToolCallAliases,
138-
relaySseWithImageGenCallRestore,
139142
restoreImageGenCallsInJson,
140143
} from "../responses-image-gen-repair";
144+
import { composeSsePayloadRewrites, relaySseWithPayloadRewrite } from "../sse-payload-rewrite";
141145
import type { EffectiveSubagentRoster, SpawnAgentSurface } from "../../codex/catalog";
142146

143147
import { buildToolBridgeMaps, collabSurface, injectDeveloperMessage, multiAgentGuidanceText } from "./collaboration";
@@ -1665,13 +1669,19 @@ export async function handleResponses(
16651669
// win32 must keep the pure native relay (Bun#32111 JS-sink segfault); elsewhere a JS pull
16661670
// relay is established practice (relayWithAbort, relaySseWithHeartbeat) and lets a
16671671
// mid-stream reset end with a clean response.failed terminal instead of a raw socket error.
1668-
const restoredBody = relaySseWithImageGenCallRestore(nativeBody, imageGenCallAliases);
1669-
const repairedBody = hasResponsesItemIdRepair(repairConfig)
1670-
? relaySseWithResponsesItemIdRepair(restoredBody, repairConfig!)
1671-
: restoredBody;
1672+
// Compose opt-in payload rewrites into one parse/stringify pass (image-gen restore first).
1673+
const payloadRewrites = [
1674+
createImageGenCallRestoreRewrite(imageGenCallAliases),
1675+
hasResponsesItemIdRepair(repairConfig)
1676+
? createResponsesItemIdPayloadRewrite(repairConfig!)
1677+
: undefined,
1678+
].filter((rewrite): rewrite is NonNullable<typeof rewrite> => rewrite !== undefined);
1679+
const rewrittenBody = payloadRewrites.length > 0
1680+
? relaySseWithPayloadRewrite(nativeBody, composeSsePayloadRewrites(...payloadRewrites))
1681+
: nativeBody;
16721682
const clientBody = process.platform === "win32" && !needsClientRewrite
16731683
? nativeBody
1674-
: relaySseWithFailedTail(repairedBody, upstream);
1684+
: relaySseWithFailedTail(rewrittenBody, upstream);
16751685
return markNativePassthroughSseResponse(new Response(clientBody, {
16761686
status: upstreamResponse.status,
16771687
headers,

src/server/sse-payload-rewrite.ts

Lines changed: 116 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,116 @@
1+
/**
2+
* Shared client-facing SSE payload rewrite shell.
3+
*
4+
* Multiple opt-in transforms (image-gen namespace restore, item-id repair, …) compose into one
5+
* parse/stringify pass so a tee'd stream is not re-framed twice per event.
6+
*/
7+
8+
export type SsePayloadRewrite = (payload: string) => string;
9+
10+
/** Split one complete SSE event block while retaining its original blank-line delimiter. */
11+
export function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null {
12+
const match = buffer.match(/\r?\n\r?\n/);
13+
if (!match || match.index === undefined) return null;
14+
return {
15+
block: buffer.slice(0, match.index),
16+
delimiter: match[0],
17+
rest: buffer.slice(match.index + match[0].length),
18+
};
19+
}
20+
21+
/** Join all data lines from one SSE event according to the event-stream field rules. */
22+
export function sseDataPayload(block: string): string | null {
23+
const data: string[] = [];
24+
for (const line of block.split(/\r?\n/)) {
25+
if (!line.startsWith("data:")) continue;
26+
const value = line.slice(5);
27+
data.push(value.startsWith(" ") ? value.slice(1) : value);
28+
}
29+
return data.length > 0 ? data.join("\n") : null;
30+
}
31+
32+
/** Replace an SSE event's data field while preserving non-data fields and newline style. */
33+
export function replaceSseDataPayload(block: string, payload: string): string {
34+
const newline = block.includes("\r\n") ? "\r\n" : "\n";
35+
const lines = block.split(/\r?\n/);
36+
const rewritten: string[] = [];
37+
let replaced = false;
38+
for (const line of lines) {
39+
if (!line.startsWith("data:")) {
40+
rewritten.push(line);
41+
continue;
42+
}
43+
if (!replaced) {
44+
rewritten.push(`data: ${payload}`);
45+
replaced = true;
46+
}
47+
}
48+
return replaced ? rewritten.join(newline) : block;
49+
}
50+
51+
/** Apply rewrites left-to-right; empty list is identity. */
52+
export function composeSsePayloadRewrites(...rewrites: SsePayloadRewrite[]): SsePayloadRewrite {
53+
if (rewrites.length === 0) return (payload) => payload;
54+
if (rewrites.length === 1) return rewrites[0]!;
55+
return (payload) => {
56+
let next = payload;
57+
for (const rewrite of rewrites) next = rewrite(next);
58+
return next;
59+
};
60+
}
61+
62+
/**
63+
* Relay an SSE body through a single JS pull wrapper, rewriting each event's data payload in place.
64+
* Non-data fields and framing are preserved; invalid JSON payloads are left to the rewrite callback.
65+
*/
66+
export function relaySseWithPayloadRewrite(
67+
body: ReadableStream<Uint8Array>,
68+
rewrite: SsePayloadRewrite,
69+
): ReadableStream<Uint8Array> {
70+
const reader = body.getReader();
71+
const decoder = new TextDecoder();
72+
const encoder = new TextEncoder();
73+
let buffer = "";
74+
75+
const emitProcessedBlocks = (
76+
controller: ReadableStreamDefaultController<Uint8Array>,
77+
flushFinal = false,
78+
): void => {
79+
let next: { block: string; delimiter: string; rest: string } | null;
80+
while ((next = nextSseBlock(buffer))) {
81+
buffer = next.rest;
82+
const payload = sseDataPayload(next.block);
83+
const rewrittenPayload = payload ? rewrite(payload) : undefined;
84+
const block = payload && rewrittenPayload !== undefined && rewrittenPayload !== payload
85+
? replaceSseDataPayload(next.block, rewrittenPayload)
86+
: next.block;
87+
controller.enqueue(encoder.encode(block + next.delimiter));
88+
}
89+
if (flushFinal && buffer.length > 0) {
90+
const payload = sseDataPayload(buffer);
91+
const rewrittenPayload = payload ? rewrite(payload) : undefined;
92+
const block = payload && rewrittenPayload !== undefined && rewrittenPayload !== payload
93+
? replaceSseDataPayload(buffer, rewrittenPayload)
94+
: buffer;
95+
controller.enqueue(encoder.encode(block));
96+
buffer = "";
97+
}
98+
};
99+
100+
return new ReadableStream<Uint8Array>({
101+
async pull(controller) {
102+
const { done, value } = await reader.read();
103+
if (done) {
104+
buffer += decoder.decode();
105+
emitProcessedBlocks(controller, true);
106+
controller.close();
107+
return;
108+
}
109+
buffer += decoder.decode(value, { stream: true });
110+
emitProcessedBlocks(controller);
111+
},
112+
cancel(reason) {
113+
reader.cancel(reason).catch(() => {});
114+
},
115+
});
116+
}

0 commit comments

Comments
 (0)