diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 3a2da161e..f8581ef30 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -19,6 +19,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 the Anthropic provider retrying persistent account usage/quota limits (e.g. `429 "This request would exceed your account's rate limit"`, `usage_limit_reached`) as if they were transient. Because the error text contains "rate limit", `isProviderRetryableError` matched it and the stream retry loop looped through its 2s/4s/8s backoff (then the `streamSimple` a/b/c policy re-minted the credential and ran the whole thing again) before surfacing the failure — even though the server's `retry-after` parked the account for minutes-to-hours. These errors are now recognized via `isUsageLimitError` and surfaced immediately to the credential-rotation layer, so e.g. `omp dry-balance --bench` reports a rate-limited account as failed at once instead of appearing to hang. - 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 = {