Python 88.3%
TypeScript 7.6%
Shell 4.1%
1#!/usr/bin/env python32"""Resilient HTTP client for the OpenAI, Anthropic, xAI and Gemini APIs — stdlib only.34STATUS: LIVE_VERIFIED 2026-09-18 (OpenAI/Anthropic) and 2026-09-19 (xAI grok-4.3, Gemini gemini-3.5-flash-lite) —52 minimal calls per provider via examples/shared/provider-abstraction/llm_provider.py, which is built on this client;6offline behaviour (incl. recorded xAI/Gemini error bodies) covered by tests/shared/test_resilient_client.py.78Features (see docs/architecture/resilience.md for the sourced rationale):9 * exponential backoff with FULL jitter, capped, bounded by attempts AND total elapsed time10 * per-request timeouts (connect/read) and a wall-clock deadline11 * provider-aware retryable-error classification:12 - HTTP 408/409/429/500/502/503/504/529 retryable by default13 - OpenAI: `error.type == "insufficient_quota"` and billing/spend-limit `error.code`s are NOT retryable,14 `slow_down` (429) and `server_is_overloaded` (503) are retryable (follow Retry-After)15 - Anthropic: 529 `overloaded_error` retryable; 429 WITHOUT `retry-after` = spend cap (not retryable by default)16 - xAI: bodies are `{"code": "<kebab>", "error": "<message>"}` (400 `invalid-argument` is ALSO the "incorrect API key"17 answer — never retry it), 422 bodies are a bare JSON string (serde error, not retryable), 429 has no Retry-After18 (backoff only), 5xx retryable; observed `x-ratelimit-*` headers parsed like OpenAI's (no reset headers)19 - Gemini: google.rpc.Status `{"error": {"code", "message", "status", "details": [...]}}` (sometimes array-wrapped on20 the OpenAI-compat layer); RESOURCE_EXHAUSTED 429 is retryable with the delay taken from `details[]`21 `google.rpc.RetryInfo.retryDelay` ("40s") or from the message text ("Please retry in 54.22s") — there is NO22 Retry-After header; a 429 whose QuotaFailure says `limit: 0` is a billing/tier gate (not retryable by default);23 UNAVAILABLE 503 / DEADLINE_EXCEEDED 504 / INTERNAL 500 / ABORTED 409 retryable; FAILED_PRECONDITION 400 (billing24 not enabled, region), INVALID_ARGUMENT, PERMISSION_DENIED, NOT_FOUND, UNAUTHENTICATED, ALREADY_EXISTS never retried25 - `x-should-retry: false|true` (Anthropic) always wins; `Retry-After` (seconds or HTTP-date) is honoured and capped26 * rate-limit header parsing (OpenAI + xAI x-ratelimit-*, Anthropic anthropic-ratelimit-*, Gemini: none exposed)27 * idempotency helper (facts: only OpenAI `POST /v1/agents/sessions/{id}/events` documents `Idempotency-Key`;28 neither Responses nor Messages do — so POST retries may double-bill; see docs)29 * circuit breaker (closed → open → half-open) keyed per provider/model30 * fallback chains: same-provider model fallback and multi-provider fallback31 * budget guard (USD, from usage + price table you supply)32 * stream resume helper for OpenAI background responses (`GET /v1/responses/{id}?stream=true&starting_after=N`)33 and a restart strategy for Anthropic (no server-side resume exists)3435No SDK required. Optional adapters at the bottom show how to plug the official SDKs' `max_retries=0` clients in.36Never logs secrets: headers named authorization / x-api-key are redacted in every repr.37"""38from __future__ import annotations3940import email.utils41import json42import os43import random44import re45import threading46import time47import urllib.error48import urllib.request49from dataclasses import dataclass, field50from datetime import datetime, timezone51from typing import Any, Callable, Iterator, Optional5253# --------------------------------------------------------------------------------------54# Transport layer (pluggable so tests run offline)55# --------------------------------------------------------------------------------------5657REDACTED_HEADERS = {"authorization", "x-api-key", "x-goog-api-key", "openai-organization", "openai-project"}5859PROVIDERS = ("openai", "anthropic", "xai", "gemini")60ENV_KEYS = {"openai": "OPENAI_API_KEY", "anthropic": "ANTHROPIC_API_KEY", "xai": "XAI_API_KEY", "gemini": "GEMINI_API_KEY"}61DEFAULT_BASE_URLS = {"openai": "https://api.openai.com", "anthropic": "https://api.anthropic.com",62 "xai": "https://api.x.ai", "gemini": "https://generativelanguage.googleapis.com"}636465@dataclass66class Request:67 method: str68 url: str69 headers: dict[str, str] = field(default_factory=dict)70 body: Optional[bytes] = None71 connect_timeout: float = 5.072 read_timeout: float = 600.0 # both official SDKs default to 600 s read timeout73 stream: bool = False7475 def safe_headers(self) -> dict[str, str]:76 return {k: ("***REDACTED***" if k.lower() in REDACTED_HEADERS else v) for k, v in self.headers.items()}777879@dataclass80class Response:81 status: int82 headers: dict[str, str]83 body: bytes = b""84 lines: Optional[Iterator[bytes]] = None # when stream=True85 elapsed_s: float = 0.08687 def header(self, name: str) -> Optional[str]:88 for k, v in self.headers.items():89 if k.lower() == name.lower():90 return v91 return None9293 def json(self) -> Any:94 try:95 return json.loads(self.body)96 except Exception: # noqa: BLE00197 return None9899100class TransportError(Exception):101 """Network-level failure (DNS, connection reset, timeout). Always retryable unless the deadline is hit."""102103104Transport = Callable[[Request], Response]105106107def urllib_transport(req: Request) -> Response:108 """Default transport. urllib has no pooling: for production prefer http.client keep-alive or httpx (see docs)."""109 t0 = time.monotonic()110 r = urllib.request.Request(req.url, method=req.method, headers=req.headers, data=req.body)111 try:112 # urllib exposes a single timeout: use the read timeout (connect is usually much shorter in practice)113 resp = urllib.request.urlopen(r, timeout=req.read_timeout)114 except urllib.error.HTTPError as e:115 return Response(e.code, dict(e.headers.items()), e.read(), elapsed_s=time.monotonic() - t0)116 except (urllib.error.URLError, TimeoutError, ConnectionError, OSError) as e:117 raise TransportError(str(e)) from e118 hdrs = dict(resp.headers.items())119 if req.stream:120 def gen() -> Iterator[bytes]:121 with resp:122 for raw in resp:123 yield raw124 return Response(resp.status, hdrs, b"", lines=gen(), elapsed_s=time.monotonic() - t0)125 with resp:126 return Response(resp.status, hdrs, resp.read(), elapsed_s=time.monotonic() - t0)127128129# --------------------------------------------------------------------------------------130# Error classification131# --------------------------------------------------------------------------------------132133RETRYABLE_STATUSES = {408, 409, 429, 500, 502, 503, 504, 529}134135# OpenAI: 429s that require user action (docs/guides/error-codes + rate-limits "Don't retry quota, billing…")136OPENAI_NON_RETRYABLE_CODES = {137 "insufficient_quota",138 "credit_balance_exhausted",139 "organization_spend_limit_exceeded",140 "project_spend_limit_exceeded",141 "organization_usage_limit_exceeded",142}143OPENAI_NON_RETRYABLE_TYPES = {"insufficient_quota", "invalid_request_error", "authentication_error",144 "permission_error", "not_found_error"}145ANTHROPIC_NON_RETRYABLE_TYPES = {"invalid_request_error", "authentication_error", "permission_error",146 "not_found_error", "billing_error", "request_too_large"}147# xAI inference API: {"code": "<kebab-case>", "error": "<message>"} — generated/fragments/errors/xai-errors.json (live 2026-09-19)148XAI_NON_RETRYABLE_CODES = {"invalid-argument", "unauthenticated:no-credentials", "unauthenticated", "permission-denied",149 "not-found", "method-not-allowed", "unsupported-media-type", "unprocessable-entity"}150# Gemini google.rpc.Status `error.status` — generated/fragments/errors/gemini-errors.json151GEMINI_RETRYABLE_STATUSES = {"RESOURCE_EXHAUSTED", "UNAVAILABLE", "DEADLINE_EXCEEDED", "INTERNAL", "ABORTED"}152GEMINI_NON_RETRYABLE_STATUSES = {"INVALID_ARGUMENT", "FAILED_PRECONDITION", "UNAUTHENTICATED", "PERMISSION_DENIED",153 "NOT_FOUND", "ALREADY_EXISTS", "OUT_OF_RANGE", "UNIMPLEMENTED", "CANCELLED"}154155156@dataclass157class Classification:158 retryable: bool159 reason: str160 retry_after_s: Optional[float] = None161 error_type: Optional[str] = None162 error_code: Optional[str] = None163164165def parse_retry_after(value: Optional[str]) -> Optional[float]:166 """`Retry-After` may be delta-seconds or an HTTP-date (RFC 7231)."""167 if not value:168 return None169 value = value.strip()170 if re.fullmatch(r"\d+(\.\d+)?", value):171 return float(value)172 try:173 dt = email.utils.parsedate_to_datetime(value)174 if dt.tzinfo is None:175 dt = dt.replace(tzinfo=timezone.utc)176 return max(0.0, (dt - datetime.now(timezone.utc)).total_seconds())177 except Exception: # noqa: BLE001178 return None179180181def extract_error(body: Any) -> tuple[Optional[str], Optional[str], Optional[str]]:182 """Return (type, code, message) from any of the four providers' error envelopes:183 - OpenAI / Anthropic: {"error": {"type", "code", "message"}}184 - Gemini (google.rpc.Status): {"error": {"code": 429, "message", "status": "RESOURCE_EXHAUSTED", "details": [...]}}185 → type = status, code = str(code). The OpenAI-compat layer may wrap it in a one-element array.186 - xAI inference: {"code": "invalid-argument", "error": "<message>"} → type = code = the kebab-case code187 - xAI 422 / some 400s: a bare JSON string → (None, None, <string>)188 - xAI Management API: {"code": 16, "message": "...", "details": []} (gRPC numbering) → type = "grpc", code = "16"189 """190 if isinstance(body, list) and len(body) == 1 and isinstance(body[0], dict):191 body = body[0]192 if isinstance(body, str):193 return None, None, body194 if not isinstance(body, dict):195 return None, None, None196 err = body.get("error")197 if isinstance(err, dict):198 if "status" in err and isinstance(err.get("code"), int): # google.rpc.Status199 return err.get("status"), str(err.get("code")), err.get("message")200 code = err.get("code")201 return err.get("type"), (str(code) if isinstance(code, int) else code), err.get("message")202 if isinstance(err, str): # xAI inference API203 code = body.get("code")204 return (str(code) if code is not None else None), (str(code) if code is not None else None), err205 if isinstance(body.get("code"), int) and "message" in body: # gRPC-style (xAI Management API)206 return "grpc", str(body["code"]), body.get("message")207 return None, None, None208209210_GEMINI_RETRY_IN_RE = re.compile(r"[Pp]lease retry in\s+(\d+(?:\.\d+)?)\s*(ms|s|m|h)?\b")211_GO_DURATION_RE = re.compile(r"^(\d+(?:\.\d+)?)(ms|s|m|h)?$")212213214def parse_gemini_retry_delay(body: Any) -> Optional[float]:215 """Gemini never sends Retry-After. The delay lives (a) in `error.details[]` as216 `{"@type": "type.googleapis.com/google.rpc.RetryInfo", "retryDelay": "40s"}` (protobuf Duration string), and/or217 (b) at the end of `error.message`: "... Please retry in 54.22098241s." Returns seconds, preferring RetryInfo."""218 if isinstance(body, list) and body and isinstance(body[0], dict):219 body = body[0]220 err = body.get("error") if isinstance(body, dict) else None221 if not isinstance(err, dict):222 return None223 for d in err.get("details") or []:224 if isinstance(d, dict) and str(d.get("@type", "")).endswith("google.rpc.RetryInfo"):225 rd = d.get("retryDelay")226 if isinstance(rd, dict): # JSON-encoded Duration object form {seconds, nanos}227 return float(rd.get("seconds", 0)) + float(rd.get("nanos", 0)) / 1e9228 m = _GO_DURATION_RE.match(str(rd or "").strip())229 if m:230 n, unit = float(m.group(1)), m.group(2) or "s"231 return n * {"ms": 0.001, "s": 1, "m": 60, "h": 3600}[unit]232 m = _GEMINI_RETRY_IN_RE.search(str(err.get("message") or ""))233 if m:234 n, unit = float(m.group(1)), m.group(2) or "s"235 return n * {"ms": 0.001, "s": 1, "m": 60, "h": 3600}[unit]236 return None237238239def gemini_quota_is_zero(body: Any) -> bool:240 """True when the 429 says the quota limit is 0 (free tier hitting a paid-only model/feature) — retrying is pointless;241 the fix is billing/tier (docs/gemini/rate-limits.md §7: `limit: 0` for Pro models on a free-tier key)."""242 if isinstance(body, list) and body and isinstance(body[0], dict):243 body = body[0]244 err = body.get("error") if isinstance(body, dict) else None245 if not isinstance(err, dict):246 return False247 lines = [ln for ln in str(err.get("message") or "").splitlines() if "Quota exceeded" in ln]248 return bool(lines) and all(re.search(r"limit:\s*0(?:\D|$)", ln) for ln in lines)249250251def classify(provider: str, resp: Response, *, anthropic_429_without_retry_after_retryable: bool = False,252 gemini_zero_quota_retryable: bool = False) -> Classification:253 """Decide whether an HTTP response should be retried. Provider is one of PROVIDERS."""254 st = resp.status255 body = resp.json()256 etype, ecode, emsg = extract_error(body)257 retry_after = parse_retry_after(resp.header("retry-after"))258259 # Explicit server hint (Anthropic documents x-should-retry; SDKs honour it)260 xsr = resp.header("x-should-retry")261 if xsr is not None:262 flag = xsr.strip().lower() == "true"263 return Classification(flag, f"x-should-retry={xsr.strip()}", retry_after, etype, ecode)264265 if 200 <= st < 400:266 return Classification(False, "success", None, etype, ecode)267268 if provider == "openai":269 if etype in OPENAI_NON_RETRYABLE_TYPES and st != 409:270 return Classification(False, f"openai error.type={etype}", None, etype, ecode)271 if ecode in OPENAI_NON_RETRYABLE_CODES:272 return Classification(False, f"openai error.code={ecode} (billing/quota: user action required)", None, etype, ecode)273 elif provider == "anthropic":274 if etype in ANTHROPIC_NON_RETRYABLE_TYPES:275 return Classification(False, f"anthropic error.type={etype}", None, etype, ecode)276 if st == 429 and retry_after is None and not anthropic_429_without_retry_after_retryable:277 # docs: "A tier spend-cap 429 has no retry-after header and keeps failing until access resumes"278 return Classification(False, "anthropic 429 without retry-after (spend cap suspected)", None, etype, ecode)279 elif provider == "xai":280 if st == 400:281 # xAI answers a malformed/unknown API key with 400 invalid-argument "Incorrect API key provided" — not 401282 hint = " (incorrect API key)" if emsg and "API key" in emsg else ""283 return Classification(False, f"xai 400 code={ecode}{hint}", None, etype, ecode)284 if ecode in XAI_NON_RETRYABLE_CODES or st in (401, 403, 404, 405, 415, 422):285 return Classification(False, f"xai http {st} code={ecode}", None, etype, ecode)286 if st == 429:287 # RESOURCE_EXHAUSTED (gRPC 8): team RPS/TPM or per-key qps/qpm/tpm; no Retry-After documented or observed → backoff.288 # Depleted prepaid credits with a $0 invoiced limit are "automatically rejected" (status undocumented).289 if emsg and re.search(r"credit|balance|billing", emsg, re.I):290 return Classification(False, "xai 429 mentions credits/billing (user action required)", None, etype, ecode)291 return Classification(True, "xai http 429 (rate limit, exponential backoff)", retry_after, etype, ecode)292 elif provider == "gemini":293 if st == 429 or etype == "RESOURCE_EXHAUSTED":294 delay = retry_after if retry_after is not None else parse_gemini_retry_delay(body)295 if gemini_quota_is_zero(body) and not gemini_zero_quota_retryable:296 return Classification(False, "gemini RESOURCE_EXHAUSTED with limit: 0 (free tier / paid-only model: enable billing)",297 delay, etype, ecode)298 return Classification(True, "gemini RESOURCE_EXHAUSTED (delay from RetryInfo/message; no Retry-After header)", delay, etype, ecode)299 if etype in GEMINI_NON_RETRYABLE_STATUSES:300 return Classification(False, f"gemini error.status={etype}", None, etype, ecode)301 if etype in GEMINI_RETRYABLE_STATUSES:302 return Classification(True, f"gemini error.status={etype}", retry_after, etype, ecode)303 if st == 409: # no status in body: ALREADY_EXISTS (not retryable) vs ABORTED (retryable) is undecidable → be safe304 return Classification(False, "gemini 409 without status (assume ALREADY_EXISTS)", None, etype, ecode)305306 if st in RETRYABLE_STATUSES:307 return Classification(True, f"http {st}", retry_after, etype, ecode)308 return Classification(False, f"http {st} not retryable", None, etype, ecode)309310311# --------------------------------------------------------------------------------------312# Rate-limit headers313# --------------------------------------------------------------------------------------314315_DUR_RE = re.compile(r"(\d+(?:\.\d+)?)(ms|s|m|h|d)")316317318def parse_openai_duration(value: str) -> Optional[float]:319 """OpenAI reset headers look like '1s', '6m0s', '250ms', '1h2m3.5s'. Returns seconds."""320 if not value:321 return None322 total, matched = 0.0, False323 for num, unit in _DUR_RE.findall(value):324 matched = True325 n = float(num)326 total += {"ms": n / 1000, "s": n, "m": n * 60, "h": n * 3600, "d": n * 86400}[unit]327 return total if matched else None328329330def parse_rfc3339(value: str) -> Optional[float]:331 """Anthropic reset headers are RFC 3339 timestamps. Returns seconds from now (>= 0)."""332 try:333 dt = datetime.fromisoformat(value.replace("Z", "+00:00"))334 return max(0.0, (dt - datetime.now(timezone.utc)).total_seconds())335 except Exception: # noqa: BLE001336 return None337338339@dataclass340class RateLimitInfo:341 provider: str342 requests_limit: Optional[int] = None343 requests_remaining: Optional[int] = None344 requests_reset_s: Optional[float] = None345 tokens_limit: Optional[int] = None346 tokens_remaining: Optional[int] = None347 tokens_reset_s: Optional[float] = None348 input_tokens_limit: Optional[int] = None # Anthropic only349 input_tokens_remaining: Optional[int] = None350 output_tokens_limit: Optional[int] = None # Anthropic only351 output_tokens_remaining: Optional[int] = None352 project_tokens_remaining: Optional[int] = None # OpenAI only353 retry_after_s: Optional[float] = None354 request_id: Optional[str] = None355 raw: dict[str, str] = field(default_factory=dict)356357358def _int(v: Optional[str]) -> Optional[int]:359 try:360 return int(v) if v is not None else None361 except ValueError:362 return None363364365def parse_rate_limit_headers(provider: str, headers: dict[str, str]) -> RateLimitInfo:366 h = {k.lower(): v for k, v in headers.items()}367 info = RateLimitInfo(provider=provider, retry_after_s=parse_retry_after(h.get("retry-after")))368 if provider in ("openai", "xai"):369 # xAI (observed 2026-09-19, undocumented): x-ratelimit-limit/remaining-requests (per-minute budget) and370 # x-ratelimit-limit/remaining-tokens (TPM) on inference responses; NO x-ratelimit-reset-* headers → reset stays None.371 info.requests_limit = _int(h.get("x-ratelimit-limit-requests"))372 info.requests_remaining = _int(h.get("x-ratelimit-remaining-requests"))373 info.requests_reset_s = parse_openai_duration(h.get("x-ratelimit-reset-requests", ""))374 info.tokens_limit = _int(h.get("x-ratelimit-limit-tokens"))375 info.tokens_remaining = _int(h.get("x-ratelimit-remaining-tokens"))376 info.tokens_reset_s = parse_openai_duration(h.get("x-ratelimit-reset-tokens", ""))377 info.project_tokens_remaining = _int(h.get("x-ratelimit-remaining-project-tokens"))378 info.request_id = h.get("x-request-id")379 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")}380 elif provider == "gemini":381 # Gemini exposes NO rate-limit, Retry-After or request-id headers (docs/gemini/rate-limits.md §8). Only the382 # undocumented X-Gemini-Service-Tier and Server-Timing are kept for observability; correlate with body `responseId`.383 info.raw = {k: v for k, v in h.items() if k in ("x-gemini-service-tier", "server-timing", "retry-after")}384 else:385 p = "anthropic-ratelimit-"386 info.requests_limit = _int(h.get(p + "requests-limit"))387 info.requests_remaining = _int(h.get(p + "requests-remaining"))388 info.requests_reset_s = parse_rfc3339(h.get(p + "requests-reset", ""))389 info.tokens_limit = _int(h.get(p + "tokens-limit"))390 info.tokens_remaining = _int(h.get(p + "tokens-remaining"))391 info.tokens_reset_s = parse_rfc3339(h.get(p + "tokens-reset", ""))392 info.input_tokens_limit = _int(h.get(p + "input-tokens-limit"))393 info.input_tokens_remaining = _int(h.get(p + "input-tokens-remaining"))394 info.output_tokens_limit = _int(h.get(p + "output-tokens-limit"))395 info.output_tokens_remaining = _int(h.get(p + "output-tokens-remaining"))396 info.request_id = h.get("request-id")397 info.raw = {k: v for k, v in h.items() if k.startswith(p) or k in ("retry-after", "request-id", "x-should-retry")}398 return info399400401# --------------------------------------------------------------------------------------402# Backoff403# --------------------------------------------------------------------------------------404405@dataclass406class RetryPolicy:407 max_attempts: int = 4 # 1 initial + 3 retries408 base_delay_s: float = 0.5409 max_delay_s: float = 20.0410 max_retry_after_s: float = 60.0 # cap on server hints (a 6-minute Retry-After should fail fast instead)411 max_total_s: float = 120.0 # wall-clock deadline for the whole operation412 jitter: Callable[[float], float] = staticmethod(lambda cap: random.uniform(0, cap))413 retry_on_transport_error: bool = True414 anthropic_429_without_retry_after_retryable: bool = False415 gemini_zero_quota_retryable: bool = False # 429 with `limit: 0` = paid-only feature on a free-tier key416417 def delay(self, attempt: int, retry_after_s: Optional[float]) -> float:418 """attempt is 1-based (number of failures so far). Full jitter: U(0, min(cap, base*2^(n-1))).419 A valid Retry-After is a MINIMUM: wait at least that long, plus a small jitter (OpenAI guide)."""420 if retry_after_s is not None:421 return min(retry_after_s, self.max_retry_after_s) + self.jitter(min(1.0, self.base_delay_s))422 cap = min(self.max_delay_s, self.base_delay_s * (2 ** (attempt - 1)))423 return self.jitter(cap)424425426class RetryExhausted(Exception):427 def __init__(self, message: str, last_response: Optional[Response], attempts: int, history: list[str]):428 super().__init__(message)429 self.last_response = last_response430 self.attempts = attempts431 self.history = history432433434class NonRetryableError(Exception):435 def __init__(self, resp: Response, classification: Classification):436 etype, ecode, msg = extract_error(resp.json())437 super().__init__(f"HTTP {resp.status} {classification.reason}: {etype}/{ecode}: {str(msg)[:300] if msg else msg}")438 self.response = resp439 self.classification = classification440441442# --------------------------------------------------------------------------------------443# Circuit breaker444# --------------------------------------------------------------------------------------445446class CircuitOpen(Exception):447 pass448449450class CircuitBreaker:451 """Classic 3-state breaker. Opens after `failure_threshold` consecutive failures, allows one probe after452 `recovery_timeout_s` (half-open), closes on probe success. Thread-safe."""453454 def __init__(self, failure_threshold: int = 5, recovery_timeout_s: float = 30.0, clock: Callable[[], float] = time.monotonic):455 self.failure_threshold = failure_threshold456 self.recovery_timeout_s = recovery_timeout_s457 self._clock = clock458 self._lock = threading.Lock()459 self.state = "closed"460 self.failures = 0461 self.opened_at: Optional[float] = None462463 def allow(self) -> bool:464 with self._lock:465 if self.state == "closed":466 return True467 if self.state == "open":468 if self._clock() - (self.opened_at or 0) >= self.recovery_timeout_s:469 self.state = "half-open"470 return True471 return False472 return True # half-open: one probe in flight allowed (callers should serialise)473474 def record_success(self) -> None:475 with self._lock:476 self.state, self.failures, self.opened_at = "closed", 0, None477478 def record_failure(self) -> None:479 with self._lock:480 self.failures += 1481 if self.state == "half-open" or self.failures >= self.failure_threshold:482 self.state, self.opened_at = "open", self._clock()483484485# --------------------------------------------------------------------------------------486# Budget guard487# --------------------------------------------------------------------------------------488489class BudgetExceeded(Exception):490 pass491492493class BudgetGuard:494 """Tracks estimated USD spend from `usage` objects. Prices are per 1M tokens: {model: {"input": x, "output": y, "cached_input": z}}.495 Unknown models use `default_price` so spend is never silently zero."""496497 def __init__(self, max_usd: float, prices: Optional[dict[str, dict[str, float]]] = None,498 default_price: Optional[dict[str, float]] = None):499 self.max_usd = max_usd500 self.prices = prices or {}501 self.default_price = default_price or {"input": 5.0, "output": 15.0, "cached_input": 0.5}502 self.spent_usd = 0.0503 self.calls = 0504 self._lock = threading.Lock()505506 def check(self) -> None:507 if self.spent_usd >= self.max_usd:508 raise BudgetExceeded(f"budget {self.max_usd:.4f} USD exhausted (spent {self.spent_usd:.4f})")509510 def estimate(self, model: str, usage: dict[str, Any]) -> float:511 p = self.prices.get(model, self.default_price)512 # OpenAI Responses / xAI Responses: input_tokens/output_tokens (+ input_tokens_details.cached_tokens; xAI output_tokens513 # already includes reasoning tokens)514 # OpenAI/xAI Chat Completions: prompt_tokens/completion_tokens (+ prompt_tokens_details.cached_tokens)515 # Anthropic Messages: input_tokens/output_tokens (+ cache_read_input_tokens, cache_creation_input_tokens)516 # Gemini usageMetadata: promptTokenCount (includes cached), candidatesTokenCount + thoughtsTokenCount (both billed as517 # output), cachedContentTokenCount518 inp = float(usage.get("input_tokens") or usage.get("prompt_tokens") or usage.get("promptTokenCount") or 0)519 out = float(usage.get("output_tokens") or usage.get("completion_tokens")520 or (float(usage.get("candidatesTokenCount") or 0) + float(usage.get("thoughtsTokenCount") or 0)) or 0)521 if "completion_tokens" in usage and "output_tokens" not in usage: # xAI/OpenAI chat: reasoning billed but not in completion_tokens522 out += float((usage.get("completion_tokens_details") or {}).get("reasoning_tokens") or 0)523 cached = float((usage.get("input_tokens_details") or {}).get("cached_tokens") or (usage.get("prompt_tokens_details") or {}).get("cached_tokens")524 or usage.get("cache_read_input_tokens") or usage.get("cachedContentTokenCount") or 0)525 cache_write = float(usage.get("cache_creation_input_tokens") or 0)526 uncached = max(0.0, inp - cached)527 cost = (uncached * p.get("input", 0) + cached * p.get("cached_input", 0) + out * p.get("output", 0)528 + cache_write * p.get("cache_write", p.get("input", 0) * 1.25)) / 1_000_000529 return cost530531 def record(self, model: str, usage: dict[str, Any]) -> float:532 cost = self.estimate(model, usage)533 with self._lock:534 self.spent_usd += cost535 self.calls += 1536 return cost537538539# --------------------------------------------------------------------------------------540# Idempotency541# --------------------------------------------------------------------------------------542543def new_idempotency_key() -> str:544 import uuid545 return str(uuid.uuid4())546547548def idempotency_headers(provider: str, path: str, key: str) -> dict[str, str]:549 """FACTS (2026-09-18):550 - OpenAI: the OpenAPI spec documents an `Idempotency-Key` request header ONLY on551 POST /v1/agents/sessions/{session_id}/events (operationId createAgentSessionEvents). Nothing for /v1/responses,552 /v1/chat/completions, /v1/embeddings... Sending the header there is harmless but has no documented effect.553 - Anthropic: no idempotency header on any Messages/Batches/Files endpoint. Webhook deliveries (inference hooks)554 expose `webhook-id` as an idempotency key for the RECEIVER side only.555 So for generation endpoints, idempotency is client-side: dedupe on your own key (see docs/architecture/resilience.md)."""556 if provider == "openai" and re.fullmatch(r"/v1/agents/sessions/[^/]+/events", path):557 return {"Idempotency-Key": key}558 return {}559560561# --------------------------------------------------------------------------------------562# The client563# --------------------------------------------------------------------------------------564565@dataclass566class CallResult:567 response: Response568 attempts: int569 rate_limit: RateLimitInfo570 history: list[str]571 provider: str572 model: Optional[str] = None573 est_cost_usd: float = 0.0574575576class ResilientClient:577 def __init__(self, provider: str, *, api_key: Optional[str] = None, base_url: Optional[str] = None,578 transport: Transport = urllib_transport, policy: Optional[RetryPolicy] = None,579 breaker: Optional[CircuitBreaker] = None, budget: Optional[BudgetGuard] = None,580 sleep: Callable[[float], None] = time.sleep, clock: Callable[[], float] = time.monotonic,581 anthropic_version: str = "2023-06-01", default_headers: Optional[dict[str, str]] = None,582 on_event: Optional[Callable[[str, dict[str, Any]], None]] = None):583 assert provider in PROVIDERS, f"provider must be one of {PROVIDERS}"584 self.provider = provider585 env_key = ENV_KEYS[provider]586 self._api_key = api_key or os.environ.get(env_key, "")587 self.base_url = (base_url or os.environ.get(env_key.replace("API_KEY", "BASE_URL")) or DEFAULT_BASE_URLS[provider]).rstrip("/")588 self.transport = transport589 self.policy = policy or RetryPolicy()590 self.breaker = breaker or CircuitBreaker()591 self.budget = budget592 self._sleep = sleep593 self._clock = clock594 self.anthropic_version = anthropic_version595 self.default_headers = default_headers or {}596 self.on_event = on_event or (lambda kind, data: None)597 self.last_rate_limit: Optional[RateLimitInfo] = None598599 # -- headers ---------------------------------------------------------------------600 def _auth_headers(self) -> dict[str, str]:601 if self.provider in ("openai", "xai"): # xAI: same Bearer scheme, key prefix xai-…; no version/beta headers exist602 return {"Authorization": f"Bearer {self._api_key}"}603 if self.provider == "gemini": # header, never `?key=` in the URL (leaks into logs/referrers)604 return {"x-goog-api-key": self._api_key}605 return {"x-api-key": self._api_key, "anthropic-version": self.anthropic_version}606607 def __repr__(self) -> str: # never leak the key608 return f"ResilientClient(provider={self.provider!r}, base_url={self.base_url!r}, key=***)"609610 # -- core ------------------------------------------------------------------------611 def request(self, method: str, path: str, json_body: Any = None, *, body: Optional[bytes] = None,612 headers: Optional[dict[str, str]] = None, stream: bool = False, connect_timeout: float = 5.0,613 read_timeout: float = 600.0, idempotency_key: Optional[str] = None, model: Optional[str] = None) -> CallResult:614 if self.budget:615 self.budget.check()616 hdrs = {"Content-Type": "application/json", **self.default_headers, **self._auth_headers(), **(headers or {})}617 if idempotency_key:618 hdrs.update(idempotency_headers(self.provider, path, idempotency_key))619 data = json.dumps(json_body).encode() if json_body is not None else body620 req = Request(method, self.base_url + path, hdrs, data, connect_timeout, read_timeout, stream)621 model = model or (json_body or {}).get("model") if isinstance(json_body, dict) else model622623 history: list[str] = []624 start = self._clock()625 last_resp: Optional[Response] = None626 attempt = 0627 while True:628 attempt += 1629 if not self.breaker.allow():630 raise CircuitOpen(f"circuit open for {self.provider}")631 try:632 resp = self.transport(req)633 except TransportError as e:634 self.breaker.record_failure()635 history.append(f"attempt {attempt}: transport error {e}")636 self.on_event("transport_error", {"attempt": attempt, "error": str(e)})637 if not self.policy.retry_on_transport_error:638 raise639 cls = Classification(True, "transport error")640 resp = None641 else:642 last_resp = resp643 self.last_rate_limit = parse_rate_limit_headers(self.provider, resp.headers)644 cls = classify(self.provider, resp,645 anthropic_429_without_retry_after_retryable=self.policy.anthropic_429_without_retry_after_retryable,646 gemini_zero_quota_retryable=self.policy.gemini_zero_quota_retryable)647 if 200 <= resp.status < 400:648 self.breaker.record_success()649 result = CallResult(resp, attempt, self.last_rate_limit, history, self.provider, model)650 if self.budget and not stream:651 j = resp.json()652 usage = (j.get("usage") or j.get("usageMetadata")) if isinstance(j, dict) else None # Gemini: usageMetadata653 if isinstance(usage, dict):654 result.est_cost_usd = self.budget.record(model or "unknown", usage)655 return result656 self.breaker.record_failure()657 history.append(f"attempt {attempt}: HTTP {resp.status} {cls.reason}")658 self.on_event("http_error", {"attempt": attempt, "status": resp.status, "reason": cls.reason,659 "rate_limit": self.last_rate_limit.raw})660 if not cls.retryable:661 raise NonRetryableError(resp, cls)662663 if attempt >= self.policy.max_attempts:664 raise RetryExhausted(f"gave up after {attempt} attempts", last_resp, attempt, history)665 delay = self.policy.delay(attempt, cls.retry_after_s)666 if self._clock() - start + delay > self.policy.max_total_s:667 raise RetryExhausted(f"deadline {self.policy.max_total_s}s would be exceeded", last_resp, attempt, history)668 self.on_event("backoff", {"attempt": attempt, "delay_s": round(delay, 3), "retry_after_s": cls.retry_after_s})669 self._sleep(delay)670671 # -- convenience -----------------------------------------------------------------672 def post_json(self, path: str, body: dict[str, Any], **kw: Any) -> CallResult:673 return self.request("POST", path, body, **kw)674675 def get(self, path: str, **kw: Any) -> CallResult:676 return self.request("GET", path, None, **kw)677678679# --------------------------------------------------------------------------------------680# Fallback chains681# --------------------------------------------------------------------------------------682683@dataclass684class Target:685 client: ResilientClient686 model: str687 path: str # "/v1/responses" (OpenAI, xAI) · "/v1/messages" (Anthropic) · "/v1beta/models/{model}:generateContent" (Gemini)688 build_body: Callable[[str], dict[str, Any]] # model -> request body689690 def resolved_path(self) -> str:691 return self.path.replace("{model}", self.model)692693694AUTH_OR_VALIDATION_ERROR_TYPES: tuple[str, ...] = (695 "authentication_error", "permission_error", "invalid_request_error", # OpenAI + Anthropic error.type696 "invalid-argument", "unauthenticated", "unauthenticated:no-credentials", "permission-denied", # xAI code697 "INVALID_ARGUMENT", "UNAUTHENTICATED", "PERMISSION_DENIED", "FAILED_PRECONDITION", # Gemini error.status698)699700701def call_with_fallback(targets: list[Target], *, fallback_on: tuple[type, ...] = (RetryExhausted, CircuitOpen, TransportError, NonRetryableError),702 skip_non_retryable_types: tuple[str, ...] = AUTH_OR_VALIDATION_ERROR_TYPES) -> CallResult:703 """Try each target in order. Same-provider model fallback = targets sharing a client with different models;704 multi-provider fallback = targets with different clients (bodies must be built per provider — see llm_provider.py).705 Auth/permission/validation errors are NOT a reason to fall back to another model with the same request, except706 when the error is model-specific (e.g. model_not_found), so callers can tune `skip_non_retryable_types`.707 Per-provider auth/validation error types: OpenAI/Anthropic `error.type` strings; xAI codes `invalid-argument`708 (incl. incorrect API key), `unauthenticated:no-credentials`; Gemini statuses `INVALID_ARGUMENT`, `UNAUTHENTICATED`,709 `PERMISSION_DENIED`. A Gemini 429 `limit: 0` (paid-only model on a free-tier key) IS a reason to fall back."""710 errors: list[str] = []711 for t in targets:712 try:713 return t.client.post_json(t.resolved_path(), t.build_body(t.model), model=t.model)714 except NonRetryableError as e:715 errors.append(f"{t.client.provider}/{t.model}: {e}")716 model_specific = (e.classification.error_code == "model_not_found" or e.response.status == 404717 or (t.client.provider == "gemini" and e.response.status == 429))718 if e.classification.error_type in skip_non_retryable_types and not model_specific:719 raise720 except fallback_on as e: # noqa: PERF203721 errors.append(f"{t.client.provider}/{t.model}: {type(e).__name__}: {e}")722 raise RetryExhausted("all fallback targets failed: " + " | ".join(errors), None, len(targets), errors)723724725# --------------------------------------------------------------------------------------726# Stream resume helpers727# --------------------------------------------------------------------------------------728729def iter_sse_json(lines: Iterator[bytes]) -> Iterator[dict[str, Any]]:730 """Tiny SSE → JSON iterator (full parser: examples/shared/streaming/sse_parser.py)."""731 event, data = None, []732 for raw in lines:733 line = raw.decode("utf-8", "replace").rstrip("\r\n")734 if line == "":735 if data:736 payload = "\n".join(data)737 if payload != "[DONE]":738 try:739 obj = json.loads(payload)740 if event and isinstance(obj, dict):741 obj.setdefault("event", event)742 yield obj743 except json.JSONDecodeError:744 pass745 event, data = None, []746 elif line.startswith("event:"):747 event = line[6:].strip()748 elif line.startswith("data:"):749 data.append(line[5:].lstrip(" "))750751752def openai_stream_with_resume(client: ResilientClient, body: dict[str, Any], *, max_reconnects: int = 5) -> Iterator[dict[str, Any]]:753 """Stream an OpenAI Response with `background: true, stream: true`, resuming after a drop with754 GET /v1/responses/{id}?stream=true&starting_after=<last sequence_number> (docs/guides/background).755 Requires the response to be retrievable: background responses are stored for the polling window; pass store=true756 to keep them longer. Terminal events: response.completed / response.failed / response.incomplete / response.cancelled.757 NOT applicable to xAI: its Responses API rejects `background` (400 "Argument not supported: background"), so xAI758 streams can only be restarted (see docs/architecture/resilience.md §6)."""759 body = {**body, "background": True, "stream": True}760 res = client.post_json("/v1/responses", body, stream=True, model=body.get("model"))761 response_id: Optional[str] = None762 cursor: Optional[int] = None763 reconnects = 0764 lines = res.response.lines765 terminal = {"response.completed", "response.failed", "response.incomplete", "response.cancelled"}766 while True:767 try:768 for ev in iter_sse_json(lines or iter(())):769 if "sequence_number" in ev:770 cursor = ev["sequence_number"]771 if response_id is None and isinstance(ev.get("response"), dict):772 response_id = ev["response"].get("id")773 yield ev774 if ev.get("type") in terminal:775 return776 return # stream ended without a terminal event: treat as complete (server closed)777 except (TransportError, OSError, ConnectionError) as e:778 reconnects += 1779 if response_id is None or reconnects > max_reconnects:780 raise RetryExhausted(f"stream dropped ({e}); cannot resume", None, reconnects, []) from e781 q = f"?stream=true&starting_after={cursor}" if cursor is not None else "?stream=true"782 res = client.get(f"/v1/responses/{response_id}{q}", stream=True)783 lines = res.response.lines784785786def anthropic_stream_with_restart(client: ResilientClient, body: dict[str, Any], *, max_restarts: int = 2,787 beta: Optional[str] = None) -> Iterator[dict[str, Any]]:788 """Anthropic has NO server-side stream resume. Strategy: on a mid-stream drop, restart the same request and789 re-yield from scratch, emitting a synthetic {"type":"restart","attempt":n} marker so consumers can discard the790 partial output they already rendered. In-stream `error` events (e.g. overloaded_error, which maps to HTTP 529)791 also trigger a restart. Note: assistant prefill to continue text is NOT supported on current models792 (api/errors "Prefill not supported"), so continuation must be done by re-asking, not by prefilling."""793 hdrs = {"anthropic-beta": beta} if beta else {}794 body = {**body, "stream": True}795 attempt = 0796 while True:797 attempt += 1798 res = client.post_json("/v1/messages", body, stream=True, headers=hdrs, model=body.get("model"))799 try:800 for ev in iter_sse_json(res.response.lines or iter(())):801 if ev.get("type") == "error":802 raise TransportError(f"in-stream error: {ev.get('error')}")803 yield ev804 if ev.get("type") == "message_stop":805 return806 return807 except (TransportError, OSError, ConnectionError):808 if attempt > max_restarts:809 raise810 yield {"type": "restart", "attempt": attempt}811812813# --------------------------------------------------------------------------------------814# Optional SDK adapters (import lazily; not required)815# --------------------------------------------------------------------------------------816817def openai_sdk_client_without_retries(**kw: Any) -> Any:818 """Return an `openai.OpenAI` with SDK retries disabled so THIS client owns the retry budget819 (nested retry loops multiply requests — OpenAI rate-limit guide). Requires `pip install openai`."""820 from openai import OpenAI # type: ignore821 return OpenAI(max_retries=0, **kw)822823824def anthropic_sdk_client_without_retries(**kw: Any) -> Any:825 from anthropic import Anthropic # type: ignore826 return Anthropic(max_retries=0, **kw)827828829def sdk_error_is_retryable(exc: Exception, provider: Optional[str] = None) -> bool:830 """Classify an SDK exception via status_code/headers/body.831 - openai / anthropic SDKs share class names (`status_code`, `response`, `body`); the module name tells them apart.832 - The `openai` SDK pointed at https://api.x.ai/v1 raises the same classes → pass provider="xai" explicitly.833 - google-genai raises `google.genai.errors.APIError` (`code`, `response`, `details` = the google.rpc body) → "gemini".834 - xai-sdk (gRPC) raises `grpc.RpcError` → map `e.code()` (RESOURCE_EXHAUSTED/UNAVAILABLE/DEADLINE_EXCEEDED/INTERNAL retryable)."""835 mod = type(exc).__module__836 if provider is None:837 provider = ("anthropic" if mod.startswith("anthropic") else "gemini" if mod.startswith("google") else "openai")838 code_fn = getattr(exc, "code", None)839 if mod.startswith("grpc") and callable(code_fn): # xai-sdk840 return getattr(code_fn(), "name", str(code_fn())) in ("RESOURCE_EXHAUSTED", "UNAVAILABLE", "DEADLINE_EXCEEDED", "INTERNAL", "UNKNOWN")841 status = getattr(exc, "status_code", None)842 if status is None and isinstance(code_fn, int): # google-genai APIError.code843 status = code_fn844 resp = getattr(exc, "response", None)845 headers = dict(getattr(resp, "headers", {}) or {})846 body = getattr(exc, "body", None)847 if body is None and provider == "gemini":848 body = getattr(exc, "details", None)849 if isinstance(body, dict) and "error" not in body:850 body = {"error": body}851 fake = Response(status or 0, headers, json.dumps(body).encode() if isinstance(body, (dict, list, str)) else b"")852 if status is None: # APIConnectionError / APITimeoutError853 return True854 return classify(provider, fake).retryable855