fix(ai): routed codex responses arg deltas by item_id
The Codex Responses stream runtime tracked a singleton
`currentItem`/`currentBlock`. With more than one tool call open
concurrently every `response.function_call_arguments.delta` was
appended to whichever item was added most recently, and the next
`response.output_item.done` for the earlier call overwrote the
sibling's stored arguments. On the agent loop's `task` tool this
surfaced as `tasks: Invalid input: expected array, received undefined`.
Open items are now tracked in a `Map<string, CodexOpenItem>` keyed by
`item.id`, each entry carrying its own `block` and `contentIndex`.
`response.function_call_arguments.{delta,done}`,
`response.custom_tool_call_input.{delta,done}`, and
`response.output_item.done` route through `openItemForEvent` and
operate on the matching entry's block; the legacy singleton-current
fallback only kicks in when an event omits `item_id`. A delta whose
keyed item already closed is dropped instead of leaking into a
sibling, and `toolcall_delta` / `toolcall_end` stream events emit the
right `contentIndex` for each call. Recovery sites converge on a new
`resetCodexStreamAccumulators` helper that clears the open-items map
in lockstep with `currentItem`/`currentBlock`/`nativeOutputItems`.
Regression test (`packages/ai/test/openai-codex-stream.test.ts`)
interleaves two function-call argument streams plus a stale
post-close delta and asserts per-call argument integrity and
per-call stream `contentIndex`.
Fixes #2619
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed Codex Responses stream mis-routing interleaved `function_call_arguments.delta` events when more than one tool call was open concurrently. The runtime tracked a singleton `currentItem`/`currentBlock`, so every delta — regardless of `item_id` — was appended to whichever item was most recently added, and `output_item.done` for the earlier call then overwrote a sibling's stored arguments (visible as `tasks: Invalid input: expected array, received undefined` on the `task` tool). Open items are now keyed by `item_id` with `output_index` fallback; deltas/done events route to the matching block, late deltas whose item already closed are dropped instead of corrupting a sibling, and `toolcall_*` stream events emit the right `contentIndex` per call ([#2619](https://github.com/can1357/oh-my-pi/issues/2619)).
|
||||
|
||||
## [15.13.1] - 2026-06-15
|
||||
|
||||
### Fixed
|
||||
|
||||
@@ -277,11 +277,27 @@ interface CodexRequestSetup {
|
||||
websocketFirstEventTimeoutMs: number | undefined;
|
||||
}
|
||||
|
||||
interface CodexOpenItem {
|
||||
item: CodexEventItem;
|
||||
block: CodexOutputBlock | null;
|
||||
/** Index of {@link block} in `output.content`; `-1` when no block was created for this item. */
|
||||
contentIndex: number;
|
||||
}
|
||||
|
||||
interface CodexStreamRuntime {
|
||||
eventStream: AsyncGenerator<Record<string, unknown>>;
|
||||
requestBodyForState: RequestBody;
|
||||
transport: CodexTransport;
|
||||
websocketState?: CodexWebSocketSessionState;
|
||||
/**
|
||||
* Items open on the wire — every `response.output_item.added` registers
|
||||
* here; `output_item.done` removes. Keyed by `item.id` so interleaved
|
||||
* `function_call_arguments.delta` (and sibling) events route to the matching
|
||||
* block instead of the most recently added one. A delta whose `item_id` is
|
||||
* not present is dropped rather than appended to a sibling.
|
||||
*/
|
||||
openItems: Map<string, CodexOpenItem>;
|
||||
/** Most recently added open item; fallback for events that omit `item_id`. */
|
||||
currentItem: CodexEventItem | null;
|
||||
currentBlock: CodexOutputBlock | null;
|
||||
nativeOutputItems: Array<Record<string, unknown>>;
|
||||
@@ -1061,6 +1077,7 @@ function createCodexStreamRuntime(initial: {
|
||||
requestBodyForState: initial.requestBodyForState,
|
||||
transport: initial.transport,
|
||||
websocketState: initial.websocketState,
|
||||
openItems: new Map(),
|
||||
currentItem: null,
|
||||
currentBlock: null,
|
||||
nativeOutputItems: [],
|
||||
@@ -1073,6 +1090,45 @@ function createCodexStreamRuntime(initial: {
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Wipe per-attempt accumulator state before a recovery path replays the turn.
|
||||
* Keeps {@link CodexStreamRuntime.openItems} and the legacy singleton-current
|
||||
* pointers in lockstep with {@link CodexStreamRuntime.nativeOutputItems} so a
|
||||
* stale delta from the failed attempt can't bind to a sibling on the retry.
|
||||
*/
|
||||
function resetCodexStreamAccumulators(runtime: CodexStreamRuntime): void {
|
||||
runtime.openItems.clear();
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Look up the open item a Codex stream event targets by `item_id`. When the
|
||||
* key is present but unknown the event is dropped (its item is already
|
||||
* closed); routing to a sibling is the bug we're fixing. When `item_id` is
|
||||
* absent we fall back to the most recently added open item to preserve the
|
||||
* prior singleton semantics for legacy/proxy streams that omit the key.
|
||||
*/
|
||||
function openItemForEvent(runtime: CodexStreamRuntime, rawEvent: Record<string, unknown>): CodexOpenItem | null {
|
||||
const itemId = typeof rawEvent.item_id === "string" ? rawEvent.item_id : "";
|
||||
if (itemId) return runtime.openItems.get(itemId) ?? null;
|
||||
let last: CodexOpenItem | null = null;
|
||||
for (const entry of runtime.openItems.values()) last = entry;
|
||||
return last;
|
||||
}
|
||||
|
||||
function closeCodexOpenItem(runtime: CodexStreamRuntime, itemId: string | undefined): void {
|
||||
if (!itemId) return;
|
||||
const entry = runtime.openItems.get(itemId);
|
||||
if (!entry) return;
|
||||
runtime.openItems.delete(itemId);
|
||||
if (runtime.currentItem === entry.item) {
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
}
|
||||
}
|
||||
|
||||
function resetWhitespaceToolCallArgumentsDelta(runtime: CodexStreamRuntime): void {
|
||||
runtime.whitespaceToolCallArgumentsDelta = undefined;
|
||||
}
|
||||
@@ -1194,11 +1250,23 @@ function handleCodexStreamEvent(
|
||||
const item = rawEvent.item as CodexEventItem;
|
||||
runtime.currentItem = item;
|
||||
runtime.currentBlock = createOutputBlockForItem(item);
|
||||
let contentIndex = -1;
|
||||
if (runtime.currentBlock) {
|
||||
output.content.push(runtime.currentBlock);
|
||||
contentIndex = output.content.length - 1;
|
||||
}
|
||||
// Track the open item so interleaved arg deltas route by id rather than
|
||||
// appending to whichever item was added most recently. Items without an
|
||||
// id are uncommon but still flow through the legacy singleton-current
|
||||
// path via openItemForEvent's fallback.
|
||||
const itemId = typeof (item as { id?: string }).id === "string" ? (item as { id: string }).id : "";
|
||||
if (itemId) {
|
||||
runtime.openItems.set(itemId, { item, block: runtime.currentBlock, contentIndex });
|
||||
}
|
||||
if (!runtime.currentBlock) return firstTokenTime;
|
||||
output.content.push(runtime.currentBlock);
|
||||
stream.push({
|
||||
type: getOutputBlockStartEventType(runtime.currentBlock),
|
||||
contentIndex: output.content.length - 1,
|
||||
contentIndex,
|
||||
partial: output,
|
||||
});
|
||||
return firstTokenTime;
|
||||
@@ -1242,7 +1310,7 @@ function handleCodexStreamEvent(
|
||||
|
||||
if (eventType === "response.function_call_arguments.done") {
|
||||
resetWhitespaceToolCallArgumentsDelta(runtime);
|
||||
handleToolCallArgumentsDone(runtime.currentItem, runtime.currentBlock, rawEvent);
|
||||
handleToolCallArgumentsDone(runtime, rawEvent);
|
||||
return firstTokenTime;
|
||||
}
|
||||
|
||||
@@ -1254,7 +1322,7 @@ function handleCodexStreamEvent(
|
||||
|
||||
if (eventType === "response.custom_tool_call_input.done") {
|
||||
resetWhitespaceToolCallArgumentsDelta(runtime);
|
||||
handleCustomToolCallInputDone(runtime.currentItem, runtime.currentBlock, rawEvent);
|
||||
handleCustomToolCallInputDone(runtime, rawEvent);
|
||||
return firstTokenTime;
|
||||
}
|
||||
|
||||
@@ -1416,36 +1484,37 @@ function handleToolCallArgumentsDelta(
|
||||
): CodexWhitespaceToolCallArgumentsDeltaInterruption | undefined {
|
||||
const delta = (rawEvent as { delta?: string }).delta || "";
|
||||
// Observe BEFORE the item/block guard: degenerate whitespace frames can keep
|
||||
// arriving after the item closed (currentBlock detached) and still count as
|
||||
// arriving after the item closed (entry detached) and still count as
|
||||
// progress for the idle watchdogs — dropping them unobserved would reopen
|
||||
// the infinite-loop hole the breaker exists for.
|
||||
const interruption = observeWhitespaceToolCallArgumentsDelta(runtime, rawEvent, delta);
|
||||
if (interruption) return interruption;
|
||||
const currentItem = runtime.currentItem;
|
||||
const currentBlock = runtime.currentBlock;
|
||||
if (currentItem?.type !== "function_call" || currentBlock?.type !== "toolCall") return undefined;
|
||||
currentBlock.partialJson += delta;
|
||||
const throttled = parseStreamingJsonThrottled(currentBlock.partialJson, currentBlock.lastParseLen ?? 0);
|
||||
// Route to the entry the event keys to; a delta whose item already closed
|
||||
// is dropped instead of leaking into a sibling tool call (#2619).
|
||||
const entry = openItemForEvent(runtime, rawEvent);
|
||||
if (!entry) return undefined;
|
||||
if (entry.item.type !== "function_call" || entry.block?.type !== "toolCall") return undefined;
|
||||
const block = entry.block;
|
||||
block.partialJson += delta;
|
||||
const throttled = parseStreamingJsonThrottled(block.partialJson, block.lastParseLen ?? 0);
|
||||
if (throttled) {
|
||||
currentBlock.arguments = throttled.value;
|
||||
currentBlock.lastParseLen = throttled.parsedLen;
|
||||
block.arguments = throttled.value;
|
||||
block.lastParseLen = throttled.parsedLen;
|
||||
}
|
||||
stream.push({ type: "toolcall_delta", contentIndex: output.content.length - 1, delta, partial: output });
|
||||
stream.push({ type: "toolcall_delta", contentIndex: entry.contentIndex, delta, partial: output });
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function handleToolCallArgumentsDone(
|
||||
currentItem: CodexEventItem | null,
|
||||
currentBlock: CodexOutputBlock | null,
|
||||
rawEvent: Record<string, unknown>,
|
||||
): void {
|
||||
if (currentItem?.type !== "function_call" || currentBlock?.type !== "toolCall") return;
|
||||
function handleToolCallArgumentsDone(runtime: CodexStreamRuntime, rawEvent: Record<string, unknown>): void {
|
||||
const entry = openItemForEvent(runtime, rawEvent);
|
||||
if (!entry || entry.item.type !== "function_call" || entry.block?.type !== "toolCall") return;
|
||||
const args = (rawEvent as { arguments?: string }).arguments;
|
||||
if (typeof args === "string") {
|
||||
currentBlock.partialJson = args;
|
||||
currentBlock.arguments = parseStreamingJson(currentBlock.partialJson);
|
||||
delete (currentBlock as { partialJson?: string }).partialJson;
|
||||
delete (currentBlock as { lastParseLen?: number }).lastParseLen;
|
||||
const block = entry.block;
|
||||
block.partialJson = args;
|
||||
block.arguments = parseStreamingJson(block.partialJson);
|
||||
delete (block as { partialJson?: string }).partialJson;
|
||||
delete (block as { lastParseLen?: number }).lastParseLen;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1459,25 +1528,23 @@ function handleCustomToolCallInputDelta(
|
||||
// Observe BEFORE the item/block guard — see handleToolCallArgumentsDelta.
|
||||
const interruption = observeWhitespaceToolCallArgumentsDelta(runtime, rawEvent, delta);
|
||||
if (interruption) return interruption;
|
||||
const currentItem = runtime.currentItem;
|
||||
const currentBlock = runtime.currentBlock;
|
||||
if (currentItem?.type !== "custom_tool_call" || currentBlock?.type !== "toolCall") return undefined;
|
||||
currentBlock.partialJson += delta;
|
||||
(currentBlock.arguments as { input?: string }).input = currentBlock.partialJson;
|
||||
stream.push({ type: "toolcall_delta", contentIndex: output.content.length - 1, delta, partial: output });
|
||||
const entry = openItemForEvent(runtime, rawEvent);
|
||||
if (!entry) return undefined;
|
||||
if (entry.item.type !== "custom_tool_call" || entry.block?.type !== "toolCall") return undefined;
|
||||
const block = entry.block;
|
||||
block.partialJson += delta;
|
||||
(block.arguments as { input?: string }).input = block.partialJson;
|
||||
stream.push({ type: "toolcall_delta", contentIndex: entry.contentIndex, delta, partial: output });
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function handleCustomToolCallInputDone(
|
||||
currentItem: CodexEventItem | null,
|
||||
currentBlock: CodexOutputBlock | null,
|
||||
rawEvent: Record<string, unknown>,
|
||||
): void {
|
||||
if (currentItem?.type !== "custom_tool_call" || currentBlock?.type !== "toolCall") return;
|
||||
function handleCustomToolCallInputDone(runtime: CodexStreamRuntime, rawEvent: Record<string, unknown>): void {
|
||||
const entry = openItemForEvent(runtime, rawEvent);
|
||||
if (!entry || entry.item.type !== "custom_tool_call" || entry.block?.type !== "toolCall") return;
|
||||
const input = (rawEvent as { input?: string }).input;
|
||||
if (typeof input === "string") {
|
||||
currentBlock.partialJson = input;
|
||||
currentBlock.arguments = { input };
|
||||
entry.block.partialJson = input;
|
||||
entry.block.arguments = { input };
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1493,32 +1560,40 @@ function handleOutputItemDone(
|
||||
const item = structuredCloneJSON(rawItem) as CodexEventItem;
|
||||
runtime.nativeOutputItems.push(item as unknown as Record<string, unknown>);
|
||||
|
||||
if (item.type === "reasoning" && runtime.currentBlock?.type === "thinking") {
|
||||
runtime.currentBlock.thinking = item.summary?.map(summary => summary.text).join("\n\n") || "";
|
||||
runtime.currentBlock.thinkingSignature = JSON.stringify(item);
|
||||
// Match the finalization to the OPEN ITEM that started this block, not the
|
||||
// singleton current — interleaved items can finish out of order, so the
|
||||
// most-recently-added block may belong to a sibling (#2619).
|
||||
const itemId = typeof (item as { id?: string }).id === "string" ? (item as { id: string }).id : "";
|
||||
const entry = itemId ? runtime.openItems.get(itemId) : null;
|
||||
const block = entry?.block ?? null;
|
||||
const contentIndex = entry?.contentIndex ?? output.content.length - 1;
|
||||
|
||||
if (item.type === "reasoning" && block?.type === "thinking") {
|
||||
block.thinking = item.summary?.map(summary => summary.text).join("\n\n") || "";
|
||||
block.thinkingSignature = JSON.stringify(item);
|
||||
stream.push({
|
||||
type: "thinking_end",
|
||||
contentIndex: output.content.length - 1,
|
||||
content: runtime.currentBlock.thinking,
|
||||
contentIndex,
|
||||
content: block.thinking,
|
||||
partial: output,
|
||||
});
|
||||
runtime.currentBlock = null;
|
||||
closeCodexOpenItem(runtime, itemId);
|
||||
return;
|
||||
}
|
||||
|
||||
if (item.type === "message" && runtime.currentBlock?.type === "text") {
|
||||
runtime.currentBlock.text = item.content
|
||||
if (item.type === "message" && block?.type === "text") {
|
||||
block.text = item.content
|
||||
.map(content => (content.type === "output_text" ? content.text : content.refusal))
|
||||
.join("");
|
||||
const phase = item.phase === "commentary" || item.phase === "final_answer" ? item.phase : undefined;
|
||||
runtime.currentBlock.textSignature = encodeTextSignatureV1(item.id, phase);
|
||||
block.textSignature = encodeTextSignatureV1(item.id, phase);
|
||||
stream.push({
|
||||
type: "text_end",
|
||||
contentIndex: output.content.length - 1,
|
||||
content: runtime.currentBlock.text,
|
||||
contentIndex,
|
||||
content: block.text,
|
||||
partial: output,
|
||||
});
|
||||
runtime.currentBlock = null;
|
||||
closeCodexOpenItem(runtime, itemId);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1529,26 +1604,25 @@ function handleOutputItemDone(
|
||||
name: item.name,
|
||||
arguments: parseStreamingJson(item.arguments || "{}"),
|
||||
};
|
||||
if (runtime.currentBlock?.type === "toolCall") {
|
||||
if (block?.type === "toolCall") {
|
||||
// Persist the authoritative final args on the stored block; the throttled
|
||||
// delta parser may have left currentBlock.arguments stale (often `{}`).
|
||||
runtime.currentBlock.arguments = toolCall.arguments;
|
||||
delete (runtime.currentBlock as { partialJson?: string }).partialJson;
|
||||
delete (runtime.currentBlock as { lastParseLen?: number }).lastParseLen;
|
||||
// Detach so a late/duplicate arguments.delta cannot append to the
|
||||
// finished block or trip the whitespace-loop guard against it.
|
||||
runtime.currentBlock = null;
|
||||
// delta parser may have left block.arguments stale (often `{}`).
|
||||
block.arguments = toolCall.arguments;
|
||||
delete (block as { partialJson?: string }).partialJson;
|
||||
delete (block as { lastParseLen?: number }).lastParseLen;
|
||||
}
|
||||
// Detach so a late/duplicate arguments.delta cannot append to the
|
||||
// finished block or trip the whitespace-loop guard against it.
|
||||
closeCodexOpenItem(runtime, itemId);
|
||||
runtime.canSafelyReplayWebsocketOverSse = false;
|
||||
stream.push({ type: "toolcall_end", contentIndex: output.content.length - 1, toolCall, partial: output });
|
||||
stream.push({ type: "toolcall_end", contentIndex, toolCall, partial: output });
|
||||
return;
|
||||
}
|
||||
|
||||
if (item.type === "custom_tool_call") {
|
||||
const rawInput =
|
||||
runtime.currentBlock?.type === "toolCall" && runtime.currentBlock.partialJson
|
||||
? runtime.currentBlock.partialJson
|
||||
: (item.input ?? "");
|
||||
const partial =
|
||||
block?.type === "toolCall" ? (block as ToolCall & { partialJson?: string }).partialJson : undefined;
|
||||
const rawInput = partial && partial.length > 0 ? partial : (item.input ?? "");
|
||||
const toolCall: ToolCall = {
|
||||
type: "toolCall",
|
||||
id: encodeResponsesToolCallId(item.call_id, item.id),
|
||||
@@ -1556,13 +1630,13 @@ function handleOutputItemDone(
|
||||
arguments: { input: rawInput },
|
||||
customWireName: item.name,
|
||||
};
|
||||
if (runtime.currentBlock?.type === "toolCall") {
|
||||
runtime.currentBlock.arguments = { input: rawInput };
|
||||
delete (runtime.currentBlock as { partialJson?: string }).partialJson;
|
||||
runtime.currentBlock = null;
|
||||
if (block?.type === "toolCall") {
|
||||
block.arguments = { input: rawInput };
|
||||
delete (block as { partialJson?: string }).partialJson;
|
||||
}
|
||||
closeCodexOpenItem(runtime, itemId);
|
||||
runtime.canSafelyReplayWebsocketOverSse = false;
|
||||
stream.push({ type: "toolcall_end", contentIndex: output.content.length - 1, toolCall, partial: output });
|
||||
stream.push({ type: "toolcall_end", contentIndex, toolCall, partial: output });
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1695,6 +1769,8 @@ function dropTrailingDegenerateToolCall(output: AssistantMessage, runtime: Codex
|
||||
if (block && block.type === "toolCall" && output.content[output.content.length - 1] === block) {
|
||||
output.content.pop();
|
||||
}
|
||||
const itemId = runtime.currentItem && (runtime.currentItem as { id?: string }).id;
|
||||
if (itemId) runtime.openItems.delete(itemId);
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
}
|
||||
@@ -1742,10 +1818,8 @@ async function tryRecoverCodexWhitespaceToolCallLoop(
|
||||
transport: runtime.transport,
|
||||
});
|
||||
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
runtime.sawTerminalEvent = false;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetWhitespaceToolCallArgumentsDelta(runtime);
|
||||
resetOutputState(context.output);
|
||||
context.firstTokenTime = undefined;
|
||||
@@ -1802,9 +1876,7 @@ async function tryReconnectCodexWebSocketOnConnectionLimit(
|
||||
if (context.output.content.length > 0) {
|
||||
// Content already emitted to the caller — cannot safely continue on a new WS.
|
||||
// Reset and replay the full request over SSE.
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
resetOutputState(context.output);
|
||||
context.firstTokenTime = undefined;
|
||||
recordCodexWebSocketFailure(websocketState, true);
|
||||
@@ -1817,9 +1889,7 @@ async function tryReconnectCodexWebSocketOnConnectionLimit(
|
||||
// over websocket, bounded by the shared retry budget: an account-scoped
|
||||
// limit can reject every fresh connection, and an unbounded loop would
|
||||
// hammer the endpoint with zero backoff.
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
context.firstTokenTime = undefined;
|
||||
if (runtime.websocketStreamRetries >= getCodexWebSocketRetryBudget()) {
|
||||
recordCodexWebSocketFailure(websocketState, true);
|
||||
@@ -1871,10 +1941,8 @@ async function tryRecoverCodexPreviousResponseNotFound(
|
||||
runtime.providerRetryAttempt += 1;
|
||||
resetCodexWebSocketAppendState(websocketState);
|
||||
resetCodexSessionMetadata(websocketState);
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
runtime.sawTerminalEvent = false;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetOutputState(context.output);
|
||||
context.firstTokenTime = undefined;
|
||||
|
||||
@@ -1921,9 +1989,7 @@ async function tryReplayWebsocketFailureOverSse(
|
||||
// Full re-send on a fresh socket: clear accumulator state from the failed
|
||||
// attempt. Content is empty here, but blockless native items (e.g.
|
||||
// web_search_call) may already have accumulated.
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
context.firstTokenTime = undefined;
|
||||
await scheduler.wait(getCodexWebSocketRetryDelayMs(runtime.websocketStreamRetries), {
|
||||
signal: context.requestSetup.requestSignal,
|
||||
@@ -1932,9 +1998,7 @@ async function tryReplayWebsocketFailureOverSse(
|
||||
return true;
|
||||
}
|
||||
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
resetOutputState(context.output);
|
||||
context.firstTokenTime = undefined;
|
||||
|
||||
@@ -1970,10 +2034,8 @@ async function tryRetryCodexProviderError(
|
||||
transport: runtime.transport,
|
||||
});
|
||||
|
||||
runtime.currentItem = null;
|
||||
runtime.currentBlock = null;
|
||||
resetCodexStreamAccumulators(runtime);
|
||||
runtime.sawTerminalEvent = false;
|
||||
runtime.nativeOutputItems.length = 0;
|
||||
resetOutputState(context.output);
|
||||
context.firstTokenTime = undefined;
|
||||
await scheduler.wait(CODEX_RETRY_DELAY_MS * runtime.providerRetryAttempt, {
|
||||
|
||||
@@ -380,6 +380,114 @@ describe("openai-codex streaming", () => {
|
||||
expect("lastParseLen" in toolCall).toBe(false);
|
||||
});
|
||||
|
||||
it("routes interleaved function-call argument deltas to the matching open item", async () => {
|
||||
const tempDir = TempDir.createSync("@pi-codex-stream-");
|
||||
setAgentDir(tempDir.path());
|
||||
const token = createCodexTestToken();
|
||||
const context = createCodexTestContext();
|
||||
// Two function calls are opened concurrently and the server interleaves
|
||||
// `function_call_arguments.delta` events by `item_id`. With the old
|
||||
// singleton current-block, every delta went to whichever item was added
|
||||
// most recently; the `task` call ended up with `arguments = {}` and the
|
||||
// sibling received the `task` payload (issue #2619). Each call must
|
||||
// retain its own arguments and emit `toolcall_*` events against its own
|
||||
// content index.
|
||||
const taskArgs = '{"ops":[{"op":"start","task":"X"}]}';
|
||||
const otherArgs = '{"input":"hello"}';
|
||||
const events: Array<Record<string, unknown>> = [
|
||||
{
|
||||
type: "response.output_item.added",
|
||||
output_index: 0,
|
||||
item: { type: "function_call", id: "fc_task", call_id: "call_task", name: "task", arguments: "" },
|
||||
},
|
||||
{
|
||||
type: "response.output_item.added",
|
||||
output_index: 1,
|
||||
item: { type: "function_call", id: "fc_other", call_id: "call_other", name: "other", arguments: "" },
|
||||
},
|
||||
{
|
||||
type: "response.function_call_arguments.delta",
|
||||
item_id: "fc_task",
|
||||
output_index: 0,
|
||||
delta: taskArgs.slice(0, 12),
|
||||
},
|
||||
{
|
||||
type: "response.function_call_arguments.delta",
|
||||
item_id: "fc_other",
|
||||
output_index: 1,
|
||||
delta: otherArgs.slice(0, 10),
|
||||
},
|
||||
{
|
||||
type: "response.function_call_arguments.delta",
|
||||
item_id: "fc_task",
|
||||
output_index: 0,
|
||||
delta: taskArgs.slice(12),
|
||||
},
|
||||
{
|
||||
type: "response.function_call_arguments.delta",
|
||||
item_id: "fc_other",
|
||||
output_index: 1,
|
||||
delta: otherArgs.slice(10),
|
||||
},
|
||||
// Stale delta for fc_task arriving after fc_other finishes must be dropped,
|
||||
// not appended to fc_other.
|
||||
{
|
||||
type: "response.output_item.done",
|
||||
output_index: 1,
|
||||
item: { type: "function_call", id: "fc_other", call_id: "call_other", name: "other", arguments: otherArgs },
|
||||
},
|
||||
{
|
||||
type: "response.function_call_arguments.delta",
|
||||
item_id: "fc_other",
|
||||
output_index: 1,
|
||||
delta: "STALE",
|
||||
},
|
||||
{
|
||||
type: "response.output_item.done",
|
||||
output_index: 0,
|
||||
item: { type: "function_call", id: "fc_task", call_id: "call_task", name: "task", arguments: taskArgs },
|
||||
},
|
||||
{
|
||||
type: "response.completed",
|
||||
response: {
|
||||
id: "resp_1",
|
||||
status: "completed",
|
||||
usage: { input_tokens: 5, output_tokens: 3, total_tokens: 8, input_tokens_details: { cached_tokens: 0 } },
|
||||
},
|
||||
},
|
||||
];
|
||||
const sse = `${events.map(e => `data: ${JSON.stringify(e)}`).join("\n\n")}\n\n`;
|
||||
const fetchMock: FetchImpl = (async () =>
|
||||
new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" } })) as FetchImpl;
|
||||
|
||||
const toolcallEnds: Array<{ contentIndex: number; name: string; argumentsJson: string }> = [];
|
||||
const model = { ...createCodexTestModel("https://chatgpt.com/backend-api"), preferWebsockets: false };
|
||||
const aem = streamOpenAICodexResponses(model, context, { apiKey: token, fetch: fetchMock });
|
||||
(async () => {
|
||||
for await (const event of aem) {
|
||||
if (event.type !== "toolcall_end") continue;
|
||||
toolcallEnds.push({
|
||||
contentIndex: event.contentIndex,
|
||||
name: event.toolCall.name,
|
||||
argumentsJson: JSON.stringify(event.toolCall.arguments),
|
||||
});
|
||||
}
|
||||
})();
|
||||
const result = await aem.result();
|
||||
|
||||
const calls = result.content.filter(c => c.type === "toolCall");
|
||||
expect(calls).toHaveLength(2);
|
||||
const byName = new Map(calls.map(c => [c.name, c] as const));
|
||||
expect(byName.get("task")?.arguments).toEqual({ ops: [{ op: "start", task: "X" }] });
|
||||
expect(byName.get("other")?.arguments).toEqual({ input: "hello" });
|
||||
// `task` is the FIRST opened block (index 0); a stale delta after fc_other
|
||||
// closed must NOT have appended "STALE" anywhere.
|
||||
expect(JSON.stringify(result.content)).not.toContain("STALE");
|
||||
// Stream events must address each tool call by its own content index.
|
||||
expect(toolcallEnds.find(e => e.name === "task")?.contentIndex).toBe(0);
|
||||
expect(toolcallEnds.find(e => e.name === "other")?.contentIndex).toBe(1);
|
||||
});
|
||||
|
||||
it("waits for caller abort when SSE streams only no-progress status events", async () => {
|
||||
const tempDir = TempDir.createSync("@pi-codex-stream-");
|
||||
setAgentDir(tempDir.path());
|
||||
|
||||
Reference in New Issue
Block a user