#!/usr/bin/env python3 """Resilient HTTP client for the OpenAI, Anthropic, xAI and Gemini APIs — stdlib only. STATUS: LIVE_VERIFIED 2026-09-18 (OpenAI/Anthropic) and 2026-09-19 (xAI grok-4.3, Gemini gemini-3.5-flash-lite) — 2 minimal calls per provider via examples/shared/provider-abstraction/llm_provider.py, which is built on this client; offline behaviour (incl. recorded xAI/Gemini error bodies) covered by tests/shared/test_resilient_client.py. Features (see docs/architecture/resilience.md for the sourced rationale): * exponential backoff with FULL jitter, capped, bounded by attempts AND total elapsed time * per-request timeouts (connect/read) and a wall-clock deadline * provider-aware retryable-error classification: - HTTP 408/409/429/500/502/503/504/529 retryable by default - OpenAI: `error.type == "insufficient_quota"` and billing/spend-limit `error.code`s are NOT retryable, `slow_down` (429) and `server_is_overloaded` (503) are retryable (follow Retry-After) - Anthropic: 529 `overloaded_error` retryable; 429 WITHOUT `retry-after` = spend cap (not retryable by default) - xAI: bodies are `{"code": "", "error": ""}` (400 `invalid-argument` is ALSO the "incorrect API key" answer — never retry it), 422 bodies are a bare JSON string (serde error, not retryable), 429 has no Retry-After (backoff only), 5xx retryable; observed `x-ratelimit-*` headers parsed like OpenAI's (no reset headers) - Gemini: google.rpc.Status `{"error": {"code", "message", "status", "details": [...]}}` (sometimes array-wrapped on the OpenAI-compat layer); RESOURCE_EXHAUSTED 429 is retryable with the delay taken from `details[]` `google.rpc.RetryInfo.retryDelay` ("40s") or from the message text ("Please retry in 54.22s") — there is NO Retry-After header; a 429 whose QuotaFailure says `limit: 0` is a billing/tier gate (not retryable by default); UNAVAILABLE 503 / DEADLINE_EXCEEDED 504 / INTERNAL 500 / ABORTED 409 retryable; FAILED_PRECONDITION 400 (billing not enabled, region), INVALID_ARGUMENT, PERMISSION_DENIED, NOT_FOUND, UNAUTHENTICATED, ALREADY_EXISTS never retried - `x-should-retry: false|true` (Anthropic) always wins; `Retry-After` (seconds or HTTP-date) is honoured and capped * rate-limit header parsing (OpenAI + xAI x-ratelimit-*, Anthropic anthropic-ratelimit-*, Gemini: none exposed) * idempotency helper (facts: only OpenAI `POST /v1/agents/sessions/{id}/events` documents `Idempotency-Key`; neither Responses nor Messages do — so POST retries may double-bill; see docs) * circuit breaker (closed → open → half-open) keyed per provider/model * fallback chains: same-provider model fallback and multi-provider fallback * budget guard (USD, from usage + price table you supply) * stream resume helper for OpenAI background responses (`GET /v1/responses/{id}?stream=true&starting_after=N`) and a restart strategy for Anthropic (no server-side resume exists) No SDK required. Optional adapters at the bottom show how to plug the official SDKs' `max_retries=0` clients in. Never logs secrets: headers named authorization / x-api-key are redacted in every repr. """ from __future__ import annotations import email.utils import json import os import random import re import threading import time import urllib.error import urllib.request from dataclasses import dataclass, field from datetime import datetime, timezone from typing import Any, Callable, Iterator, Optional # -------------------------------------------------------------------------------------- # Transport layer (pluggable so tests run offline) # -------------------------------------------------------------------------------------- REDACTED_HEADERS = {"authorization", "x-api-key", "x-goog-api-key", "openai-organization", "openai-project"} PROVIDERS = ("openai", "anthropic", "xai", "gemini") ENV_KEYS = {"openai": "OPENAI_API_KEY", "anthropic": "ANTHROPIC_API_KEY", "xai": "XAI_API_KEY", "gemini": "GEMINI_API_KEY"} DEFAULT_BASE_URLS = {"openai": "https://api.openai.com", "anthropic": "https://api.anthropic.com", "xai": "https://api.x.ai", "gemini": "https://generativelanguage.googleapis.com"} @dataclass class Request: method: str url: str headers: dict[str, str] = field(default_factory=dict) body: Optional[bytes] = None connect_timeout: float = 5.0 read_timeout: float = 600.0 # both official SDKs default to 600 s read timeout stream: bool = False def safe_headers(self) -> dict[str, str]: return {k: ("***REDACTED***" if k.lower() in REDACTED_HEADERS else v) for k, v in self.headers.items()} @dataclass class Response: status: int headers: dict[str, str] body: bytes = b"" lines: Optional[Iterator[bytes]] = None # when stream=True elapsed_s: float = 0.0 def header(self, name: str) -> Optional[str]: for k, v in self.headers.items(): if k.lower() == name.lower(): return v return None def json(self) -> Any: try: return json.loads(self.body) except Exception: # noqa: BLE001 return None class TransportError(Exception): """Network-level failure (DNS, connection reset, timeout). Always retryable unless the deadline is hit.""" Transport = Callable[[Request], Response] def urllib_transport(req: Request) -> Response: """Default transport. urllib has no pooling: for production prefer http.client keep-alive or httpx (see docs).""" t0 = time.monotonic() r = urllib.request.Request(req.url, method=req.method, headers=req.headers, data=req.body) try: # urllib exposes a single timeout: use the read timeout (connect is usually much shorter in practice) resp = urllib.request.urlopen(r, timeout=req.read_timeout) except urllib.error.HTTPError as e: return Response(e.code, dict(e.headers.items()), e.read(), elapsed_s=time.monotonic() - t0) except (urllib.error.URLError, TimeoutError, ConnectionError, OSError) as e: raise TransportError(str(e)) from e hdrs = dict(resp.headers.items()) if req.stream: def gen() -> Iterator[bytes]: with resp: for raw in resp: yield raw return Response(resp.status, hdrs, b"", lines=gen(), elapsed_s=time.monotonic() - t0) with resp: return Response(resp.status, hdrs, resp.read(), elapsed_s=time.monotonic() - t0) # -------------------------------------------------------------------------------------- # Error classification # -------------------------------------------------------------------------------------- RETRYABLE_STATUSES = {408, 409, 429, 500, 502, 503, 504, 529} # OpenAI: 429s that require user action (docs/guides/error-codes + rate-limits "Don't retry quota, billing…") OPENAI_NON_RETRYABLE_CODES = { "insufficient_quota", "credit_balance_exhausted", "organization_spend_limit_exceeded", "project_spend_limit_exceeded", "organization_usage_limit_exceeded", } OPENAI_NON_RETRYABLE_TYPES = {"insufficient_quota", "invalid_request_error", "authentication_error", "permission_error", "not_found_error"} ANTHROPIC_NON_RETRYABLE_TYPES = {"invalid_request_error", "authentication_error", "permission_error", "not_found_error", "billing_error", "request_too_large"} # xAI inference API: {"code": "", "error": ""} — generated/fragments/errors/xai-errors.json (live 2026-09-19) XAI_NON_RETRYABLE_CODES = {"invalid-argument", "unauthenticated:no-credentials", "unauthenticated", "permission-denied", "not-found", "method-not-allowed", "unsupported-media-type", "unprocessable-entity"} # Gemini google.rpc.Status `error.status` — generated/fragments/errors/gemini-errors.json GEMINI_RETRYABLE_STATUSES = {"RESOURCE_EXHAUSTED", "UNAVAILABLE", "DEADLINE_EXCEEDED", "INTERNAL", "ABORTED"} GEMINI_NON_RETRYABLE_STATUSES = {"INVALID_ARGUMENT", "FAILED_PRECONDITION", "UNAUTHENTICATED", "PERMISSION_DENIED", "NOT_FOUND", "ALREADY_EXISTS", "OUT_OF_RANGE", "UNIMPLEMENTED", "CANCELLED"} @dataclass class Classification: retryable: bool reason: str retry_after_s: Optional[float] = None error_type: Optional[str] = None error_code: Optional[str] = None def parse_retry_after(value: Optional[str]) -> Optional[float]: """`Retry-After` may be delta-seconds or an HTTP-date (RFC 7231).""" if not value: return None value = value.strip() if re.fullmatch(r"\d+(\.\d+)?", value): return float(value) try: dt = email.utils.parsedate_to_datetime(value) if dt.tzinfo is None: dt = dt.replace(tzinfo=timezone.utc) return max(0.0, (dt - datetime.now(timezone.utc)).total_seconds()) except Exception: # noqa: BLE001 return None def extract_error(body: Any) -> tuple[Optional[str], Optional[str], Optional[str]]: """Return (type, code, message) from any of the four providers' error envelopes: - OpenAI / Anthropic: {"error": {"type", "code", "message"}} - Gemini (google.rpc.Status): {"error": {"code": 429, "message", "status": "RESOURCE_EXHAUSTED", "details": [...]}} → type = status, code = str(code). The OpenAI-compat layer may wrap it in a one-element array. - xAI inference: {"code": "invalid-argument", "error": ""} → type = code = the kebab-case code - xAI 422 / some 400s: a bare JSON string → (None, None, ) - xAI Management API: {"code": 16, "message": "...", "details": []} (gRPC numbering) → type = "grpc", code = "16" """ if isinstance(body, list) and len(body) == 1 and isinstance(body[0], dict): body = body[0] if isinstance(body, str): return None, None, body if not isinstance(body, dict): return None, None, None err = body.get("error") if isinstance(err, dict): if "status" in err and isinstance(err.get("code"), int): # google.rpc.Status return err.get("status"), str(err.get("code")), err.get("message") code = err.get("code") return err.get("type"), (str(code) if isinstance(code, int) else code), err.get("message") if isinstance(err, str): # xAI inference API code = body.get("code") return (str(code) if code is not None else None), (str(code) if code is not None else None), err if isinstance(body.get("code"), int) and "message" in body: # gRPC-style (xAI Management API) return "grpc", str(body["code"]), body.get("message") return None, None, None _GEMINI_RETRY_IN_RE = re.compile(r"[Pp]lease retry in\s+(\d+(?:\.\d+)?)\s*(ms|s|m|h)?\b") _GO_DURATION_RE = re.compile(r"^(\d+(?:\.\d+)?)(ms|s|m|h)?$") def parse_gemini_retry_delay(body: Any) -> Optional[float]: """Gemini never sends Retry-After. The delay lives (a) in `error.details[]` as `{"@type": "type.googleapis.com/google.rpc.RetryInfo", "retryDelay": "40s"}` (protobuf Duration string), and/or (b) at the end of `error.message`: "... Please retry in 54.22098241s." Returns seconds, preferring RetryInfo.""" if isinstance(body, list) and body and isinstance(body[0], dict): body = body[0] err = body.get("error") if isinstance(body, dict) else None if not isinstance(err, dict): return None for d in err.get("details") or []: if isinstance(d, dict) and str(d.get("@type", "")).endswith("google.rpc.RetryInfo"): rd = d.get("retryDelay") if isinstance(rd, dict): # JSON-encoded Duration object form {seconds, nanos} return float(rd.get("seconds", 0)) + float(rd.get("nanos", 0)) / 1e9 m = _GO_DURATION_RE.match(str(rd or "").strip()) if m: n, unit = float(m.group(1)), m.group(2) or "s" return n * {"ms": 0.001, "s": 1, "m": 60, "h": 3600}[unit] m = _GEMINI_RETRY_IN_RE.search(str(err.get("message") or "")) if m: n, unit = float(m.group(1)), m.group(2) or "s" return n * {"ms": 0.001, "s": 1, "m": 60, "h": 3600}[unit] return None def gemini_quota_is_zero(body: Any) -> bool: """True when the 429 says the quota limit is 0 (free tier hitting a paid-only model/feature) — retrying is pointless; the fix is billing/tier (docs/gemini/rate-limits.md §7: `limit: 0` for Pro models on a free-tier key).""" if isinstance(body, list) and body and isinstance(body[0], dict): body = body[0] err = body.get("error") if isinstance(body, dict) else None if not isinstance(err, dict): return False lines = [ln for ln in str(err.get("message") or "").splitlines() if "Quota exceeded" in ln] return bool(lines) and all(re.search(r"limit:\s*0(?:\D|$)", ln) for ln in lines) def classify(provider: str, resp: Response, *, anthropic_429_without_retry_after_retryable: bool = False, gemini_zero_quota_retryable: bool = False) -> Classification: """Decide whether an HTTP response should be retried. Provider is one of PROVIDERS.""" st = resp.status body = resp.json() etype, ecode, emsg = extract_error(body) retry_after = parse_retry_after(resp.header("retry-after")) # Explicit server hint (Anthropic documents x-should-retry; SDKs honour it) xsr = resp.header("x-should-retry") if xsr is not None: flag = xsr.strip().lower() == "true" return Classification(flag, f"x-should-retry={xsr.strip()}", retry_after, etype, ecode) if 200 <= st < 400: return Classification(False, "success", None, etype, ecode) if provider == "openai": if etype in OPENAI_NON_RETRYABLE_TYPES and st != 409: return Classification(False, f"openai error.type={etype}", None, etype, ecode) if ecode in OPENAI_NON_RETRYABLE_CODES: return Classification(False, f"openai error.code={ecode} (billing/quota: user action required)", None, etype, ecode) elif provider == "anthropic": if etype in ANTHROPIC_NON_RETRYABLE_TYPES: return Classification(False, f"anthropic error.type={etype}", None, etype, ecode) if st == 429 and retry_after is None and not anthropic_429_without_retry_after_retryable: # docs: "A tier spend-cap 429 has no retry-after header and keeps failing until access resumes" return Classification(False, "anthropic 429 without retry-after (spend cap suspected)", None, etype, ecode) elif provider == "xai": if st == 400: # xAI answers a malformed/unknown API key with 400 invalid-argument "Incorrect API key provided" — not 401 hint = " (incorrect API key)" if emsg and "API key" in emsg else "" return Classification(False, f"xai 400 code={ecode}{hint}", None, etype, ecode) if ecode in XAI_NON_RETRYABLE_CODES or st in (401, 403, 404, 405, 415, 422): return Classification(False, f"xai http {st} code={ecode}", None, etype, ecode) if st == 429: # RESOURCE_EXHAUSTED (gRPC 8): team RPS/TPM or per-key qps/qpm/tpm; no Retry-After documented or observed → backoff. # Depleted prepaid credits with a $0 invoiced limit are "automatically rejected" (status undocumented). if emsg and re.search(r"credit|balance|billing", emsg, re.I): return Classification(False, "xai 429 mentions credits/billing (user action required)", None, etype, ecode) return Classification(True, "xai http 429 (rate limit, exponential backoff)", retry_after, etype, ecode) elif provider == "gemini": if st == 429 or etype == "RESOURCE_EXHAUSTED": delay = retry_after if retry_after is not None else parse_gemini_retry_delay(body) if gemini_quota_is_zero(body) and not gemini_zero_quota_retryable: return Classification(False, "gemini RESOURCE_EXHAUSTED with limit: 0 (free tier / paid-only model: enable billing)", delay, etype, ecode) return Classification(True, "gemini RESOURCE_EXHAUSTED (delay from RetryInfo/message; no Retry-After header)", delay, etype, ecode) if etype in GEMINI_NON_RETRYABLE_STATUSES: return Classification(False, f"gemini error.status={etype}", None, etype, ecode) if etype in GEMINI_RETRYABLE_STATUSES: return Classification(True, f"gemini error.status={etype}", retry_after, etype, ecode) if st == 409: # no status in body: ALREADY_EXISTS (not retryable) vs ABORTED (retryable) is undecidable → be safe return Classification(False, "gemini 409 without status (assume ALREADY_EXISTS)", None, etype, ecode) if st in RETRYABLE_STATUSES: return Classification(True, f"http {st}", retry_after, etype, ecode) return Classification(False, f"http {st} not retryable", None, etype, ecode) # -------------------------------------------------------------------------------------- # Rate-limit headers # -------------------------------------------------------------------------------------- _DUR_RE = re.compile(r"(\d+(?:\.\d+)?)(ms|s|m|h|d)") def parse_openai_duration(value: str) -> Optional[float]: """OpenAI reset headers look like '1s', '6m0s', '250ms', '1h2m3.5s'. Returns seconds.""" if not value: return None total, matched = 0.0, False for num, unit in _DUR_RE.findall(value): matched = True n = float(num) total += {"ms": n / 1000, "s": n, "m": n * 60, "h": n * 3600, "d": n * 86400}[unit] return total if matched else None def parse_rfc3339(value: str) -> Optional[float]: """Anthropic reset headers are RFC 3339 timestamps. Returns seconds from now (>= 0).""" try: dt = datetime.fromisoformat(value.replace("Z", "+00:00")) return max(0.0, (dt - datetime.now(timezone.utc)).total_seconds()) except Exception: # noqa: BLE001 return None @dataclass class RateLimitInfo: provider: str requests_limit: Optional[int] = None requests_remaining: Optional[int] = None requests_reset_s: Optional[float] = None tokens_limit: Optional[int] = None tokens_remaining: Optional[int] = None tokens_reset_s: Optional[float] = None input_tokens_limit: Optional[int] = None # Anthropic only input_tokens_remaining: Optional[int] = None output_tokens_limit: Optional[int] = None # Anthropic only output_tokens_remaining: Optional[int] = None project_tokens_remaining: Optional[int] = None # OpenAI only retry_after_s: Optional[float] = None request_id: Optional[str] = None raw: dict[str, str] = field(default_factory=dict) def _int(v: Optional[str]) -> Optional[int]: try: return int(v) if v is not None else None except ValueError: return None def parse_rate_limit_headers(provider: str, headers: dict[str, str]) -> RateLimitInfo: h = {k.lower(): v for k, v in headers.items()} info = RateLimitInfo(provider=provider, retry_after_s=parse_retry_after(h.get("retry-after"))) if provider in ("openai", "xai"): # xAI (observed 2026-09-19, undocumented): x-ratelimit-limit/remaining-requests (per-minute budget) and # x-ratelimit-limit/remaining-tokens (TPM) on inference responses; NO x-ratelimit-reset-* headers → reset stays None. info.requests_limit = _int(h.get("x-ratelimit-limit-requests")) info.requests_remaining = _int(h.get("x-ratelimit-remaining-requests")) info.requests_reset_s = parse_openai_duration(h.get("x-ratelimit-reset-requests", "")) info.tokens_limit = _int(h.get("x-ratelimit-limit-tokens")) info.tokens_remaining = _int(h.get("x-ratelimit-remaining-tokens")) info.tokens_reset_s = parse_openai_duration(h.get("x-ratelimit-reset-tokens", "")) info.project_tokens_remaining = _int(h.get("x-ratelimit-remaining-project-tokens")) info.request_id = h.get("x-request-id") info.raw = {k: v for k, v in h.items() if k.startswith("x-ratelimit") or k in ("retry-after", "x-request-id", "x-zero-data-retention")} elif provider == "gemini": # Gemini exposes NO rate-limit, Retry-After or request-id headers (docs/gemini/rate-limits.md §8). Only the # undocumented X-Gemini-Service-Tier and Server-Timing are kept for observability; correlate with body `responseId`. info.raw = {k: v for k, v in h.items() if k in ("x-gemini-service-tier", "server-timing", "retry-after")} else: p = "anthropic-ratelimit-" info.requests_limit = _int(h.get(p + "requests-limit")) info.requests_remaining = _int(h.get(p + "requests-remaining")) info.requests_reset_s = parse_rfc3339(h.get(p + "requests-reset", "")) info.tokens_limit = _int(h.get(p + "tokens-limit")) info.tokens_remaining = _int(h.get(p + "tokens-remaining")) info.tokens_reset_s = parse_rfc3339(h.get(p + "tokens-reset", "")) info.input_tokens_limit = _int(h.get(p + "input-tokens-limit")) info.input_tokens_remaining = _int(h.get(p + "input-tokens-remaining")) info.output_tokens_limit = _int(h.get(p + "output-tokens-limit")) info.output_tokens_remaining = _int(h.get(p + "output-tokens-remaining")) info.request_id = h.get("request-id") info.raw = {k: v for k, v in h.items() if k.startswith(p) or k in ("retry-after", "request-id", "x-should-retry")} return info # -------------------------------------------------------------------------------------- # Backoff # -------------------------------------------------------------------------------------- @dataclass class RetryPolicy: max_attempts: int = 4 # 1 initial + 3 retries base_delay_s: float = 0.5 max_delay_s: float = 20.0 max_retry_after_s: float = 60.0 # cap on server hints (a 6-minute Retry-After should fail fast instead) max_total_s: float = 120.0 # wall-clock deadline for the whole operation jitter: Callable[[float], float] = staticmethod(lambda cap: random.uniform(0, cap)) retry_on_transport_error: bool = True anthropic_429_without_retry_after_retryable: bool = False gemini_zero_quota_retryable: bool = False # 429 with `limit: 0` = paid-only feature on a free-tier key def delay(self, attempt: int, retry_after_s: Optional[float]) -> float: """attempt is 1-based (number of failures so far). Full jitter: U(0, min(cap, base*2^(n-1))). A valid Retry-After is a MINIMUM: wait at least that long, plus a small jitter (OpenAI guide).""" if retry_after_s is not None: return min(retry_after_s, self.max_retry_after_s) + self.jitter(min(1.0, self.base_delay_s)) cap = min(self.max_delay_s, self.base_delay_s * (2 ** (attempt - 1))) return self.jitter(cap) class RetryExhausted(Exception): def __init__(self, message: str, last_response: Optional[Response], attempts: int, history: list[str]): super().__init__(message) self.last_response = last_response self.attempts = attempts self.history = history class NonRetryableError(Exception): def __init__(self, resp: Response, classification: Classification): etype, ecode, msg = extract_error(resp.json()) super().__init__(f"HTTP {resp.status} {classification.reason}: {etype}/{ecode}: {str(msg)[:300] if msg else msg}") self.response = resp self.classification = classification # -------------------------------------------------------------------------------------- # Circuit breaker # -------------------------------------------------------------------------------------- class CircuitOpen(Exception): pass class CircuitBreaker: """Classic 3-state breaker. Opens after `failure_threshold` consecutive failures, allows one probe after `recovery_timeout_s` (half-open), closes on probe success. Thread-safe.""" def __init__(self, failure_threshold: int = 5, recovery_timeout_s: float = 30.0, clock: Callable[[], float] = time.monotonic): self.failure_threshold = failure_threshold self.recovery_timeout_s = recovery_timeout_s self._clock = clock self._lock = threading.Lock() self.state = "closed" self.failures = 0 self.opened_at: Optional[float] = None def allow(self) -> bool: with self._lock: if self.state == "closed": return True if self.state == "open": if self._clock() - (self.opened_at or 0) >= self.recovery_timeout_s: self.state = "half-open" return True return False return True # half-open: one probe in flight allowed (callers should serialise) def record_success(self) -> None: with self._lock: self.state, self.failures, self.opened_at = "closed", 0, None def record_failure(self) -> None: with self._lock: self.failures += 1 if self.state == "half-open" or self.failures >= self.failure_threshold: self.state, self.opened_at = "open", self._clock() # -------------------------------------------------------------------------------------- # Budget guard # -------------------------------------------------------------------------------------- class BudgetExceeded(Exception): pass class BudgetGuard: """Tracks estimated USD spend from `usage` objects. Prices are per 1M tokens: {model: {"input": x, "output": y, "cached_input": z}}. Unknown models use `default_price` so spend is never silently zero.""" def __init__(self, max_usd: float, prices: Optional[dict[str, dict[str, float]]] = None, default_price: Optional[dict[str, float]] = None): self.max_usd = max_usd self.prices = prices or {} self.default_price = default_price or {"input": 5.0, "output": 15.0, "cached_input": 0.5} self.spent_usd = 0.0 self.calls = 0 self._lock = threading.Lock() def check(self) -> None: if self.spent_usd >= self.max_usd: raise BudgetExceeded(f"budget {self.max_usd:.4f} USD exhausted (spent {self.spent_usd:.4f})") def estimate(self, model: str, usage: dict[str, Any]) -> float: p = self.prices.get(model, self.default_price) # OpenAI Responses / xAI Responses: input_tokens/output_tokens (+ input_tokens_details.cached_tokens; xAI output_tokens # already includes reasoning tokens) # OpenAI/xAI Chat Completions: prompt_tokens/completion_tokens (+ prompt_tokens_details.cached_tokens) # Anthropic Messages: input_tokens/output_tokens (+ cache_read_input_tokens, cache_creation_input_tokens) # Gemini usageMetadata: promptTokenCount (includes cached), candidatesTokenCount + thoughtsTokenCount (both billed as # output), cachedContentTokenCount inp = float(usage.get("input_tokens") or usage.get("prompt_tokens") or usage.get("promptTokenCount") or 0) out = float(usage.get("output_tokens") or usage.get("completion_tokens") or (float(usage.get("candidatesTokenCount") or 0) + float(usage.get("thoughtsTokenCount") or 0)) or 0) if "completion_tokens" in usage and "output_tokens" not in usage: # xAI/OpenAI chat: reasoning billed but not in completion_tokens out += float((usage.get("completion_tokens_details") or {}).get("reasoning_tokens") or 0) cached = float((usage.get("input_tokens_details") or {}).get("cached_tokens") or (usage.get("prompt_tokens_details") or {}).get("cached_tokens") or usage.get("cache_read_input_tokens") or usage.get("cachedContentTokenCount") or 0) cache_write = float(usage.get("cache_creation_input_tokens") or 0) uncached = max(0.0, inp - cached) cost = (uncached * p.get("input", 0) + cached * p.get("cached_input", 0) + out * p.get("output", 0) + cache_write * p.get("cache_write", p.get("input", 0) * 1.25)) / 1_000_000 return cost def record(self, model: str, usage: dict[str, Any]) -> float: cost = self.estimate(model, usage) with self._lock: self.spent_usd += cost self.calls += 1 return cost # -------------------------------------------------------------------------------------- # Idempotency # -------------------------------------------------------------------------------------- def new_idempotency_key() -> str: import uuid return str(uuid.uuid4()) def idempotency_headers(provider: str, path: str, key: str) -> dict[str, str]: """FACTS (2026-09-18): - OpenAI: the OpenAPI spec documents an `Idempotency-Key` request header ONLY on POST /v1/agents/sessions/{session_id}/events (operationId createAgentSessionEvents). Nothing for /v1/responses, /v1/chat/completions, /v1/embeddings... Sending the header there is harmless but has no documented effect. - Anthropic: no idempotency header on any Messages/Batches/Files endpoint. Webhook deliveries (inference hooks) expose `webhook-id` as an idempotency key for the RECEIVER side only. So for generation endpoints, idempotency is client-side: dedupe on your own key (see docs/architecture/resilience.md).""" if provider == "openai" and re.fullmatch(r"/v1/agents/sessions/[^/]+/events", path): return {"Idempotency-Key": key} return {} # -------------------------------------------------------------------------------------- # The client # -------------------------------------------------------------------------------------- @dataclass class CallResult: response: Response attempts: int rate_limit: RateLimitInfo history: list[str] provider: str model: Optional[str] = None est_cost_usd: float = 0.0 class ResilientClient: def __init__(self, provider: str, *, api_key: Optional[str] = None, base_url: Optional[str] = None, transport: Transport = urllib_transport, policy: Optional[RetryPolicy] = None, breaker: Optional[CircuitBreaker] = None, budget: Optional[BudgetGuard] = None, sleep: Callable[[float], None] = time.sleep, clock: Callable[[], float] = time.monotonic, anthropic_version: str = "2023-06-01", default_headers: Optional[dict[str, str]] = None, on_event: Optional[Callable[[str, dict[str, Any]], None]] = None): assert provider in PROVIDERS, f"provider must be one of {PROVIDERS}" self.provider = provider env_key = ENV_KEYS[provider] self._api_key = api_key or os.environ.get(env_key, "") self.base_url = (base_url or os.environ.get(env_key.replace("API_KEY", "BASE_URL")) or DEFAULT_BASE_URLS[provider]).rstrip("/") self.transport = transport self.policy = policy or RetryPolicy() self.breaker = breaker or CircuitBreaker() self.budget = budget self._sleep = sleep self._clock = clock self.anthropic_version = anthropic_version self.default_headers = default_headers or {} self.on_event = on_event or (lambda kind, data: None) self.last_rate_limit: Optional[RateLimitInfo] = None # -- headers --------------------------------------------------------------------- def _auth_headers(self) -> dict[str, str]: if self.provider in ("openai", "xai"): # xAI: same Bearer scheme, key prefix xai-…; no version/beta headers exist return {"Authorization": f"Bearer {self._api_key}"} if self.provider == "gemini": # header, never `?key=` in the URL (leaks into logs/referrers) return {"x-goog-api-key": self._api_key} return {"x-api-key": self._api_key, "anthropic-version": self.anthropic_version} def __repr__(self) -> str: # never leak the key return f"ResilientClient(provider={self.provider!r}, base_url={self.base_url!r}, key=***)" # -- core ------------------------------------------------------------------------ def request(self, method: str, path: str, json_body: Any = None, *, body: Optional[bytes] = None, headers: Optional[dict[str, str]] = None, stream: bool = False, connect_timeout: float = 5.0, read_timeout: float = 600.0, idempotency_key: Optional[str] = None, model: Optional[str] = None) -> CallResult: if self.budget: self.budget.check() hdrs = {"Content-Type": "application/json", **self.default_headers, **self._auth_headers(), **(headers or {})} if idempotency_key: hdrs.update(idempotency_headers(self.provider, path, idempotency_key)) data = json.dumps(json_body).encode() if json_body is not None else body req = Request(method, self.base_url + path, hdrs, data, connect_timeout, read_timeout, stream) model = model or (json_body or {}).get("model") if isinstance(json_body, dict) else model history: list[str] = [] start = self._clock() last_resp: Optional[Response] = None attempt = 0 while True: attempt += 1 if not self.breaker.allow(): raise CircuitOpen(f"circuit open for {self.provider}") try: resp = self.transport(req) except TransportError as e: self.breaker.record_failure() history.append(f"attempt {attempt}: transport error {e}") self.on_event("transport_error", {"attempt": attempt, "error": str(e)}) if not self.policy.retry_on_transport_error: raise cls = Classification(True, "transport error") resp = None else: last_resp = resp self.last_rate_limit = parse_rate_limit_headers(self.provider, resp.headers) cls = classify(self.provider, resp, anthropic_429_without_retry_after_retryable=self.policy.anthropic_429_without_retry_after_retryable, gemini_zero_quota_retryable=self.policy.gemini_zero_quota_retryable) if 200 <= resp.status < 400: self.breaker.record_success() result = CallResult(resp, attempt, self.last_rate_limit, history, self.provider, model) if self.budget and not stream: j = resp.json() usage = (j.get("usage") or j.get("usageMetadata")) if isinstance(j, dict) else None # Gemini: usageMetadata if isinstance(usage, dict): result.est_cost_usd = self.budget.record(model or "unknown", usage) return result self.breaker.record_failure() history.append(f"attempt {attempt}: HTTP {resp.status} {cls.reason}") self.on_event("http_error", {"attempt": attempt, "status": resp.status, "reason": cls.reason, "rate_limit": self.last_rate_limit.raw}) if not cls.retryable: raise NonRetryableError(resp, cls) if attempt >= self.policy.max_attempts: raise RetryExhausted(f"gave up after {attempt} attempts", last_resp, attempt, history) delay = self.policy.delay(attempt, cls.retry_after_s) if self._clock() - start + delay > self.policy.max_total_s: raise RetryExhausted(f"deadline {self.policy.max_total_s}s would be exceeded", last_resp, attempt, history) self.on_event("backoff", {"attempt": attempt, "delay_s": round(delay, 3), "retry_after_s": cls.retry_after_s}) self._sleep(delay) # -- convenience ----------------------------------------------------------------- def post_json(self, path: str, body: dict[str, Any], **kw: Any) -> CallResult: return self.request("POST", path, body, **kw) def get(self, path: str, **kw: Any) -> CallResult: return self.request("GET", path, None, **kw) # -------------------------------------------------------------------------------------- # Fallback chains # -------------------------------------------------------------------------------------- @dataclass class Target: client: ResilientClient model: str path: str # "/v1/responses" (OpenAI, xAI) · "/v1/messages" (Anthropic) · "/v1beta/models/{model}:generateContent" (Gemini) build_body: Callable[[str], dict[str, Any]] # model -> request body def resolved_path(self) -> str: return self.path.replace("{model}", self.model) AUTH_OR_VALIDATION_ERROR_TYPES: tuple[str, ...] = ( "authentication_error", "permission_error", "invalid_request_error", # OpenAI + Anthropic error.type "invalid-argument", "unauthenticated", "unauthenticated:no-credentials", "permission-denied", # xAI code "INVALID_ARGUMENT", "UNAUTHENTICATED", "PERMISSION_DENIED", "FAILED_PRECONDITION", # Gemini error.status ) def call_with_fallback(targets: list[Target], *, fallback_on: tuple[type, ...] = (RetryExhausted, CircuitOpen, TransportError, NonRetryableError), skip_non_retryable_types: tuple[str, ...] = AUTH_OR_VALIDATION_ERROR_TYPES) -> CallResult: """Try each target in order. Same-provider model fallback = targets sharing a client with different models; multi-provider fallback = targets with different clients (bodies must be built per provider — see llm_provider.py). Auth/permission/validation errors are NOT a reason to fall back to another model with the same request, except when the error is model-specific (e.g. model_not_found), so callers can tune `skip_non_retryable_types`. Per-provider auth/validation error types: OpenAI/Anthropic `error.type` strings; xAI codes `invalid-argument` (incl. incorrect API key), `unauthenticated:no-credentials`; Gemini statuses `INVALID_ARGUMENT`, `UNAUTHENTICATED`, `PERMISSION_DENIED`. A Gemini 429 `limit: 0` (paid-only model on a free-tier key) IS a reason to fall back.""" errors: list[str] = [] for t in targets: try: return t.client.post_json(t.resolved_path(), t.build_body(t.model), model=t.model) except NonRetryableError as e: errors.append(f"{t.client.provider}/{t.model}: {e}") model_specific = (e.classification.error_code == "model_not_found" or e.response.status == 404 or (t.client.provider == "gemini" and e.response.status == 429)) if e.classification.error_type in skip_non_retryable_types and not model_specific: raise except fallback_on as e: # noqa: PERF203 errors.append(f"{t.client.provider}/{t.model}: {type(e).__name__}: {e}") raise RetryExhausted("all fallback targets failed: " + " | ".join(errors), None, len(targets), errors) # -------------------------------------------------------------------------------------- # Stream resume helpers # -------------------------------------------------------------------------------------- def iter_sse_json(lines: Iterator[bytes]) -> Iterator[dict[str, Any]]: """Tiny SSE → JSON iterator (full parser: examples/shared/streaming/sse_parser.py).""" event, data = None, [] for raw in lines: line = raw.decode("utf-8", "replace").rstrip("\r\n") if line == "": if data: payload = "\n".join(data) if payload != "[DONE]": try: obj = json.loads(payload) if event and isinstance(obj, dict): obj.setdefault("event", event) yield obj except json.JSONDecodeError: pass event, data = None, [] elif line.startswith("event:"): event = line[6:].strip() elif line.startswith("data:"): data.append(line[5:].lstrip(" ")) def openai_stream_with_resume(client: ResilientClient, body: dict[str, Any], *, max_reconnects: int = 5) -> Iterator[dict[str, Any]]: """Stream an OpenAI Response with `background: true, stream: true`, resuming after a drop with GET /v1/responses/{id}?stream=true&starting_after= (docs/guides/background). Requires the response to be retrievable: background responses are stored for the polling window; pass store=true to keep them longer. Terminal events: response.completed / response.failed / response.incomplete / response.cancelled. NOT applicable to xAI: its Responses API rejects `background` (400 "Argument not supported: background"), so xAI streams can only be restarted (see docs/architecture/resilience.md §6).""" body = {**body, "background": True, "stream": True} res = client.post_json("/v1/responses", body, stream=True, model=body.get("model")) response_id: Optional[str] = None cursor: Optional[int] = None reconnects = 0 lines = res.response.lines terminal = {"response.completed", "response.failed", "response.incomplete", "response.cancelled"} while True: try: for ev in iter_sse_json(lines or iter(())): if "sequence_number" in ev: cursor = ev["sequence_number"] if response_id is None and isinstance(ev.get("response"), dict): response_id = ev["response"].get("id") yield ev if ev.get("type") in terminal: return return # stream ended without a terminal event: treat as complete (server closed) except (TransportError, OSError, ConnectionError) as e: reconnects += 1 if response_id is None or reconnects > max_reconnects: raise RetryExhausted(f"stream dropped ({e}); cannot resume", None, reconnects, []) from e q = f"?stream=true&starting_after={cursor}" if cursor is not None else "?stream=true" res = client.get(f"/v1/responses/{response_id}{q}", stream=True) lines = res.response.lines def anthropic_stream_with_restart(client: ResilientClient, body: dict[str, Any], *, max_restarts: int = 2, beta: Optional[str] = None) -> Iterator[dict[str, Any]]: """Anthropic has NO server-side stream resume. Strategy: on a mid-stream drop, restart the same request and re-yield from scratch, emitting a synthetic {"type":"restart","attempt":n} marker so consumers can discard the partial output they already rendered. In-stream `error` events (e.g. overloaded_error, which maps to HTTP 529) also trigger a restart. Note: assistant prefill to continue text is NOT supported on current models (api/errors "Prefill not supported"), so continuation must be done by re-asking, not by prefilling.""" hdrs = {"anthropic-beta": beta} if beta else {} body = {**body, "stream": True} attempt = 0 while True: attempt += 1 res = client.post_json("/v1/messages", body, stream=True, headers=hdrs, model=body.get("model")) try: for ev in iter_sse_json(res.response.lines or iter(())): if ev.get("type") == "error": raise TransportError(f"in-stream error: {ev.get('error')}") yield ev if ev.get("type") == "message_stop": return return except (TransportError, OSError, ConnectionError): if attempt > max_restarts: raise yield {"type": "restart", "attempt": attempt} # -------------------------------------------------------------------------------------- # Optional SDK adapters (import lazily; not required) # -------------------------------------------------------------------------------------- def openai_sdk_client_without_retries(**kw: Any) -> Any: """Return an `openai.OpenAI` with SDK retries disabled so THIS client owns the retry budget (nested retry loops multiply requests — OpenAI rate-limit guide). Requires `pip install openai`.""" from openai import OpenAI # type: ignore return OpenAI(max_retries=0, **kw) def anthropic_sdk_client_without_retries(**kw: Any) -> Any: from anthropic import Anthropic # type: ignore return Anthropic(max_retries=0, **kw) def sdk_error_is_retryable(exc: Exception, provider: Optional[str] = None) -> bool: """Classify an SDK exception via status_code/headers/body. - openai / anthropic SDKs share class names (`status_code`, `response`, `body`); the module name tells them apart. - The `openai` SDK pointed at https://api.x.ai/v1 raises the same classes → pass provider="xai" explicitly. - google-genai raises `google.genai.errors.APIError` (`code`, `response`, `details` = the google.rpc body) → "gemini". - xai-sdk (gRPC) raises `grpc.RpcError` → map `e.code()` (RESOURCE_EXHAUSTED/UNAVAILABLE/DEADLINE_EXCEEDED/INTERNAL retryable).""" mod = type(exc).__module__ if provider is None: provider = ("anthropic" if mod.startswith("anthropic") else "gemini" if mod.startswith("google") else "openai") code_fn = getattr(exc, "code", None) if mod.startswith("grpc") and callable(code_fn): # xai-sdk return getattr(code_fn(), "name", str(code_fn())) in ("RESOURCE_EXHAUSTED", "UNAVAILABLE", "DEADLINE_EXCEEDED", "INTERNAL", "UNKNOWN") status = getattr(exc, "status_code", None) if status is None and isinstance(code_fn, int): # google-genai APIError.code status = code_fn resp = getattr(exc, "response", None) headers = dict(getattr(resp, "headers", {}) or {}) body = getattr(exc, "body", None) if body is None and provider == "gemini": body = getattr(exc, "details", None) if isinstance(body, dict) and "error" not in body: body = {"error": body} fake = Response(status or 0, headers, json.dumps(body).encode() if isinstance(body, (dict, list, str)) else b"") if status is None: # APIConnectionError / APITimeoutError return True return classify(provider, fake).retryable