Skip to content

Commit 3c32ac8

Browse files
committed
address review feedback for deepseek responses repair
Make response ids unique, mint missing item ids under rewriteNonCanonicalIds, share rewrite/trailer state through one factory, pass trailers through the Windows eager path, and add focused regression tests. Also link the PR to #938.
1 parent c86af23 commit 3c32ac8

5 files changed

Lines changed: 102 additions & 30 deletions

File tree

src/server/relay-eager.ts

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,8 @@ export type EagerRelayHooks = {
5858
};
5959

6060
export type EagerRelayOptions = {
61+
/** Optional trailer emitted after rewrite flush (e.g. data: [DONE]). */
62+
trailer?: string | (() => string | undefined);
6163
/** Bounded client queue in bytes; producer pauses above it. Default 8 MiB. */
6264
maxQueueBytes?: number;
6365
/** Transient-budget owner for the inline-rewrite frame buffer. */
@@ -232,6 +234,12 @@ export function relaySseEagerBounded(
232234
queuedBytes += tail.byteLength;
233235
try { controllerRef?.enqueue(tail); } catch { /* client already gone */ }
234236
}
237+
const trailer = typeof opts?.trailer === "function" ? opts.trailer() : opts?.trailer;
238+
if (trailer && !cancelled) {
239+
const trailerBytes = new TextEncoder().encode(trailer);
240+
queuedBytes += trailerBytes.byteLength;
241+
try { controllerRef?.enqueue(trailerBytes); } catch { /* client already gone */ }
242+
}
235243
}
236244
if (!hooks.sawTerminal() && !cancelled && !upstream.signal.aborted) {
237245
syntheticKind = "incomplete";

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

Lines changed: 33 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,9 @@ function mapRawId(
8686
}
8787

8888
if (!rawId) {
89-
if (!state.repairMissingTerminalIds) return null;
89+
// When non-canonical rewrite is active, mint a request-local id even if the upstream
90+
// omitted the terminal id, so Codex lifecycle correlation stays stable.
91+
if (!state.repairMissingTerminalIds && !state.rewriteNonCanonicalIds) return null;
9092
const minted = mintCanonicalId(type, state.scope, outputIndex);
9193
state.budget?.chargeRetained(new TextEncoder().encode(JSON.stringify([outputIndex, minted])).byteLength, { kind: "item_ids" });
9294
state.outputIds[type].set(outputIndex, minted);
@@ -151,7 +153,8 @@ function mapResponseId(state: ResponsesItemIdRepairState, rawId: string | undefi
151153
if (rawId.startsWith("resp_")) return rawId;
152154
const existing = state.responseIdMap.get(rawId);
153155
if (existing) return existing;
154-
const minted = `resp_ocx_${state.scope}`;
156+
// Keep each distinct upstream response id unique within the request-local scope.
157+
const minted = `resp_ocx_${state.scope}_${state.responseIdMap.size}`;
155158
state.responseIdMap.set(rawId, minted);
156159
return minted;
157160
}
@@ -220,7 +223,7 @@ function rewriteOutputItem(
220223
if (!mapped) return { item, changed: false };
221224
const currentId = typeof item.id === "string" ? item.id : undefined;
222225
if (currentId === mapped) return { item, changed: false };
223-
if (currentId === undefined && !state.repairMissingTerminalIds) return { item, changed: false };
226+
if (currentId === undefined && !state.repairMissingTerminalIds && !state.rewriteNonCanonicalIds) return { item, changed: false };
224227
return { item: { ...item, id: mapped }, changed: true };
225228
}
226229

@@ -254,7 +257,7 @@ function rewriteItemIdField(
254257
}
255258
if (!mapped) return { event, changed: false };
256259
if (currentId === mapped) return { event, changed: false };
257-
if (currentId === undefined && !state.repairMissingTerminalIds) return { event, changed: false };
260+
if (currentId === undefined && !state.repairMissingTerminalIds && !state.rewriteNonCanonicalIds) return { event, changed: false };
258261
const next: Record<string, unknown> = { ...event, item_id: mapped };
259262
if (state.rewriteNonCanonicalIds && "logprobs" in next) delete next.logprobs;
260263
return { event: next, changed: true };
@@ -403,36 +406,45 @@ function repairEventPayload(
403406
* rewriteNonCanonicalIds rewrites DeepSeek-style UUID item/response ids, converts
404407
* reasoning_text streams into Codex-friendly encrypted_content reasoning items, and
405408
* ensures a terminal [DONE] trailer so CLI/TUI turns finish.
409+
*
410+
* The rewrite and trailer share one request-local state object via closure so
411+
* composition cannot lose the trailer state.
406412
*/
407-
export function createResponsesItemIdPayloadRewrite(
413+
export interface ResponsesItemIdRepairHandlers {
414+
rewrite: SsePayloadRewrite;
415+
trailer: () => string | undefined;
416+
}
417+
418+
export function createResponsesItemIdRepairHandlers(
408419
config: ResponsesItemIdRepairConfig,
409420
budget?: TranslatorBudget,
410-
): SsePayloadRewrite {
421+
): ResponsesItemIdRepairHandlers {
411422
const state = createRepairState(config, budget);
412-
const rewrite: SsePayloadRewrite = (payload) => repairEventPayload(payload, state);
413-
(rewrite as SsePayloadRewrite & { __state?: ResponsesItemIdRepairState }).__state = state;
414-
return rewrite;
423+
return {
424+
rewrite: (payload) => repairEventPayload(payload, state),
425+
trailer: () => {
426+
if (state.sawCompleted && !state.sawDoneTrailer) return "data: [DONE]\n\n";
427+
return undefined;
428+
},
429+
};
415430
}
416431

417-
export function createResponsesItemIdDoneTrailer(
418-
rewrite: SsePayloadRewrite,
419-
): () => string | undefined {
420-
return () => {
421-
const state = (rewrite as SsePayloadRewrite & { __state?: ResponsesItemIdRepairState }).__state;
422-
if (!state) return undefined;
423-
if (state.sawCompleted && !state.sawDoneTrailer) return "data: [DONE]\n\n";
424-
return undefined;
425-
};
432+
/** Backward-compatible helper used by existing tests and callers. */
433+
export function createResponsesItemIdPayloadRewrite(
434+
config: ResponsesItemIdRepairConfig,
435+
budget?: TranslatorBudget,
436+
): SsePayloadRewrite {
437+
return createResponsesItemIdRepairHandlers(config, budget).rewrite;
426438
}
427439

428440
export function relaySseWithResponsesItemIdRepair(
429441
body: ReadableStream<Uint8Array>,
430442
config: ResponsesItemIdRepairConfig,
431443
budget: TranslatorBudget,
432444
): ReadableStream<Uint8Array> {
433-
const rewrite = createResponsesItemIdPayloadRewrite(config, budget);
434-
return relaySseWithPayloadRewrite(body, rewrite, budget, {
435-
trailer: createResponsesItemIdDoneTrailer(rewrite),
445+
const handlers = createResponsesItemIdRepairHandlers(config, budget);
446+
return relaySseWithPayloadRewrite(body, handlers.rewrite, budget, {
447+
trailer: handlers.trailer,
436448
});
437449
}
438450

src/server/responses/core.ts

Lines changed: 17 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -147,7 +147,7 @@ import { relaySseEagerBounded } from "../relay-eager";
147147
import { isWin32EagerRewrite, selectEagerPath } from "../../lib/bun-stream-caps";
148148
import { cancelBodyOnAbort } from "../../lib/abort";
149149
import {
150-
createResponsesItemIdPayloadRewrite,
150+
createResponsesItemIdRepairHandlers,
151151
hasResponsesItemIdRepair,
152152
} from "../responses-item-id-repair";
153153
import {
@@ -1818,12 +1818,13 @@ async function handleResponsesInner(
18181818
if (isEventStream && upstreamResponse.body) {
18191819
const repairConfig = route.provider.responsesItemIdRepair;
18201820
const needsClientRewrite = imageGenCallAliases.size > 0 || hasResponsesItemIdRepair(repairConfig);
1821+
const itemIdRepair = hasResponsesItemIdRepair(repairConfig)
1822+
? createResponsesItemIdRepairHandlers(repairConfig!, translatorBudget)
1823+
: undefined;
18211824
// Compose opt-in payload rewrites into one parse/stringify pass (image-gen restore first).
18221825
const payloadRewrites = [
18231826
createImageGenCallRestoreRewrite(imageGenCallAliases),
1824-
hasResponsesItemIdRepair(repairConfig)
1825-
? createResponsesItemIdPayloadRewrite(repairConfig!, translatorBudget)
1826-
: undefined,
1827+
itemIdRepair?.rewrite,
18271828
].filter((rewrite): rewrite is NonNullable<typeof rewrite> => rewrite !== undefined);
18281829
// #864: win32 rewrite traffic must never enter the tee()+JS-pull chain
18291830
// (Bun#32111 JS-sink segfault — text frames pass, the terminal block is
@@ -1887,7 +1888,12 @@ async function handleResponsesInner(
18871888
},
18881889
onClientCancel: () => options.onNativePassthroughCancel?.(),
18891890
onDone: () => unregisterTurn(turnAc),
1890-
}, win32EagerRewrite ? { rewriteBudget: translatorBudget } : undefined);
1891+
}, win32EagerRewrite
1892+
? {
1893+
rewriteBudget: translatorBudget,
1894+
...(itemIdRepair ? { trailer: itemIdRepair.trailer } : {}),
1895+
}
1896+
: undefined);
18911897
// selectEagerPath admits only no-rewrite traffic on both eligible platforms;
18921898
// win32 rewrite traffic reaches this relay too, but with the payload rewrite
18931899
// applied inline — never via an image/item-id JS pull wrapper (#32111, #864).
@@ -1961,7 +1967,12 @@ async function handleResponsesInner(
19611967
// relay is established practice (relayWithAbort, relaySseWithHeartbeat) and lets a
19621968
// mid-stream reset end with a clean response.failed terminal instead of a raw socket error.
19631969
const rewrittenBody = payloadRewrites.length > 0
1964-
? relaySseWithPayloadRewrite(nativeBody, composeSsePayloadRewrites(...payloadRewrites), translatorBudget)
1970+
? relaySseWithPayloadRewrite(
1971+
nativeBody,
1972+
composeSsePayloadRewrites(...payloadRewrites),
1973+
translatorBudget,
1974+
itemIdRepair ? { trailer: itemIdRepair.trailer } : undefined,
1975+
)
19651976
: nativeBody;
19661977
const clientBody = process.platform === "win32" && !needsClientRewrite
19671978
? nativeBody

src/server/sse-payload-rewrite.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -141,7 +141,7 @@ export function relaySseWithPayloadRewrite(
141141
while ((next = nextSseBlock(buffer))) {
142142
replaceBuffer(next.rest);
143143
const payload = sseDataPayload(next.block);
144-
if (payload) {
144+
if (payload !== null) {
145145
const rewrittenPayload = rewrite(payload);
146146
if (rewrittenPayload === null) continue;
147147
const block = rewrittenPayload !== payload
@@ -156,7 +156,7 @@ export function relaySseWithPayloadRewrite(
156156
}
157157
if (flushFinal && buffer.length > 0) {
158158
const payload = sseDataPayload(buffer);
159-
if (payload) {
159+
if (payload !== null) {
160160
const rewrittenPayload = rewrite(payload);
161161
if (rewrittenPayload !== null) {
162162
const block = rewrittenPayload !== payload

tests/responses-item-id-repair-deepseek.test.ts

Lines changed: 42 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -83,12 +83,53 @@ describe("DeepSeek Responses item-id repair", () => {
8383
expect(reasoningAdded.id).toMatch(/^rs_ocx_[0-9a-f]+_0$/);
8484
expect(messageAdded.id).toMatch(/^msg_ocx_[0-9a-f]+_1$/);
8585
expect(textDelta.item_id).toBe(messageAdded.id);
86-
expect(completed.id).toMatch(/^resp_ocx_[0-9a-f]+$/);
86+
expect(completed.id).toMatch(/^resp_ocx_[0-9a-f]+_\d+$/);
8787
expect(completed.output[0].id).toBe(reasoningAdded.id);
8888
expect(String(completed.output[0].encrypted_content)).toMatch(/^ocxr1:/);
8989
expect(completed.output[1].id).toBe(messageAdded.id);
9090
});
9191

92+
93+
test("mints unique response ids for distinct upstream response ids", async () => {
94+
const upstream = [
95+
`event: response.completed\ndata: {"type":"response.completed","response":{"id":"11111111-1111-1111-1111-111111111111","status":"completed","output":[]}}\n\n`,
96+
`event: response.completed\ndata: {"type":"response.completed","response":{"id":"22222222-2222-2222-2222-222222222222","status":"completed","output":[]}}\n\n`,
97+
].join("");
98+
const events = await parseSse(await readAll(relaySseWithResponsesItemIdRepair(streamFromText(upstream), {
99+
rewriteNonCanonicalIds: true,
100+
})));
101+
const first = (events[0].response as { id: string }).id;
102+
const second = (events[1].response as { id: string }).id;
103+
expect(first).toMatch(/^resp_ocx_[0-9a-f]+_0$/);
104+
expect(second).toMatch(/^resp_ocx_[0-9a-f]+_1$/);
105+
expect(first).not.toBe(second);
106+
});
107+
108+
test("mints reasoning ids when rewriteNonCanonicalIds is enabled without repairMissingTerminalIds", async () => {
109+
const upstream = [
110+
`event: response.output_item.added\ndata: {"type":"response.output_item.added","output_index":0,"item":{"type":"reasoning","summary":[]}}\n\n`,
111+
`event: response.output_item.done\ndata: {"type":"response.output_item.done","output_index":0,"item":{"type":"reasoning","summary":[]}}\n\n`,
112+
].join("");
113+
const events = await parseSse(await readAll(relaySseWithResponsesItemIdRepair(streamFromText(upstream), {
114+
rewriteNonCanonicalIds: true,
115+
})));
116+
const added = events[0].item as Record<string, unknown>;
117+
const done = events[1].item as Record<string, unknown>;
118+
expect(added.id).toMatch(/^rs_ocx_[0-9a-f]+_0$/);
119+
expect(done.id).toBe(added.id);
120+
});
121+
122+
test("does not duplicate an upstream [DONE] trailer", async () => {
123+
const upstream = [
124+
`event: response.completed\ndata: {"type":"response.completed","response":{"id":"58786d76-18be-43d9-b6f5-8922796fbe28","status":"completed","output":[]}}\n\n`,
125+
`data: [DONE]\n\n`,
126+
].join("");
127+
const repaired = await readAll(relaySseWithResponsesItemIdRepair(streamFromText(upstream), {
128+
rewriteNonCanonicalIds: true,
129+
}));
130+
expect(repaired.match(/data: \[DONE\]/g)).toHaveLength(1);
131+
});
132+
92133
test("opt-in flag is recognized", () => {
93134
expect(hasResponsesItemIdRepair({ rewriteNonCanonicalIds: true })).toBe(true);
94135
});

0 commit comments

Comments
 (0)