diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 1763b8067..8f69050a7 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -16,6 +16,8 @@ - Fixed tool call ID normalization for Anthropic-compatible models - Fixed Anthropic Messages replay sanitizing malformed tool-call IDs, including aborted native tool calls with empty IDs, so retries no longer send invalid `tool_use.id` / `tool_result.tool_use_id` pairs. +- Fixed the Codex Responses WebSocket transport attributing a prior turn's output to the current one on a reused connection: a trailing/duplicate frame from a cleanly-completed previous response that slipped past the queue drain could be consumed as this request's terminal (ending the turn with empty output) or as a stale tool call. Frames are now keyed by `response.id` — a frame carrying the previous response's id is dropped, and one carrying a third id (or a regressed `sequence_number`) fails closed so the turn retries instead of mixing two responses' streams. Idless frames (deltas, the rate-limit/metadata preamble, `response.created`-less streams) still pass through, matching upstream codex-rs. +- Fixed `transformMessages` pulling an earlier, orphaned tool result onto a later tool call that reused the same id (left behind when compaction folded the originating `tool_use` into a summary). The pending-call flush now pairs each call with a result positioned *after* its assistant turn, so a reused id surfaces its own output rather than a prior turn's. ### Fixed diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 0164e4903..b7838fd18 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -205,6 +205,17 @@ function isCodexStreamProgressEvent(event: unknown): boolean { return typeof type === "string" && CODEX_ADDITIONAL_PROGRESS_EVENT_TYPES.has(type); } +function extractCodexFrameResponseId(frame: Record): string | undefined { + const response = (frame as { response?: { id?: unknown } }).response; + const id = response?.id; + return typeof id === "string" && id.length > 0 ? id : undefined; +} + +function extractCodexFrameSequenceNumber(frame: Record): number | undefined { + const raw = (frame as { sequence_number?: unknown }).sequence_number; + return typeof raw === "number" && Number.isFinite(raw) ? Math.trunc(raw) : undefined; +} + type CodexWebSocketTimeoutDetails = { lastEventAt: number; lastEventType?: string; @@ -2403,6 +2414,12 @@ class CodexWebSocketConnection { #lastInboundAt = 0; /** Wall-clock of the last heartbeat ping we issued; 0 if none yet. */ #lastPingAt = 0; + /** + * Most recent `response.id` accepted on this socket, retained across + * requests. Lets the next request drop a trailing/duplicate frame from the + * previous (cleanly-completed) response that outlived the queue drain. + */ + #lastSeenResponseId?: string; constructor(url: string, headers: Record, options: CodexWebSocketConnectionOptions) { this.#url = url; @@ -2650,6 +2667,11 @@ class CodexWebSocketConnection { let lastProgressEventType: string | undefined; let lastEventAt = lastProgressAt; let lastEventType: string | undefined; + // Cross-request frame guard: lock onto this response's id and reject + // frames belonging to another response interleaved on the reused socket. + let activeResponseId: string | undefined; + let lastSequence: number | undefined; + const priorResponseId = this.#lastSeenResponseId; while (true) { let timeoutMs: number | undefined; let timeoutReason: string; @@ -2691,8 +2713,51 @@ class CodexWebSocketConnection { if (next === null) { throw new CodexWebSocketTransportError(`websocket closed before response completion`); } - sawFirstEvent = true; const eventType = typeof next.type === "string" ? next.type : ""; + // Cross-request frame guard. The socket is reused across turns. Upstream + // codex-rs leans on the protocol guarantee that nothing follows a + // response's terminal event, but our queue can still surface a trailing + // or duplicate frame from a cleanly-completed prior response after + // #dropStaleFrames() drained the queue at send time. Attaching such a + // frame to THIS turn misattributes an earlier turn's output (a stale + // `response.completed` ends the turn early; a stale item makes the model + // see an unrelated call). Only lifecycle events (created/completed/ + // failed/incomplete) carry a `response.id` — exactly the harmful ones — + // so key the guard on it and let idless frames (deltas, the rate-limit/ + // metadata preamble, created-less streams) pass through, matching + // upstream rather than gating on `response.created`. + const frameResponseId = extractCodexFrameResponseId(next); + const frameSequence = extractCodexFrameSequenceNumber(next); + if (frameResponseId !== undefined) { + if (activeResponseId === undefined) { + if (priorResponseId !== undefined && frameResponseId === priorResponseId) { + // Trailing/duplicate frame of the previous response that + // outlived the drain. Drop without locking or advancing the + // first-event clocks so our own response can still start. + continue; + } + activeResponseId = frameResponseId; + } else if (frameResponseId !== activeResponseId) { + // A different response is interleaving on the socket; the idless + // deltas that follow are indistinguishable, so fail closed + // (retryable) instead of risking misattribution. + this.close("stale-frame"); + throw new CodexWebSocketTransportError( + `websocket frame for response ${frameResponseId} interleaved into active response ${activeResponseId}`, + ); + } + this.#lastSeenResponseId = frameResponseId; + } + if (frameSequence !== undefined) { + if (activeResponseId !== undefined && lastSequence !== undefined && frameSequence < lastSequence) { + this.close("stale-frame"); + throw new CodexWebSocketTransportError( + `websocket sequence_number ${frameSequence} regressed below ${lastSequence} within response ${activeResponseId}`, + ); + } + lastSequence = frameSequence; + } + sawFirstEvent = true; lastEventAt = Date.now(); lastEventType = eventType || undefined; if (isCodexStreamProgressEvent(next)) { diff --git a/packages/ai/src/providers/transform-messages.ts b/packages/ai/src/providers/transform-messages.ts index b563360da..6ba7c018d 100644 --- a/packages/ai/src/providers/transform-messages.ts +++ b/packages/ai/src/providers/transform-messages.ts @@ -397,12 +397,33 @@ export function transformMessages( maxNormalizedToolCallIdLength, duplicateToolCallIdSuffixPrefix, ); - const realToolResultsById = new Map(); - for (const msg of transformed) { - if (msg.role === "toolResult" && !realToolResultsById.has(msg.toolCallId)) { - realToolResultsById.set(msg.toolCallId, msg); + // All real tool results, keyed by id, in document order. One id can map to + // more than one result: compaction can fold an assistant `tool_use` into a + // summary string while its `tool_result` survives, and a later turn may reuse + // the id. `takeRealToolResult` pulls the earliest unconsumed result positioned + // AFTER the call's assistant turn, so an orphaned earlier result is never + // pulled forward onto a later call (which would surface a prior turn's output). + type IndexedToolResult = { index: number; msg: ToolResultMessage; consumed: boolean }; + const realToolResultsById = new Map(); + for (let index = 0; index < transformed.length; index++) { + const msg = transformed[index]; + if (msg.role === "toolResult") { + const entry: IndexedToolResult = { index, msg, consumed: false }; + const entries = realToolResultsById.get(msg.toolCallId); + if (entries) entries.push(entry); + else realToolResultsById.set(msg.toolCallId, [entry]); } } + const takeRealToolResult = (id: string, afterIndex: number): ToolResultMessage | undefined => { + const entries = realToolResultsById.get(id); + if (!entries) return undefined; + for (const entry of entries) { + if (entry.consumed || entry.index <= afterIndex) continue; + entry.consumed = true; + return entry.msg; + } + return undefined; + }; // Anthropic rejects `tool_result` blocks whose `tool_use_id` does not appear in a prior // `tool_use` block. After handoff/compaction folds an assistant turn into a summary @@ -421,8 +442,12 @@ export function transformMessages( // followed by exactly one corresponding tool result. const result: Message[] = []; let pendingToolCalls: ToolCall[] = []; + // Index of the assistant turn that declared `pendingToolCalls`; a pulled + // result must be positioned after it (see `takeRealToolResult`). + let pendingToolCallsStartIndex = -1; let pendingAbortedToolCalls = new Map(); let pendingAbortedTimestamp: number | undefined; + let pendingAbortedStartIndex = -1; // Track which tool calls already have an emitted result so delayed/duplicate // toolResult messages cannot create a second provider-visible result. const toolCallStatus = new Map(); @@ -431,7 +456,7 @@ export function transformMessages( if (pendingToolCalls.length === 0) return; for (const tc of pendingToolCalls) { if (toolCallStatus.has(tc.id)) continue; - const realToolResult = realToolResultsById.get(tc.id); + const realToolResult = takeRealToolResult(tc.id, pendingToolCallsStartIndex); if (realToolResult) { result.push(realToolResult); toolCallStatus.set(tc.id, ToolCallStatus.Resolved); @@ -454,7 +479,7 @@ export function transformMessages( if (pendingAbortedTimestamp === undefined) return; for (const tc of pendingAbortedToolCalls.values()) { if (toolCallStatus.has(tc.id)) continue; - const realToolResult = realToolResultsById.get(tc.id); + const realToolResult = takeRealToolResult(tc.id, pendingAbortedStartIndex); if (realToolResult) { result.push(realToolResult); toolCallStatus.set(tc.id, ToolCallStatus.Resolved); @@ -510,11 +535,13 @@ export function transformMessages( result.push(msg); pendingAbortedToolCalls = new Map(toolCalls.map(toolCall => [toolCall.id, toolCall] as const)); pendingAbortedTimestamp = assistantMsg.timestamp; + pendingAbortedStartIndex = i; continue; } if (toolCalls.length > 0) { pendingToolCalls = toolCalls; + pendingToolCallsStartIndex = i; } result.push(msg); diff --git a/packages/ai/test/duplicate-tool-results.test.ts b/packages/ai/test/duplicate-tool-results.test.ts index 7745abd9e..1ec6e5c32 100644 --- a/packages/ai/test/duplicate-tool-results.test.ts +++ b/packages/ai/test/duplicate-tool-results.test.ts @@ -195,6 +195,43 @@ describe("Duplicate Tool Results Regression", () => { expect((toolResults[0] as ToolResultMessage).content).toEqual([{ type: "text", text: "todo updated" }]); }); + it("routes a reused tool-call id to its own result, never an earlier orphaned one", () => { + // Compaction folded the assistant turn that originally issued `sharedId` + // into a summary string, but its tool result survived as an orphan. A + // later turn reuses the same id, and a developer note sits between that + // call and its real result — forcing a pending-call flush before the real + // result is reached. The flush must pull THIS turn's result, not the + // earlier orphan's output (regression: a tool call returning an earlier, + // unrelated command's output). + const sharedId = "toolu_shared_reuse_1"; + const messages: Message[] = [ + { role: "user", content: "first request", timestamp: 1 }, + // Orphaned result: its originating tool_use was compacted away. + makeEvalToolResult(sharedId, "OUTPUT FROM EARLIER COMMAND", 2), + { role: "user", content: "second request", timestamp: 3 }, + makeEvalAssistantMessage(sharedId, 4), + { role: "developer", content: "guidance between call and result", timestamp: 5 }, + makeEvalToolResult(sharedId, "OUTPUT FROM CURRENT COMMAND", 6), + ]; + + const transformed = transformMessages(messages, model); + + const results = getToolResults(transformed).filter(result => result.toolCallId === sharedId); + expect(results).toHaveLength(1); + expect(results[0]!.content).toEqual([{ type: "text", text: "OUTPUT FROM CURRENT COMMAND" }]); + + // The surviving result must land immediately after the reusing assistant turn. + const assistantIdx = transformed.findIndex( + message => + message.role === "assistant" && + message.content.some(block => block.type === "toolCall" && block.id === sharedId), + ); + expect(transformed[assistantIdx + 1]?.role).toBe("toolResult"); + expect((transformed[assistantIdx + 1] as ToolResultMessage).content).toEqual([ + { type: "text", text: "OUTPUT FROM CURRENT COMMAND" }, + ]); + }); + it("should not duplicate tool results for aborted messages when results already exist", () => { const toolCallId = "toolu_aborted_test_123"; diff --git a/packages/ai/test/openai-codex-stream.test.ts b/packages/ai/test/openai-codex-stream.test.ts index 66f1fc54a..1f0c2750e 100644 --- a/packages/ai/test/openai-codex-stream.test.ts +++ b/packages/ai/test/openai-codex-stream.test.ts @@ -2261,6 +2261,102 @@ describe("openai-codex streaming", () => { }); }); + it("drops a stale terminal frame from the prior response leaking onto a reused websocket", async () => { + const tempDir = TempDir.createSync("@pi-codex-stale-frame-"); + setAgentDir(tempDir.path()); + const token = createCodexTestToken(); + const fetchMock = vi.fn(async () => { + throw new Error("SSE fallback should not be called"); + }); + + // On the reused connection's second request, a trailing/duplicate + // `response.completed` from the previous response slips past the queue + // drain and arrives before this request's own frames. The transport must + // drop it (its `response.id` is the prior response's) rather than consume + // it as request 2's terminal — which would end the turn with empty output + // or, worse, attribute the prior turn's output to this one. + class StaleFrameWebSocket extends MockWebSocket { + #sendCount = 0; + + constructor(url: string, options?: { headers?: WsHeaders }) { + super(url, options); + this.scheduleOpen(); + } + + send(): void { + this.#sendCount += 1; + if (this.#sendCount === 1) { + this.emitCodexResponse({ + messageId: "msg_1", + responseId: "resp_1", + text: "First answer", + terminalType: "response.completed", + includeCreated: true, + }); + return; + } + this.sendJson({ + type: "response.completed", + response: { id: "resp_1", status: "completed", usage: DEFAULT_USAGE }, + }); + this.emitCodexResponse({ + messageId: "msg_2", + responseId: "resp_2", + text: "Second answer", + terminalType: "response.completed", + includeCreated: true, + }); + } + } + + global.WebSocket = StaleFrameWebSocket as unknown as typeof WebSocket; + const model: Model<"openai-codex-responses"> = buildModel({ + id: "gpt-5.3-codex-spark", + name: "GPT-5.3 Codex Spark", + api: "openai-codex-responses", + provider: "openai-codex", + baseUrl: "https://chatgpt.com/backend-api", + reasoning: true, + preferWebsockets: true, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 128000, + maxTokens: 128000, + }); + const providerSessionState = new Map(); + const firstContext: Context = { + systemPrompt: ["You are a helpful assistant."], + messages: [{ role: "user", content: "First question", timestamp: Date.now() }], + }; + const first = await streamOpenAICodexResponses(model, firstContext, { + fetch: fetchMock as FetchImpl, + apiKey: token, + sessionId: "ws-stale-frame-session", + providerSessionState, + }).result(); + const secondContext: Context = { + systemPrompt: ["You are a helpful assistant."], + messages: [ + ...firstContext.messages, + first, + { role: "user", content: "Second question", timestamp: Date.now() }, + ], + }; + const second = await streamOpenAICodexResponses(model, secondContext, { + fetch: fetchMock as FetchImpl, + apiKey: token, + sessionId: "ws-stale-frame-session", + providerSessionState, + }).result(); + + const secondText = second.content + .filter((block): block is { type: "text"; text: string } => block.type === "text") + .map(block => block.text) + .join(""); + expect(secondText).toBe("Second answer"); + expect(fetchMock).not.toHaveBeenCalled(); + }); + it("applies onPayload to the final chained websocket frame", async () => { const tempDir = TempDir.createSync("@pi-codex-ws-payload-hook-"); setAgentDir(tempDir.path());