Python 88.3%
TypeScript 7.6%
Shell 4.1%
1/**2 * Resilient fetch-based client for the OpenAI, Anthropic, xAI and Gemini APIs (no SDK required).3 *4 * STATUS: LIVE_VERIFIED 2026-09-18 (OpenAI/Anthropic) and 2026-09-19 (xAI grok-4.3, Gemini gemini-3.5-flash-lite) — 2 minimal5 * calls per provider through examples/shared/provider-abstraction/llmProvider.ts, which is built on this client; offline6 * behaviour (incl. recorded xAI/Gemini error bodies) exercised by `node --experimental-strip-types resilientClient.ts --selftest`.7 *8 * Run: node --env-file=.env --experimental-strip-types examples/shared/resilient-client/resilientClient.ts --selftest9 *10 * Mirrors resilient_client.py — see docs/architecture/resilience.md for the sourced rules:11 * - exponential backoff with full jitter, capped, bounded by attempts + wall-clock deadline12 * - per-request timeouts via AbortController (connect+read combined; fetch has no separate connect timeout)13 * - provider-aware retryable classification (408/409/429/500/502/503/504/529; OpenAI insufficient_quota & billing14 * codes not retryable; Anthropic 429 without retry-after = spend cap; xAI `{code, error}` bodies with 400 = incorrect15 * API key and bare-string 422s never retried, 429 backoff without Retry-After; Gemini google.rpc `error.status`16 * (RESOURCE_EXHAUSTED retryable using RetryInfo.retryDelay / "Please retry in Ns" from the message, `limit: 0` = tier gate,17 * FAILED_PRECONDITION never retried); `x-should-retry` wins; Retry-After honoured/capped)18 * - rate-limit header parsing (OpenAI + xAI x-ratelimit-*, Anthropic anthropic-ratelimit-*, Gemini: none exposed)19 * - circuit breaker, fallback chains, budget guard, OpenAI stream resume (starting_after), Anthropic stream restart20 *21 * Connection pooling: Node's global fetch (undici) keeps connections alive by default; tune with a custom22 * `undici.Agent({ connections, keepAliveTimeout })` and `setGlobalDispatcher` if you need to cap concurrency.23 * Never logs secrets: `safeHeaders()` redacts authorization / x-api-key / x-goog-api-key.24 */2526export type Provider = "openai" | "anthropic" | "xai" | "gemini";27export const PROVIDERS: readonly Provider[] = ["openai", "anthropic", "xai", "gemini"];28export const ENV_KEYS: Record<Provider, string> = { openai: "OPENAI_API_KEY", anthropic: "ANTHROPIC_API_KEY", xai: "XAI_API_KEY", gemini: "GEMINI_API_KEY" };29export const DEFAULT_BASE_URLS: Record<Provider, string> = {30 openai: "https://api.openai.com", anthropic: "https://api.anthropic.com", xai: "https://api.x.ai", gemini: "https://generativelanguage.googleapis.com",31};3233export interface Req {34 method: string;35 url: string;36 headers: Record<string, string>;37 body?: string;38 timeoutMs: number;39 stream: boolean;40}4142export interface Resp {43 status: number;44 headers: Record<string, string>;45 text: string; // empty when streaming46 body?: ReadableStream<Uint8Array> | null; // when stream=true47}4849export class TransportError extends Error {}5051export type Transport = (req: Req) => Promise<Resp>;5253const REDACTED = new Set(["authorization", "x-api-key", "x-goog-api-key", "openai-organization", "openai-project"]);54export function safeHeaders(h: Record<string, string>): Record<string, string> {55 return Object.fromEntries(Object.entries(h).map(([k, v]) => [k, REDACTED.has(k.toLowerCase()) ? "***REDACTED***" : v]));56}5758export const fetchTransport: Transport = async (req) => {59 const ctrl = new AbortController();60 // unref: a pending timeout must never keep the process alive once the stream has been consumed61 const timer = setTimeout(() => ctrl.abort(new Error(`timeout after ${req.timeoutMs}ms`)), req.timeoutMs);62 (timer as any).unref?.();63 try {64 const r = await fetch(req.url, { method: req.method, headers: req.headers, body: req.body, signal: ctrl.signal });65 const headers: Record<string, string> = {};66 r.headers.forEach((v, k) => (headers[k.toLowerCase()] = v));67 if (req.stream && r.ok && r.body) {68 // the timeout covers the WHOLE stream; clear it when the body ends (identity transform lets us observe EOF)69 const body = r.body.pipeThrough(new TransformStream<Uint8Array, Uint8Array>({ flush() { clearTimeout(timer); } }));70 return { status: r.status, headers, text: "", body };71 }72 return { status: r.status, headers, text: await r.text() };73 } catch (e: any) {74 throw new TransportError(e?.message ?? String(e));75 } finally {76 if (!req.stream) clearTimeout(timer);77 }78};7980// ---------------------------------------------------------------- classification81export const RETRYABLE_STATUSES = new Set([408, 409, 429, 500, 502, 503, 504, 529]);82export const OPENAI_NON_RETRYABLE_CODES = new Set([83 "insufficient_quota", "credit_balance_exhausted", "organization_spend_limit_exceeded",84 "project_spend_limit_exceeded", "organization_usage_limit_exceeded",85]);86const OPENAI_NON_RETRYABLE_TYPES = new Set(["insufficient_quota", "invalid_request_error", "authentication_error", "permission_error", "not_found_error"]);87const ANTHROPIC_NON_RETRYABLE_TYPES = new Set(["invalid_request_error", "authentication_error", "permission_error", "not_found_error", "billing_error", "request_too_large"]);88// xAI inference API codes (generated/fragments/errors/xai-errors.json, live 2026-09-19)89export const XAI_NON_RETRYABLE_CODES = new Set(["invalid-argument", "unauthenticated:no-credentials", "unauthenticated", "permission-denied", "not-found", "method-not-allowed", "unsupported-media-type", "unprocessable-entity"]);90// Gemini google.rpc.Status `error.status` (generated/fragments/errors/gemini-errors.json)91export const GEMINI_RETRYABLE_STATUSES = new Set(["RESOURCE_EXHAUSTED", "UNAVAILABLE", "DEADLINE_EXCEEDED", "INTERNAL", "ABORTED"]);92export const GEMINI_NON_RETRYABLE_STATUSES = new Set(["INVALID_ARGUMENT", "FAILED_PRECONDITION", "UNAUTHENTICATED", "PERMISSION_DENIED", "NOT_FOUND", "ALREADY_EXISTS", "OUT_OF_RANGE", "UNIMPLEMENTED", "CANCELLED"]);9394export interface Classification { retryable: boolean; reason: string; retryAfterS?: number; errorType?: string; errorCode?: string }9596export function parseRetryAfter(v?: string | null): number | undefined {97 if (!v) return undefined;98 const s = v.trim();99 if (/^\d+(\.\d+)?$/.test(s)) return parseFloat(s);100 const t = Date.parse(s);101 return Number.isNaN(t) ? undefined : Math.max(0, (t - Date.now()) / 1000);102}103104function parseJson(text: string): any { try { return JSON.parse(text); } catch { return undefined; } }105const unwrap = (body: any) => (Array.isArray(body) && body.length === 1 && body[0] && typeof body[0] === "object" ? body[0] : body);106/**107 * (type, code, message) from any of the four envelopes:108 * OpenAI/Anthropic {error:{type,code,message}} · Gemini google.rpc.Status {error:{code:429,message,status,details[]}} (type=status,109 * code=String(code); may be array-wrapped on the OpenAI-compat layer) · xAI {code:"invalid-argument", error:"msg"} (type=code=code) ·110 * xAI bare-string 422 bodies (message only) · xAI Management API gRPC {code:16,message,details} (type="grpc").111 */112export function extractError(body: any): { type?: string; code?: string; message?: string } {113 body = unwrap(body);114 if (typeof body === "string") return { message: body };115 if (!body || typeof body !== "object") return {};116 const e = body.error;117 if (e && typeof e === "object") {118 if ("status" in e && typeof e.code === "number") return { type: e.status, code: String(e.code), message: e.message };119 return { type: e.type, code: typeof e.code === "number" ? String(e.code) : e.code, message: e.message };120 }121 if (typeof e === "string") { const c = body.code !== undefined ? String(body.code) : undefined; return { type: c, code: c, message: e }; }122 if (typeof body.code === "number" && "message" in body) return { type: "grpc", code: String(body.code), message: body.message };123 return {};124}125126const DUR = { ms: 0.001, s: 1, m: 60, h: 3600 } as const;127/** Gemini never sends Retry-After: read `details[]` google.rpc.RetryInfo.retryDelay ("40s") or "Please retry in 54.22s" in the message. */128export function parseGeminiRetryDelay(body: any): number | undefined {129 const e = unwrap(body)?.error;130 if (!e || typeof e !== "object") return undefined;131 for (const d of e.details ?? []) {132 if (d && typeof d === "object" && String(d["@type"] ?? "").endsWith("google.rpc.RetryInfo")) {133 const rd = d.retryDelay;134 if (rd && typeof rd === "object") return Number(rd.seconds ?? 0) + Number(rd.nanos ?? 0) / 1e9;135 const m = /^(\d+(?:\.\d+)?)(ms|s|m|h)?$/.exec(String(rd ?? "").trim());136 if (m) return parseFloat(m[1]) * DUR[(m[2] ?? "s") as keyof typeof DUR];137 }138 }139 const m = /[Pp]lease retry in\s+(\d+(?:\.\d+)?)\s*(ms|s|m|h)?\b/.exec(String(e.message ?? ""));140 return m ? parseFloat(m[1]) * DUR[(m[2] ?? "s") as keyof typeof DUR] : undefined;141}142/** 429 whose every "Quota exceeded" line says `limit: 0` = paid-only model/feature on a free-tier key → retrying is pointless. */143export function geminiQuotaIsZero(body: any): boolean {144 const e = unwrap(body)?.error;145 if (!e || typeof e !== "object") return false;146 const lines = String(e.message ?? "").split("\n").filter((l) => l.includes("Quota exceeded"));147 return lines.length > 0 && lines.every((l) => /limit:\s*0(?:\D|$)/.test(l));148}149150export function classify(provider: Provider, resp: Resp, opts: { anthropic429WithoutRetryAfterRetryable?: boolean; geminiZeroQuotaRetryable?: boolean } = {}): Classification {151 const body = parseJson(resp.text);152 const { type, code, message } = extractError(body);153 const retryAfterS = parseRetryAfter(resp.headers["retry-after"]);154 const xsr = resp.headers["x-should-retry"];155 if (xsr !== undefined) return { retryable: xsr.trim().toLowerCase() === "true", reason: `x-should-retry=${xsr.trim()}`, retryAfterS, errorType: type, errorCode: code };156 const st = resp.status;157 if (st >= 200 && st < 400) return { retryable: false, reason: "success", errorType: type, errorCode: code };158 if (provider === "openai") {159 if (type && OPENAI_NON_RETRYABLE_TYPES.has(type) && st !== 409) return { retryable: false, reason: `openai error.type=${type}`, errorType: type, errorCode: code };160 if (code && OPENAI_NON_RETRYABLE_CODES.has(code)) return { retryable: false, reason: `openai error.code=${code} (billing/quota: user action required)`, errorType: type, errorCode: code };161 } else if (provider === "anthropic") {162 if (type && ANTHROPIC_NON_RETRYABLE_TYPES.has(type)) return { retryable: false, reason: `anthropic error.type=${type}`, errorType: type, errorCode: code };163 if (st === 429 && retryAfterS === undefined && !opts.anthropic429WithoutRetryAfterRetryable)164 return { retryable: false, reason: "anthropic 429 without retry-after (spend cap suspected)", errorType: type, errorCode: code };165 } else if (provider === "xai") {166 // 400 invalid-argument is ALSO xAI's answer to an incorrect API key (not 401) — never retry a 400167 if (st === 400) return { retryable: false, reason: `xai 400 code=${code}${message?.includes("API key") ? " (incorrect API key)" : ""}`, errorType: type, errorCode: code };168 if ((code && XAI_NON_RETRYABLE_CODES.has(code)) || [401, 403, 404, 405, 415, 422].includes(st)) return { retryable: false, reason: `xai http ${st} code=${code}`, errorType: type, errorCode: code };169 if (st === 429) {170 if (message && /credit|balance|billing/i.test(message)) return { retryable: false, reason: "xai 429 mentions credits/billing (user action required)", errorType: type, errorCode: code };171 return { retryable: true, reason: "xai http 429 (rate limit, exponential backoff)", retryAfterS, errorType: type, errorCode: code }; // no Retry-After documented/observed172 }173 } else if (provider === "gemini") {174 if (st === 429 || type === "RESOURCE_EXHAUSTED") {175 const delay = retryAfterS ?? parseGeminiRetryDelay(body);176 if (geminiQuotaIsZero(body) && !opts.geminiZeroQuotaRetryable) return { retryable: false, reason: "gemini RESOURCE_EXHAUSTED with limit: 0 (free tier / paid-only model: enable billing)", retryAfterS: delay, errorType: type, errorCode: code };177 return { retryable: true, reason: "gemini RESOURCE_EXHAUSTED (delay from RetryInfo/message; no Retry-After header)", retryAfterS: delay, errorType: type, errorCode: code };178 }179 if (type && GEMINI_NON_RETRYABLE_STATUSES.has(type)) return { retryable: false, reason: `gemini error.status=${type}`, errorType: type, errorCode: code };180 if (type && GEMINI_RETRYABLE_STATUSES.has(type)) return { retryable: true, reason: `gemini error.status=${type}`, retryAfterS, errorType: type, errorCode: code };181 if (st === 409) return { retryable: false, reason: "gemini 409 without status (assume ALREADY_EXISTS)", errorType: type, errorCode: code };182 }183 if (RETRYABLE_STATUSES.has(st)) return { retryable: true, reason: `http ${st}`, retryAfterS, errorType: type, errorCode: code };184 return { retryable: false, reason: `http ${st} not retryable`, errorType: type, errorCode: code };185}186187// ---------------------------------------------------------------- rate-limit headers188export function parseOpenAIDuration(v?: string): number | undefined {189 if (!v) return undefined;190 let total = 0, matched = false;191 for (const m of v.matchAll(/(\d+(?:\.\d+)?)(ms|s|m|h|d)/g)) {192 matched = true;193 const n = parseFloat(m[1]);194 total += { ms: n / 1000, s: n, m: n * 60, h: n * 3600, d: n * 86400 }[m[2] as "ms" | "s" | "m" | "h" | "d"];195 }196 return matched ? total : undefined;197}198export function parseRfc3339(v?: string): number | undefined {199 if (!v) return undefined;200 const t = Date.parse(v);201 return Number.isNaN(t) ? undefined : Math.max(0, (t - Date.now()) / 1000);202}203const int = (v?: string) => (v !== undefined && /^-?\d+$/.test(v) ? parseInt(v, 10) : undefined);204205export interface RateLimitInfo {206 provider: Provider; requestsLimit?: number; requestsRemaining?: number; requestsResetS?: number;207 tokensLimit?: number; tokensRemaining?: number; tokensResetS?: number;208 inputTokensLimit?: number; inputTokensRemaining?: number; outputTokensLimit?: number; outputTokensRemaining?: number;209 projectTokensRemaining?: number; retryAfterS?: number; requestId?: string; raw: Record<string, string>;210}211export function parseRateLimitHeaders(provider: Provider, headers: Record<string, string>): RateLimitInfo {212 const h: Record<string, string> = {};213 for (const [k, v] of Object.entries(headers)) h[k.toLowerCase()] = v;214 const info: RateLimitInfo = { provider, retryAfterS: parseRetryAfter(h["retry-after"]), raw: {} };215 if (provider === "openai" || provider === "xai") { // xAI (observed, undocumented): same x-ratelimit-limit/remaining-* names, no reset headers216 info.requestsLimit = int(h["x-ratelimit-limit-requests"]); info.requestsRemaining = int(h["x-ratelimit-remaining-requests"]);217 info.requestsResetS = parseOpenAIDuration(h["x-ratelimit-reset-requests"]);218 info.tokensLimit = int(h["x-ratelimit-limit-tokens"]); info.tokensRemaining = int(h["x-ratelimit-remaining-tokens"]);219 info.tokensResetS = parseOpenAIDuration(h["x-ratelimit-reset-tokens"]);220 info.projectTokensRemaining = int(h["x-ratelimit-remaining-project-tokens"]); info.requestId = h["x-request-id"];221 info.raw = Object.fromEntries(Object.entries(h).filter(([k]) => k.startsWith("x-ratelimit") || ["retry-after", "x-request-id", "x-zero-data-retention"].includes(k)));222 } else if (provider === "gemini") { // no rate-limit / Retry-After / request-id headers at all; keep the observability ones223 info.raw = Object.fromEntries(Object.entries(h).filter(([k]) => ["x-gemini-service-tier", "server-timing", "retry-after"].includes(k)));224 } else {225 const p = "anthropic-ratelimit-";226 info.requestsLimit = int(h[p + "requests-limit"]); info.requestsRemaining = int(h[p + "requests-remaining"]); info.requestsResetS = parseRfc3339(h[p + "requests-reset"]);227 info.tokensLimit = int(h[p + "tokens-limit"]); info.tokensRemaining = int(h[p + "tokens-remaining"]); info.tokensResetS = parseRfc3339(h[p + "tokens-reset"]);228 info.inputTokensLimit = int(h[p + "input-tokens-limit"]); info.inputTokensRemaining = int(h[p + "input-tokens-remaining"]);229 info.outputTokensLimit = int(h[p + "output-tokens-limit"]); info.outputTokensRemaining = int(h[p + "output-tokens-remaining"]);230 info.requestId = h["request-id"];231 info.raw = Object.fromEntries(Object.entries(h).filter(([k]) => k.startsWith(p) || ["retry-after", "request-id", "x-should-retry"].includes(k)));232 }233 return info;234}235236// ---------------------------------------------------------------- backoff237export interface RetryPolicy {238 maxAttempts: number; baseDelayS: number; maxDelayS: number; maxRetryAfterS: number; maxTotalS: number;239 jitter: (cap: number) => number; retryOnTransportError: boolean; anthropic429WithoutRetryAfterRetryable: boolean; geminiZeroQuotaRetryable?: boolean;240}241export const defaultPolicy = (): RetryPolicy => ({242 maxAttempts: 4, baseDelayS: 0.5, maxDelayS: 20, maxRetryAfterS: 60, maxTotalS: 120,243 jitter: (cap) => Math.random() * cap, retryOnTransportError: true, anthropic429WithoutRetryAfterRetryable: false, geminiZeroQuotaRetryable: false,244});245export function computeDelay(p: RetryPolicy, attempt: number, retryAfterS?: number): number {246 if (retryAfterS !== undefined) return Math.min(retryAfterS, p.maxRetryAfterS) + p.jitter(Math.min(1, p.baseDelayS));247 return p.jitter(Math.min(p.maxDelayS, p.baseDelayS * 2 ** (attempt - 1)));248}249250export class RetryExhausted extends Error {251 lastResponse: Resp | undefined; attempts: number; history: string[];252 constructor(msg: string, lastResponse: Resp | undefined, attempts: number, history: string[]) { super(msg); this.lastResponse = lastResponse; this.attempts = attempts; this.history = history; }253}254export class NonRetryableError extends Error {255 response: Resp; classification: Classification;256 constructor(response: Resp, classification: Classification) {257 const { type, code, message } = extractError(parseJson(response.text));258 super(`HTTP ${response.status} ${classification.reason}: ${type}/${code}: ${String(message ?? "").slice(0, 300)}`);259 this.response = response; this.classification = classification;260 }261}262export class CircuitOpen extends Error {}263export class BudgetExceeded extends Error {}264265// ---------------------------------------------------------------- circuit breaker266export class CircuitBreaker {267 state: "closed" | "open" | "half-open" = "closed"; failures = 0; openedAt?: number;268 failureThreshold: number; recoveryTimeoutS: number; private clock: () => number;269 constructor(failureThreshold = 5, recoveryTimeoutS = 30, clock: () => number = () => performance.now() / 1000) {270 this.failureThreshold = failureThreshold; this.recoveryTimeoutS = recoveryTimeoutS; this.clock = clock;271 }272 allow(): boolean {273 if (this.state === "closed") return true;274 if (this.state === "open") {275 if (this.clock() - (this.openedAt ?? 0) >= this.recoveryTimeoutS) { this.state = "half-open"; return true; }276 return false;277 }278 return true;279 }280 recordSuccess() { this.state = "closed"; this.failures = 0; this.openedAt = undefined; }281 recordFailure() { this.failures++; if (this.state === "half-open" || this.failures >= this.failureThreshold) { this.state = "open"; this.openedAt = this.clock(); } }282}283284// ---------------------------------------------------------------- budget guard285export type Price = { input: number; output: number; cached_input?: number; cache_write?: number }; // USD per 1M tokens286export class BudgetGuard {287 spentUsd = 0; calls = 0; maxUsd: number; prices: Record<string, Price>; defaultPrice: Price;288 constructor(maxUsd: number, prices: Record<string, Price> = {}, defaultPrice: Price = { input: 5, output: 15, cached_input: 0.5 }) {289 this.maxUsd = maxUsd; this.prices = prices; this.defaultPrice = defaultPrice;290 }291 check() { if (this.spentUsd >= this.maxUsd) throw new BudgetExceeded(`budget ${this.maxUsd} USD exhausted (spent ${this.spentUsd.toFixed(6)})`); }292 estimate(model: string, usage: any): number {293 const p = this.prices[model] ?? this.defaultPrice;294 // OpenAI/xAI Responses, OpenAI/xAI Chat, Anthropic Messages, Gemini usageMetadata (candidates + thoughts are both billed output)295 const inp = Number(usage?.input_tokens ?? usage?.prompt_tokens ?? usage?.promptTokenCount ?? 0);296 let out = Number(usage?.output_tokens ?? usage?.completion_tokens ?? (usage?.candidatesTokenCount !== undefined || usage?.thoughtsTokenCount !== undefined297 ? Number(usage?.candidatesTokenCount ?? 0) + Number(usage?.thoughtsTokenCount ?? 0) : 0));298 if (usage?.completion_tokens !== undefined && usage?.output_tokens === undefined) out += Number(usage?.completion_tokens_details?.reasoning_tokens ?? 0); // chat: reasoning billed separately299 const cached = Number(usage?.input_tokens_details?.cached_tokens ?? usage?.prompt_tokens_details?.cached_tokens ?? usage?.cache_read_input_tokens ?? usage?.cachedContentTokenCount ?? 0);300 const cacheWrite = Number(usage?.cache_creation_input_tokens ?? 0);301 const uncached = Math.max(0, inp - cached);302 return (uncached * p.input + cached * (p.cached_input ?? 0) + out * p.output + cacheWrite * (p.cache_write ?? p.input * 1.25)) / 1e6;303 }304 record(model: string, usage: any): number { const c = this.estimate(model, usage); this.spentUsd += c; this.calls++; return c; }305}306307// ---------------------------------------------------------------- idempotency (facts, see .py twin)308export function idempotencyHeaders(provider: Provider, path: string, key: string): Record<string, string> {309 // Only documented spot: OpenAI POST /v1/agents/sessions/{id}/events (Idempotency-Key). None for Responses/Messages.310 return provider === "openai" && /^\/v1\/agents\/sessions\/[^/]+\/events$/.test(path) ? { "Idempotency-Key": key } : {};311}312313// ---------------------------------------------------------------- client314export interface CallResult { response: Resp; attempts: number; rateLimit: RateLimitInfo; history: string[]; provider: Provider; model?: string; estCostUsd: number }315316export interface ClientOptions {317 apiKey?: string; baseUrl?: string; transport?: Transport; policy?: RetryPolicy; breaker?: CircuitBreaker; budget?: BudgetGuard;318 sleep?: (ms: number) => Promise<void>; clock?: () => number; anthropicVersion?: string; defaultHeaders?: Record<string, string>;319 onEvent?: (kind: string, data: Record<string, unknown>) => void;320}321322export class ResilientClient {323 readonly provider: Provider; readonly baseUrl: string; private apiKey: string; transport: Transport; policy: RetryPolicy;324 breaker: CircuitBreaker; budget?: BudgetGuard; private sleep: (ms: number) => Promise<void>; private clock: () => number;325 anthropicVersion: string; defaultHeaders: Record<string, string>; onEvent: (k: string, d: Record<string, unknown>) => void;326 lastRateLimit?: RateLimitInfo;327328 constructor(provider: Provider, o: ClientOptions = {}) {329 if (!PROVIDERS.includes(provider)) throw new Error(`provider must be one of ${PROVIDERS.join(", ")}`);330 this.provider = provider;331 const envKey = ENV_KEYS[provider];332 this.apiKey = o.apiKey ?? process.env[envKey] ?? "";333 this.baseUrl = (o.baseUrl ?? process.env[envKey.replace("API_KEY", "BASE_URL")] ?? DEFAULT_BASE_URLS[provider]).replace(/\/$/, "");334 this.transport = o.transport ?? fetchTransport; this.policy = o.policy ?? defaultPolicy(); this.breaker = o.breaker ?? new CircuitBreaker();335 this.budget = o.budget; this.sleep = o.sleep ?? ((ms) => new Promise((r) => setTimeout(r, ms))); this.clock = o.clock ?? (() => performance.now() / 1000);336 this.anthropicVersion = o.anthropicVersion ?? "2023-06-01"; this.defaultHeaders = o.defaultHeaders ?? {}; this.onEvent = o.onEvent ?? (() => {});337 }338 toString() { return `ResilientClient(${this.provider}, ${this.baseUrl}, key=***)`; }339340 private authHeaders(): Record<string, string> {341 if (this.provider === "openai" || this.provider === "xai") return { Authorization: `Bearer ${this.apiKey}` }; // xAI: same scheme, no version/beta headers342 if (this.provider === "gemini") return { "x-goog-api-key": this.apiKey }; // header, never `?key=` in the URL343 return { "x-api-key": this.apiKey, "anthropic-version": this.anthropicVersion };344 }345346 async request(method: string, path: string, jsonBody?: any, o: { headers?: Record<string, string>; stream?: boolean; timeoutMs?: number; idempotencyKey?: string; model?: string; rawBody?: string } = {}): Promise<CallResult> {347 this.budget?.check();348 const headers = { "Content-Type": "application/json", ...this.defaultHeaders, ...this.authHeaders(), ...(o.headers ?? {}) };349 if (o.idempotencyKey) Object.assign(headers, idempotencyHeaders(this.provider, path, o.idempotencyKey));350 const req: Req = { method, url: this.baseUrl + path, headers, body: jsonBody !== undefined ? JSON.stringify(jsonBody) : o.rawBody, timeoutMs: o.timeoutMs ?? 600_000, stream: !!o.stream };351 const model = o.model ?? jsonBody?.model;352 const history: string[] = []; const start = this.clock(); let last: Resp | undefined; let attempt = 0;353 for (;;) {354 attempt++;355 if (!this.breaker.allow()) throw new CircuitOpen(`circuit open for ${this.provider}`);356 let cls: Classification;357 try {358 const resp = await this.transport(req);359 last = resp;360 this.lastRateLimit = parseRateLimitHeaders(this.provider, resp.headers);361 cls = classify(this.provider, resp, { anthropic429WithoutRetryAfterRetryable: this.policy.anthropic429WithoutRetryAfterRetryable, geminiZeroQuotaRetryable: this.policy.geminiZeroQuotaRetryable });362 if (resp.status >= 200 && resp.status < 400) {363 this.breaker.recordSuccess();364 const result: CallResult = { response: resp, attempts: attempt, rateLimit: this.lastRateLimit, history, provider: this.provider, model, estCostUsd: 0 };365 if (this.budget && !o.stream) { const j = parseJson(resp.text); const usage = j?.usage ?? j?.usageMetadata; if (usage) result.estCostUsd = this.budget.record(model ?? "unknown", usage); }366 return result;367 }368 this.breaker.recordFailure();369 history.push(`attempt ${attempt}: HTTP ${resp.status} ${cls.reason}`);370 this.onEvent("http_error", { attempt, status: resp.status, reason: cls.reason, rateLimit: this.lastRateLimit.raw });371 if (!cls.retryable) throw new NonRetryableError(resp, cls);372 } catch (e) {373 if (!(e instanceof TransportError)) throw e;374 this.breaker.recordFailure();375 history.push(`attempt ${attempt}: transport error ${e.message}`);376 this.onEvent("transport_error", { attempt, error: e.message });377 if (!this.policy.retryOnTransportError) throw e;378 cls = { retryable: true, reason: "transport error" };379 }380 if (attempt >= this.policy.maxAttempts) throw new RetryExhausted(`gave up after ${attempt} attempts`, last, attempt, history);381 const delay = computeDelay(this.policy, attempt, cls.retryAfterS);382 if (this.clock() - start + delay > this.policy.maxTotalS) throw new RetryExhausted(`deadline ${this.policy.maxTotalS}s would be exceeded`, last, attempt, history);383 this.onEvent("backoff", { attempt, delayS: delay, retryAfterS: cls.retryAfterS });384 await this.sleep(delay * 1000);385 }386 }387 postJson(path: string, body: any, o?: Parameters<ResilientClient["request"]>[3]) { return this.request("POST", path, body, o); }388 get(path: string, o?: Parameters<ResilientClient["request"]>[3]) { return this.request("GET", path, undefined, o); }389}390391// ---------------------------------------------------------------- fallback392/** `path` may contain `{model}` (Gemini: "/v1beta/models/{model}:generateContent"). */393export interface Target { client: ResilientClient; model: string; path: string; buildBody: (model: string) => any }394export const AUTH_OR_VALIDATION_ERROR_TYPES = new Set([395 "authentication_error", "permission_error", "invalid_request_error", // OpenAI + Anthropic error.type396 "invalid-argument", "unauthenticated", "unauthenticated:no-credentials", "permission-denied", // xAI code397 "INVALID_ARGUMENT", "UNAUTHENTICATED", "PERMISSION_DENIED", "FAILED_PRECONDITION", // Gemini error.status398]);399export async function callWithFallback(targets: Target[], skipTypes = AUTH_OR_VALIDATION_ERROR_TYPES): Promise<CallResult> {400 const errors: string[] = [];401 for (const t of targets) {402 try {403 return await t.client.postJson(t.path.replace("{model}", t.model), t.buildBody(t.model), { model: t.model });404 } catch (e: any) {405 errors.push(`${t.client.provider}/${t.model}: ${e?.message ?? e}`);406 const modelSpecific = e instanceof NonRetryableError && (e.classification.errorCode === "model_not_found" || e.response.status === 404407 || (t.client.provider === "gemini" && e.response.status === 429)); // Gemini `limit: 0` = this model is paid-only for this key408 if (e instanceof NonRetryableError && e.classification.errorType && skipTypes.has(e.classification.errorType) && !modelSpecific) throw e;409 if (!(e instanceof NonRetryableError || e instanceof RetryExhausted || e instanceof CircuitOpen || e instanceof TransportError)) throw e;410 }411 }412 throw new RetryExhausted("all fallback targets failed: " + errors.join(" | "), undefined, targets.length, errors);413}414415// ---------------------------------------------------------------- SSE iteration + stream resume416export async function* iterSseJson(body: ReadableStream<Uint8Array> | null | undefined): AsyncGenerator<any> {417 if (!body) return;418 const reader = body.getReader(); const dec = new TextDecoder(); let buf = ""; let event: string | undefined; let data: string[] = [];419 const flush = () => { const payload = data.join("\n"); event = undefined; data = []; if (!payload || payload === "[DONE]") return undefined; try { const o = JSON.parse(payload); if (event && o && typeof o === "object" && !("event" in o)) o.event = event; return o; } catch { return undefined; } };420 for (;;) {421 const { value, done } = await reader.read();422 buf += done ? "" : dec.decode(value, { stream: true });423 let idx;424 while ((idx = buf.indexOf("\n")) >= 0) {425 const line = buf.slice(0, idx).replace(/\r$/, ""); buf = buf.slice(idx + 1);426 if (line === "") { const o = flush(); if (o !== undefined) yield o; }427 else if (line.startsWith("event:")) event = line.slice(6).trim();428 else if (line.startsWith("data:")) data.push(line.slice(5).replace(/^ /, ""));429 }430 if (done) { if (data.length) { const o = flush(); if (o !== undefined) yield o; } return; }431 }432}433434/** OpenAI background+stream with resume via GET /v1/responses/{id}?stream=true&starting_after=<sequence_number>. */435export async function* openaiStreamWithResume(client: ResilientClient, body: any, maxReconnects = 5): AsyncGenerator<any> {436 const terminal = new Set(["response.completed", "response.failed", "response.incomplete", "response.cancelled"]);437 let res = await client.postJson("/v1/responses", { ...body, background: true, stream: true }, { stream: true, model: body.model });438 let responseId: string | undefined; let cursor: number | undefined; let reconnects = 0;439 for (;;) {440 try {441 for await (const ev of iterSseJson(res.response.body)) {442 if (typeof ev.sequence_number === "number") cursor = ev.sequence_number;443 if (!responseId && ev.response?.id) responseId = ev.response.id;444 yield ev;445 if (terminal.has(ev.type)) return;446 }447 return;448 } catch (e: any) {449 if (!responseId || ++reconnects > maxReconnects) throw new RetryExhausted(`stream dropped (${e?.message}); cannot resume`, undefined, reconnects, []);450 res = await client.get(`/v1/responses/${responseId}?stream=true${cursor !== undefined ? `&starting_after=${cursor}` : ""}`, { stream: true });451 }452 }453}454455/** Anthropic: no server-side resume → restart the request and emit {type:"restart"} so consumers discard partial output. */456export async function* anthropicStreamWithRestart(client: ResilientClient, body: any, maxRestarts = 2, beta?: string): AsyncGenerator<any> {457 let attempt = 0;458 for (;;) {459 attempt++;460 const res = await client.postJson("/v1/messages", { ...body, stream: true }, { stream: true, model: body.model, headers: beta ? { "anthropic-beta": beta } : undefined });461 try {462 for await (const ev of iterSseJson(res.response.body)) {463 if (ev.type === "error") throw new TransportError(`in-stream error: ${JSON.stringify(ev.error)}`);464 yield ev;465 if (ev.type === "message_stop") return;466 }467 return;468 } catch (e) {469 if (attempt > maxRestarts) throw e;470 yield { type: "restart", attempt };471 }472 }473}474475// ---------------------------------------------------------------- offline self-test (mock transport, no network)476if (process.argv.includes("--selftest")) {477 const mk = (status: number, body?: any, headers: Record<string, string> = {}): Resp => ({ status, headers, text: body === undefined ? "" : JSON.stringify(body) });478 const scripted = (script: (Resp | Error)[]) => { const reqs: Req[] = []; const t: Transport = async (r) => { reqs.push(r); const it = script.shift()!; if (it instanceof Error) throw it; return it; }; return { t, reqs }; };479 const assert = (c: unknown, m: string) => { if (!c) { console.error("FAIL", m); process.exit(1); } };480 const noJitter = { ...defaultPolicy(), baseDelayS: 0.1, jitter: (cap: number) => cap };481482 // classification483 assert(!classify("openai", mk(429, { error: { type: "insufficient_quota", code: "insufficient_quota" } })).retryable, "insufficient_quota not retryable");484 assert(classify("openai", mk(429, { error: { type: "rate_limit_error", code: "slow_down" } }, { "retry-after": "7" })).retryAfterS === 7, "slow_down retry-after");485 assert(classify("anthropic", mk(529, { error: { type: "overloaded_error" } })).retryable, "529 retryable");486 assert(!classify("anthropic", mk(429, { error: { type: "rate_limit_error" } })).retryable, "spend cap 429 not retryable");487 assert(!classify("anthropic", mk(500, {}, { "x-should-retry": "false" })).retryable, "x-should-retry false");488 assert(parseOpenAIDuration("6m0s") === 360 && parseOpenAIDuration("250ms") === 0.25, "duration parse");489 assert(parseRateLimitHeaders("openai", { "x-ratelimit-remaining-tokens": "149984" }).tokensRemaining === 149984, "oa headers");490 assert(computeDelay({ ...noJitter }, 3, undefined) === 0.4 && computeDelay({ ...noJitter, jitter: () => 0 }, 1, 600) === 60, "delay");491492 // retry loop493 const s1 = scripted([mk(500, { error: {} }), mk(429, { error: { type: "rate_limit_error" } }, { "retry-after": "2" }), mk(200, { usage: { input_tokens: 10, output_tokens: 2 } })]);494 const sleeps: number[] = [];495 const c1 = new ResilientClient("openai", { apiKey: "k", transport: s1.t, policy: noJitter, sleep: async (ms) => { sleeps.push(ms); }, budget: new BudgetGuard(1, { m: { input: 1, output: 1 } }) });496 const r1 = await c1.postJson("/v1/responses", { model: "m", input: "x" });497 assert(r1.attempts === 3 && Math.round(sleeps[0]) === 100 && Math.round(sleeps[1]) === 2100, `retry loop ${sleeps}`);498 assert(Math.abs(r1.estCostUsd - 12 / 1e6) < 1e-12, "budget");499 assert(safeHeaders(s1.reqs[0].headers).Authorization === "***REDACTED***", "redaction");500501 // exhaustion + breaker + fallback502 const s2 = scripted(Array(4).fill(mk(503, { error: { type: "service_unavailable_error", code: "server_is_overloaded" } })));503 const c2 = new ResilientClient("openai", { apiKey: "k", transport: s2.t, policy: noJitter, sleep: async () => {} });504 let exhausted = false; try { await c2.get("/v1/models"); } catch (e) { exhausted = e instanceof RetryExhausted; } assert(exhausted, "exhausted");505 const b = new CircuitBreaker(2, 10, () => 0); b.recordFailure(); b.recordFailure(); assert(b.state === "open" && !b.allow(), "breaker opens");506 const s3 = scripted([mk(200, { id: "msg" })]);507 const c3 = new ResilientClient("anthropic", { apiKey: "k", transport: s3.t });508 const openBreaker = new CircuitBreaker(1, 1000); openBreaker.recordFailure(); // simulate a provider whose circuit is open509 const cOpen = new ResilientClient("openai", { apiKey: "k", transport: scripted([]).t, breaker: openBreaker });510 const s2b = scripted([mk(404, { error: { type: "invalid_request_error", code: "model_not_found" } })]);511 const c2b = new ResilientClient("openai", { apiKey: "k", transport: s2b.t, policy: noJitter, sleep: async () => {} });512 const fb = await callWithFallback([513 { client: cOpen, model: "gpt-5.4", path: "/v1/responses", buildBody: (m) => ({ model: m }) }, // CircuitOpen → next514 { client: c2b, model: "gpt-5.4-nano", path: "/v1/responses", buildBody: (m) => ({ model: m }) }, // model_not_found → next515 { client: c3, model: "claude-haiku-4-5-20251001", path: "/v1/messages", buildBody: (m) => ({ model: m, max_tokens: 8, messages: [] }) },516 ]).catch((e) => { assert(false, `fallback threw ${e}`); throw e; });517 assert(fb.provider === "anthropic", "fallback to anthropic");518 assert(s3.reqs[0].headers["anthropic-version"] === "2023-06-01" && s3.reqs[0].headers["x-api-key"] === "k", "anthropic headers");519520 // SSE iteration + resume521 const enc = new TextEncoder();522 // pull-based mock so already-delivered events are consumed before the simulated connection drop523 const streamOf = (events: any[], dropAfter?: number) => { let i = 0; return new ReadableStream<Uint8Array>({ pull(ctl) {524 if (dropAfter !== undefined && i >= dropAfter) { ctl.error(new Error("dropped")); return; }525 if (i >= events.length) { ctl.close(); return; }526 const e = events[i++]; ctl.enqueue(enc.encode(`event: ${e.type}\ndata: ${JSON.stringify(e)}\n\n`));527 } }); };528 const s4 = scripted([529 { status: 200, headers: {}, text: "", body: streamOf([{ type: "response.created", sequence_number: 0, response: { id: "resp_9" } }, { type: "response.output_text.delta", sequence_number: 1, delta: "OK" }, { type: "x" }], 2) },530 { status: 200, headers: {}, text: "", body: streamOf([{ type: "response.output_text.delta", sequence_number: 2, delta: "!" }, { type: "response.completed", sequence_number: 3, response: {} }]) },531 ]);532 const c4 = new ResilientClient("openai", { apiKey: "k", transport: s4.t });533 const got: string[] = []; for await (const ev of openaiStreamWithResume(c4, { model: "m", input: "x" })) got.push(ev.type);534 assert(got.join(",") === "response.created,response.output_text.delta,response.output_text.delta,response.completed", `resume events ${got}`);535 assert(s4.reqs[1].method === "GET" && s4.reqs[1].url.endsWith("/v1/responses/resp_9?stream=true&starting_after=1"), "resume URL");536 const s5 = scripted([537 { status: 200, headers: {}, text: "", body: streamOf([{ type: "message_start" }, { type: "error", error: { type: "overloaded_error" } }]) },538 { status: 200, headers: {}, text: "", body: streamOf([{ type: "message_start" }, { type: "message_stop" }]) },539 ]);540 const c5 = new ResilientClient("anthropic", { apiKey: "k", transport: s5.t });541 const got5: string[] = []; for await (const ev of anthropicStreamWithRestart(c5, { model: "m", max_tokens: 1, messages: [] })) got5.push(ev.type);542 assert(got5.join(",") === "message_start,restart,message_start,message_stop", `restart ${got5}`);543 assert(idempotencyHeaders("openai", "/v1/agents/sessions/s/events", "k")["Idempotency-Key"] === "k" && !idempotencyHeaders("openai", "/v1/responses", "k")["Idempotency-Key"], "idempotency");544545 // xAI + Gemini classification (bodies from generated/fragments/errors/{xai,gemini}-errors.json)546 const xaiBadKey = classify("xai", mk(400, { code: "invalid-argument", error: "Incorrect API key provided. You can obtain an API key from https://console.x.ai." }));547 assert(!xaiBadKey.retryable && xaiBadKey.reason.includes("incorrect API key") && xaiBadKey.errorCode === "invalid-argument", "xai 400 bad key");548 assert(!classify("xai", mk(422, "Failed to deserialize the JSON body into the target type: text: invalid type: integer `123`, expected a string at line 1 column 33")).retryable, "xai 422 string");549 assert(classify("xai", mk(429, { code: "resource-exhausted", error: "rate limit" })).retryable && classify("xai", mk(500, {})).retryable, "xai 429/500");550 assert(!classify("xai", mk(404, { code: "not-found", error: "The model grok-2-image does not exist or your team … does not have access to it." })).retryable, "xai 404");551 const gem429 = { error: { code: 429, message: "You exceeded your current quota…\n* Quota exceeded for metric: generativelanguage.googleapis.com/generate_content_free_tier_requests, limit: 15, model: gemini-3.5-flash-lite\nPlease retry in 50.868302469s.", status: "RESOURCE_EXHAUSTED", details: [{ "@type": "type.googleapis.com/google.rpc.QuotaFailure", violations: [{ quotaMetric: "generativelanguage.googleapis.com/generate_content_free_tier_requests", quotaId: "GenerateRequestsPerMinutePerProjectPerModel-FreeTier", quotaDimensions: { model: "gemini-3.5-flash-lite", location: "global" } }] }] } };552 const g1 = classify("gemini", mk(429, gem429)); assert(g1.retryable && Math.abs((g1.retryAfterS ?? 0) - 50.868302469) < 1e-9 && g1.errorType === "RESOURCE_EXHAUSTED" && g1.errorCode === "429", `gemini 429 ${JSON.stringify(g1)}`);553 const gemRetryInfo = { error: { code: 429, message: "x Please retry in 54.22s.", status: "RESOURCE_EXHAUSTED", details: [{ "@type": "type.googleapis.com/google.rpc.RetryInfo", retryDelay: "40s" }] } };554 assert(parseGeminiRetryDelay(gemRetryInfo) === 40, "RetryInfo wins over message");555 const gemZero = { error: { code: 429, message: "…\n* Quota exceeded for metric: generativelanguage.googleapis.com/generate_content_free_tier_requests, limit: 0, model: gemini-3.1-pro\nPlease retry in 50.868302469s.", status: "RESOURCE_EXHAUSTED" } };556 assert(!classify("gemini", mk(429, gemZero)).retryable && classify("gemini", mk(429, gemZero), { geminiZeroQuotaRetryable: true }).retryable, "gemini limit 0");557 assert(!classify("gemini", mk(400, { error: { code: 400, message: "Precondition check failed.", status: "FAILED_PRECONDITION" } })).retryable, "FAILED_PRECONDITION");558 assert(classify("gemini", mk(503, { error: { code: 503, message: "The model is overloaded. Please try again later.", status: "UNAVAILABLE" } })).retryable, "UNAVAILABLE");559 assert(classify("gemini", mk(504, { error: { code: 504, message: "Deadline expired", status: "DEADLINE_EXCEEDED" } })).retryable, "DEADLINE_EXCEEDED");560 const wrapped = classify("gemini", mk(400, [{ error: { code: 400, message: "Missing or invalid Authorization header.", status: "INVALID_ARGUMENT" } }]));561 assert(!wrapped.retryable && wrapped.errorType === "INVALID_ARGUMENT", "array-wrapped gemini error");562 assert(parseRateLimitHeaders("xai", { "x-ratelimit-limit-requests": "1800", "x-ratelimit-remaining-tokens": "10000000", "x-request-id": "0ef8" }).tokensRemaining === 10000000, "xai headers");563 assert(parseRateLimitHeaders("gemini", { "x-gemini-service-tier": "standard" }).tokensRemaining === undefined, "gemini no rl headers");564 const sx = scripted([mk(200, { id: "r" })]); const cx = new ResilientClient("xai", { apiKey: "xk", transport: sx.t }); await cx.postJson("/v1/responses", { model: "grok-4.3" });565 assert(sx.reqs[0].headers.Authorization === "Bearer xk" && sx.reqs[0].url.startsWith("https://api.x.ai/"), "xai auth/base");566 const sg = scripted([mk(200, { candidates: [], usageMetadata: { promptTokenCount: 10, candidatesTokenCount: 2, thoughtsTokenCount: 3 } })]);567 const cg = new ResilientClient("gemini", { apiKey: "gk", transport: sg.t, budget: new BudgetGuard(1, { g: { input: 1, output: 1 } }) });568 const rg = await cg.postJson("/v1beta/models/g:generateContent", { contents: [] }, { model: "g" });569 assert(sg.reqs[0].headers["x-goog-api-key"] === "gk" && !sg.reqs[0].url.includes("key=") && Math.abs(rg.estCostUsd - 15 / 1e6) < 1e-12, "gemini auth/budget (usageMetadata)");570 assert(safeHeaders(sg.reqs[0].headers)["x-goog-api-key"] === "***REDACTED***", "gemini redaction");571 // multi-provider fallback: gemini free-tier limit 0 → xai → anthropic572 const sgz = scripted([mk(429, gemZero)]); const cgz = new ResilientClient("gemini", { apiKey: "gk", transport: sgz.t, policy: noJitter, sleep: async () => {} });573 const sxz = scripted([mk(400, { code: "invalid-argument", error: "Incorrect API key provided." })]); const cxz = new ResilientClient("xai", { apiKey: "bad", transport: sxz.t });574 let authStops = false; try { await callWithFallback([{ client: cxz, model: "grok-4.3", path: "/v1/responses", buildBody: (m) => ({ model: m }) }, { client: c3, model: "m", path: "/v1/messages", buildBody: (m) => ({ model: m }) }]); } catch (e) { authStops = e instanceof NonRetryableError; }575 assert(authStops, "xai bad key stops fallback");576 const s3b = scripted([mk(200, { id: "msg2" })]); const c3b = new ResilientClient("anthropic", { apiKey: "k", transport: s3b.t });577 const fb2 = await callWithFallback([{ client: cgz, model: "gemini-3.1-pro-preview", path: "/v1beta/models/{model}:generateContent", buildBody: () => ({ contents: [] }) }, { client: c3b, model: "claude-haiku-4-5-20251001", path: "/v1/messages", buildBody: (m) => ({ model: m, max_tokens: 8, messages: [] }) }]);578 assert(fb2.provider === "anthropic" && sgz.reqs[0].url.endsWith("/v1beta/models/gemini-3.1-pro-preview:generateContent"), "gemini limit 0 → fallback, {model} substituted");579 console.log("resilientClient.ts selftest: all assertions passed");580}581