fix(ai/providers): prevented cross-turn frame and tool result misattribution
- Implement ID-based guards for Codex WebSocket frames to reject stale or unauthorized interleaved frames from previous turns. - Check sequence numbers within Codex responses to detect and handle out-of-order frame delivery. - Restrict tool result consumption to occurrences appearing after the associated tool call, preventing the reuse of orphaned results from earlier conversation turns.
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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, unknown>): 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<string, unknown>): 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<string, string>, 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)) {
|
||||
|
||||
@@ -397,12 +397,33 @@ export function transformMessages<TApi extends Api>(
|
||||
maxNormalizedToolCallIdLength,
|
||||
duplicateToolCallIdSuffixPrefix,
|
||||
);
|
||||
const realToolResultsById = new Map<string, ToolResultMessage>();
|
||||
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<string, IndexedToolResult[]>();
|
||||
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<TApi extends Api>(
|
||||
// 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<string, ToolCall>();
|
||||
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<string, ToolCallStatus>();
|
||||
@@ -431,7 +456,7 @@ export function transformMessages<TApi extends Api>(
|
||||
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<TApi extends Api>(
|
||||
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<TApi extends Api>(
|
||||
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);
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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<string, ProviderSessionState>();
|
||||
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());
|
||||
|
||||
Reference in New Issue
Block a user