Skip to content
Open
1 change: 1 addition & 0 deletions packages/agent/.changes/res-1271-retry-owner.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- Removed the unused `maxRetryDelayMs` agent option; retry delays are owned by the session retry loop.
4 changes: 0 additions & 4 deletions packages/agent/src/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -112,7 +112,6 @@ export interface AgentOptions {
sessionId?: string;
thinkingBudgets?: ThinkingBudgets;
transport?: Transport;
maxRetryDelayMs?: number;
toolExecution?: ToolExecutionMode;
}

Expand Down Expand Up @@ -216,7 +215,6 @@ export class Agent {
public sessionId?: string;
public thinkingBudgets?: ThinkingBudgets;
public transport: Transport;
public maxRetryDelayMs?: number;
public toolExecution: ToolExecutionMode;

constructor(options: AgentOptions = {}) {
Expand All @@ -237,7 +235,6 @@ export class Agent {
this.sessionId = options.sessionId;
this.thinkingBudgets = options.thinkingBudgets;
this.transport = options.transport ?? "auto";
this.maxRetryDelayMs = options.maxRetryDelayMs;
this.toolExecution = options.toolExecution ?? "parallel";
}

Expand Down Expand Up @@ -470,7 +467,6 @@ export class Agent {
onResponse: this.onResponse,
transport: this.transport,
thinkingBudgets: this.thinkingBudgets,
maxRetryDelayMs: this.maxRetryDelayMs,
toolExecution: this.toolExecution,
beforeToolCall: this.beforeToolCall,
afterToolCall: this.afterToolCall,
Expand Down
2 changes: 0 additions & 2 deletions packages/agent/src/proxy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,6 @@ type ProxySerializableStreamOptions = Pick<
| "metadata"
| "transport"
| "thinkingBudgets"
| "maxRetryDelayMs"
>;

export interface ProxyStreamOptions extends ProxySerializableStreamOptions {
Expand Down Expand Up @@ -96,7 +95,6 @@ function buildProxyRequestOptions(options: ProxyStreamOptions): ProxySerializabl
metadata: options.metadata,
transport: options.transport,
thinkingBudgets: options.thinkingBudgets,
maxRetryDelayMs: options.maxRetryDelayMs,
};
}

Expand Down
4 changes: 4 additions & 0 deletions packages/ai/.changes/res-1271-retry-owner.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
- Removed SDK-internal provider retries (OpenAI, Azure, Anthropic, Bedrock, Codex): providers make a single attempt and report a structured failure so the agent's own retry loop owns every retry.
- Added the server-requested Retry-After delay to structured stream-failure diagnostics.
- Added structured stream-failure diagnostics to the OpenAI-completions and Codex providers, including friendly messages for Codex nested usage-limit error payloads.
- Removed the `maxRetries` and `maxRetryDelayMs` stream options.
1 change: 1 addition & 0 deletions packages/ai/src/providers/amazon-bedrock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,7 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream", BedrockOpt

const config: BedrockRuntimeClientConfig = {
profile: options.profile,
maxAttempts: 1,
};
const configuredRegion = getConfiguredBedrockRegion(options);
const hasConfiguredProfile = hasConfiguredBedrockProfile();
Expand Down
5 changes: 4 additions & 1 deletion packages/ai/src/providers/anthropic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -521,7 +521,6 @@ export const streamAnthropic: StreamFunction<"anthropic-messages", AnthropicOpti
const requestOptions = {
...(options?.signal ? { signal: options.signal } : {}),
...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}),
...(options?.maxRetries !== undefined ? { maxRetries: options.maxRetries } : {}),
};
const response = await client.messages.create({ ...params, stream: true }, requestOptions).asResponse();
await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model);
Expand Down Expand Up @@ -861,6 +860,7 @@ function createClient(

if (model.provider === "cloudflare-ai-gateway") {
const client = new Anthropic({
maxRetries: 0,
apiKey: null,
authToken: null,
baseURL: resolveCloudflareBaseUrl(model),
Expand All @@ -884,6 +884,7 @@ function createClient(

if (model.provider === "github-copilot") {
const client = new Anthropic({
maxRetries: 0,
apiKey: null,
authToken: apiKey,
baseURL: model.baseUrl,
Expand All @@ -905,6 +906,7 @@ function createClient(

if (isOAuthToken(apiKey)) {
const client = new Anthropic({
maxRetries: 0,
apiKey: null,
authToken: apiKey,
baseURL: model.baseUrl,
Expand All @@ -926,6 +928,7 @@ function createClient(
}

const client = new Anthropic({
maxRetries: 0,
apiKey,
baseURL: model.baseUrl,
dangerouslyAllowBrowser: true,
Expand Down
2 changes: 1 addition & 1 deletion packages/ai/src/providers/azure-openai-responses.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,6 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses"
const requestOptions = {
...(options?.signal ? { signal: options.signal } : {}),
...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}),
...(options?.maxRetries !== undefined ? { maxRetries: options.maxRetries } : {}),
};
const { data: openaiStream, response } = await client.responses.create(params, requestOptions).withResponse();
await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model);
Expand Down Expand Up @@ -231,6 +230,7 @@ function createClient(model: Model<"azure-openai-responses">, apiKey: string, op
dangerouslyAllowBrowser: true,
defaultHeaders: headers,
baseURL: baseUrl,
maxRetries: 0,
});
}

Expand Down
189 changes: 94 additions & 95 deletions packages/ai/src/providers/openai-codex-responses.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,13 +40,12 @@ import {
} from "../utils/diagnostics.js";
import { AssistantMessageEventStream } from "../utils/event-stream.js";
import { headersToRecord } from "../utils/headers.js";
import { parseRetryAfterMs, recordStreamFailure } from "../utils/stream-failure.js";
import { convertResponsesMessages, convertResponsesTools, processResponsesStream } from "./openai-responses-shared.js";
import { buildBaseOptions } from "./simple-options.js";

const DEFAULT_CODEX_BASE_URL = "https://chatgpt.com/backend-api";
const JWT_CLAIM_PATH = "https://api.openai.com/auth" as const;
const MAX_RETRIES = 3;
const BASE_DELAY_MS = 1000;
const CODEX_TOOL_CALL_PROVIDERS = new Set(["openai", "openai-codex", "opencode"]);
const WEBSOCKET_MESSAGE_TOO_BIG_CLOSE_CODE = 1009;

Expand Down Expand Up @@ -87,27 +86,6 @@ interface RequestBody {
[key: string]: unknown;
}

function isRetryableError(status: number, errorText: string): boolean {
if (status === 429 || status === 500 || status === 502 || status === 503 || status === 504) {
return true;
}
return /rate.?limit|overloaded|service.?unavailable|upstream.?connect|connection.?refused/i.test(errorText);
}

function sleep(ms: number, signal?: AbortSignal): Promise<void> {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(new Error("Request was aborted"));
return;
}
const timeout = setTimeout(resolve, ms);
signal?.addEventListener("abort", () => {
clearTimeout(timeout);
reject(new Error("Request was aborted"));
});
});
}

export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses", OpenAICodexResponsesOptions> = (
model: Model<"openai-codex-responses">,
context: Context,
Expand Down Expand Up @@ -211,61 +189,31 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses"
}
}

let response: Response | undefined;
let lastError: Error | undefined;
if (options?.signal?.aborted) {
throw new Error("Request was aborted");
}

for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) {
if (options?.signal?.aborted) {
let response: Response;
try {
response = await fetch(resolveCodexUrl(model.baseUrl), {
method: "POST",
headers: sseHeaders,
body: bodyJson,
signal: options?.signal,
});
} catch (error) {
if (
options?.signal?.aborted ||
(error instanceof Error && (error.name === "AbortError" || error.message === "Request was aborted"))
) {
throw new Error("Request was aborted");
}

try {
response = await fetch(resolveCodexUrl(model.baseUrl), {
method: "POST",
headers: sseHeaders,
body: bodyJson,
signal: options?.signal,
});
await options?.onResponse?.(
{ status: response.status, headers: headersToRecord(response.headers) },
model,
);

if (response.ok) {
break;
}

const errorText = await response.text();
if (attempt < MAX_RETRIES && isRetryableError(response.status, errorText)) {
const delayMs = BASE_DELAY_MS * 2 ** attempt;
await sleep(delayMs, options?.signal);
continue;
}

const fakeResponse = new Response(errorText, {
status: response.status,
statusText: response.statusText,
});
const info = await parseErrorResponse(fakeResponse);
throw new Error(info.friendlyMessage || info.message);
} catch (error) {
if (error instanceof Error) {
if (error.name === "AbortError" || error.message === "Request was aborted") {
throw new Error("Request was aborted");
}
}
lastError = error instanceof Error ? error : new Error(String(error));
if (attempt < MAX_RETRIES && !lastError.message.includes("usage limit")) {
const delayMs = BASE_DELAY_MS * 2 ** attempt;
await sleep(delayMs, options?.signal);
continue;
}
throw lastError;
}
throw error;
}
await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model);

if (!response?.ok) {
throw lastError ?? new Error("Failed after retries");
if (!response.ok) {
throw await parseErrorResponse(response);
}

if (!response.body) {
Expand All @@ -288,6 +236,7 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses"
}
output.stopReason = options?.signal?.aborted ? "aborted" : "error";
output.errorMessage = error instanceof Error ? error.message : String(error);
recordStreamFailure(model, output, error);
stream.push({ type: "error", reason: output.stopReason, error: output });
stream.end();
}
Expand Down Expand Up @@ -437,12 +386,25 @@ async function processStream(

class CodexApiError extends Error {
readonly code?: string;
readonly status?: number;
readonly retryAfterMs?: number;
readonly payload?: Record<string, unknown>;

constructor(message: string, options?: { code?: string; payload?: Record<string, unknown>; cause?: unknown }) {
constructor(
message: string,
options?: {
code?: string;
status?: number;
retryAfterMs?: number;
payload?: Record<string, unknown>;
cause?: unknown;
},
) {
super(message);
this.name = "CodexApiError";
this.code = options?.code;
this.status = options?.status;
this.retryAfterMs = options?.retryAfterMs;
this.payload = options?.payload;
this.cause = options?.cause;
}
Expand All @@ -469,12 +431,24 @@ async function* mapCodexEvents(events: AsyncIterable<Record<string, unknown>>):
if (!type) continue;

if (type === "error") {
const code = (event as { code?: string }).code || "";
const message = (event as { message?: string }).message || "";
throw new CodexApiError(`Codex error: ${message || code || JSON.stringify(event)}`, {
code: code || undefined,
payload: event,
});
// Errors arrive flat ({ code, message }) or nested under event.error.
const flatCode = typeof event.code === "string" ? event.code : "";
const flatMessage = typeof event.message === "string" ? event.message : "";
const nested = event.error && typeof event.error === "object" ? (event.error as CodexErrorPayload) : undefined;
const statusCode = (event as { status_code?: unknown }).status_code;
const status = typeof statusCode === "number" ? statusCode : undefined;
const code = flatCode || nested?.code || nested?.type || undefined;
const usageLimit = nested ? codexUsageLimitMessage(nested, status) : undefined;
const message = flatMessage || nested?.message || "";
throw new CodexApiError(
usageLimit?.friendlyMessage ?? `Codex error: ${message || code || JSON.stringify(event)}`,
{
code,
status,
retryAfterMs: usageLimit?.retryAfterMs,
payload: event,
},
);
}

if (type === "response.failed") {
Expand Down Expand Up @@ -1183,33 +1157,58 @@ async function processWebSocketStream(
}
}

async function parseErrorResponse(response: Response): Promise<{ message: string; friendlyMessage?: string }> {
interface CodexErrorPayload {
code?: string;
type?: string;
message?: string;
plan_type?: string;
resets_at?: number;
}

function codexUsageLimitMessage(
err: CodexErrorPayload,
status?: number,
): { friendlyMessage: string; retryAfterMs?: number } | undefined {
const code = err.code || err.type || "";
if (!/usage_limit_reached|usage_not_included|rate_limit_exceeded/i.test(code) && status !== 429) {
return undefined;
}
const plan = err.plan_type ? ` (${err.plan_type.toLowerCase()} plan)` : "";
const retryAfterMs = err.resets_at !== undefined ? Math.max(0, err.resets_at * 1000 - Date.now()) : undefined;
const when = retryAfterMs !== undefined ? ` Try again in ~${Math.round(retryAfterMs / 60000)} min.` : "";
return {
friendlyMessage: `You have hit your ChatGPT usage limit${plan}.${when}`.trim(),
retryAfterMs,
};
}

async function parseErrorResponse(response: Response): Promise<CodexApiError> {
const raw = await response.text();
let message = raw || response.statusText || "Request failed";
let friendlyMessage: string | undefined;
let code: string | undefined;
let retryAfterMs = parseRetryAfterMs(response.headers);

try {
const parsed = JSON.parse(raw) as {
error?: { code?: string; type?: string; message?: string; plan_type?: string; resets_at?: number };
};
const parsed = JSON.parse(raw) as { error?: CodexErrorPayload };
const err = parsed?.error;
if (err) {
const code = err.code || err.type || "";
if (/usage_limit_reached|usage_not_included|rate_limit_exceeded/i.test(code) || response.status === 429) {
const plan = err.plan_type ? ` (${err.plan_type.toLowerCase()} plan)` : "";
const mins = err.resets_at
? Math.max(0, Math.round((err.resets_at * 1000 - Date.now()) / 60000))
: undefined;
const when = mins !== undefined ? ` Try again in ~${mins} min.` : "";
friendlyMessage = `You have hit your ChatGPT usage limit${plan}.${when}`.trim();
code = err.code || err.type || undefined;
const usageLimit = codexUsageLimitMessage(err, response.status);
if (usageLimit) {
message = usageLimit.friendlyMessage;
// Neither server delay (Retry-After header, resets_at body) may undercut the other.
if (usageLimit.retryAfterMs !== undefined) {
retryAfterMs = Math.max(retryAfterMs ?? 0, usageLimit.retryAfterMs);
}
} else {
message = err.message || message;
}
message = err.message || friendlyMessage || message;
}
} catch {
// Unparseable error body: fall back to the raw message.
}

return { message, friendlyMessage };
return new CodexApiError(message, { code, status: response.status, retryAfterMs });
}

function extractAccountId(token: string): string {
Expand Down
Loading
Loading