diff --git a/packages/ai/src/providers/amazon-bedrock.ts b/packages/ai/src/providers/amazon-bedrock.ts index 5c2ad689f..12553360f 100644 --- a/packages/ai/src/providers/amazon-bedrock.ts +++ b/packages/ai/src/providers/amazon-bedrock.ts @@ -41,6 +41,7 @@ import type { ToolResultMessage, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump, withHttpStatus } from "../utils/http-inspector"; import { parseStreamingJson } from "../utils/json-parse"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; import { transformMessages } from "./transform-messages"; @@ -94,6 +95,7 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream"> = ( }; const blocks = output.content as Block[]; + let rawRequestDump: RawHttpRequestDump | undefined; const config: BedrockRuntimeClientConfig = { region: options.region, @@ -133,6 +135,14 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream"> = ( additionalModelRequestFields: buildAdditionalModelRequestFields(model, options), }; options?.onPayload?.(commandInput); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: `https://bedrock-runtime.${config.region}.amazonaws.com/model/${model.id}/converse-stream`, + body: commandInput, + }; const command = new ConverseStreamCommand(commandInput); const response = await client.send(command, { abortSignal: options.signal }); @@ -160,7 +170,7 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream"> = ( } else if (item.modelStreamErrorException) { throw new Error(`Model stream error: ${item.modelStreamErrorException.message}`); } else if (item.validationException) { - throw new Error(`Validation error: ${item.validationException.message}`); + throw withHttpStatus(new Error(`Validation error: ${item.validationException.message}`), 400); } else if (item.throttlingException) { throw new Error(`Throttling error: ${item.throttlingException.message}`); } else if (item.serviceUnavailableException) { @@ -186,7 +196,11 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream"> = ( delete (block as Block).partialJson; } output.stopReason = options.signal?.aborted ? "aborted" : "error"; - output.errorMessage = error instanceof Error ? error.message : JSON.stringify(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + error instanceof Error ? error.message : JSON.stringify(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index ef9fbc3c2..82bc87dfb 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -4,6 +4,7 @@ import type { MessageCreateParamsStreaming, MessageParam, } from "@anthropic-ai/sdk/resources/messages"; +import { abortableSleep } from "@oh-my-pi/pi-utils"; import { calculateCost } from "../models"; import { getEnvApiKey, OUTPUT_FALLBACK_BUFFER } from "../stream"; import type { @@ -25,6 +26,7 @@ import type { ToolResultMessage, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { parseStreamingJson } from "../utils/json-parse"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; @@ -290,6 +292,19 @@ function mergeHeaders(...headerSources: (Record | undefined)[]): return merged; } +const PROVIDER_MAX_RETRIES = 3; +const PROVIDER_BASE_DELAY_MS = 2000; + +/** + * Check if an error from the Anthropic SDK is a rate-limit or transient error + * that the SDK itself didn't retry (e.g. z.ai returns non-429 status with rate limit in body). + */ +function isProviderRetryableError(error: unknown): boolean { + if (!(error instanceof Error)) return false; + const msg = error.message; + return /rate.?limit|too many requests|overloaded|service.?unavailable|1302/i.test(msg); +} + export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( model: Model<"anthropic-messages">, context: Context, @@ -318,6 +333,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { const apiKey = options?.apiKey ?? getEnvApiKey(model.provider) ?? ""; @@ -342,163 +358,206 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( }); const params = buildParams(model, context, isOAuthToken, options); options?.onPayload?.(params); - const anthropicStream = client.messages.stream({ ...params, stream: true }, { signal: options?.signal }); - stream.push({ type: "start", partial: output }); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: `${model.baseUrl ?? "https://api.anthropic.com"}/v1/messages`, + body: params, + }; type Block = (ThinkingContent | TextContent | (ToolCall & { partialJson: string })) & { index: number }; const blocks = output.content as Block[]; + stream.push({ type: "start", partial: output }); + // Retry loop for rate-limit errors from proxies (e.g. z.ai) that the SDK doesn't handle. + // These errors surface when iterating the stream, so we retry the full stream creation. + // Only retry if no content blocks have been emitted yet (safe to restart). + let providerRetryAttempt = 0; + let started = false; + do { + const anthropicStream = client.messages.stream({ ...params, stream: true }, { signal: options?.signal }); - for await (const event of anthropicStream) { - if (event.type === "message_start") { - // Capture initial token usage from message_start event - // This ensures we have input token counts even if the stream is aborted early - output.usage.input = event.message.usage.input_tokens || 0; - output.usage.output = event.message.usage.output_tokens || 0; - output.usage.cacheRead = event.message.usage.cache_read_input_tokens || 0; - output.usage.cacheWrite = event.message.usage.cache_creation_input_tokens || 0; - // Anthropic doesn't provide total_tokens, compute from components - output.usage.totalTokens = - output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; - calculateCost(model, output.usage); - } else if (event.type === "content_block_start") { - if (!firstTokenTime) firstTokenTime = Date.now(); - if (event.content_block.type === "text") { - const block: Block = { - type: "text", - text: "", - index: event.index, - }; - output.content.push(block); - stream.push({ type: "text_start", contentIndex: output.content.length - 1, partial: output }); - } else if (event.content_block.type === "thinking") { - const block: Block = { - type: "thinking", - thinking: "", - thinkingSignature: "", - index: event.index, - }; - output.content.push(block); - stream.push({ type: "thinking_start", contentIndex: output.content.length - 1, partial: output }); - } else if (event.content_block.type === "tool_use") { - const block: Block = { - type: "toolCall", - id: event.content_block.id, - name: isOAuthToken ? fromClaudeCodeName(event.content_block.name) : event.content_block.name, - arguments: (event.content_block.input as Record) ?? {}, - partialJson: "", - index: event.index, - }; - output.content.push(block); - stream.push({ type: "toolcall_start", contentIndex: output.content.length - 1, partial: output }); - } - } else if (event.type === "content_block_delta") { - if (event.delta.type === "text_delta") { - const index = blocks.findIndex(b => b.index === event.index); - const block = blocks[index]; - if (block && block.type === "text") { - block.text += event.delta.text; - stream.push({ - type: "text_delta", - contentIndex: index, - delta: event.delta.text, - partial: output, - }); - } - } else if (event.delta.type === "thinking_delta") { - const index = blocks.findIndex(b => b.index === event.index); - const block = blocks[index]; - if (block && block.type === "thinking") { - block.thinking += event.delta.thinking; - stream.push({ - type: "thinking_delta", - contentIndex: index, - delta: event.delta.thinking, - partial: output, - }); - } - } else if (event.delta.type === "input_json_delta") { - const index = blocks.findIndex(b => b.index === event.index); - const block = blocks[index]; - if (block && block.type === "toolCall") { - block.partialJson += event.delta.partial_json; - block.arguments = parseStreamingJson(block.partialJson); - stream.push({ - type: "toolcall_delta", - contentIndex: index, - delta: event.delta.partial_json, - partial: output, - }); - } - } else if (event.delta.type === "signature_delta") { - const index = blocks.findIndex(b => b.index === event.index); - const block = blocks[index]; - if (block && block.type === "thinking") { - block.thinkingSignature = block.thinkingSignature || ""; - block.thinkingSignature += event.delta.signature; + try { + for await (const event of anthropicStream) { + started = true; + if (event.type === "message_start") { + // Capture initial token usage from message_start event + // This ensures we have input token counts even if the stream is aborted early + output.usage.input = event.message.usage.input_tokens || 0; + output.usage.output = event.message.usage.output_tokens || 0; + output.usage.cacheRead = event.message.usage.cache_read_input_tokens || 0; + output.usage.cacheWrite = event.message.usage.cache_creation_input_tokens || 0; + // Anthropic doesn't provide total_tokens, compute from components + output.usage.totalTokens = + output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; + calculateCost(model, output.usage); + } else if (event.type === "content_block_start") { + if (!firstTokenTime) firstTokenTime = Date.now(); + if (event.content_block.type === "text") { + const block: Block = { + type: "text", + text: "", + index: event.index, + }; + output.content.push(block); + stream.push({ type: "text_start", contentIndex: output.content.length - 1, partial: output }); + } else if (event.content_block.type === "thinking") { + const block: Block = { + type: "thinking", + thinking: "", + thinkingSignature: "", + index: event.index, + }; + output.content.push(block); + stream.push({ + type: "thinking_start", + contentIndex: output.content.length - 1, + partial: output, + }); + } else if (event.content_block.type === "tool_use") { + const block: Block = { + type: "toolCall", + id: event.content_block.id, + name: isOAuthToken ? fromClaudeCodeName(event.content_block.name) : event.content_block.name, + arguments: (event.content_block.input as Record) ?? {}, + partialJson: "", + index: event.index, + }; + output.content.push(block); + stream.push({ + type: "toolcall_start", + contentIndex: output.content.length - 1, + partial: output, + }); + } + } else if (event.type === "content_block_delta") { + if (event.delta.type === "text_delta") { + const index = blocks.findIndex(b => b.index === event.index); + const block = blocks[index]; + if (block && block.type === "text") { + block.text += event.delta.text; + stream.push({ + type: "text_delta", + contentIndex: index, + delta: event.delta.text, + partial: output, + }); + } + } else if (event.delta.type === "thinking_delta") { + const index = blocks.findIndex(b => b.index === event.index); + const block = blocks[index]; + if (block && block.type === "thinking") { + block.thinking += event.delta.thinking; + stream.push({ + type: "thinking_delta", + contentIndex: index, + delta: event.delta.thinking, + partial: output, + }); + } + } else if (event.delta.type === "input_json_delta") { + const index = blocks.findIndex(b => b.index === event.index); + const block = blocks[index]; + if (block && block.type === "toolCall") { + block.partialJson += event.delta.partial_json; + block.arguments = parseStreamingJson(block.partialJson); + stream.push({ + type: "toolcall_delta", + contentIndex: index, + delta: event.delta.partial_json, + partial: output, + }); + } + } else if (event.delta.type === "signature_delta") { + const index = blocks.findIndex(b => b.index === event.index); + const block = blocks[index]; + if (block && block.type === "thinking") { + block.thinkingSignature = block.thinkingSignature || ""; + block.thinkingSignature += event.delta.signature; + } + } + } else if (event.type === "content_block_stop") { + const index = blocks.findIndex(b => b.index === event.index); + const block = blocks[index]; + if (block) { + delete (block as { index?: number }).index; + if (block.type === "text") { + stream.push({ + type: "text_end", + contentIndex: index, + content: block.text, + partial: output, + }); + } else if (block.type === "thinking") { + stream.push({ + type: "thinking_end", + contentIndex: index, + content: block.thinking, + partial: output, + }); + } else if (block.type === "toolCall") { + block.arguments = parseStreamingJson(block.partialJson); + delete (block as { partialJson?: string }).partialJson; + stream.push({ + type: "toolcall_end", + contentIndex: index, + toolCall: block, + partial: output, + }); + } + } + } else if (event.type === "message_delta") { + if (event.delta.stop_reason) { + output.stopReason = mapStopReason(event.delta.stop_reason); + } + // Only update usage fields if present (not null). + // Preserves input_tokens from message_start when proxies omit it in message_delta. + if (event.usage.input_tokens != null) { + output.usage.input = event.usage.input_tokens; + } + if (event.usage.output_tokens != null) { + output.usage.output = event.usage.output_tokens; + } + if (event.usage.cache_read_input_tokens != null) { + output.usage.cacheRead = event.usage.cache_read_input_tokens; + } + if (event.usage.cache_creation_input_tokens != null) { + output.usage.cacheWrite = event.usage.cache_creation_input_tokens; + } + // Anthropic doesn't provide total_tokens, compute from components + output.usage.totalTokens = + output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; + calculateCost(model, output.usage); } } - } else if (event.type === "content_block_stop") { - const index = blocks.findIndex(b => b.index === event.index); - const block = blocks[index]; - if (block) { - delete (block as any).index; - if (block.type === "text") { - stream.push({ - type: "text_end", - contentIndex: index, - content: block.text, - partial: output, - }); - } else if (block.type === "thinking") { - stream.push({ - type: "thinking_end", - contentIndex: index, - content: block.thinking, - partial: output, - }); - } else if (block.type === "toolCall") { - block.arguments = parseStreamingJson(block.partialJson); - delete (block as any).partialJson; - stream.push({ - type: "toolcall_end", - contentIndex: index, - toolCall: block, - partial: output, - }); - } + + if (options?.signal?.aborted) { + throw new Error("Request was aborted"); } - } else if (event.type === "message_delta") { - if (event.delta.stop_reason) { - output.stopReason = mapStopReason(event.delta.stop_reason); + + if (output.stopReason === "aborted" || output.stopReason === "error") { + throw new Error("An unknown error occurred"); } - // Only update usage fields if present (not null). - // Preserves input_tokens from message_start when proxies omit it in message_delta. - if (event.usage.input_tokens != null) { - output.usage.input = event.usage.input_tokens; + break; // Stream completed successfully + } catch (streamError) { + // Only retry if: not aborted, no content emitted yet, retries left, and error is retryable + if ( + options?.signal?.aborted || + firstTokenTime !== undefined || + providerRetryAttempt >= PROVIDER_MAX_RETRIES || + !isProviderRetryableError(streamError) + ) { + throw streamError; } - if (event.usage.output_tokens != null) { - output.usage.output = event.usage.output_tokens; - } - if (event.usage.cache_read_input_tokens != null) { - output.usage.cacheRead = event.usage.cache_read_input_tokens; - } - if (event.usage.cache_creation_input_tokens != null) { - output.usage.cacheWrite = event.usage.cache_creation_input_tokens; - } - // Anthropic doesn't provide total_tokens, compute from components - output.usage.totalTokens = - output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; - calculateCost(model, output.usage); + providerRetryAttempt++; + const delayMs = PROVIDER_BASE_DELAY_MS * 2 ** (providerRetryAttempt - 1); + await abortableSleep(delayMs, options?.signal); + // Reset output state for clean retry + output.content.length = 0; + output.stopReason = "stop"; } - } - - if (options?.signal?.aborted) { - throw new Error("Request was aborted"); - } - - if (output.stopReason === "aborted" || output.stopReason === "error") { - throw new Error("An unknown error occurred"); - } + } while (!started); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; @@ -507,7 +566,11 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( } catch (error) { for (const block of output.content) delete (block as any).index; output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/azure-openai-responses.ts b/packages/ai/src/providers/azure-openai-responses.ts index a6e92d012..9238ac107 100644 --- a/packages/ai/src/providers/azure-openai-responses.ts +++ b/packages/ai/src/providers/azure-openai-responses.ts @@ -30,6 +30,7 @@ import type { ToolChoice, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { parseStreamingJson } from "../utils/json-parse"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; @@ -103,6 +104,7 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { // Create Azure OpenAI client @@ -110,6 +112,14 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" const client = createClient(model, apiKey, options); const params = buildParams(model, context, options, deploymentName); options?.onPayload?.(params); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: `${resolveAzureConfig(model, options).baseUrl}/responses`, + body: params, + }; const openaiStream = await client.responses.create( params, options?.signal ? { signal: options.signal } : undefined, @@ -351,7 +361,11 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" } catch (error) { for (const block of output.content) delete (block as { index?: number }).index; output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/google-gemini-cli.ts b/packages/ai/src/providers/google-gemini-cli.ts index 8cf459ed9..e3d72a273 100644 --- a/packages/ai/src/providers/google-gemini-cli.ts +++ b/packages/ai/src/providers/google-gemini-cli.ts @@ -19,6 +19,7 @@ import type { ToolCall, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump, withHttpStatus } from "../utils/http-inspector"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; import { convertMessages, @@ -316,6 +317,7 @@ export const streamGoogleGeminiCli: StreamFunction<"google-gemini-cli"> = ( stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { // apiKey is JSON-encoded: { token, projectId } @@ -356,6 +358,14 @@ export const streamGoogleGeminiCli: StreamFunction<"google-gemini-cli"> = ( ...(options?.headers ?? {}), }; const requestBodyJson = JSON.stringify(requestBody); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + body: requestBody, + headers: requestHeaders, + }; // Fetch with retry logic for rate limits and transient errors let response: Response | undefined; @@ -416,7 +426,10 @@ export const streamGoogleGeminiCli: StreamFunction<"google-gemini-cli"> = ( } // Not retryable or budget exceeded - throw new Error(`Cloud Code Assist API error (${response.status}): ${extractErrorMessage(errorText)}`); + throw withHttpStatus( + new Error(`Cloud Code Assist API error (${response.status}): ${extractErrorMessage(errorText)}`), + response.status, + ); } catch (error) { // Check for abort - fetch throws AbortError, our code throws "Request was aborted" if (error instanceof Error) { @@ -693,7 +706,10 @@ export const streamGoogleGeminiCli: StreamFunction<"google-gemini-cli"> = ( if (!currentResponse.ok) { const retryErrorText = await currentResponse.text(); - throw new Error(`Cloud Code Assist API error (${currentResponse.status}): ${retryErrorText}`); + throw withHttpStatus( + new Error(`Cloud Code Assist API error (${currentResponse.status}): ${retryErrorText}`), + currentResponse.status, + ); } } @@ -731,7 +747,11 @@ export const streamGoogleGeminiCli: StreamFunction<"google-gemini-cli"> = ( } } output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = error instanceof Error ? error.message : JSON.stringify(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + error instanceof Error ? error.message : JSON.stringify(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/google-vertex.ts b/packages/ai/src/providers/google-vertex.ts index 55b897e48..c13d9bbfc 100644 --- a/packages/ai/src/providers/google-vertex.ts +++ b/packages/ai/src/providers/google-vertex.ts @@ -19,6 +19,7 @@ import type { ToolCall, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; import type { GoogleThinkingLevel } from "./google-gemini-cli"; @@ -83,6 +84,7 @@ export const streamGoogleVertex: StreamFunction<"google-vertex"> = ( stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { const project = resolveProject(options); @@ -90,6 +92,14 @@ export const streamGoogleVertex: StreamFunction<"google-vertex"> = ( const client = createClient(model, project, location); const params = buildParams(model, context, options); options?.onPayload?.(params); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: `https://${location}-aiplatform.googleapis.com/v1/projects/${project}/locations/${location}/publishers/google/models/${model.id}:streamGenerateContent`, + body: params, + }; const googleStream = await client.models.generateContentStream(params); stream.push({ type: "start", partial: output }); @@ -275,7 +285,11 @@ export const streamGoogleVertex: StreamFunction<"google-vertex"> = ( } } output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/google.ts b/packages/ai/src/providers/google.ts index 6073c977d..d095b2cdf 100644 --- a/packages/ai/src/providers/google.ts +++ b/packages/ai/src/providers/google.ts @@ -18,6 +18,7 @@ import type { ToolCall, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; import type { GoogleThinkingLevel } from "./google-gemini-cli"; @@ -73,12 +74,21 @@ export const streamGoogle: StreamFunction<"google-generative-ai"> = ( stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { const apiKey = options?.apiKey || getEnvApiKey(model.provider) || ""; const client = createClient(model, apiKey); const params = buildParams(model, context, options); options?.onPayload?.(params); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: model.baseUrl ? `${model.baseUrl}/models/${model.id}:streamGenerateContent` : undefined, + body: params, + }; const googleStream = await client.models.generateContentStream(params); stream.push({ type: "start", partial: output }); @@ -261,7 +271,11 @@ export const streamGoogle: StreamFunction<"google-generative-ai"> = ( } } output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 218f94b9b..c8857350a 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -28,6 +28,7 @@ import type { ToolChoice, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { parseStreamingJson } from "../utils/json-parse"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; @@ -312,6 +313,7 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" }; let websocketState: CodexWebSocketSessionState | undefined; let usingWebsocket = false; + let rawRequestDump: RawHttpRequestDump | undefined; try { const apiKey = options?.apiKey || getEnvApiKey(model.provider) || ""; @@ -368,6 +370,14 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" const reasoningEffort = transformedBody.reasoning?.effort ?? null; const requestHeaders = { ...(model.headers ?? {}), ...(options?.headers ?? {}) }; + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url, + body: transformedBody, + }; const providerSessionState = getCodexProviderSessionState(options?.providerSessionState); const sessionKey = getCodexWebSocketSessionKey(options?.sessionId, model, accountId, baseUrl); const publicSessionKey = getCodexPublicSessionKey(options?.sessionId, model, baseUrl); @@ -795,7 +805,11 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" resetCodexSessionMetadata(websocketState); } output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); @@ -1302,7 +1316,8 @@ async function openCodexSseEventStream( if (!response.ok) { const info = await parseCodexError(response); const error = new Error(info.friendlyMessage || info.message); - (error as { headers?: Headers }).headers = response.headers; + (error as { headers?: Headers; status?: number }).headers = response.headers; + (error as { headers?: Headers; status?: number }).status = response.status; throw error; } if (!response.body) { diff --git a/packages/ai/src/providers/openai-completions.ts b/packages/ai/src/providers/openai-completions.ts index 958127640..c07cfed4d 100644 --- a/packages/ai/src/providers/openai-completions.ts +++ b/packages/ai/src/providers/openai-completions.ts @@ -28,6 +28,7 @@ import type { ToolResultMessage, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { parseStreamingJson } from "../utils/json-parse"; import { getKimiCommonHeaders } from "../utils/oauth/kimi"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; @@ -150,12 +151,21 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = ( stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { const apiKey = options?.apiKey || getEnvApiKey(model.provider) || ""; const client = await createClient(model, context, apiKey, options?.headers); const params = buildParams(model, context, options); options?.onPayload?.(params); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: `${model.baseUrl ?? "https://api.openai.com/v1"}/chat/completions`, + body: params, + }; const openaiStream = await client.chat.completions.create(params, { signal: options?.signal }); stream.push({ type: "start", partial: output }); @@ -434,7 +444,11 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = ( } catch (error) { for (const block of output.content) delete (block as any).index; output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); // Some providers via OpenRouter include extra details here. const rawMetadata = (error as { error?: { metadata?: { raw?: string } } })?.error?.metadata?.raw; if (rawMetadata) output.errorMessage += `\n${rawMetadata}`; diff --git a/packages/ai/src/providers/openai-responses.ts b/packages/ai/src/providers/openai-responses.ts index e319912e7..f66d118ec 100644 --- a/packages/ai/src/providers/openai-responses.ts +++ b/packages/ai/src/providers/openai-responses.ts @@ -29,6 +29,7 @@ import type { ToolChoice, } from "../types"; import { AssistantMessageEventStream } from "../utils/event-stream"; +import { appendRawHttpRequestDumpFor400, type RawHttpRequestDump } from "../utils/http-inspector"; import { parseStreamingJson } from "../utils/json-parse"; import { formatErrorMessageWithRetryAfter } from "../utils/retry-after"; import { sanitizeSurrogates } from "../utils/sanitize-unicode"; @@ -109,6 +110,7 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( stopReason: "stop", timestamp: Date.now(), }; + let rawRequestDump: RawHttpRequestDump | undefined; try { // Create OpenAI client @@ -116,6 +118,14 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( const client = createClient(model, context, apiKey, options?.headers); const params = buildParams(model, context, options); options?.onPayload?.(params); + rawRequestDump = { + provider: model.provider, + api: output.api, + model: model.id, + method: "POST", + url: `${model.baseUrl ?? "https://api.openai.com/v1"}/responses`, + body: params, + }; const openaiStream = await client.responses.create( params, options?.signal ? { signal: options.signal } : undefined, @@ -357,7 +367,11 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( } catch (error) { for (const block of output.content) delete (block as any).index; output.stopReason = options?.signal?.aborted ? "aborted" : "error"; - output.errorMessage = formatErrorMessageWithRetryAfter(error); + output.errorMessage = await appendRawHttpRequestDumpFor400( + formatErrorMessageWithRetryAfter(error), + error, + rawRequestDump, + ); output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); diff --git a/packages/ai/src/providers/transform-messages.ts b/packages/ai/src/providers/transform-messages.ts index 8806312e6..718fddb1d 100644 --- a/packages/ai/src/providers/transform-messages.ts +++ b/packages/ai/src/providers/transform-messages.ts @@ -95,6 +95,7 @@ export function transformMessages( const result: Message[] = []; let pendingToolCalls: ToolCall[] = []; let existingToolResultIds = new Set(); + const skippedToolCallIds = new Set(); for (let i = 0; i < transformed.length; i++) { const msg = transformed[i]; @@ -118,13 +119,20 @@ export function transformMessages( existingToolResultIds = new Set(); } - // Skip errored/aborted assistant messages entirely. - // These are incomplete turns that shouldn't be replayed: - // - May have partial content (reasoning without message, incomplete tool calls) - // - Replaying them can cause API errors (e.g., OpenAI "reasoning without following item") - // - The model should retry from the last valid state + // For errored/aborted assistant messages, convert tool calls to text summaries + // so the model retains awareness of what it attempted. This avoids orphaned + // tool_use blocks while preserving context for the retry. const assistantMsg = msg as AssistantMessage; if (assistantMsg.stopReason === "error" || assistantMsg.stopReason === "aborted") { + const patchedContent = assistantMsg.content.flatMap(block => { + if (block.type !== "toolCall") return block; + const tc = block as ToolCall; + skippedToolCallIds.add(tc.id); + return { type: "text" as const, text: `[Tool call aborted: ${tc.name}]` }; + }); + if (patchedContent.length > 0) { + result.push({ ...assistantMsg, content: patchedContent }); + } continue; } @@ -137,6 +145,8 @@ export function transformMessages( result.push(msg); } else if (msg.role === "toolResult") { + // Skip tool results whose corresponding assistant tool_use was from a skipped (aborted/errored) message + if (skippedToolCallIds.has(msg.toolCallId)) continue; existingToolResultIds.add(msg.toolCallId); result.push(msg); } else if (msg.role === "user") { diff --git a/packages/ai/src/utils/http-inspector.ts b/packages/ai/src/utils/http-inspector.ts new file mode 100644 index 000000000..aa3f131ee --- /dev/null +++ b/packages/ai/src/utils/http-inspector.ts @@ -0,0 +1,108 @@ +import * as os from "node:os"; +import * as path from "node:path"; + +export type RawHttpRequestDump = { + provider: string; + api: string; + model: string; + method?: string; + url?: string; + headers?: Record; + body?: unknown; +}; + +type ErrorWithStatus = { + status?: unknown; + statusCode?: unknown; + response?: { status?: unknown }; + cause?: unknown; +}; + +const SENSITIVE_HEADERS = ["authorization", "x-api-key", "api-key", "cookie", "set-cookie", "proxy-authorization"]; + +export async function appendRawHttpRequestDumpFor400( + message: string, + error: unknown, + dump: RawHttpRequestDump | undefined, +): Promise { + if (!dump || getStatusCode(error) !== 400) { + return message; + } + + const sanitizedDump = sanitizeDump(dump); + const fileName = `${Date.now()}-${Bun.hash(JSON.stringify(sanitizedDump)).toString(36)}.json`; + const filePath = path.join(os.homedir(), ".omp", "logs", "http-400-requests", fileName); + + try { + await Bun.write(filePath, `${JSON.stringify(sanitizedDump, null, 2)}\n`); + return `${message}\nraw-http-request=${filePath}`; + } catch (writeError) { + const writeMessage = writeError instanceof Error ? writeError.message : String(writeError); + return `${message}\nraw-http-request-save-failed=${writeMessage}`; + } +} + +export function withHttpStatus(error: unknown, status: number): Error { + const wrapped = error instanceof Error ? error : new Error(String(error)); + (wrapped as ErrorWithStatus).status = status; + return wrapped; +} + +function getStatusCode(error: unknown): number | undefined { + if (!error || typeof error !== "object") { + return undefined; + } + + const typedError = error as ErrorWithStatus; + const directStatus = toStatusCode(typedError.status) ?? toStatusCode(typedError.statusCode); + if (directStatus !== undefined) { + return directStatus; + } + + const responseStatus = toStatusCode(typedError.response?.status); + if (responseStatus !== undefined) { + return responseStatus; + } + + if (typedError.cause) { + return getStatusCode(typedError.cause); + } + + return undefined; +} + +function toStatusCode(value: unknown): number | undefined { + if (typeof value === "number" && Number.isFinite(value)) { + return value; + } + if (typeof value === "string") { + const parsed = Number(value); + if (Number.isFinite(parsed)) { + return parsed; + } + } + return undefined; +} + +function sanitizeDump(dump: RawHttpRequestDump): RawHttpRequestDump { + return { + ...dump, + headers: redactHeaders(dump.headers), + }; +} + +function redactHeaders(headers: Record | undefined): Record | undefined { + if (!headers) { + return undefined; + } + + const redacted: Record = {}; + for (const [key, value] of Object.entries(headers)) { + if (SENSITIVE_HEADERS.includes(key.toLowerCase())) { + redacted[key] = "[redacted]"; + continue; + } + redacted[key] = value; + } + return redacted; +} diff --git a/packages/ai/test/github-copilot-anthropic-auth.test.ts b/packages/ai/test/github-copilot-anthropic-auth.test.ts index e37caa06a..342cf5385 100644 --- a/packages/ai/test/github-copilot-anthropic-auth.test.ts +++ b/packages/ai/test/github-copilot-anthropic-auth.test.ts @@ -73,22 +73,6 @@ describe("Anthropic Copilot auth config", () => { expect(beta).toContain("interleaved-thinking-2025-05-14"); }); - it("does not include interleaved-thinking beta when disabled", () => { - const model = makeCopilotClaudeModel(); - const options = buildAnthropicClientOptions({ - model, - apiKey: "ghu_test", - extraBetas: [], - stream: true, - dynamicHeaders: {}, - }); - - const beta = options.defaultHeaders["anthropic-beta"]; - if (beta) { - expect(beta).not.toContain("interleaved-thinking-2025-05-14"); - } - }); - it("does not include fine-grained-tool-streaming beta for Copilot", () => { const model = makeCopilotClaudeModel(); const options = buildAnthropicClientOptions({ diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 78e4fb55c..f317e5075 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -551,7 +551,9 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} // Use provided or create AuthStorage and ModelRegistry const authStorage = options.authStorage ?? (await discoverAuthStorage(agentDir)); const modelRegistry = options.modelRegistry ?? new ModelRegistry(authStorage); - await modelRegistry.refresh(); + if (!options.modelRegistry) { + await modelRegistry.refresh(); + } time("discoverModels"); const settings = options.settings ?? (await Settings.init({ cwd, agentDir }));