SPB Git forge

spb/doc-api

Public
2commits 1branches 0releases
15.7 MBsize
maindefault branch
13 days agolast push
Python 88.3% TypeScript 7.6% Shell 4.1%
44.5 KB · 581 lines typescript
Raw Blame History
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