Merge remote-tracking branch 'origin/farm/b74e4820/hashline-replace-drops-body-lines'

This commit is contained in:
can1357
2026-06-08 03:11:51 +02:00
3 changed files with 153 additions and 53 deletions
+1
View File
@@ -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
@@ -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<number, OpenFunctionCall>();
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);
@@ -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<string, unknown>).output_index)).toEqual([0, 1]);
expect(argumentDeltas.map(f => (f.data as Record<string, unknown>).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<string, unknown>).output_index)).toEqual([0, 1]);
expect(argumentDone.map(f => (f.data as Record<string, unknown>).arguments)).toEqual([
'{"input":"first"}',
'{"path":"second"}',
]);
const doneItems = frames
.filter(f => f.event === "response.output_item.done")
.map(f => (f.data as Record<string, unknown>).item as Record<string, unknown>)
.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 = {