From 65bb1c8c4d24cda93e348bc3873bdb57ed9aa46c Mon Sep 17 00:00:00 2001 From: roboomp Date: Mon, 8 Jun 2026 01:02:22 +0000 Subject: [PATCH] fix(ai): routed responses tool deltas by content index Codex review on #2082 found that the Responses compatibility encoder still kept a singleton open function-call item. When MiniMax object-argument chunks are flushed late by contentIndex, a later parallel toolcall_start could close the first item and make the late delta append to the second item instead. Keep OpenFunctionCall state in a map keyed by contentIndex, allocate Responses output indexes when items open, and close each function item from its matching toolcall_end. Late deltas now use event.contentIndex, preserving deferred object-argument flushes for the original tool call even after later parallel starts. Added an encodeStream regression that starts two parallel calls, emits the first call's arguments only after the second start, and verifies both argument delta/done events and output_item.done payloads stay attached to their own output indexes. Fixes #2080 --- packages/ai/CHANGELOG.md | 1 + .../src/providers/openai-responses-server.ts | 135 +++++++++++------- .../auth-gateway-openai-responses.test.ts | 70 +++++++++ 3 files changed, 153 insertions(+), 53 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 0df633642..64871ad49 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -15,6 +15,7 @@ - Fixed duplicate upstream `tool_call_id` values collapsing distinct tool calls during message transformation, preserving one call/result pairing per emitted tool call before provider replay and keeping generated duplicate IDs distinct after OpenAI/Mistral wire-length caps. ([#2055](https://github.com/can1357/oh-my-pi/issues/2055)) - Fixed MiniMax-compatible OpenAI-completions hosts losing tool-call argument content when `function.arguments` is streamed as an object across more than one delta. The accumulator added in #1776 wrote `block.partialArgs = rawArgs` per chunk, so every chunk but the last was overwritten — for an `edit` call this surfaced as a tail-slice of the patch text being applied (e.g. a single-line `replace 91..91:` body extending the deletion across the surrounding rows). Chunks are now shallow-merged; for shared string keys, `startsWith` distinguishes cumulative restatements (take the latest) from per-chunk-delta fragments (concatenate). Per-chunk `toolcall_delta` emission for the object branch is suppressed (the previous code emitted `JSON.stringify(rawArgs)` per chunk, which fed downstream concat consumers — `packages/agent/src/proxy.ts`, `openai-chat-server`, `openai-responses-server`, `anthropic-messages-server` — an invalid sequence like `{"input":"a"}{"input":"b"}`); the merged object is flushed instead as a single concat-safe delta in `finishToolCallBlock` before `toolcall_end`, so accumulators reconstruct the args correctly. The single-chunk shape covered by the existing #1776 regression test stays correct end-to-end. ([#2080](https://github.com/can1357/oh-my-pi/issues/2080)) +- Fixed the OpenAI Responses compatibility server misrouting late `toolcall_delta` events for earlier parallel tool calls after a later `toolcall_start`. The encoder now keeps OpenFunctionCall state by content index, allocates output indexes at item start, and closes each tool item by its own `toolcall_end`, preserving deferred MiniMax object-argument flushes for the matching call. ([#2080](https://github.com/can1357/oh-my-pi/issues/2080)) ## [15.10.1] - 2026-06-07 diff --git a/packages/ai/src/providers/openai-responses-server.ts b/packages/ai/src/providers/openai-responses-server.ts index 7fe0c3bf5..de462a831 100644 --- a/packages/ai/src/providers/openai-responses-server.ts +++ b/packages/ai/src/providers/openai-responses-server.ts @@ -698,6 +698,7 @@ interface OpenFunctionCall { kind: "function_call"; itemId: string; outputIndex: number; + contentIndex: number; callId: string; name: string; argsText: string; @@ -729,7 +730,9 @@ export function encodeStream( let createdAt = Math.floor(Date.now() / 1000); let outputIndex = 0; const state: { open: OpenItem | null } = { open: null }; + const openFunctionCalls = new Map(); const finishedItems: OutputItem[] = []; + const allocateOutputIndex = (): number => outputIndex++; const responseSnapshot = (status: ResponseStatus, output: OutputItem[] | []) => ({ id: responseId, @@ -742,6 +745,7 @@ export function encodeStream( }); const openMessage = (): OpenMessage => { + const itemOutputIndex = allocateOutputIndex(); const itemId = makeMsgId(); const item = { type: "message" as const, @@ -750,11 +754,11 @@ export function encodeStream( role: "assistant" as const, content: [] as Array<{ type: "output_text"; text: string; annotations: never[] }>, }; - emit("response.output_item.added", { output_index: outputIndex, item }); + emit("response.output_item.added", { output_index: itemOutputIndex, item }); const next: OpenMessage = { kind: "message", itemId, - outputIndex, + outputIndex: itemOutputIndex, contentIndex: 0, currentPartText: "", content: [], @@ -764,6 +768,7 @@ export function encodeStream( }; const openReasoning = (partial: AssistantMessage, contentIndex: number): OpenReasoning => { + const itemOutputIndex = allocateOutputIndex(); const part = partial.content[contentIndex]; const itemId = part && part.type === "thinking" ? reasoningItemId(part) : makeReasoningId(); const item = { @@ -771,22 +776,23 @@ export function encodeStream( id: itemId, summary: [] as Array<{ type: "summary_text"; text: string }>, }; - emit("response.output_item.added", { output_index: outputIndex, item }); + emit("response.output_item.added", { output_index: itemOutputIndex, item }); // Open the summary part. Real OpenAI streams summary text in the // canonical `reasoning_summary_*` lifecycle; pi-ai's own decoder // reads `summary[].text` from the eventual `output_item.done`. emit("response.reasoning_summary_part.added", { item_id: itemId, - output_index: outputIndex, + output_index: itemOutputIndex, summary_index: 0, part: { type: "summary_text", text: "" }, }); - const next: OpenReasoning = { kind: "reasoning", itemId, outputIndex, reasoningText: "" }; + const next: OpenReasoning = { kind: "reasoning", itemId, outputIndex: itemOutputIndex, reasoningText: "" }; state.open = next; return next; }; const openToolCall = (partial: AssistantMessage, contentIndex: number): OpenFunctionCall => { + const itemOutputIndex = allocateOutputIndex(); const part = partial.content[contentIndex]; const tc = part && part.type === "toolCall" ? part : undefined; const customWireName: string | undefined = @@ -814,20 +820,65 @@ export function encodeStream( arguments: "", status: "in_progress", }; - emit("response.output_item.added", { output_index: outputIndex, item }); + emit("response.output_item.added", { output_index: itemOutputIndex, item }); const next: OpenFunctionCall = { kind: "function_call", itemId, - outputIndex, + outputIndex: itemOutputIndex, + contentIndex, callId, name, argsText: "", ...(isCustom ? { customWireName } : {}), }; + openFunctionCalls.set(contentIndex, next); state.open = next; return next; }; + const closeFunctionCall = (call: OpenFunctionCall): void => { + const text = call.argsText ?? ""; + if (call.customWireName) { + const item = { + type: "custom_tool_call", + id: call.itemId, + call_id: call.callId ?? "", + name: call.customWireName, + input: text, + status: "completed", + }; + emit("response.output_item.done", { output_index: call.outputIndex, item }); + finishedItems.push({ + type: "custom_tool_call", + id: call.itemId, + call_id: call.callId ?? "", + name: call.customWireName, + input: text, + status: "completed", + }); + } else { + const item = { + type: "function_call", + id: call.itemId, + call_id: call.callId ?? "", + name: call.name ?? "", + arguments: text, + status: "completed", + }; + emit("response.output_item.done", { output_index: call.outputIndex, item }); + finishedItems.push({ + type: "function_call", + id: call.itemId, + call_id: call.callId ?? "", + name: call.name ?? "", + arguments: text, + status: "completed", + }); + } + openFunctionCalls.delete(call.contentIndex); + if (state.open === call) state.open = null; + }; + const closeOpen = () => { if (!state.open) return; if (state.open.kind === "message") { @@ -846,6 +897,7 @@ export function encodeStream( status: "completed", content: state.open.content, }); + state.open = null; } else if (state.open.kind === "reasoning") { const summary = [{ type: "summary_text" as const, text: state.open.reasoningText ?? "" }]; const item = { @@ -859,50 +911,23 @@ export function encodeStream( id: state.open.itemId, summary, }); + state.open = null; } else { - const text = state.open.argsText ?? ""; - if (state.open.customWireName) { - const item = { - type: "custom_tool_call", - id: state.open.itemId, - call_id: state.open.callId ?? "", - name: state.open.customWireName, - input: text, - status: "completed", - }; - emit("response.output_item.done", { output_index: state.open.outputIndex, item }); - finishedItems.push({ - type: "custom_tool_call", - id: state.open.itemId, - call_id: state.open.callId ?? "", - name: state.open.customWireName, - input: text, - status: "completed", - }); - } else { - const item = { - type: "function_call", - id: state.open.itemId, - call_id: state.open.callId ?? "", - name: state.open.name ?? "", - arguments: text, - status: "completed", - }; - emit("response.output_item.done", { output_index: state.open.outputIndex, item }); - finishedItems.push({ - type: "function_call", - id: state.open.itemId, - call_id: state.open.callId ?? "", - name: state.open.name ?? "", - arguments: text, - status: "completed", - }); - } + closeFunctionCall(state.open); } - outputIndex++; - state.open = null; }; + const closeOpenFunctionCalls = (): void => { + for (const call of [...openFunctionCalls.values()]) { + closeFunctionCall(call); + } + }; + + const functionCallForEvent = (contentIndex: number): OpenFunctionCall | undefined => { + const byIndex = openFunctionCalls.get(contentIndex); + if (byIndex) return byIndex; + return state.open?.kind === "function_call" ? state.open : undefined; + }; try { let finalMessage: AssistantMessage | null = null; let failureMessage: AssistantMessage | null = null; @@ -941,6 +966,7 @@ export function encodeStream( cur = state.open; cur.currentPartText = ""; } else { + closeOpenFunctionCalls(); if (state.open) closeOpen(); cur = openMessage(); } @@ -992,6 +1018,7 @@ export function encodeStream( break; } case "thinking_start": { + closeOpenFunctionCalls(); if (state.open) closeOpen(); openReasoning(ev.partial, ev.contentIndex); break; @@ -1029,13 +1056,13 @@ export function encodeStream( break; } case "toolcall_start": { - if (state.open) closeOpen(); + if (state.open && state.open.kind !== "function_call") closeOpen(); openToolCall(ev.partial, ev.contentIndex); break; } case "toolcall_delta": { - if (state.open?.kind !== "function_call") break; - const cur: OpenFunctionCall = state.open; + const cur = functionCallForEvent(ev.contentIndex); + if (!cur) break; cur.argsText += ev.delta; if (cur.customWireName) { emit("response.custom_tool_call_input.delta", { @@ -1053,8 +1080,8 @@ export function encodeStream( break; } case "toolcall_end": { - if (state.open?.kind !== "function_call") break; - const cur: OpenFunctionCall = state.open; + const cur = functionCallForEvent(ev.contentIndex); + if (!cur) break; // Promote possibly-late info from the canonical ToolCall. const tc = ev.toolCall; if (tc.customWireName && !cur.customWireName) cur.customWireName = tc.customWireName; @@ -1087,7 +1114,7 @@ export function encodeStream( name: cur.name, }); } - closeOpen(); + closeFunctionCall(cur); break; } case "done": { @@ -1102,6 +1129,7 @@ export function encodeStream( } if (failureMessage) { + closeOpenFunctionCalls(); if (state.open) closeOpen(); controller.enqueue( encoder.encode( @@ -1120,6 +1148,7 @@ export function encodeStream( return; } + closeOpenFunctionCalls(); if (state.open) closeOpen(); const message = finalMessage ?? ((await events.result().catch(() => null)) as AssistantMessage | null); diff --git a/packages/ai/test/auth-gateway-openai-responses.test.ts b/packages/ai/test/auth-gateway-openai-responses.test.ts index c6cbd8c9b..091928abf 100644 --- a/packages/ai/test/auth-gateway-openai-responses.test.ts +++ b/packages/ai/test/auth-gateway-openai-responses.test.ts @@ -495,6 +495,76 @@ describe("openai-responses encodeStream", () => { expect(output[2]!.id).not.toBe(output[2]!.call_id); }); + it("routes late tool-call deltas by contentIndex after later parallel starts", async () => { + const stream = new AssistantMessageEventStream(); + const base: AssistantMessage = { + role: "assistant", + api: "openai-responses", + provider: "openai", + model: "gpt-5", + content: [], + usage: zeroUsage(), + stopReason: "toolUse", + timestamp: 1_700_000_000_000, + }; + const callA = { type: "toolCall" as const, id: "call_a", name: "edit", arguments: {} }; + const callB = { type: "toolCall" as const, id: "call_b", name: "read", arguments: {} }; + const partialA: AssistantMessage = { ...base, content: [callA] }; + const partialBoth: AssistantMessage = { ...base, content: [callA, callB] }; + const finalMessage: AssistantMessage = { + ...base, + content: [ + { ...callA, arguments: { input: "first" } }, + { ...callB, arguments: { path: "second" } }, + ], + }; + + queueMicrotask(() => { + stream.push({ type: "start", partial: base }); + stream.push({ type: "toolcall_start", contentIndex: 0, partial: partialA }); + stream.push({ type: "toolcall_start", contentIndex: 1, partial: partialBoth }); + stream.push({ type: "toolcall_delta", contentIndex: 0, delta: '{"input":"first"}', partial: partialBoth }); + stream.push({ + type: "toolcall_end", + contentIndex: 0, + toolCall: { ...callA, arguments: { input: "first" } }, + partial: partialBoth, + }); + stream.push({ type: "toolcall_delta", contentIndex: 1, delta: '{"path":"second"}', partial: partialBoth }); + stream.push({ + type: "toolcall_end", + contentIndex: 1, + toolCall: { ...callB, arguments: { path: "second" } }, + partial: partialBoth, + }); + stream.push({ type: "done", reason: "toolUse", message: finalMessage }); + }); + + const raw = await collectStream(encodeStream(stream, "gpt-5-requested")); + const frames = parseSse(raw); + const argumentDeltas = frames.filter(f => f.event === "response.function_call_arguments.delta"); + expect(argumentDeltas.map(f => (f.data as Record).output_index)).toEqual([0, 1]); + expect(argumentDeltas.map(f => (f.data as Record).delta)).toEqual([ + '{"input":"first"}', + '{"path":"second"}', + ]); + + const argumentDone = frames.filter(f => f.event === "response.function_call_arguments.done"); + expect(argumentDone.map(f => (f.data as Record).output_index)).toEqual([0, 1]); + expect(argumentDone.map(f => (f.data as Record).arguments)).toEqual([ + '{"input":"first"}', + '{"path":"second"}', + ]); + + const doneItems = frames + .filter(f => f.event === "response.output_item.done") + .map(f => (f.data as Record).item as Record) + .filter(item => item.type === "function_call"); + + expect(doneItems).toHaveLength(2); + expect(doneItems[0]).toMatchObject({ call_id: "call_a", name: "edit", arguments: '{"input":"first"}' }); + expect(doneItems[1]).toMatchObject({ call_id: "call_b", name: "read", arguments: '{"path":"second"}' }); + }); it("emits response.incomplete for length-limited streams", async () => { const stream = new AssistantMessageEventStream(); const message: AssistantMessage = {