diff --git a/packages/ai/src/providers/azure-openai-responses.ts b/packages/ai/src/providers/azure-openai-responses.ts index 711f02b33..1ce6d0102 100644 --- a/packages/ai/src/providers/azure-openai-responses.ts +++ b/packages/ai/src/providers/azure-openai-responses.ts @@ -28,9 +28,10 @@ import { iterateWithIdleTimeout, } from "../utils/idle-iterator"; import { sanitizeSchemaForOpenAIResponses, toolWireSchema } from "../utils/schema"; +import { createSdkStreamRequestOptions } from "../utils/sdk-stream-timeout"; import { notifyRawSseEvent } from "../utils/sse-debug"; import { mapToOpenAIResponsesToolChoice } from "../utils/tool-choice"; -import { normalizeOpenAIResponsesPromptCacheKey, supportsDeveloperRole } from "./openai-responses"; +import { getOpenAIResponsesCacheSessionId, supportsDeveloperRole } from "./openai-responses"; import { appendResponsesToolResultMessages, applyCommonResponsesSamplingParams, @@ -156,10 +157,7 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" } let openaiStream: AsyncIterable; try { - const requestOptions = - requestTimeoutMs === undefined - ? { signal: requestSignal } - : { signal: requestSignal, timeout: requestTimeoutMs }; + const requestOptions = createSdkStreamRequestOptions(requestSignal, requestTimeoutMs); openaiStream = await client.responses.create(params, requestOptions); } catch (error) { if (error instanceof OpenAIConnectionTimeoutError && !abortTracker.wasCallerAbort()) { @@ -306,7 +304,10 @@ function buildParams( model: deploymentName, input: messages, stream: true, - prompt_cache_key: normalizeOpenAIResponsesPromptCacheKey(options?.promptCacheKey ?? options?.sessionId), + prompt_cache_key: getOpenAIResponsesCacheSessionId(options), + // Encrypted reasoning replay (applyResponsesReasoningParams) requires + // stateless responses, matching the openai provider. + store: false, }; applyCommonResponsesSamplingParams(params, options, model); @@ -332,6 +333,7 @@ function convertMessages( const messages: ResponseInput = []; const transformedMessages = transformMessages(context.messages, model, normalizeResponsesToolCallIdForTransform); const knownCallIds = new Set(); + const customCallIds = new Set(); const systemPrompts = normalizeSystemPrompts(context.systemPrompt); if (systemPrompts.length > 0) { @@ -351,11 +353,18 @@ function convertMessages( content: msg.role === "developer" && typeof msg.content === "string" ? msg.content.toWellFormed() : content, }); } else if (msg.role === "assistant") { - const outputItems = convertResponsesAssistantMessage(msg as AssistantMessage, model, msgIndex, knownCallIds); + const outputItems = convertResponsesAssistantMessage( + msg as AssistantMessage, + model, + msgIndex, + knownCallIds, + true, + customCallIds, + ); if (outputItems.length === 0) continue; messages.push(...outputItems); } else if (msg.role === "toolResult") { - appendResponsesToolResultMessages(messages, msg, model, strictResponsesPairing, knownCallIds); + appendResponsesToolResultMessages(messages, msg, model, strictResponsesPairing, knownCallIds, customCallIds); } msgIndex++; } diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 17ff4a411..021044bdc 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -136,7 +136,6 @@ const CODEX_RETRYABLE_EVENT_MESSAGE = const CODEX_PROVIDER_SESSION_STATE_KEY = "openai-codex-responses"; const X_CODEX_TURN_STATE_HEADER = "x-codex-turn-state"; const X_MODELS_ETAG_HEADER = "x-models-etag"; -const X_REASONING_INCLUDED_HEADER = "x-reasoning-included"; /** Connection-level websocket failures that should immediately fall back to SSE without retrying. */ const CODEX_WEBSOCKET_FATAL_PATTERNS = ["websocket error:", "websocket closed before open", "connection timeout"]; /** Max total time to spend retrying 429s with server-provided delays (5 minutes). */ @@ -196,7 +195,6 @@ type CodexWebSocketSessionState = { canAppend: boolean; turnState?: string; modelsEtag?: string; - reasoningIncluded?: boolean; connection?: CodexWebSocketConnection; lastTransport?: CodexTransport; fallbackCount: number; @@ -383,6 +381,7 @@ function isCodexWebSocketRetryableStreamError(error: unknown): boolean { message.includes("websocket ping failed") || message.includes("websocket pong timeout") || message.includes("websocket message queue exceeded") || + message.includes("websocket request already in progress") || message.includes("idle timeout waiting for websocket") || message.includes("timeout waiting for first websocket event") || message.includes("syntaxerror") || @@ -434,11 +433,6 @@ function updateCodexSessionMetadataFromHeaders( if (modelsEtag && modelsEtag.length > 0) { state.modelsEtag = modelsEtag; } - const reasoningIncluded = resolvedHeaders.get(X_REASONING_INCLUDED_HEADER); - if (reasoningIncluded !== null) { - const normalized = reasoningIncluded.trim().toLowerCase(); - state.reasoningIncluded = normalized.length === 0 ? true : normalized !== "false"; - } } function extractCodexWebSocketHandshakeHeaders(socket: Bun.WebSocket, openEvent?: Event): Headers | undefined { @@ -714,9 +708,9 @@ async function buildTransformedCodexRequestBody( prompt_cache_key: promptCacheKey, }; - if (options?.maxTokens) { - params.max_output_tokens = options.maxTokens; - } + // `maxTokens` is intentionally not forwarded: transformRequestBody strips + // `max_output_tokens`/`max_completion_tokens` (the Codex backend rejects + // caller-supplied output caps). if (options?.temperature !== undefined) { params.temperature = options.temperature; } @@ -766,7 +760,7 @@ async function buildTransformedCodexRequestBody( const developerMessages = systemPrompts.slice(1); const codexOptions: CodexRequestOptions = { reasoningEffort: options?.reasoning, - reasoningSummary: options?.reasoningSummary ?? "auto", + reasoningSummary: options?.reasoningSummary === undefined ? "auto" : options.reasoningSummary, textVerbosity: options?.textVerbosity, include: options?.include, }; @@ -1294,12 +1288,16 @@ function handleToolCallArgumentsDelta( stream: AssistantMessageEventStream, output: AssistantMessage, ): 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 + // 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; - const delta = (rawEvent as { delta?: string }).delta || ""; - const interruption = observeWhitespaceToolCallArgumentsDelta(runtime, rawEvent, delta); - if (interruption) return interruption; currentBlock.partialJson += delta; const throttled = parseStreamingJsonThrottled(currentBlock.partialJson, currentBlock.lastParseLen ?? 0); if (throttled) { @@ -1331,12 +1329,13 @@ function handleCustomToolCallInputDelta( stream: AssistantMessageEventStream, output: AssistantMessage, ): CodexWhitespaceToolCallArgumentsDeltaInterruption | undefined { + const delta = (rawEvent as { delta?: string }).delta || ""; + // 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; - const delta = (rawEvent as { delta?: string }).delta || ""; - const interruption = observeWhitespaceToolCallArgumentsDelta(runtime, rawEvent, delta); - if (interruption) return interruption; currentBlock.partialJson += delta; (currentBlock.arguments as { input?: string }).input = currentBlock.partialJson; stream.push({ type: "toolcall_delta", contentIndex: output.content.length - 1, delta, partial: output }); @@ -1363,7 +1362,9 @@ function handleOutputItemDone( runtime: CodexStreamRuntime, rawEvent: Record, ): void { - const item = structuredCloneJSON(rawEvent.item) as CodexEventItem; + const rawItem = rawEvent.item; + if (!rawItem || typeof rawItem !== "object") return; + const item = structuredCloneJSON(rawItem) as CodexEventItem; runtime.nativeOutputItems.push(item as unknown as Record); if (item.type === "reasoning" && runtime.currentBlock?.type === "thinking") { @@ -1485,8 +1486,11 @@ function handleResponseCompleted( if (typeof response?.id === "string" && response.id.length > 0) { state.lastResponseId = response.id; state.lastResponseItems = stripInputItemIds(structuredCloneJSON(runtime.nativeOutputItems)); + state.canAppend = rawEvent.type === "response.done" || rawEvent.type === "response.completed"; + } else { + // Without a response id the append baseline cannot be trusted. + state.canAppend = false; } - state.canAppend = rawEvent.type === "response.done" || rawEvent.type === "response.completed"; } calculateCost(model, output.usage); @@ -1541,7 +1545,9 @@ function dropTrailingDegenerateToolCall(output: AssistantMessage, runtime: Codex * scratch — bounded by {@link CODEX_WHITESPACE_LOOP_RETRY_LIMIT}. Sampling * nondeterminism usually breaks the loop on a fresh attempt; once the budget is * exhausted the original error is surfaced (now without the junk tool call - * polluting the message). + * polluting the message). Replay is refused once a toolcall_end was already + * delivered to the consumer (`canSafelyReplayWebsocketOverSse`) — it would + * re-emit the same tool calls. */ async function tryRecoverCodexWhitespaceToolCallLoop( context: CodexStreamProcessingContext, @@ -1554,7 +1560,11 @@ async function tryRecoverCodexWhitespaceToolCallLoop( // Drop the half-built degenerate tool call whether or not we retry, so it // never reaches the caller's message. dropTrailingDegenerateToolCall(context.output, runtime); - if (runtime.whitespaceLoopRetries >= CODEX_WHITESPACE_LOOP_RETRY_LIMIT || context.options?.signal?.aborted) { + if ( + runtime.whitespaceLoopRetries >= CODEX_WHITESPACE_LOOP_RETRY_LIMIT || + !runtime.canSafelyReplayWebsocketOverSse || + context.options?.signal?.aborted + ) { return false; } @@ -1641,8 +1651,24 @@ async function tryReconnectCodexWebSocketOnConnectionLimit( return true; } - // No content emitted yet — reconnect over websocket. + // No content emitted yet — clear accumulator state from the failed attempt + // (blockless native items can exist even with empty content) and reconnect + // 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; + context.firstTokenTime = undefined; + if (runtime.websocketStreamRetries >= getCodexWebSocketRetryBudget()) { + recordCodexWebSocketFailure(websocketState, true); + await reopenCodexSseRuntimeStream(context, runtime, websocketState); + return true; + } runtime.websocketStreamRetries += 1; + await scheduler.wait(getCodexWebSocketRetryDelayMs(runtime.websocketStreamRetries), { + signal: context.requestSetup.requestSignal, + }); await reopenCodexWebSocketRuntimeStream(context, runtime, websocketState); return true; } @@ -1718,6 +1744,13 @@ async function tryReplayWebsocketFailureOverSse( if (!activateFallback) { runtime.websocketStreamRetries += 1; + // 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; + context.firstTokenTime = undefined; await scheduler.wait(getCodexWebSocketRetryDelayMs(runtime.websocketStreamRetries), { signal: context.requestSetup.requestSignal, }); @@ -1725,13 +1758,11 @@ async function tryReplayWebsocketFailureOverSse( return true; } - if (replayingBufferedOutputOverSse) { - runtime.currentItem = null; - runtime.currentBlock = null; - runtime.nativeOutputItems.length = 0; - resetOutputState(context.output); - context.firstTokenTime = undefined; - } + runtime.currentItem = null; + runtime.currentBlock = null; + runtime.nativeOutputItems.length = 0; + resetOutputState(context.output); + context.firstTokenTime = undefined; await reopenCodexSseRuntimeStream(context, runtime, state); return true; @@ -1906,8 +1937,19 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" }, startTime, } satisfies CodexStreamProcessingContext); - const failure = await handleCodexStreamFailure(failureContext, error); - stream.push({ type: "error", reason: failure.stopReason as "error" | "aborted", error: failure }); + try { + const failure = await handleCodexStreamFailure(failureContext, error); + stream.push({ type: "error", reason: failure.stopReason as "error" | "aborted", error: failure }); + } catch (failureError) { + // Last resort — the failure handler itself threw (exotic error object or + // request-dump formatting). Never leave the stream un-ended. + logger.error("Codex stream failure handler threw", { + error: failureError instanceof Error ? failureError.message : String(failureError), + }); + output.stopReason = "error"; + output.errorMessage ??= error instanceof Error ? error.message : String(error); + stream.push({ type: "error", reason: "error", error: output }); + } stream.end(); } })(); @@ -2022,13 +2064,18 @@ function resetCodexWebSocketAppendState(state: CodexWebSocketSessionState): void function resetCodexSessionMetadata(state: CodexWebSocketSessionState): void { state.turnState = undefined; state.modelsEtag = undefined; - state.reasoningIncluded = undefined; } function recordCodexWebSocketFailure(state: CodexWebSocketSessionState, activateFallback: boolean): void { resetCodexWebSocketAppendState(state); - state.connection?.close("fallback"); - state.connection = undefined; + // Never tear down a CONNECTING socket: it belongs to a concurrent caller's + // in-flight handshake (prewarm/request race); closing it would reject that + // caller with a fatal "websocket closed before open" and disable websockets + // for the whole session. + if (state.connection && !state.connection.isConnecting()) { + state.connection.close("fallback"); + state.connection = undefined; + } state.lastFallbackAt = Date.now(); if (activateFallback && !state.disableWebsocket) { state.disableWebsocket = true; @@ -2436,6 +2483,9 @@ class CodexWebSocketConnection { if (this.#activeRequest) { throw createCodexWebSocketTransportError("websocket request already in progress"); } + if (signal?.aborted) { + throw createCodexWebSocketTransportError("request was aborted"); + } this.#activeRequest = true; this.#streamObserver = onSseEvent; // Drain any non-error frames left over from a prior request before sending. @@ -2453,13 +2503,7 @@ class CodexWebSocketConnection { this.close("aborted"); this.#push(createCodexWebSocketTransportError("request was aborted")); }; - if (signal) { - if (signal.aborted) { - onAbort(); - } else { - signal.addEventListener("abort", onAbort, { once: true }); - } - } + if (signal) signal.addEventListener("abort", onAbort, { once: true }); try { const debugSession = isRequestDebugEnabled() @@ -2477,8 +2521,13 @@ class CodexWebSocketConnection { const requestPayload = JSON.stringify(request); notifyCodexWebSocketOutbound(onSseEvent, request, requestPayload); + // Re-check liveness: the debug-session await above can outlive the socket. + const socket = this.#socket; + if (!socket || socket.readyState !== WebSocket.OPEN) { + throw createCodexWebSocketTransportError("websocket connection is unavailable"); + } try { - this.#socket.send(requestPayload); + socket.send(requestPayload); } catch (error) { throw createCodexWebSocketTransportError( `websocket send failed: ${error instanceof Error ? error.message : String(error)}`, @@ -2760,13 +2809,16 @@ async function getOrCreateCodexWebSocketConnection( // CONNECTING socket rejects the concurrent caller (prewarm racing the first // request) with a fatal "websocket closed before open", which would disable // websockets for the entire session. - const pending = state.connection; - if (pending && !pending.isOpen() && pending.isConnecting()) { + // Bounded re-join: a fresh handshake may have been started by yet another + // caller while we awaited the previous one. + for (let joinAttempt = 0; joinAttempt < 3; joinAttempt += 1) { + const pending = state.connection; + if (!pending || pending.isOpen() || !pending.isConnecting()) break; try { await pending.connect(signal); } catch { - // The handshake owner surfaces its own failure; fall through and - // re-evaluate (state.connection may have been replaced or cleared). + // The handshake owner surfaces its own failure; re-evaluate below + // (state.connection may have been replaced or cleared). } } if (state.connection?.isOpen()) { @@ -2893,6 +2945,7 @@ function createCodexHeaders( } else { headers.delete(OPENAI_HEADERS.CONVERSATION_ID); headers.delete(OPENAI_HEADERS.SESSION_ID); + headers.delete("x-client-request-id"); } if (state?.turnState) { headers.set(X_CODEX_TURN_STATE_HEADER, state.turnState); @@ -2931,6 +2984,7 @@ function redactHeaders(headers: Headers): Record { lower.includes("account") || lower.includes("session") || lower.includes("conversation") || + lower === "x-client-request-id" || lower === "cookie" ) { redacted[key] = "[redacted]"; @@ -3010,11 +3064,13 @@ function convertMessages(model: Model<"openai-codex-responses">, context: Contex if (msg.role === "assistant") { const assistantMsg = msg as AssistantMessage; - const providerPayload = getOpenAIResponsesHistoryPayload( - assistantMsg.providerPayload, - model.provider, - assistantMsg.provider, - ); + // Native items are model-bound (reasoning carries encrypted content + // minted by the producing model); after a mid-session model switch fall + // back to block re-encode, which strips foreign signatures. + const providerPayload = + assistantMsg.api === model.api && assistantMsg.model === model.id + ? getOpenAIResponsesHistoryPayload(assistantMsg.providerPayload, model.provider, assistantMsg.provider) + : undefined; const historyItems = providerPayload?.items as Array | undefined; if (historyItems) { for (const item of historyItems) { @@ -3167,7 +3223,9 @@ function isRetryableCodexFailureEvent(rawEvent: Record): boolea } function createCodexProviderStreamError(rawEvent: Record): CodexProviderStreamError { - const code = getString(rawEvent.code) ?? ""; + const response = asRecord(rawEvent.response); + const nestedError = asRecord(rawEvent.error) ?? (response ? asRecord(response.error) : null); + const code = getString(rawEvent.code) ?? getString(nestedError?.code) ?? getString(nestedError?.type) ?? ""; const message = getString(rawEvent.message) ?? ""; const formattedMessage = typeof rawEvent.type === "string" && rawEvent.type === "error" diff --git a/packages/ai/src/providers/openai-codex/request-transformer.ts b/packages/ai/src/providers/openai-codex/request-transformer.ts index 91abe9dc5..342708b2e 100644 --- a/packages/ai/src/providers/openai-codex/request-transformer.ts +++ b/packages/ai/src/providers/openai-codex/request-transformer.ts @@ -105,8 +105,8 @@ function orphanFunctionOutputToMessage(item: InputItem, callId: string): InputIt * Repair both halves of unpaired tool exchanges so the Responses input grammar * stays valid — the API rejects either orphan with a 400: * - * - `function_call_output` with no matching `function_call` → folded into an - * assistant message (`400 No tool call found for function call output …`). + * - `function_call_output` / `custom_tool_call_output` with no matching call → + * folded into an assistant message (`400 No tool call found for … output`). * Regression of #472 / #1351. * - `function_call` / `custom_tool_call` with no matching `*_output` → a * placeholder output is synthesized immediately after the call @@ -131,7 +131,11 @@ function repairToolCallPairs(input: InputItem[]): InputItem[] { for (const item of input) { const callId = typeof item.call_id === "string" ? item.call_id : undefined; - if (item.type === "function_call_output" && callId !== undefined && !callIds.has(callId)) { + if ( + (item.type === "function_call_output" || item.type === "custom_tool_call_output") && + callId !== undefined && + !callIds.has(callId) + ) { repaired.push(orphanFunctionOutputToMessage(item, callId)); continue; } diff --git a/packages/ai/src/providers/openai-responses-shared.ts b/packages/ai/src/providers/openai-responses-shared.ts index ac2d53936..2652d1821 100644 --- a/packages/ai/src/providers/openai-responses-shared.ts +++ b/packages/ai/src/providers/openai-responses-shared.ts @@ -1,4 +1,4 @@ -import { structuredCloneJSON } from "@oh-my-pi/pi-utils"; +import { logger, structuredCloneJSON } from "@oh-my-pi/pi-utils"; import type OpenAI from "openai"; import type { ResponseCustomToolCall, @@ -311,6 +311,7 @@ export function convertResponsesAssistantMessage( customCallIds?: Set, ): ResponseInput { const outputItems: ResponseInput = []; + let unsignedTextBlocks = 0; const isDifferentModel = assistantMsg.model !== model.id && assistantMsg.provider === model.provider && assistantMsg.api === model.api; @@ -320,7 +321,12 @@ export function convertResponsesAssistantMessage( continue; } if (block.thinkingSignature) { - outputItems.push(JSON.parse(block.thinkingSignature) as ResponseReasoningItem); + try { + outputItems.push(JSON.parse(block.thinkingSignature) as ResponseReasoningItem); + } catch { + // Legacy/corrupt persisted signature — skip the reasoning item + // rather than failing the whole request build. + } } continue; } @@ -329,7 +335,10 @@ export function convertResponsesAssistantMessage( const parsedSignature = parseTextSignature(block.textSignature); let msgId = parsedSignature?.id; if (!msgId) { - msgId = `msg_${msgIndex}`; + // Distinct ids per unsigned block: several text blocks in one message + // (cross-provider replay downgrades thinking → text) must not share an id. + msgId = unsignedTextBlocks === 0 ? `msg_${msgIndex}` : `msg_${msgIndex}_${unsignedTextBlocks}`; + unsignedTextBlocks += 1; } else if (msgId.length > 64) { msgId = `msg_${Bun.hash(msgId).toString(36)}`; } @@ -394,10 +403,6 @@ export function appendResponsesToolResultMessages( const hasImages = toolResult.content.some((block): block is ImageContent => block.type === "image"); const omittedImages = hasImages && !supportsImages; const normalized = normalizeResponsesToolCallId(toolResult.toolCallId); - if (strictResponsesPairing && !knownCallIds.has(normalized.callId)) { - return; - } - const output = ( omittedImages ? joinTextWithImagePlaceholder(textResult, true) @@ -405,6 +410,19 @@ export function appendResponsesToolResultMessages( ? textResult : "(see attached image)" ).toWellFormed(); + if (strictResponsesPairing && !knownCallIds.has(normalized.callId)) { + // Strict backends (Azure, Copilot) reject unpaired outputs outright, but + // silently dropping the result loses information the model needs. Fold it + // into an assistant note instead (same shape as repairOrphanResponsesToolOutputs). + const limit = 16_000; + const noteText = output.length > limit ? `${output.slice(0, limit)}\n...[truncated]` : output; + messages.push({ + type: "message", + role: "assistant", + content: `[Orphan ${toolResult.toolName || "tool"} result; call_id=${normalized.callId}]: ${noteText}`, + } as ResponseInput[number]); + return; + } if (customCallIds?.has(normalized.callId)) { messages.push({ type: "custom_tool_call_output", @@ -733,9 +751,15 @@ export async function processResponsesStream( : item.content?.[0]?.type === "reasoning_text" ? (item.content[0].text ?? "") : ""; - const reasoningBlock = output.content.find( - b => b.type === "thinking" && (b as ThinkingContent).itemId === item.id, - ) as ThinkingContent | undefined; + // Prefer the routed entry; the bare itemId find misroutes when ids are + // absent (`undefined === undefined` matches the FIRST thinking block) and + // misses entirely when the done-event id drifts from the added-event id. + const reasoningBlock = + entry?.block.type === "thinking" + ? entry.block + : (output.content.find(b => b.type === "thinking" && (b as ThinkingContent).itemId === item.id) as + | ThinkingContent + | undefined); if (reasoningBlock) { reasoningBlock.thinking = thinking; reasoningBlock.thinkingSignature = JSON.stringify(item); @@ -747,18 +771,25 @@ export async function processResponsesStream( }); } closeOpenItem(event.output_index, item.id, entry); - } else if (item.type === "message" && entry?.block.type === "text") { - const block = entry.block; - block.text = item.content + } else if (item.type === "message") { + const block = entry?.block.type === "text" ? entry.block : undefined; + const text = item.content .map(part => (part.type === "output_text" ? (part.text ?? "") : (part.refusal ?? ""))) .join(""); - block.textSignature = encodeTextSignatureV1(item.id, item.phase ?? undefined); - stream.push({ - type: "text_end", - contentIndex: contentIndexOf(block), - content: block.text, - partial: output, - }); + const textSignature = encodeTextSignatureV1(item.id, item.phase ?? undefined); + let contentIndex: number; + if (block) { + block.text = text; + block.textSignature = textSignature; + contentIndex = contentIndexOf(block); + } else { + // `output_item.added` never arrived (lossy proxy) — synthesize the + // block so the final message still carries the authoritative text. + const synthesized: TextContent = { type: "text", text, textSignature }; + output.content.push(synthesized); + contentIndex = output.content.length - 1; + } + stream.push({ type: "text_end", contentIndex, content: text, partial: output }); closeOpenItem(event.output_index, item.id, entry); } else if (item.type === "function_call") { const block = entry?.block.type === "toolCall" ? entry.block : undefined; @@ -773,6 +804,7 @@ export async function processResponsesStream( name: item.name, arguments: args, }; + let contentIndex: number; if (block) { // Persist the authoritative final args on the stored block. The // throttled delta parser may have skipped the last partial parse, @@ -782,8 +814,14 @@ export async function processResponsesStream( delete (block as { partialJson?: string }).partialJson; delete (block as { lastParseLen?: number }).lastParseLen; delete (block as { argumentsDone?: boolean }).argumentsDone; + contentIndex = contentIndexOf(block); + } else { + // `output_item.added` never arrived (lossy proxy) — synthesize the + // block so the final message carries the call the consumer was told + // completed (the agent loop executes tools from message.content). + output.content.push(toolCall); + contentIndex = output.content.length - 1; } - const contentIndex = block ? contentIndexOf(block) : output.content.length - 1; closeOpenItem(event.output_index, item.id, entry, item.call_id); stream.push({ type: "toolcall_end", contentIndex, toolCall, partial: output }); } else if (item.type === "custom_tool_call") { @@ -796,19 +834,39 @@ export async function processResponsesStream( arguments: { input: rawInput }, customWireName: item.name, }; + let contentIndex: number; if (block) { // Persist the final input on the stored block and drop the transient // accumulation buffer, mirroring the function_call branch above. block.arguments = { input: rawInput }; delete (block as { partialJson?: string }).partialJson; delete (block as { lastParseLen?: number }).lastParseLen; + contentIndex = contentIndexOf(block); + } else { + output.content.push(toolCall); + contentIndex = output.content.length - 1; } - const contentIndex = block ? contentIndexOf(block) : output.content.length - 1; closeOpenItem(event.output_index, item.id, entry, item.call_id); stream.push({ type: "toolcall_end", contentIndex, toolCall, partial: output }); } } else if (event.type === "response.completed" || event.type === "response.incomplete") { const response = event.response; + // Finalize any toolCall block whose output_item.done never arrived: the + // throttled delta parser may have left block.arguments stale, and the + // toolUse override below would hand the agent incomplete arguments. + for (const open of openItemsInOrder) { + if (open.block.type !== "toolCall") continue; + const block = open.block; + if (block.partialJson && !block.argumentsDone) { + block.arguments = + open.item.type === "custom_tool_call" + ? { input: block.partialJson } + : parseStreamingJson(block.partialJson); + } + delete (block as { partialJson?: string }).partialJson; + delete (block as { lastParseLen?: number }).lastParseLen; + delete (block as { argumentsDone?: boolean }).argumentsDone; + } if (response?.id) { output.responseId = response.id; } @@ -828,12 +886,19 @@ export async function processResponsesStream( : "Unknown error (no error details in response)"; throw new Error(message); } + if (response?.status === "incomplete" && response.incomplete_details?.reason === "content_filter") { + // A content-filtered turn is a failure, not a token-cap truncation — + // mapping it to "length" would route the agent loop into "shorten your + // output" recovery against a filtered prompt. + throw new Error("incomplete: content_filter"); + } if (output.content.some(block => block.type === "toolCall") && output.stopReason === "stop") { output.stopReason = "toolUse"; } } else if (event.type === "error") { throw new Error(`Error Code ${event.code}: ${event.message}`); } else if (event.type === "response.failed") { + populateResponsesUsageFromResponse(output, event.response?.usage); const error = event.response?.error ?? (event.response as any)?.status_details?.error; const details = event.response?.incomplete_details; const message = error @@ -860,8 +925,12 @@ export function mapOpenAIResponsesStopReason(status: OpenAI.Responses.ResponseSt case "queued": return "stop"; default: { + // Compile-time exhaustiveness; at runtime a brand-new status from the + // server must degrade gracefully instead of failing a fully-streamed + // response. const exhaustive: never = status; - throw new Error(`Unhandled stop reason: ${exhaustive}`); + logger.warn("Unhandled OpenAI Responses stop reason", { status: exhaustive }); + return "stop"; } } } @@ -967,7 +1036,9 @@ export function applyResponsesReasoningParams

0 ? { reasoningTokens } : {}), cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }; + if (premiumRequests !== undefined) { + output.usage.premiumRequests = premiumRequests; + } } diff --git a/packages/ai/src/providers/openai-responses.ts b/packages/ai/src/providers/openai-responses.ts index a9a7c579d..d385d0e06 100644 --- a/packages/ai/src/providers/openai-responses.ts +++ b/packages/ai/src/providers/openai-responses.ts @@ -271,6 +271,12 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( const { data, response, request_id } = await client.responses .create(params, requestOptions) .withResponse(); + // Disarm the first-event watchdog as soon as headers arrive — a slow + // onResponse callback must not abort an already-connected stream. + if (requestTimeout !== undefined) { + clearTimeout(requestTimeout); + requestTimeout = undefined; + } await notifyProviderResponse(options, response, model, request_id); return data; } catch (error) { @@ -309,7 +315,6 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( nativeOutputItems.push(structuredCloneJSON(item) as unknown as Record); }, }); - if (premiumRequestsTotal !== undefined) output.usage.premiumRequests = premiumRequestsTotal; const firstEventTimeoutError = abortTracker.getLocalAbortReason(); if (firstEventTimeoutError) { @@ -464,7 +469,9 @@ function buildParams( instructions: systemInstructions, stream: true, prompt_cache_key: promptCacheKey, - prompt_cache_retention: promptCacheKey ? getPromptCacheRetention(model.baseUrl, cacheRetention) : undefined, + prompt_cache_retention: promptCacheKey + ? getPromptCacheRetention(resolvedBaseUrl ?? model.baseUrl, cacheRetention) + : undefined, store: false, stream_options: model.provider === "openai" ? { include_obfuscation: false } : undefined, }; @@ -587,9 +594,13 @@ function convertConversationMessages( messages.push({ role: "user", content }); } else if (msg.role === "assistant") { const assistantMsg = msg as AssistantMessage; - const providerPayload = shouldReplayNativeHistory - ? getOpenAIResponsesHistoryPayload(assistantMsg.providerPayload, model.provider, assistantMsg.provider) - : undefined; + // Native items are model-bound (reasoning carries encrypted content minted + // by the producing model); after a mid-session model switch fall back to + // block re-encode, which strips foreign signatures. + const providerPayload = + shouldReplayNativeHistory && assistantMsg.api === model.api && assistantMsg.model === model.id + ? getOpenAIResponsesHistoryPayload(assistantMsg.providerPayload, model.provider, assistantMsg.provider) + : undefined; const historyItems = providerPayload?.items; if (historyItems) { const sanitizedHistoryItems = sanitizeOpenAIResponsesHistoryItemsForReplay(filterReasoning(historyItems)); diff --git a/packages/ai/src/providers/transform-messages.ts b/packages/ai/src/providers/transform-messages.ts index a129c1b2f..34346d99c 100644 --- a/packages/ai/src/providers/transform-messages.ts +++ b/packages/ai/src/providers/transform-messages.ts @@ -19,9 +19,24 @@ const enum ToolCallStatus { const MAX_TOOL_CALL_ID_LENGTH = 64; function appendDuplicateSuffix(originalId: string, suffix: string, maxLength: number): string { - if (originalId.length + suffix.length <= maxLength) return `${originalId}${suffix}`; + // Responses-family ids are composites (`callId|itemId`): the wire call_id is + // the FIRST segment (normalizeResponsesToolCallId splits on `|`), so the + // suffix must land on every segment or the duplicate collapses back onto the + // original call_id at encode time. The length budget applies per segment, + // matching the per-segment caps of the provider normalizers. + if (originalId.includes("|")) { + return originalId + .split("|") + .map(segment => appendSegmentDuplicateSuffix(segment, suffix, maxLength)) + .join("|"); + } + return appendSegmentDuplicateSuffix(originalId, suffix, maxLength); +} + +function appendSegmentDuplicateSuffix(segment: string, suffix: string, maxLength: number): string { + if (segment.length + suffix.length <= maxLength) return `${segment}${suffix}`; const prefixBudget = Math.max(0, maxLength - suffix.length); - return `${originalId.slice(0, prefixBudget)}${suffix}`; + return `${segment.slice(0, prefixBudget)}${suffix}`; } type PendingToolResultRewrite = { replacementId: string } | undefined;