diff --git a/.env-sample b/.env-sample index ee85daa..dc6611d 100644 --- a/.env-sample +++ b/.env-sample @@ -92,6 +92,7 @@ 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) # ============================== # 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/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 5db660f..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() @@ -120,7 +134,13 @@ class ProwlarrScraper(BaseScraper): continue try: - all_results.extend(await response.json()) + 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}" 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: