Merge pull request #400 from g0ldyy/development

feat: dev to main
This commit is contained in:
Goldy
2025-12-21 18:42:48 +01:00
committed by GitHub
7 changed files with 297 additions and 112 deletions
+5 -3
View File
@@ -82,16 +82,17 @@ BYPASS_PROXY_URL=http://warp:1080 # To bypass scraper IP blacklists
SCRAPE_JACKETT=False # Context mode: live, background, both, false
JACKETT_URL=http://127.0.0.1:9117
JACKETT_API_KEY=XXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX
JACKETT_INDEXERS='["EXAMPLE1_CHANGETHIS", "EXAMPLE2_CHANGETHIS"]' # Get names from https://github.com/Jackett/Jackett/tree/master/src/Jackett.Common/Definitions
JACKETT_INDEXERS=[] # Leave empty to automatically use all configured/healthy indexers. Or specify a list of indexer IDs to use (e.g. '["oxtorrent", "torrent9"]').
# Prowlarr Configuration
SCRAPE_PROWLARR=False # Context mode: live, background, both, false
PROWLARR_URL=http://127.0.0.1:9696
PROWLARR_API_KEY=XXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX
PROWLARR_INDEXERS='["EXAMPLE1_CHANGETHIS", "EXAMPLE2_CHANGETHIS"]'
PROWLARR_INDEXERS=[] # Leave empty to automatically use all configured/healthy indexers. Or specify a list of indexer IDs.
# Shared Settings
INDEXER_MANAGER_TIMEOUT=60 # Max time to get search results (seconds) - Shared by both
INDEXER_MANAGER_TIMEOUT=30 # Max time to get search results (seconds) - Shared by both
INDEXER_MANAGER_UPDATE_INTERVAL=900 # Time in seconds between indexer updates (default: 900s / 15m)
# ============================== #
# Torrent Settings #
@@ -189,6 +190,7 @@ TORBOX_API_KEY=XXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX
SCRAPE_YGGTORRENT=False
YGGTORRENT_USERNAME=username
YGGTORRENT_PASSWORD=password
YGGTORRENT_PASSKEY=passkey
YGGTORRENT_MAX_CONCURRENT_PAGES=5
# ============================== #
+10
View File
@@ -19,6 +19,7 @@ from comet.core.logger import logger
from comet.core.models import settings
from comet.services.anime import anime_mapper
from comet.services.bandwidth import bandwidth_monitor
from comet.services.indexer_manager import indexer_manager
from comet.services.torrent_manager import (add_torrent_queue,
torrent_update_queue)
from comet.services.trackers import download_best_trackers
@@ -66,9 +67,18 @@ async def lifespan(app: FastAPI):
if settings.BACKGROUND_SCRAPER_ENABLED:
background_scraper_task = asyncio.create_task(background_scraper.start())
# Start indexer manager
indexer_manager_task = asyncio.create_task(indexer_manager.run())
try:
yield
finally:
indexer_manager_task.cancel()
try:
await indexer_manager_task
except asyncio.CancelledError:
pass
if background_scraper_task:
await background_scraper.stop()
background_scraper_task.cancel()
+17 -3
View File
@@ -196,7 +196,12 @@ def log_startup_info(settings):
jackett_info = ""
if settings.is_any_context_enabled(settings.SCRAPE_JACKETT):
jackett_info = f" - {settings.JACKETT_URL} - Indexers: {', '.join(settings.JACKETT_INDEXERS)}"
indexers = (
", ".join(settings.JACKETT_INDEXERS)
if settings.JACKETT_INDEXERS
else "All Configured/Healthy"
)
jackett_info = f" - {settings.JACKETT_URL} - Indexers: {indexers}"
logger.log(
"COMET",
f"Jackett Scraper: {settings.format_scraper_mode(settings.SCRAPE_JACKETT)}{jackett_info}",
@@ -204,13 +209,22 @@ def log_startup_info(settings):
prowlarr_info = ""
if settings.is_any_context_enabled(settings.SCRAPE_PROWLARR):
prowlarr_info = f" - {settings.PROWLARR_URL} - Indexers: {', '.join(settings.PROWLARR_INDEXERS)}"
indexers = (
", ".join(settings.PROWLARR_INDEXERS)
if settings.PROWLARR_INDEXERS
else "All Configured/Healthy"
)
prowlarr_info = f" - {settings.PROWLARR_URL} - Indexers: {indexers}"
logger.log(
"COMET",
f"Prowlarr Scraper: {settings.format_scraper_mode(settings.SCRAPE_PROWLARR)}{prowlarr_info}",
)
logger.log("COMET", f"Indexer Manager Timeout: {settings.INDEXER_MANAGER_TIMEOUT}s")
logger.log(
"COMET",
f"Indexer Manager Update Interval: {settings.INDEXER_MANAGER_UPDATE_INTERVAL}s",
)
logger.log("COMET", f"Get Torrent Timeout: {settings.GET_TORRENT_TIMEOUT}s")
logger.log("COMET", f"Magnet Resolve Timeout: {settings.MAGNET_RESOLVE_TIMEOUT}s")
logger.log(
@@ -328,7 +342,7 @@ def log_startup_info(settings):
)
yggtorrent_info = (
f" - Username: {settings.YGGTORRENT_USERNAME} - Password: {settings.YGGTORRENT_PASSWORD}"
f" - Username: {settings.YGGTORRENT_USERNAME} - Password: {settings.YGGTORRENT_PASSWORD} - Passkey: {settings.YGGTORRENT_PASSKEY}"
if settings.is_any_context_enabled(settings.SCRAPE_YGGTORRENT)
else ""
)
+2
View File
@@ -55,6 +55,7 @@ class AppSettings(BaseSettings):
INDEXER_MANAGER_MODE: Union[bool, str] = "both"
INDEXER_MANAGER_TIMEOUT: Optional[int] = 30
INDEXER_MANAGER_INDEXERS: List[str] = []
INDEXER_MANAGER_UPDATE_INTERVAL: Optional[int] = 900
SCRAPE_JACKETT: Union[bool, str] = False
JACKETT_URL: Optional[str] = "http://127.0.0.1:9117"
JACKETT_API_KEY: Optional[str] = None
@@ -102,6 +103,7 @@ class AppSettings(BaseSettings):
SCRAPE_YGGTORRENT: Union[bool, str] = False
YGGTORRENT_USERNAME: Optional[str] = None
YGGTORRENT_PASSWORD: Optional[str] = None
YGGTORRENT_PASSKEY: Optional[str] = None
YGGTORRENT_MAX_CONCURRENT_PAGES: Optional[int] = 5
CUSTOM_HEADER_HTML: Optional[str] = None
PROXY_DEBRID_STREAM: Optional[bool] = False
+1 -18
View File
@@ -95,28 +95,11 @@ class ProwlarrScraper(BaseScraper):
)
try:
indexers = [indexer.lower() for indexer in settings.PROWLARR_INDEXERS]
get_indexers = await self.session.get(
f"{self.url}/api/v1/indexer",
headers={"X-Api-Key": settings.PROWLARR_API_KEY},
timeout=INDEXER_TIMEOUT,
)
get_indexers = await get_indexers.json()
indexers_id = []
for indexer in get_indexers:
if (
indexer["name"].lower() in indexers
or indexer["definitionName"].lower() in indexers
):
indexers_id.append(indexer["id"])
tasks = []
for query in queries:
tasks.append(
self.session.get(
f"{self.url}/api/v1/search?query={query}&indexerIds={'&indexerIds='.join(str(indexer_id) for indexer_id in indexers_id)}&type=search",
f"{self.url}/api/v1/search?query={query}&indexerIds={'&indexerIds='.join(str(indexer_id) for indexer_id in settings.PROWLARR_INDEXERS)}&type=search",
headers={"X-Api-Key": settings.PROWLARR_API_KEY},
timeout=INDEXER_TIMEOUT,
)
+66 -88
View File
@@ -9,9 +9,11 @@ from comet.core.logger import logger
from comet.core.models import settings
from comet.scrapers.base import BaseScraper
from comet.scrapers.models import ScrapeRequest
from comet.services.torrent_manager import extract_torrent_metadata
from comet.utils.formatting import size_to_bytes
YGG_URL = "https://www.yggtorrent.org"
TRACKER_URL = "http://tracker.p2p-world.net:8080"
LOGIN_PAGE = "/auth/login"
LOGIN_PROCESS_PAGE = "/auth/process_login"
@@ -22,45 +24,18 @@ NAME_PATTERN = re.compile(r'<a[^>]+href="([^"]+)"[^>]*>(.*)', re.DOTALL)
class YGGTorrentScraper(BaseScraper):
_domain = None
_session = None
_lock = asyncio.Lock()
def __init__(self, manager, session: aiohttp.ClientSession):
super().__init__(manager, session)
@classmethod
async def _get_domain(cls):
if cls._domain:
return cls._domain
try:
async with requests.AsyncSession(impersonate="chrome") as session:
response = await session.get("https://ygg.re", allow_redirects=False)
if "location" in response.headers:
location = response.headers["location"]
domain = urlparse(location).netloc
logger.info(f"Resolved YGG domain to: {domain}")
cls._domain = domain
return domain
elif response.status_code == 200:
cls._domain = urlparse(response.url).netloc
return cls._domain
except Exception as e:
logger.error(f"Error resolving YGG domain: {e}")
return None
return None
@classmethod
async def _ensure_session(cls):
async with cls._lock:
if cls._session:
domain = await cls._get_domain()
if not domain:
return False
try:
response = await cls._session.get(f"https://{domain}/")
response = await cls._session.get(f"{YGG_URL}/")
if response.status_code == 200 and "Déconnexion" in response.text:
return True
except Exception:
@@ -70,15 +45,13 @@ class YGGTorrentScraper(BaseScraper):
await cls._session.close()
cls._session = None
domain = await cls._get_domain()
if not domain:
return False
domain = urlparse(YGG_URL).netloc
session = requests.AsyncSession(impersonate="chrome")
session.cookies.set("account_created", "true", domain=domain)
try:
response = await session.get(f"https://{domain}{LOGIN_PAGE}")
response = await session.get(f"{YGG_URL}{LOGIN_PAGE}")
if response.status_code != 200:
logger.error(f"Failed to get login page: {response.status_code}")
await session.close()
@@ -90,7 +63,7 @@ class YGGTorrentScraper(BaseScraper):
}
response = await session.post(
f"https://{domain}{LOGIN_PROCESS_PAGE}", data=payload
f"{YGG_URL}{LOGIN_PROCESS_PAGE}", data=payload
)
if response.status_code != 200:
@@ -98,7 +71,7 @@ class YGGTorrentScraper(BaseScraper):
await session.close()
return False
response = await session.get(f"https://{domain}/")
response = await session.get(f"{YGG_URL}/")
if "Déconnexion" in response.text or "logout" in response.text.lower():
logger.info("Successfully logged in to YGGTorrent.")
cls._session = session
@@ -112,53 +85,63 @@ class YGGTorrentScraper(BaseScraper):
await session.close()
return False
async def _download_torrent(self, url):
async def _process_torrent(self, url, title, seeders, size):
try:
response = await self._session.get(url)
if response.status_code == 200:
return response.content
logger.warning(f"Failed to download torrent file: {response.status_code}")
except Exception as e:
logger.warning(f"Exception downloading torrent file: {e}")
return None
async def _process_torrent(self, domain, torrent_id, title, seeders, size):
download_url = f"https://{domain}/engine/download_torrent?id={torrent_id}"
torrent_content = await self._download_torrent(download_url)
results = []
if torrent_content:
try:
metadata = extract_torrent_metadata(torrent_content)
if metadata:
for file in metadata["files"]:
results.append(
{
"title": file["name"],
"infoHash": metadata["info_hash"].lower(),
"fileIndex": file["index"],
"seeders": seeders,
"size": file["size"],
"tracker": "YGGTorrent",
"sources": metadata["announce_list"],
}
)
else:
logger.warning(
f"Failed to extract metadata for torrent {torrent_id}"
)
except Exception as e:
if response.status_code != 200:
logger.warning(
f"Exception extracting metadata for torrent {torrent_id}: {e}"
f"Failed to fetch torrent page {url}: {response.status_code}"
)
else:
logger.warning(f"Failed to download torrent content for {torrent_id}")
return []
return results
html_content = response.text
async def _scrape_page(self, domain, query, offset, semaphore):
url = f"https://{domain}/engine/search?name={query}&do=search&page={offset}&category=2145"
hash_match = re.search(
r"Hash\s*</td>\s*<td[^>]*>(.*?)</td>",
html_content,
re.IGNORECASE | re.DOTALL,
)
if not hash_match:
logger.warning(f"Could not find hash in page {url}")
return []
info_hash = (
hash_match.group(1).strip().replace("Tester", "").strip().lower()
)
info_hash = re.sub(r"<[^>]+>", "", info_hash).strip()
if not info_hash or len(info_hash) != 40:
logger.warning(f"Invalid hash found: {info_hash} in {url}")
return []
if not settings.YGGTORRENT_PASSKEY:
logger.warning(
"YGGTORRENT_PASSKEY not set, cannot construct source URL."
)
return []
source = f"{TRACKER_URL}/{settings.YGGTORRENT_PASSKEY}/announce"
return [
{
"title": title,
"infoHash": info_hash,
"fileIndex": None,
"seeders": seeders,
"size": size,
"tracker": "YGGTorrent",
"sources": [source],
}
]
except Exception as e:
logger.warning(f"Exception processing torrent {url}: {e}")
return []
async def _scrape_page(self, query, offset, semaphore):
url = f"{YGG_URL}/engine/search?name={query}&do=search&page={offset}&category=2145"
async with semaphore:
try:
@@ -193,10 +176,11 @@ class YGGTorrentScraper(BaseScraper):
name_match = NAME_PATTERN.search(name_html)
href = name_match.group(1)
id_match = re.search(r"/torrent/.*/(\d+)-", href)
torrent_id = id_match.group(1)
title = name_match.group(2)
if href.startswith("/"):
href = f"{YGG_URL}{href}"
title = name_match.group(2).split("</a>")[0].strip()
seeders = int(tds[7])
@@ -205,9 +189,7 @@ class YGGTorrentScraper(BaseScraper):
clean_size = re.sub(r"(\d)([a-zA-Z])", r"\1 \2", clean_size)
size = size_to_bytes(clean_size)
tasks.append(
self._process_torrent(domain, torrent_id, title, seeders, size)
)
tasks.append(self._process_torrent(href, title, seeders, size))
results = await asyncio.gather(*tasks)
return [r for res in results for r in res], total_results
@@ -220,16 +202,12 @@ class YGGTorrentScraper(BaseScraper):
semaphore = asyncio.Semaphore(limit)
# Fetch page 0 first to get total results
results, total_results = await self._scrape_page(
self._domain, request.title, 0, semaphore
)
results, total_results = await self._scrape_page(request.title, 0, semaphore)
if total_results > 50:
tasks = []
for offset in range(50, total_results, 50):
tasks.append(
self._scrape_page(self._domain, request.title, offset, semaphore)
)
tasks.append(self._scrape_page(request.title, offset, semaphore))
if tasks:
pages_results = await asyncio.gather(*tasks)
+196
View File
@@ -0,0 +1,196 @@
import asyncio
import xml.etree.ElementTree as ET
from datetime import datetime, timezone
from typing import Optional
import aiohttp
from comet.core.constants import INDEXER_TIMEOUT
from comet.core.logger import logger
from comet.core.models import settings
class IndexerManager:
def __init__(self):
self.session: Optional[aiohttp.ClientSession] = None
self.refresh_interval = settings.INDEXER_MANAGER_UPDATE_INTERVAL
self.original_jackett_config = settings.JACKETT_INDEXERS.copy()
self.original_prowlarr_config = settings.PROWLARR_INDEXERS.copy()
async def get_session(self):
if self.session is None or self.session.closed:
self.session = aiohttp.ClientSession()
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:
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}")
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}")
async def run(self):
while True:
await self.update_jackett()
await self.update_prowlarr()
await asyncio.sleep(self.refresh_interval)
indexer_manager = IndexerManager()