From 7e5eee8718a20d667a6b9714ffa05a299f90810a Mon Sep 17 00:00:00 2001 From: serhiizghama Date: Sat, 4 Jul 2026 03:20:49 +0700 Subject: [PATCH 1/4] fix(responses): unique output_index for message + tool call When assistant text streamed before a tool call, the message output item opened at outputItemIndex without advancing it, so the following tool call reused the same output_index. The message's done events also read the mutated counter, drifting from its added index. Reserve a dedicated index for the message and order the final output array by output_index. --- .../tools/convert-streaming-to-responses.ts | 94 ++++++++++++------- 1 file changed, 59 insertions(+), 35 deletions(-) diff --git a/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts b/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts index 3ed140c592..9db38bdbe5 100644 --- a/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts +++ b/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts @@ -9,6 +9,8 @@ interface StreamingState { model: string; createdAt: number; outputItemIndex: number; + messageOutputIndex: number; + reasoningOutputIndex: number; contentPartStarted: boolean; outputItemStarted: boolean; messageId: string; @@ -65,6 +67,8 @@ export function createStreamingState( model, createdAt: Math.floor(Date.now() / 1000), outputItemIndex: 0, + messageOutputIndex: 0, + reasoningOutputIndex: 0, contentPartStarted: false, outputItemStarted: false, messageId: `msg_${shortid(24)}`, @@ -256,10 +260,11 @@ export function processStreamChunk( if (delta.reasoning) { if (!state.reasoningStarted) { state.reasoningStarted = true; + state.reasoningOutputIndex = state.outputItemIndex; events.push( emitEvent(state, "response.output_item.added", { type: "response.output_item.added", - output_index: state.outputItemIndex, + output_index: state.reasoningOutputIndex, item: { type: "reasoning", id: state.reasoningId, @@ -330,11 +335,14 @@ export function processStreamChunk( if (state.reasoningStarted) { state.outputItemIndex++; } + // Claim this slot and advance so a later tool call gets its own + // output_index instead of reusing the message's. + state.messageOutputIndex = state.outputItemIndex++; events.push( emitEvent(state, "response.output_item.added", { type: "response.output_item.added", - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, item: { type: "message", id: state.messageId, @@ -352,7 +360,7 @@ export function processStreamChunk( emitEvent(state, "response.content_part.added", { type: "response.content_part.added", item_id: state.messageId, - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, content_index: 0, part: { type: "output_text", text: "", annotations: [] }, }), @@ -364,7 +372,7 @@ export function processStreamChunk( emitEvent(state, "response.output_text.delta", { type: "response.output_text.delta", item_id: state.messageId, - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, content_index: 0, delta: delta.content, }), @@ -420,7 +428,7 @@ export function createCompletionEvents(state: StreamingState): SSEEvent[] { emitEvent(state, "response.output_text.done", { type: "response.output_text.done", item_id: state.messageId, - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, content_index: 0, text: state.fullContent.join(""), }), @@ -429,7 +437,7 @@ export function createCompletionEvents(state: StreamingState): SSEEvent[] { emitEvent(state, "response.content_part.done", { type: "response.content_part.done", item_id: state.messageId, - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, content_index: 0, part: { type: "output_text", @@ -445,7 +453,7 @@ export function createCompletionEvents(state: StreamingState): SSEEvent[] { events.push( emitEvent(state, "response.output_item.done", { type: "response.output_item.done", - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, item: { type: "message", id: state.messageId, @@ -487,53 +495,69 @@ export function createCompletionEvents(state: StreamingState): SSEEvent[] { status = "incomplete"; } - // Build final output array - const output: Record[] = []; + // Build final output array in output_index order so it matches the + // streaming events — a message streamed before a tool call keeps its + // lower index instead of always being listed last. + const output: { index: number; item: Record }[] = []; if (state.reasoningStarted) { output.push({ - type: "reasoning", - id: state.reasoningId, - summary: [ - { - type: "summary_text", - text: state.fullReasoning.join(""), - }, - ], + index: state.reasoningOutputIndex, + item: { + type: "reasoning", + id: state.reasoningId, + summary: [ + { + type: "summary_text", + text: state.fullReasoning.join(""), + }, + ], + }, }); } for (const tc of state.toolCalls.values()) { output.push({ - type: "function_call", - id: tc.id, - call_id: tc.callId, - name: tc.name, - arguments: tc.arguments, - status: "completed", + index: tc.outputIndex, + item: { + type: "function_call", + id: tc.id, + call_id: tc.callId, + name: tc.name, + arguments: tc.arguments, + status: "completed", + }, }); } if (state.fullContent.length > 0) { output.push({ - type: "message", - id: state.messageId, - role: "assistant", - content: [ - { - type: "output_text", - text: state.fullContent.join(""), - annotations: [], - }, - ], - status: "completed", + index: state.messageOutputIndex, + item: { + type: "message", + id: state.messageId, + role: "assistant", + content: [ + { + type: "output_text", + text: state.fullContent.join(""), + annotations: [], + }, + ], + status: "completed", + }, }); } + output.sort((a, b) => a.index - b.index); + events.push( emitEvent(state, "response.completed", { type: "response.completed", - response: buildResponsePayload(state, { status, output }), + response: buildResponsePayload(state, { + status, + output: output.map((o) => o.item), + }), }), ); From 2476be4ba09d143f81d3af01f4b1e16999a83345 Mon Sep 17 00:00:00 2001 From: serhiizghama Date: Sat, 4 Jul 2026 03:20:53 +0700 Subject: [PATCH 2/4] test(responses): cover message/tool output_index ordering --- apps/gateway/src/responses/responses.spec.ts | 62 ++++++++++++++++++++ 1 file changed, 62 insertions(+) diff --git a/apps/gateway/src/responses/responses.spec.ts b/apps/gateway/src/responses/responses.spec.ts index 67ca77fe13..50e79dae51 100644 --- a/apps/gateway/src/responses/responses.spec.ts +++ b/apps/gateway/src/responses/responses.spec.ts @@ -633,6 +633,68 @@ describe("streaming conversion", () => { expect(JSON.parse(fcDone!.data).item.name).toBe("get_weather"); }); + it("gives the message and a following tool call distinct output_index values", () => { + const state = createStreamingState("gpt-4o-mini"); + const contentEvents = processStreamChunk( + { choices: [{ delta: { content: "Let me check" } }] }, + state, + ); + const toolEvents = processStreamChunk( + { + choices: [ + { + delta: { + tool_calls: [ + { + index: 0, + id: "call_1", + function: { name: "get_weather", arguments: "{}" }, + }, + ], + }, + }, + ], + }, + state, + ); + + const msgAdded = JSON.parse( + contentEvents.find( + (e) => + e.event === "response.output_item.added" && + JSON.parse(e.data).item.type === "message", + )!.data, + ); + const fcAdded = JSON.parse( + toolEvents.find( + (e) => + e.event === "response.output_item.added" && + JSON.parse(e.data).item.type === "function_call", + )!.data, + ); + expect(msgAdded.output_index).not.toBe(fcAdded.output_index); + + const events = createCompletionEvents(state); + const msgDone = JSON.parse( + events.find( + (e) => + e.event === "response.output_item.done" && + JSON.parse(e.data).item.type === "message", + )!.data, + ); + // The message keeps the same output_index across added and done. + expect(msgDone.output_index).toBe(msgAdded.output_index); + + // The final output array is ordered by output_index (message before tool). + const completed = JSON.parse( + events.find((e) => e.event === "response.completed")!.data, + ); + const types = completed.response.output.map( + (o: Record) => o.type, + ); + expect(types).toEqual(["message", "function_call"]); + }); + it("maps length finish_reason to incomplete status in streaming", () => { const state = createStreamingState("gpt-4o-mini"); processStreamChunk( From e69c307619e193ddb0527819b4f1f54f029cc844 Mon Sep 17 00:00:00 2001 From: serhiizghama Date: Tue, 14 Jul 2026 13:26:40 +0700 Subject: [PATCH 3/4] fix(gateway): claim reasoning output slot on start The close-reasoning branch ran on every tool-call chunk while no message had started, so reasoning followed by multi-chunk tool calls inflated the shared index and a later tool call got an output_index past its final response.output position. Claim the reasoning slot once when reasoning starts and drop the deferred increments. --- apps/gateway/src/responses/responses.spec.ts | 123 ++++++++++++++++++ .../tools/convert-streaming-to-responses.ts | 13 +- 2 files changed, 126 insertions(+), 10 deletions(-) diff --git a/apps/gateway/src/responses/responses.spec.ts b/apps/gateway/src/responses/responses.spec.ts index 3051f9335b..4f7f740c17 100644 --- a/apps/gateway/src/responses/responses.spec.ts +++ b/apps/gateway/src/responses/responses.spec.ts @@ -841,6 +841,129 @@ describe("streaming conversion", () => { expect(types).toEqual(["message", "function_call"]); }); + it("gives reasoning and a following message distinct output_index values", () => { + const state = createStreamingState("gpt-4o-mini"); + const reasoningEvents = processStreamChunk( + { choices: [{ delta: { reasoning: "let me think" } }] }, + state, + ); + const contentEvents = processStreamChunk( + { choices: [{ delta: { content: "Here is the answer" } }] }, + state, + ); + + const reasoningAdded = JSON.parse( + reasoningEvents.find( + (e) => + e.event === "response.output_item.added" && + JSON.parse(e.data).item.type === "reasoning", + )!.data, + ); + const msgAdded = JSON.parse( + contentEvents.find( + (e) => + e.event === "response.output_item.added" && + JSON.parse(e.data).item.type === "message", + )!.data, + ); + expect(reasoningAdded.output_index).not.toBe(msgAdded.output_index); + + const events = createCompletionEvents(state); + const completed = JSON.parse( + events.find((e) => e.event === "response.completed")!.data, + ); + const types = completed.response.output.map( + (o: Record) => o.type, + ); + expect(types).toEqual(["reasoning", "message"]); + }); + + it("keeps tool-call output_index aligned when reasoning precedes multi-chunk tool calls", () => { + const state = createStreamingState("gpt-4o-mini"); + processStreamChunk( + { choices: [{ delta: { reasoning: "thinking" } }] }, + state, + ); + // First tool call opens. + processStreamChunk( + { + choices: [ + { + delta: { + tool_calls: [ + { + index: 0, + id: "call_a", + function: { name: "get_weather", arguments: "" }, + }, + ], + }, + }, + ], + }, + state, + ); + // Extra argument chunk for the same tool call — previously this + // bumped the shared index on every chunk, inflating later slots. + processStreamChunk( + { + choices: [ + { + delta: { + tool_calls: [ + { index: 0, function: { arguments: '{"city":"NYC"}' } }, + ], + }, + }, + ], + }, + state, + ); + // Second tool call opens. + const secondToolEvents = processStreamChunk( + { + choices: [ + { + delta: { + tool_calls: [ + { + index: 1, + id: "call_b", + function: { name: "get_time", arguments: "{}" }, + }, + ], + }, + }, + ], + }, + state, + ); + + const secondAdded = JSON.parse( + secondToolEvents.find( + (e) => + e.event === "response.output_item.added" && + JSON.parse(e.data).item.type === "function_call", + )!.data, + ); + + const events = createCompletionEvents(state); + const completed = JSON.parse( + events.find((e) => e.event === "response.completed")!.data, + ); + const output = completed.response.output as Record[]; + + // The streamed output_index for the second tool call must match its + // position in the final, index-sorted response.output array. + const finalPos = output.findIndex( + (o) => o.type === "function_call" && o.name === "get_time", + ); + expect(secondAdded.output_index).toBe(finalPos); + + const types = output.map((o) => o.type); + expect(types).toEqual(["reasoning", "function_call", "function_call"]); + }); + it("maps length finish_reason to incomplete status in streaming", () => { const state = createStreamingState("gpt-4o-mini"); processStreamChunk( diff --git a/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts b/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts index 02654141db..632157025b 100644 --- a/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts +++ b/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts @@ -285,7 +285,9 @@ export function processStreamChunk( if (delta.reasoning) { if (!state.reasoningStarted) { state.reasoningStarted = true; - state.reasoningOutputIndex = state.outputItemIndex; + // Claim the reasoning slot immediately so later tool calls and the + // message get their own indices instead of reusing this one. + state.reasoningOutputIndex = state.outputItemIndex++; events.push( emitEvent(state, "response.output_item.added", { type: "response.output_item.added", @@ -303,11 +305,6 @@ export function processStreamChunk( // Handle tool_calls delta if (delta.tool_calls) { - // If reasoning was streamed but tool calls arrive (no content), close reasoning index - if (state.reasoningStarted && !state.outputItemStarted) { - state.outputItemIndex++; - } - for (const tc of delta.tool_calls) { const existing = state.toolCalls.get(tc.index); if (!existing) { @@ -356,10 +353,6 @@ export function processStreamChunk( if (delta.content) { if (!state.outputItemStarted) { state.outputItemStarted = true; - // If reasoning was streamed, close it first - if (state.reasoningStarted) { - state.outputItemIndex++; - } // Claim this slot and advance so a later tool call gets its own // output_index instead of reusing the message's. state.messageOutputIndex = state.outputItemIndex++; From a74a3e4d9bfdd4947bdb81b86b58a19db142e7a7 Mon Sep 17 00:00:00 2001 From: Luca Steeb Date: Tue, 14 Jul 2026 20:06:00 +0100 Subject: [PATCH 4/4] fix(gateway): annotation events use message output_index Co-Authored-By: Claude Fable 5 --- apps/gateway/src/responses/responses.spec.ts | 44 +++++++++++++++++++ .../tools/convert-streaming-to-responses.ts | 2 +- 2 files changed, 45 insertions(+), 1 deletion(-) diff --git a/apps/gateway/src/responses/responses.spec.ts b/apps/gateway/src/responses/responses.spec.ts index 8657fb36c6..78875e2fc8 100644 --- a/apps/gateway/src/responses/responses.spec.ts +++ b/apps/gateway/src/responses/responses.spec.ts @@ -964,6 +964,50 @@ describe("streaming conversion", () => { expect(types).toEqual(["reasoning", "function_call", "function_call"]); }); + it("emits annotation events with the message's output_index", () => { + const state = createStreamingState("gpt-4o-mini"); + const contentEvents = processStreamChunk( + { choices: [{ delta: { content: "According to the docs" } }] }, + state, + ); + const annotationEvents = processStreamChunk( + { + choices: [ + { + delta: { + annotations: [ + { + type: "url_citation", + url_citation: { + url: "https://example.com", + title: "Example", + start_index: 0, + end_index: 10, + }, + }, + ], + }, + }, + ], + }, + state, + ); + + const msgAdded = JSON.parse( + contentEvents.find( + (e) => + e.event === "response.output_item.added" && + JSON.parse(e.data).item.type === "message", + )!.data, + ); + const annAdded = JSON.parse( + annotationEvents.find( + (e) => e.event === "response.output_text.annotation.added", + )!.data, + ); + expect(annAdded.output_index).toBe(msgAdded.output_index); + }); + it("maps length finish_reason to incomplete status in streaming", () => { const state = createStreamingState("gpt-4o-mini"); processStreamChunk( diff --git a/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts b/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts index 632157025b..c2048f8463 100644 --- a/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts +++ b/apps/gateway/src/responses/tools/convert-streaming-to-responses.ts @@ -409,7 +409,7 @@ export function processStreamChunk( emitEvent(state, "response.output_text.annotation.added", { type: "response.output_text.annotation.added", item_id: state.messageId, - output_index: state.outputItemIndex, + output_index: state.messageOutputIndex, content_index: 0, annotation_index: annotationIndex, annotation,