SPB Git forge

spb/doc-api

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