diff --git a/.env-sample b/.env-sample index a28a6d6..db979a4 100644 --- a/.env-sample +++ b/.env-sample @@ -103,6 +103,9 @@ PROXY_ETHOS=always # always: use proxy for everything; on_failure: retry with pr # Scraper specific proxies can be set using [SCRAPERNAME]_PROXY_URL # Example: TORRENTIO_PROXY_URL=http://... +RATELIMIT_MAX_RETRIES=3 # Maximum number of retries for 429 Too Many Requests errors. Set to 0 to disable retries. +RATELIMIT_RETRY_BASE_DELAY=1.0 # Base delay in seconds for exponential backoff (e.g., 1.0 -> 1s, 2s, 4s, 8s...) + # ============================== # # Jackett & Prowlarr Settings # # ============================== # diff --git a/comet/core/logger.py b/comet/core/logger.py index 0a6cbee..5aa55e0 100644 --- a/comet/core/logger.py +++ b/comet/core/logger.py @@ -224,6 +224,10 @@ def log_startup_info(settings): "COMET", f"Global Proxy: {settings.GLOBAL_PROXY_URL} - Ethos: {settings.PROXY_ETHOS}", ) + logger.log( + "COMET", + f"Rate Limit Manager: Max Retries={settings.RATELIMIT_MAX_RETRIES} - Base Delay={settings.RATELIMIT_RETRY_BASE_DELAY}s", + ) jackett_info = "" if settings.is_any_context_enabled(settings.SCRAPE_JACKETT): diff --git a/comet/core/models.py b/comet/core/models.py index b7b5a29..3f2d765 100644 --- a/comet/core/models.py +++ b/comet/core/models.py @@ -131,6 +131,8 @@ class AppSettings(BaseSettings): TMDB_READ_ACCESS_TOKEN: Optional[str] = None GLOBAL_PROXY_URL: Optional[str] = None PROXY_ETHOS: Optional[str] = "always" + RATELIMIT_MAX_RETRIES: Optional[int] = 3 + RATELIMIT_RETRY_BASE_DELAY: Optional[float] = 1.0 @field_validator("INDEXER_MANAGER_TYPE") def set_indexer_manager_type(cls, v, values): diff --git a/comet/utils/network_manager.py b/comet/utils/network_manager.py index a35aff2..176d551 100644 --- a/comet/utils/network_manager.py +++ b/comet/utils/network_manager.py @@ -1,3 +1,4 @@ +import asyncio import os import socket from typing import Dict, Optional @@ -145,24 +146,58 @@ class _RequestContextManager: raise e async def _attempt_request(self, use_proxy, proxy): - if self.wrapper.impersonate: - # Use curl_cffi - curl_proxy = self.wrapper._resolved_proxy_url if proxy else None - session = await self.wrapper._get_curl_session() - raw_response = await session.request( - self.method, self.url, proxy=curl_proxy, **self.kwargs - ) - self.response = ResponseWrapper(raw_response, "curl") - return self.response - else: - # Use aiohttp - session = await self.wrapper._get_aiohttp_session() - self.aiohttp_cm = session.request( - self.method, self.url, proxy=proxy, **self.kwargs - ) - raw_response = await self.aiohttp_cm.__aenter__() - self.response = ResponseWrapper(raw_response, "aiohttp") - return self.response + max_retries = max(0, settings.RATELIMIT_MAX_RETRIES) + base_delay = settings.RATELIMIT_RETRY_BASE_DELAY + for attempt in range(max_retries + 1): + if self.wrapper.impersonate: + # Use curl_cffi + curl_proxy = self.wrapper._resolved_proxy_url if proxy else None + session = await self.wrapper._get_curl_session() + raw_response = await session.request( + self.method, self.url, proxy=curl_proxy, **self.kwargs + ) + self.response = ResponseWrapper(raw_response, "curl") + else: + # Use aiohttp + session = await self.wrapper._get_aiohttp_session() + self.aiohttp_cm = session.request( + self.method, self.url, proxy=proxy, **self.kwargs + ) + raw_response = await self.aiohttp_cm.__aenter__() + self.response = ResponseWrapper(raw_response, "aiohttp") + + if self.response.status != 429: + return self.response + + # Handle 429 Too Many Requests + if attempt < max_retries: + retry_after = self.response.headers.get("Retry-After") + try: + delay = float(retry_after) if retry_after else None + except (ValueError, TypeError): + delay = None + + if delay is None: + delay = base_delay * (2**attempt) + + # Enforce a minimum delay if the server indicates 0 or very small retry-after + delay = max(delay, base_delay) + + logger.warning( + f"[{self.wrapper.scraper_name}] Received 429 Too Many Requests. Retrying in {delay}s... (Attempt {attempt + 1}/{max_retries})" + ) + + # Cleanup aiohttp context manager for the failed attempt + if not self.wrapper.impersonate and self.aiohttp_cm: + await self.aiohttp_cm.__aexit__(None, None, None) + self.aiohttp_cm = None + + await asyncio.sleep(delay) + else: + logger.error( + f"[{self.wrapper.scraper_name}] Max retries ({max_retries}) exceeded for 429 Too Many Requests." + ) + return self.response async def __aexit__(self, exc_type, exc, tb): if self.aiohttp_cm: