fix(ai): improved error handling and provider robustness
- Added raw HTTP request dumps to error messages for 400 status codes. - Implemented client-side retry for Anthropic streaming on transient errors. - Preserved context by converting aborted assistant tool calls to text. - Optimized model registry refresh in coding agent sessions.
This commit is contained in:
@@ -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 });
|
||||
|
||||
@@ -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<string, string> | 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<string, any>) ?? {},
|
||||
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<string, unknown>) ?? {},
|
||||
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 });
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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}`;
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -95,6 +95,7 @@ export function transformMessages<TApi extends Api>(
|
||||
const result: Message[] = [];
|
||||
let pendingToolCalls: ToolCall[] = [];
|
||||
let existingToolResultIds = new Set<string>();
|
||||
const skippedToolCallIds = new Set<string>();
|
||||
|
||||
for (let i = 0; i < transformed.length; i++) {
|
||||
const msg = transformed[i];
|
||||
@@ -118,13 +119,20 @@ export function transformMessages<TApi extends Api>(
|
||||
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<TApi extends Api>(
|
||||
|
||||
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") {
|
||||
|
||||
@@ -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<string, string>;
|
||||
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<string> {
|
||||
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<string, string> | undefined): Record<string, string> | undefined {
|
||||
if (!headers) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
const redacted: Record<string, string> = {};
|
||||
for (const [key, value] of Object.entries(headers)) {
|
||||
if (SENSITIVE_HEADERS.includes(key.toLowerCase())) {
|
||||
redacted[key] = "[redacted]";
|
||||
continue;
|
||||
}
|
||||
redacted[key] = value;
|
||||
}
|
||||
return redacted;
|
||||
}
|
||||
@@ -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({
|
||||
|
||||
@@ -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 }));
|
||||
|
||||
Reference in New Issue
Block a user