fix(ai): tracked unkeyed current open item via currentEntry

The keyed maps only see items whose `output_item.added` carries `item.id`
or `output_index`. A fully keyless add never reached either map, and the
unkeyed fallback was scanning those maps instead of the actual current
item — so `function_call_arguments.delta` / `output_item.done` for a
keyless tool call landed on null and the stored block kept `{}`. Mixed
streams also picked an older `output_index` entry before a later
id-only current item for the same reason.

`CodexStreamRuntime.currentEntry` now always points at the most recently
added `output_item.added` (whether or not it has keys). `openItemForEvent`
returns `currentEntry` when both `item_id` and `output_index` are absent,
and `closeCodexOpenItem` clears `currentEntry` (alongside the legacy
`currentItem` / `currentBlock` mirrors) when its item closes. The keyed
maps are unchanged so deliberate drop-on-mismatch for keyed events still
holds.

Regression tests cover the fully keyless function-call stream and the
mixed id-only vs `output_index`-only ordering. Existing reasoning/message
flow keeps singleton semantics through the same fallback.

Fixes #2619
This commit is contained in:
roboomp
2026-06-15 06:28:35 +00:00
parent 69ae68db8a
commit da94f27edd
2 changed files with 127 additions and 21 deletions
@@ -302,7 +302,14 @@ interface CodexStreamRuntime {
* call items omit `id`; these still carry `output_index` on deltas/done.
*/
openItemsByOutputIndex: Map<number, CodexOpenItem>;
/** Most recently added open item; fallback for events that omit all keys. */
/**
* Most recently added open item for events that omit both `item_id` and
* `output_index`. Always tracks the latest `output_item.added`, including
* fully keyless items that never make it into the keyed maps; cleared when
* its item closes.
*/
currentEntry: CodexOpenItem | null;
/** Convenience mirrors of {@link currentEntry} for legacy singleton handlers. */
currentItem: CodexEventItem | null;
currentBlock: CodexOutputBlock | null;
nativeOutputItems: Array<Record<string, unknown>>;
@@ -1084,6 +1091,7 @@ function createCodexStreamRuntime(initial: {
websocketState: initial.websocketState,
openItems: new Map(),
openItemsByOutputIndex: new Map(),
currentEntry: null,
currentItem: null,
currentBlock: null,
nativeOutputItems: [],
@@ -1105,6 +1113,7 @@ function createCodexStreamRuntime(initial: {
function resetCodexStreamAccumulators(runtime: CodexStreamRuntime): void {
runtime.openItems.clear();
runtime.openItemsByOutputIndex.clear();
runtime.currentEntry = null;
runtime.currentItem = null;
runtime.currentBlock = null;
runtime.nativeOutputItems.length = 0;
@@ -1115,25 +1124,23 @@ function resetCodexStreamAccumulators(runtime: CodexStreamRuntime): void {
* uniquely identifies a response item; `output_index` covers idless function
* call items. A keyed event whose target is already closed is dropped instead
* of being routed to a sibling. Only streams that omit both keys fall back to
* the most recently added item, preserving legacy/proxy singleton semantics.
* {@link CodexStreamRuntime.currentEntry} — the most recently added item,
* including fully keyless ones that never reached the keyed maps.
*/
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;
const outputIndex = readOptionalInteger(rawEvent.output_index);
if (outputIndex !== undefined) return runtime.openItemsByOutputIndex.get(outputIndex) ?? null;
let last: CodexOpenItem | null = null;
for (const entry of runtime.openItemsByOutputIndex.values()) last = entry;
if (last) return last;
for (const entry of runtime.openItems.values()) last = entry;
return last;
return runtime.currentEntry;
}
function closeCodexOpenItem(runtime: CodexStreamRuntime, entry: CodexOpenItem | null | undefined): void {
if (!entry) return;
if (entry.itemId) runtime.openItems.delete(entry.itemId);
if (entry.outputIndex !== undefined) runtime.openItemsByOutputIndex.delete(entry.outputIndex);
if (runtime.currentItem === entry.item) {
if (runtime.currentEntry === entry) {
runtime.currentEntry = null;
runtime.currentItem = null;
runtime.currentBlock = null;
}
@@ -1272,6 +1279,7 @@ function handleCodexStreamEvent(
const itemId = typeof (item as { id?: string }).id === "string" ? (item as { id: string }).id : undefined;
const outputIndex = readOptionalInteger(rawEvent.output_index);
const entry: CodexOpenItem = { item, block: runtime.currentBlock, contentIndex, itemId, outputIndex };
runtime.currentEntry = entry;
if (itemId) runtime.openItems.set(itemId, entry);
if (outputIndex !== undefined) runtime.openItemsByOutputIndex.set(outputIndex, entry);
if (!runtime.currentBlock) return firstTokenTime;
@@ -1782,19 +1790,7 @@ function dropTrailingDegenerateToolCall(output: AssistantMessage, runtime: Codex
if (block && block.type === "toolCall" && output.content[output.content.length - 1] === block) {
output.content.pop();
}
let entry =
runtime.currentItem && (runtime.currentItem as { id?: string }).id
? runtime.openItems.get((runtime.currentItem as { id: string }).id)
: undefined;
if (!entry) {
for (const open of runtime.openItemsByOutputIndex.values()) {
if (open.item === runtime.currentItem || open.block === block) {
entry = open;
break;
}
}
}
closeCodexOpenItem(runtime, entry);
closeCodexOpenItem(runtime, runtime.currentEntry);
}
/**
@@ -558,6 +558,116 @@ describe("openai-codex streaming", () => {
expect(toolcallEnds.find(e => e.name === "apply_patch")?.contentIndex).toBe(1);
});
it("routes fully keyless deltas/done to the latest open item via currentEntry", async () => {
const tempDir = TempDir.createSync("@pi-codex-stream-");
setAgentDir(tempDir.path());
const token = createCodexTestToken();
const context = createCodexTestContext();
// Pathological legacy/proxy stream: `output_item.added` carries no `id`
// AND no `output_index`, so neither keyed map ever receives the item.
// `function_call_arguments.delta` / `output_item.done` likewise lack
// both keys. The runtime must still route them via `currentEntry`
// (the latest live `output_item.added`) instead of dropping.
const taskArgs = '{"tasks":[{"assignment":"keyless"}]}';
const events: Array<Record<string, unknown>> = [
{
type: "response.output_item.added",
item: { type: "function_call", call_id: "call_keyless", name: "task", arguments: "" },
},
{ type: "response.function_call_arguments.delta", delta: taskArgs.slice(0, 12) },
{ type: "response.function_call_arguments.delta", delta: taskArgs.slice(12) },
{
type: "response.output_item.done",
item: { type: "function_call", call_id: "call_keyless", name: "task", arguments: taskArgs },
},
{
type: "response.completed",
response: {
id: "resp_keyless",
status: "completed",
usage: {
input_tokens: 1,
output_tokens: 1,
total_tokens: 2,
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 model = { ...createCodexTestModel("https://chatgpt.com/backend-api"), preferWebsockets: false };
const result = await streamOpenAICodexResponses(model, context, {
apiKey: token,
fetch: fetchMock as FetchImpl,
}).result();
const call = result.content.find(c => c.type === "toolCall");
expect(call?.name).toBe("task");
expect(call?.arguments).toEqual({ tasks: [{ assignment: "keyless" }] });
});
it("prefers a later id-only current item over an older output_index entry on unkeyed events", async () => {
const tempDir = TempDir.createSync("@pi-codex-stream-");
setAgentDir(tempDir.path());
const token = createCodexTestToken();
const context = createCodexTestContext();
// Mixed key shapes: the first call is output_index-keyed only, the
// second is id-only and is now the latest open item. An unkeyed delta
// must address the second call (currentEntry), not whatever the
// keyed-map iteration happens to surface first.
const idOnlyArgs = '{"input":"id-only-current"}';
const events: Array<Record<string, unknown>> = [
{
type: "response.output_item.added",
output_index: 0,
item: { type: "function_call", call_id: "call_old", name: "older", arguments: "" },
},
{
type: "response.output_item.added",
item: { type: "function_call", id: "fc_id_only", call_id: "call_new", name: "newer", arguments: "" },
},
// Keyless delta + done for the newer call — must route to fc_id_only.
{ type: "response.function_call_arguments.delta", delta: idOnlyArgs },
{
type: "response.output_item.done",
item: { type: "function_call", id: "fc_id_only", call_id: "call_new", name: "newer", arguments: idOnlyArgs },
},
// Close the older one explicitly with its key so the test verifies isolation.
{
type: "response.output_item.done",
output_index: 0,
item: { type: "function_call", call_id: "call_old", name: "older", arguments: "{}" },
},
{
type: "response.completed",
response: {
id: "resp_mixed",
status: "completed",
usage: {
input_tokens: 1,
output_tokens: 1,
total_tokens: 2,
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 model = { ...createCodexTestModel("https://chatgpt.com/backend-api"), preferWebsockets: false };
const result = await streamOpenAICodexResponses(model, context, {
apiKey: token,
fetch: fetchMock as FetchImpl,
}).result();
const calls = result.content.filter(c => c.type === "toolCall");
const byName = new Map(calls.map(c => [c.name, c] as const));
expect(byName.get("newer")?.arguments).toEqual({ input: "id-only-current" });
expect(byName.get("older")?.arguments).toEqual({});
});
it("waits for caller abort when SSE streams only no-progress status events", async () => {
const tempDir = TempDir.createSync("@pi-codex-stream-");
setAgentDir(tempDir.path());