fix(openai-responses-providers): patched response flow and pairing logic

- Fixed Azure and OpenAI response flows by using stream request options, cache metadata, and replay checks.
- Fixed Codex websocket flow by clearing runtime state on reconnect failures and OPEN-state sends.
- Fixed response parsing robustness by skipping malformed signatures and repairing orphaned outputs.
- Fixed tool/result handling by mapping unknown call IDs, splitting composite IDs, and truncating duplicates.
This commit is contained in:
can1357
2026-06-09 23:17:26 +02:00
parent 647eb1e47e
commit 228c1ddd8d
6 changed files with 269 additions and 94 deletions
@@ -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<ResponseStreamEvent>;
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<string>();
const customCallIds = new Set<string>();
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++;
}
@@ -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<string, unknown>,
): 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<string, unknown>);
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<string, string> {
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<ResponseInput[number]> | undefined;
if (historyItems) {
for (const item of historyItems) {
@@ -3167,7 +3223,9 @@ function isRetryableCodexFailureEvent(rawEvent: Record<string, unknown>): boolea
}
function createCodexProviderStreamError(rawEvent: Record<string, unknown>): 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"
@@ -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;
}
@@ -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<TApi extends Api>(
customCallIds?: Set<string>,
): 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<TApi extends Api>(
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<TApi extends Api>(
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<TApi extends Api>(
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<TApi extends Api>(
? 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<TApi extends Api>(
: 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<TApi extends Api>(
});
}
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<TApi extends Api>(
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<TApi extends Api>(
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<TApi extends Api>(
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<TApi extends Api>(
: "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<P extends OpenAI.Responses.Respons
// multi-turn conversations when store is false (items aren't persisted server-side, so
// we must include the full content). See: https://github.com/can1357/oh-my-pi/issues/41
if (includeEncryptedReasoning) {
params.include = ["reasoning.encrypted_content"];
const include = params.include ?? [];
if (!include.includes("reasoning.encrypted_content")) include.push("reasoning.encrypted_content");
params.include = include;
}
if (options?.reasoning || options?.reasoningSummary !== undefined) {
@@ -1022,6 +1093,10 @@ export function populateResponsesUsageFromResponse(
if (!usage) return;
const cachedTokens = usage.input_tokens_details?.cached_tokens || 0;
const reasoningTokens = usage.output_tokens_details?.reasoning_tokens || 0;
// Wholesale replacement must not drop provider-annotated extras (Copilot
// premium-request accounting): the failed/cancelled paths throw right after
// this call with no later chance to re-apply.
const premiumRequests = output.usage.premiumRequests;
output.usage = {
input: (usage.input_tokens || 0) - cachedTokens,
output: usage.output_tokens || 0,
@@ -1031,4 +1106,7 @@ export function populateResponsesUsageFromResponse(
...(reasoningTokens > 0 ? { reasoningTokens } : {}),
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
if (premiumRequests !== undefined) {
output.usage.premiumRequests = premiumRequests;
}
}
+16 -5
View File
@@ -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<unknown>(item) as unknown as Record<string, unknown>);
},
});
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));
@@ -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;