Merge pull request #403 from g0ldyy/development

Development
This commit is contained in:
Goldy
2025-12-21 21:51:01 +01:00
committed by GitHub
9 changed files with 234 additions and 156 deletions
+3 -2
View File
@@ -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)
# ============================== #
+1 -1
View File
@@ -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"]
+4
View File
@@ -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",
+1
View File
@@ -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
+4 -1
View File
@@ -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
+1 -1
View File
@@ -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:
+13
View File
@@ -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()
+49 -6
View File
@@ -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}"
+158 -145
View File
@@ -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: