Spaces:
Running
Running
| from __future__ import annotations | |
| import asyncio | |
| from email.utils import parsedate_to_datetime | |
| import logging | |
| import random | |
| from datetime import datetime, timezone | |
| from typing import Any, Awaitable, Callable | |
| import httpx | |
| from app.config import MODEL_VERSION | |
| logger = logging.getLogger(__name__) | |
| class ProviderError(RuntimeError): | |
| def __init__(self, message: str, status_code: int | None = None): | |
| super().__init__(message) | |
| self.status_code = status_code | |
| class ResilientHTTP: | |
| def __init__(self, timeout: float = 20.0, retries: int = 3): | |
| self.timeout = timeout | |
| self.retries = max(1, retries) | |
| self._client: httpx.AsyncClient | None = None | |
| def _get_client(self) -> httpx.AsyncClient: | |
| if self._client is None: | |
| timeout = httpx.Timeout( | |
| timeout=self.timeout, | |
| connect=min(self.timeout, 10.0), | |
| read=self.timeout, | |
| write=min(self.timeout, 10.0), | |
| pool=min(self.timeout, 10.0), | |
| ) | |
| self._client = httpx.AsyncClient( | |
| timeout=timeout, | |
| follow_redirects=True, | |
| limits=httpx.Limits(max_connections=8, max_keepalive_connections=4), | |
| headers={"User-Agent": f"SafeBetAI/{MODEL_VERSION}"}, | |
| ) | |
| return self._client | |
| async def aclose(self) -> None: | |
| if self._client is not None: | |
| await self._client.aclose() | |
| self._client = None | |
| def _retry_delay(response: httpx.Response, attempt: int) -> float: | |
| retry_after = response.headers.get("retry-after") | |
| if retry_after: | |
| try: | |
| return min(max(float(retry_after), 0.0), 65.0) | |
| except ValueError: | |
| try: | |
| when = parsedate_to_datetime(retry_after) | |
| if when.tzinfo is None: | |
| when = when.replace(tzinfo=timezone.utc) | |
| seconds = (when - datetime.now(timezone.utc)).total_seconds() | |
| return min(max(seconds, 0.0), 65.0) | |
| except Exception: | |
| pass | |
| return min(1.2 * (2 ** attempt) + random.uniform(0.05, 0.55), 12.0) | |
| def _safe_error_detail( | |
| response: httpx.Response, | |
| params: dict[str, Any] | None, | |
| headers: dict[str, str] | None, | |
| ) -> str: | |
| detail = response.text[:500].replace("\n", " ") | |
| for mapping in (params or {}, headers or {}): | |
| for key, value in mapping.items(): | |
| key_lower = str(key).lower() | |
| if not any(marker in key_lower for marker in ("key", "token", "auth", "secret")): | |
| continue | |
| secret = str(value) | |
| if len(secret) >= 4: | |
| detail = detail.replace(secret, "[redacted]") | |
| return detail | |
| async def get_json( | |
| self, | |
| url: str, | |
| *, | |
| params: dict[str, Any] | None = None, | |
| headers: dict[str, str] | None = None, | |
| allow_status: set[int] | None = None, | |
| before_attempt: Callable[[], Awaitable[None]] | None = None, | |
| ) -> tuple[Any, httpx.Headers]: | |
| allow_status = allow_status or set() | |
| last_exc: Exception | None = None | |
| client = self._get_client() | |
| for attempt in range(self.retries): | |
| try: | |
| if before_attempt is not None: | |
| await before_attempt() | |
| response = await client.get(url, params=params, headers=headers) | |
| if response.status_code in allow_status: | |
| return None, response.headers | |
| if response.status_code in (408, 425, 429, 500, 502, 503, 504): | |
| if attempt < self.retries - 1: | |
| delay = self._retry_delay(response, attempt) | |
| logger.warning( | |
| "HTTP %s em %s; retry %d/%d em %.1fs", | |
| response.status_code, | |
| url, | |
| attempt + 1, | |
| self.retries - 1, | |
| delay, | |
| ) | |
| await asyncio.sleep(delay) | |
| continue | |
| if response.status_code >= 400: | |
| detail = self._safe_error_detail(response, params, headers) | |
| raise ProviderError( | |
| f"HTTP {response.status_code} em {url}: {detail}", | |
| response.status_code, | |
| ) | |
| try: | |
| return response.json(), response.headers | |
| except ValueError as exc: | |
| raise ProviderError( | |
| f"JSON inválido recebido de {url}: {exc}", | |
| response.status_code, | |
| ) from exc | |
| except ProviderError: | |
| raise | |
| except (httpx.TimeoutException, httpx.TransportError) as exc: | |
| last_exc = exc | |
| if attempt < self.retries - 1: | |
| delay = min(1.0 * (2 ** attempt) + random.uniform(0.05, 0.55), 8.0) | |
| logger.warning( | |
| "Falha de rede em %s; retry %d/%d em %.1fs: %s", | |
| url, | |
| attempt + 1, | |
| self.retries - 1, | |
| delay, | |
| type(exc).__name__, | |
| ) | |
| await asyncio.sleep(delay) | |
| continue | |
| break | |
| raise ProviderError(f"Falha de rede em {url}: {last_exc}") | |