mirror of
https://github.com/g0ldyy/comet.git
synced 2026-01-12 01:16:12 +01:00
Merge pull request #449 from g0ldyy/feat/better-bitmagnet-scraper
feat: update BitMagnet scraper to use IMDb ID and media type for queries
This commit is contained in:
@@ -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 #
|
||||
# ============================== #
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -55,15 +55,30 @@ class BitmagnetScraper(BaseScraper):
|
||||
continue
|
||||
return torrents
|
||||
|
||||
async def scrape_bitmagnet_page(self, query, offset, limit):
|
||||
async def scrape_bitmagnet_page(
|
||||
self, imdb_id, scrape_type, offset, limit, season=None, episode=None
|
||||
):
|
||||
try:
|
||||
params = {"t": "search", "q": query, "offset": offset, "limit": limit}
|
||||
params = {
|
||||
"t": scrape_type,
|
||||
"imdbid": imdb_id,
|
||||
"offset": offset,
|
||||
"limit": limit,
|
||||
}
|
||||
if season:
|
||||
params["season"] = season
|
||||
if episode:
|
||||
params["ep"] = episode
|
||||
async with self.session.get(
|
||||
f"{self.url}/torznab/api", params=params
|
||||
) as response:
|
||||
data_text = await response.text()
|
||||
if not data_text.strip():
|
||||
return []
|
||||
root = ET.fromstring(data_text)
|
||||
return self.parse_bitmagnet_items(root)
|
||||
except ET.ParseError:
|
||||
return []
|
||||
except Exception as e:
|
||||
logger.warning(f"Error scraping BitMagnet page offset={offset}: {e}")
|
||||
return []
|
||||
@@ -71,7 +86,10 @@ class BitmagnetScraper(BaseScraper):
|
||||
async def scrape(self, request: ScrapeRequest):
|
||||
torrents = []
|
||||
limit = 100
|
||||
query = request.title
|
||||
imdb_id = request.media_only_id
|
||||
scrape_type = "movie" if request.media_type == "movie" else "tvsearch"
|
||||
season = request.season
|
||||
episode = request.episode
|
||||
|
||||
batch_size = settings.BITMAGNET_MAX_CONCURRENT_PAGES
|
||||
offset = 0
|
||||
@@ -85,7 +103,11 @@ class BitmagnetScraper(BaseScraper):
|
||||
current_offset = offset + (i * limit)
|
||||
if current_offset >= settings.BITMAGNET_MAX_OFFSET:
|
||||
break
|
||||
tasks.append(self.scrape_bitmagnet_page(query, current_offset, limit))
|
||||
tasks.append(
|
||||
self.scrape_bitmagnet_page(
|
||||
imdb_id, scrape_type, current_offset, limit, season, episode
|
||||
)
|
||||
)
|
||||
|
||||
if not tasks:
|
||||
break
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user