diff --git a/.env-sample b/.env-sample index 1848473..dc6611d 100644 --- a/.env-sample +++ b/.env-sample @@ -92,13 +92,14 @@ PROWLARR_INDEXERS=[] # Leave empty to automatically use all configured/healthy i # Shared Settings INDEXER_MANAGER_TIMEOUT=30 # Max time to get search results (seconds) - Shared by both +INDEXER_MANAGER_WAIT_TIMEOUT=30 # Max time to wait for the indexer manager to initialize (seconds) INDEXER_MANAGER_UPDATE_INTERVAL=900 # Time in seconds between indexer updates (default: 900s / 15m) # ============================== # # Torrent Settings # # ============================== # -GET_TORRENT_TIMEOUT=5 # Max time to obtain torrent info hash (seconds) -DOWNLOAD_TORRENT_FILES=False # Enable torrent file retrieval (instead of magnet link only) +GET_TORRENT_TIMEOUT=5 # Max time to download .torrent file (seconds) +DOWNLOAD_TORRENT_FILES=False # Enable torrent file retrieval from magnet link MAGNET_RESOLVE_TIMEOUT=60 # Max time to resolve a magnet link (seconds) # ============================== # diff --git a/comet/api/endpoints/stream.py b/comet/api/endpoints/stream.py index 2a4e880..382b3b8 100644 --- a/comet/api/endpoints/stream.py +++ b/comet/api/endpoints/stream.py @@ -466,7 +466,7 @@ async def stream( if torrent["fileIndex"] is not None: the_stream["fileIdx"] = torrent["fileIndex"] - if len(torrent["sources"]) == 0: + if not torrent["sources"]: the_stream["sources"] = trackers else: the_stream["sources"] = torrent["sources"] diff --git a/comet/core/logger.py b/comet/core/logger.py index a434750..85518bb 100644 --- a/comet/core/logger.py +++ b/comet/core/logger.py @@ -221,6 +221,10 @@ def log_startup_info(settings): ) logger.log("COMET", f"Indexer Manager Timeout: {settings.INDEXER_MANAGER_TIMEOUT}s") + logger.log( + "COMET", + f"Indexer Manager Wait Timeout: {settings.INDEXER_MANAGER_WAIT_TIMEOUT}s", + ) logger.log( "COMET", f"Indexer Manager Update Interval: {settings.INDEXER_MANAGER_UPDATE_INTERVAL}s", diff --git a/comet/core/models.py b/comet/core/models.py index 461d2b2..0450931 100644 --- a/comet/core/models.py +++ b/comet/core/models.py @@ -56,6 +56,7 @@ class AppSettings(BaseSettings): INDEXER_MANAGER_TIMEOUT: Optional[int] = 30 INDEXER_MANAGER_INDEXERS: List[str] = [] INDEXER_MANAGER_UPDATE_INTERVAL: Optional[int] = 900 + INDEXER_MANAGER_WAIT_TIMEOUT: Optional[int] = 30 SCRAPE_JACKETT: Union[bool, str] = False JACKETT_URL: Optional[str] = "http://127.0.0.1:9117" JACKETT_API_KEY: Optional[str] = None diff --git a/comet/metadata/imdb.py b/comet/metadata/imdb.py index 57cbfde..fad4167 100644 --- a/comet/metadata/imdb.py +++ b/comet/metadata/imdb.py @@ -16,5 +16,8 @@ async def get_imdb_metadata(session: aiohttp.ClientSession, id: str): year_end = int(element["yr"].split("-")[1]) if "yr" in element else None return title, year, year_end except Exception as e: - logger.warning(f"Exception while getting IMDB metadata for {id}: {e}") + additional_info = "" + if metadata: + additional_info = f"- API Response: {metadata}" + logger.warning(f"Exception while getting IMDB metadata for {id}: {e}{additional_info}") return None, None, None diff --git a/comet/scrapers/aiostreams.py b/comet/scrapers/aiostreams.py index 9ac48e3..79736f8 100644 --- a/comet/scrapers/aiostreams.py +++ b/comet/scrapers/aiostreams.py @@ -48,7 +48,7 @@ class AiostreamsScraper(BaseScraper): "seeders": torrent.get("seeders", None), "size": torrent["size"], "tracker": tracker, - "sources": torrent.get("sources", []), + "sources": torrent.get("sources") or [], } ) except Exception as e: diff --git a/comet/scrapers/jackett.py b/comet/scrapers/jackett.py index e95f370..059f1cf 100644 --- a/comet/scrapers/jackett.py +++ b/comet/scrapers/jackett.py @@ -8,6 +8,7 @@ from comet.core.logger import logger from comet.core.models import settings from comet.scrapers.base import BaseScraper from comet.scrapers.models import ScrapeRequest, ScrapeResult +from comet.services.indexer_manager import indexer_manager from comet.services.torrent_manager import (add_torrent_queue, download_torrent, extract_torrent_metadata, @@ -100,6 +101,18 @@ class JackettScraper(BaseScraper): return [] async def scrape(self, request: ScrapeRequest): + if not settings.JACKETT_INDEXERS: + try: + await asyncio.wait_for( + indexer_manager.jackett_initialized.wait(), + timeout=settings.INDEXER_MANAGER_WAIT_TIMEOUT, + ) + except asyncio.TimeoutError: + pass + + if not settings.JACKETT_INDEXERS: + logger.warning("No Jackett indexers available, skipping scrape.") + return [] torrents: List[ScrapeResult] = [] seen: Set[str] = set() diff --git a/comet/scrapers/prowlarr.py b/comet/scrapers/prowlarr.py index c82c5d1..0104bc5 100644 --- a/comet/scrapers/prowlarr.py +++ b/comet/scrapers/prowlarr.py @@ -8,6 +8,7 @@ from comet.core.logger import logger from comet.core.models import settings from comet.scrapers.base import BaseScraper from comet.scrapers.models import ScrapeRequest, ScrapeResult +from comet.services.indexer_manager import indexer_manager from comet.services.torrent_manager import (add_torrent_queue, download_torrent, extract_torrent_metadata, @@ -84,6 +85,19 @@ class ProwlarrScraper(BaseScraper): return torrents async def scrape(self, request: ScrapeRequest): + if not settings.PROWLARR_INDEXERS: + try: + await asyncio.wait_for( + indexer_manager.prowlarr_initialized.wait(), + timeout=settings.INDEXER_MANAGER_WAIT_TIMEOUT, + ) + except asyncio.TimeoutError: + pass + + if not settings.PROWLARR_INDEXERS: + logger.warning("No Prowlarr indexers available, skipping scrape.") + return [] + torrents: List[ScrapeResult] = [] seen: Set[str] = set() @@ -105,10 +119,32 @@ class ProwlarrScraper(BaseScraper): ) ) - responses = await asyncio.gather(*tasks) + responses = await asyncio.gather(*tasks, return_exceptions=True) all_results = [] for response in responses: - all_results.extend(await response.json()) + if isinstance(response, Exception): + if isinstance(response, asyncio.TimeoutError): + logger.warning( + f"Timeout while getting torrents for {request.title} with Prowlarr" + ) + else: + logger.warning( + f"Exception while getting torrents for {request.title} with Prowlarr: {response}" + ) + continue + + try: + response_json = await response.json() + if isinstance(response_json, list): + all_results.extend(response_json) + else: + logger.warning( + f"Unexpected Prowlarr response format: {response_json}" + ) + except Exception as e: + logger.warning( + f"Exception while parsing Prowlarr JSON response: {e}" + ) torrent_tasks = [] for result in all_results: @@ -120,10 +156,17 @@ class ProwlarrScraper(BaseScraper): self.process_torrent(result, request.media_only_id, request.season) ) - processed_torrents = await asyncio.gather(*torrent_tasks) - torrents = [ - t for sublist in processed_torrents for t in sublist if t["infoHash"] - ] + processed_torrents = await asyncio.gather( + *torrent_tasks, return_exceptions=True + ) + for sublist in processed_torrents: + if isinstance(sublist, Exception): + logger.warning(f"Error processing torrent with Prowlarr: {sublist}") + continue + for t in sublist: + if t["infoHash"]: + torrents.append(t) + except Exception as e: logger.warning( f"Exception while getting torrents for {request.title} with Prowlarr: {e}" diff --git a/comet/services/indexer_manager.py b/comet/services/indexer_manager.py index 3a8e4e7..4cdba1f 100644 --- a/comet/services/indexer_manager.py +++ b/comet/services/indexer_manager.py @@ -16,6 +16,8 @@ class IndexerManager: self.refresh_interval = settings.INDEXER_MANAGER_UPDATE_INTERVAL self.original_jackett_config = settings.JACKETT_INDEXERS.copy() self.original_prowlarr_config = settings.PROWLARR_INDEXERS.copy() + self.jackett_initialized = asyncio.Event() + self.prowlarr_initialized = asyncio.Event() async def get_session(self): if self.session is None or self.session.closed: @@ -23,168 +25,179 @@ class IndexerManager: return self.session async def update_jackett(self): - if ( - not settings.is_any_context_enabled(settings.SCRAPE_JACKETT) - or not settings.JACKETT_URL - or not settings.JACKETT_API_KEY - ): - return - try: - session = await self.get_session() - url = f"{settings.JACKETT_URL}/api/v2.0/indexers/!status:failing/results/torznab/api" - params = { - "apikey": settings.JACKETT_API_KEY, - "t": "indexers", - "configured": "true", - } - async with session.get( - url, params=params, timeout=INDEXER_TIMEOUT - ) as response: - if response.status != 200: + if ( + not settings.is_any_context_enabled(settings.SCRAPE_JACKETT) + or not settings.JACKETT_URL + or not settings.JACKETT_API_KEY + ): + return + + try: + session = await self.get_session() + url = f"{settings.JACKETT_URL}/api/v2.0/indexers/!status:failing/results/torznab/api" + params = { + "apikey": settings.JACKETT_API_KEY, + "t": "indexers", + "configured": "true", + } + async with session.get( + url, params=params, timeout=INDEXER_TIMEOUT + ) as response: + if response.status != 200: + logger.warning( + f"Failed to fetch Jackett indexers: {response.status}" + ) + return + + content = await response.text() + root = ET.fromstring(content) + active_ids = [] + + for indexer in root.findall("indexer"): + indexer_id = indexer.get("id") + if not indexer_id: + continue + active_ids.append(indexer_id) + + # Filter if original config exists + if self.original_jackett_config: + config_set = {x.lower() for x in self.original_jackett_config} + filtered_ids = [] + for indexer in root.findall("indexer"): + pid = indexer.get("id") + title = indexer.find("title") + name = title.text if title is not None else "" + if pid.lower() in config_set or name.lower() in config_set: + filtered_ids.append(pid) + active_ids = filtered_ids + + if sorted(settings.JACKETT_INDEXERS) != sorted(active_ids): + settings.JACKETT_INDEXERS = active_ids + logger.log( + "COMET", + f"Updated Jackett indexers ({len(active_ids)}): {', '.join(active_ids)}", + ) + + except Exception as e: + logger.warning(f"Error updating Jackett indexers: {e}") + + finally: + self.jackett_initialized.set() + + async def update_prowlarr(self): + try: + if ( + not settings.is_any_context_enabled(settings.SCRAPE_PROWLARR) + or not settings.PROWLARR_URL + or not settings.PROWLARR_API_KEY + ): + return + + try: + session = await self.get_session() + headers = {"X-Api-Key": settings.PROWLARR_API_KEY} + + indexers_task = session.get( + f"{settings.PROWLARR_URL}/api/v1/indexer", + headers=headers, + timeout=INDEXER_TIMEOUT, + ) + statuses_task = session.get( + f"{settings.PROWLARR_URL}/api/v1/indexerstatus", + headers=headers, + timeout=INDEXER_TIMEOUT, + ) + + responses = await asyncio.gather( + indexers_task, statuses_task, return_exceptions=True + ) + + if any(isinstance(r, Exception) for r in responses): + logger.warning("Failed to fetch Prowlarr indexers or statuses") + return + + resp_idx, resp_stat = responses + + if resp_idx.status != 200 or resp_stat.status != 200: logger.warning( - f"Failed to fetch Jackett indexers: {response.status}" + f"Prowlarr error: Indexers {resp_idx.status}, Status {resp_stat.status}" ) return - content = await response.text() - root = ET.fromstring(content) + indexers = await resp_idx.json() + statuses = await resp_stat.json() + + status_map = {s["indexerId"]: s for s in statuses} active_ids = [] + current_time = datetime.now(timezone.utc) - for indexer in root.findall("indexer"): - indexer_id = indexer.get("id") - if not indexer_id: + for indexer in indexers: + if not indexer.get("enable"): continue - active_ids.append(indexer_id) - # Filter if original config exists - if self.original_jackett_config: - config_set = {x.lower() for x in self.original_jackett_config} + if indexer.get("protocol") != "torrent": + continue + + idx_id = indexer.get("id") + + # Check health + status = status_map.get(idx_id, {}) + disabled_till = status.get("disabledTill") + if disabled_till: + try: + dt = datetime.fromisoformat( + disabled_till.replace("Z", "+00:00") + ) + if dt > current_time: + continue + except ValueError: + pass # Ignore parsing error, assume enabled + + active_ids.append(str(idx_id)) + + # Apply original config filter configuration + if self.original_prowlarr_config: + config_set = {x.lower() for x in self.original_prowlarr_config} filtered_ids = [] - for indexer in root.findall("indexer"): - pid = indexer.get("id") - title = indexer.find("title") - name = title.text if title is not None else "" - if pid.lower() in config_set or name.lower() in config_set: - filtered_ids.append(pid) + for indexer in indexers: + idx_id_str = str(indexer.get("id")) + if idx_id_str not in active_ids: + continue + + name = indexer.get("name", "").lower() + def_name = indexer.get("definitionName", "").lower() + + if ( + name in config_set + or def_name in config_set + or idx_id_str in config_set # support ID in config too + ): + filtered_ids.append(idx_id_str) active_ids = filtered_ids - if sorted(settings.JACKETT_INDEXERS) != sorted(active_ids): - settings.JACKETT_INDEXERS = active_ids + if sorted(settings.PROWLARR_INDEXERS) != sorted(active_ids): + settings.PROWLARR_INDEXERS = active_ids + + # Map IDs to names for logging + id_to_name = { + str(i.get("id")): i.get("name", str(i.get("id"))) + for i in indexers + } + active_names = [ + id_to_name.get(idx_id, idx_id) for idx_id in active_ids + ] + logger.log( "COMET", - f"Updated Jackett indexers ({len(active_ids)}): {', '.join(active_ids)}", + f"Updated Prowlarr indexers ({len(active_ids)}): {', '.join(active_names)}", ) - except Exception as e: - logger.warning(f"Error updating Jackett indexers: {e}") + except Exception as e: + logger.warning(f"Error updating Prowlarr indexers: {e}") - async def update_prowlarr(self): - if ( - not settings.is_any_context_enabled(settings.SCRAPE_PROWLARR) - or not settings.PROWLARR_URL - or not settings.PROWLARR_API_KEY - ): - return - - try: - session = await self.get_session() - headers = {"X-Api-Key": settings.PROWLARR_API_KEY} - - indexers_task = session.get( - f"{settings.PROWLARR_URL}/api/v1/indexer", - headers=headers, - timeout=INDEXER_TIMEOUT, - ) - statuses_task = session.get( - f"{settings.PROWLARR_URL}/api/v1/indexerstatus", - headers=headers, - timeout=INDEXER_TIMEOUT, - ) - - responses = await asyncio.gather( - indexers_task, statuses_task, return_exceptions=True - ) - - if any(isinstance(r, Exception) for r in responses): - logger.warning("Failed to fetch Prowlarr indexers or statuses") - return - - resp_idx, resp_stat = responses - - if resp_idx.status != 200 or resp_stat.status != 200: - logger.warning( - f"Prowlarr error: Indexers {resp_idx.status}, Status {resp_stat.status}" - ) - return - - indexers = await resp_idx.json() - statuses = await resp_stat.json() - - status_map = {s["indexerId"]: s for s in statuses} - active_ids = [] - current_time = datetime.now(timezone.utc) - - for indexer in indexers: - if not indexer.get("enable"): - continue - - if indexer.get("protocol") != "torrent": - continue - - idx_id = indexer.get("id") - - # Check health - status = status_map.get(idx_id, {}) - disabled_till = status.get("disabledTill") - if disabled_till: - try: - dt = datetime.fromisoformat( - disabled_till.replace("Z", "+00:00") - ) - if dt > current_time: - continue - except ValueError: - pass # Ignore parsing error, assume enabled - - active_ids.append(str(idx_id)) - - # Apply original config filter configuration - if self.original_prowlarr_config: - config_set = {x.lower() for x in self.original_prowlarr_config} - filtered_ids = [] - for indexer in indexers: - idx_id_str = str(indexer.get("id")) - if idx_id_str not in active_ids: - continue - - name = indexer.get("name", "").lower() - def_name = indexer.get("definitionName", "").lower() - - if ( - name in config_set - or def_name in config_set - or idx_id_str in config_set # support ID in config too - ): - filtered_ids.append(idx_id_str) - active_ids = filtered_ids - - if sorted(settings.PROWLARR_INDEXERS) != sorted(active_ids): - settings.PROWLARR_INDEXERS = active_ids - - # Map IDs to names for logging - id_to_name = { - str(i.get("id")): i.get("name", str(i.get("id"))) for i in indexers - } - active_names = [id_to_name.get(idx_id, idx_id) for idx_id in active_ids] - - logger.log( - "COMET", - f"Updated Prowlarr indexers ({len(active_ids)}): {', '.join(active_names)}", - ) - - except Exception as e: - logger.warning(f"Error updating Prowlarr indexers: {e}") + finally: + self.prowlarr_initialized.set() async def run(self): while True: