diff --git a/.env-sample b/.env-sample index fcfac1e..97b2af3 100644 --- a/.env-sample +++ b/.env-sample @@ -221,6 +221,8 @@ TORRENT_DISABLED_STREAM_URL=https://comet.fast # Optional URL included in the pl # Content Filtering # # ============================== # REMOVE_ADULT_CONTENT=False +DIGITAL_RELEASE_FILTER=False # Filter unreleased content +TMDB_READ_ACCESS_TOKEN= # Optional: Provide your own TMDB Read Access Token to avoid using the shared default key # ============================== # # UI Customization # diff --git a/comet/api/app.py b/comet/api/app.py index 1f20dc5..1e885c6 100644 --- a/comet/api/app.py +++ b/comet/api/app.py @@ -16,6 +16,7 @@ from comet.background_scraper.worker import background_scraper from comet.core.database import (cleanup_expired_locks, cleanup_expired_sessions, setup_database, teardown_database) +from comet.core.execution import setup_executor, shutdown_executor from comet.core.logger import logger from comet.core.models import settings from comet.services.anime import anime_mapper @@ -49,6 +50,7 @@ async def lifespan(app: FastAPI): # loop.set_debug(True) await setup_database() + setup_executor() await download_best_trackers() # Load anime ID mapping for enhanced metadata and anime detection @@ -107,6 +109,7 @@ async def lifespan(app: FastAPI): await torrent_update_queue.stop() await teardown_database() + shutdown_executor() tags_metadata = [ diff --git a/comet/api/endpoints/stream.py b/comet/api/endpoints/stream.py index 7c34911..c464e72 100644 --- a/comet/api/endpoints/stream.py +++ b/comet/api/endpoints/stream.py @@ -8,7 +8,9 @@ from fastapi import APIRouter, BackgroundTasks, Request from comet.core.config_validation import config_check from comet.core.logger import logger from comet.core.models import database, settings, trackers +from comet.debrid.exceptions import DebridAuthError from comet.debrid.manager import get_debrid_extension +from comet.metadata.filter import release_filter from comet.metadata.manager import MetadataScraper from comet.services.debrid import DebridService from comet.services.lock import DistributedLock, is_scrape_in_progress @@ -135,6 +137,9 @@ async def stream( b64config: str = None, chilllink: bool = False, ): + if media_type not in ["movie", "series"]: + return {"streams": []} + if "tmdb:" in media_id: return {"streams": []} @@ -167,8 +172,27 @@ async def stream( async with aiohttp.ClientSession(connector=connector) as session: metadata_scraper = MetadataScraper(session) - # First, check if metadata is already cached id, season, episode = parse_media_id(media_type, media_id) + + # Digital Release Filter + if settings.DIGITAL_RELEASE_FILTER: + is_released = await release_filter.check_is_released( + session, media_type, media_id, season, episode + ) + + if not is_released: + logger.log("FILTER", f"🚫 {media_id} is not released yet. Skipping.") + return { + "streams": [ + { + "name": "[🚫] Comet", + "description": "Content not digitally released yet.", + "url": "https://comet.fast", + } + ] + } + + # Check if metadata is already cached cached_metadata = await metadata_scraper.get_cached( id, season if "kitsu" not in media_id else 1, episode ) @@ -390,14 +414,25 @@ async def stream( and debrid_service != "torrent" ): logger.log("SCRAPER", "🔄 Checking availability on debrid service...") - await debrid_service_instance.get_and_cache_availability( - session, - torrent_manager.torrents, - media_id, - media_only_id, - season, - episode, - ) + try: + await debrid_service_instance.get_and_cache_availability( + session, + torrent_manager.torrents, + media_id, + media_only_id, + season, + episode, + ) + except DebridAuthError as e: + return { + "streams": [ + { + "name": "[❌] Comet", + "description": e.display_message, + "url": "https://comet.fast", + } + ] + } if debrid_service != "torrent": cached_count = sum( diff --git a/comet/core/database.py b/comet/core/database.py index 253e542..c31afa7 100644 --- a/comet/core/database.py +++ b/comet/core/database.py @@ -294,10 +294,10 @@ async def setup_database(): ) await database.execute( - f""" + """ CREATE TABLE IF NOT EXISTS bandwidth_stats ( id INTEGER PRIMARY KEY, - total_bytes {"INTEGER" if settings.DATABASE_TYPE == "sqlite" else "BIGINT"}, + total_bytes BIGINT, last_updated INTEGER ) """ @@ -361,6 +361,16 @@ async def setup_database(): """ ) + await database.execute( + """ + CREATE TABLE IF NOT EXISTS digital_release_cache ( + media_id TEXT PRIMARY KEY, + release_date BIGINT, + timestamp INTEGER + ) + """ + ) + # ============================================================================= # TORRENTS TABLE INDEXES - Most critical for performance # ============================================================================= @@ -617,6 +627,13 @@ async def setup_database(): """ ) + await database.execute( + """ + CREATE INDEX IF NOT EXISTS idx_digital_release_timestamp + ON digital_release_cache (timestamp) + """ + ) + if settings.DATABASE_TYPE == "sqlite": await database.execute("PRAGMA busy_timeout=30000") # 30 seconds timeout await database.execute("PRAGMA journal_mode=WAL") @@ -697,6 +714,14 @@ async def _run_startup_cleanup(): {"cache_ttl": settings.DEBRID_CACHE_TTL, "current_time": current_time}, ) + await database.execute( + """ + DELETE FROM digital_release_cache + WHERE timestamp + :cache_ttl < :current_time; + """, + {"cache_ttl": settings.METADATA_CACHE_TTL, "current_time": current_time}, + ) + await database.execute("DELETE FROM download_links_cache") await database.execute( diff --git a/comet/core/execution.py b/comet/core/execution.py new file mode 100644 index 0000000..52821f9 --- /dev/null +++ b/comet/core/execution.py @@ -0,0 +1,43 @@ +import atexit +import multiprocessing +import os +import signal +from concurrent.futures import ProcessPoolExecutor + +_mp_context = None +try: + _mp_context = multiprocessing.get_context("forkserver") +except ValueError: + _mp_context = multiprocessing.get_context("spawn") + +app_executor = None + + +def worker_initializer(): + signal.signal(signal.SIGINT, signal.SIG_IGN) + + +def setup_executor(max_workers: int | None = None): + global app_executor + + if max_workers is None: + cpu_count = os.cpu_count() or 1 + max_workers = min(cpu_count, 4) + + app_executor = ProcessPoolExecutor( + max_workers=max_workers, mp_context=_mp_context, initializer=worker_initializer + ) + + +def shutdown_executor(): + global app_executor + if app_executor: + app_executor.shutdown(wait=True, cancel_futures=True) + app_executor = None + + +atexit.register(shutdown_executor) + + +def get_executor(): + return app_executor diff --git a/comet/core/log_levels.py b/comet/core/log_levels.py index 81da5bd..9fe25dc 100644 --- a/comet/core/log_levels.py +++ b/comet/core/log_levels.py @@ -62,6 +62,12 @@ CUSTOM_LOG_LEVELS = { "loguru_color": "", "no": 20, }, + "FILTER": { + "color": "#FFD700", + "icon": "🛡️", + "loguru_color": "", + "no": 35, + }, } ALL_LOG_LEVELS = {**STANDARD_LOG_LEVELS, **CUSTOM_LOG_LEVELS} diff --git a/comet/core/logger.py b/comet/core/logger.py index b1fe700..11cbf69 100644 --- a/comet/core/logger.py +++ b/comet/core/logger.py @@ -390,6 +390,13 @@ def log_startup_info(settings): ) logger.log("COMET", f"Remove Adult Content: {bool(settings.REMOVE_ADULT_CONTENT)}") + logger.log( + "COMET", f"Digital Release Filter: {bool(settings.DIGITAL_RELEASE_FILTER)}" + ) + logger.log( + "COMET", + f"TMDB Read Access Token: {settings.TMDB_READ_ACCESS_TOKEN if settings.TMDB_READ_ACCESS_TOKEN else 'Shared'}", + ) logger.log("COMET", f"Custom Header HTML: {bool(settings.CUSTOM_HEADER_HTML)}") background_scraper_display = ( diff --git a/comet/core/models.py b/comet/core/models.py index 0450931..f419152 100644 --- a/comet/core/models.py +++ b/comet/core/models.py @@ -129,6 +129,8 @@ class AppSettings(BaseSettings): BACKGROUND_SCRAPER_MAX_SERIES_PER_RUN: Optional[int] = 100 ANIME_MAPPING_SOURCE: Optional[str] = "remote" ANIME_MAPPING_REFRESH_INTERVAL: Optional[int] = 86400 + DIGITAL_RELEASE_FILTER: Optional[bool] = False + TMDB_READ_ACCESS_TOKEN: Optional[str] = None @field_validator("INDEXER_MANAGER_TYPE") def set_indexer_manager_type(cls, v, values): diff --git a/comet/debrid/exceptions.py b/comet/debrid/exceptions.py new file mode 100644 index 0000000..867c553 --- /dev/null +++ b/comet/debrid/exceptions.py @@ -0,0 +1,20 @@ +class DebridError(Exception): + """Base exception for debrid-related errors.""" + + def __init__(self, message: str, display_message: str = None): + self.message = message + self.display_message = display_message or message + super().__init__(self.message) + + +class DebridAuthError(DebridError): + """Raised when debrid authentication fails (not premium, invalid API key, etc.).""" + + def __init__(self, debrid_name: str, message: str = None): + self.debrid_name = debrid_name + default_message = f"{debrid_name}: Authentication failed or not premium" + display_message = ( + message + or f"{debrid_name}: Invalid API key or no active subscription.\nPlease check your debrid account." + ) + super().__init__(default_message, display_message) diff --git a/comet/debrid/stremthru.py b/comet/debrid/stremthru.py index 43a796d..561fac2 100644 --- a/comet/debrid/stremthru.py +++ b/comet/debrid/stremthru.py @@ -4,13 +4,19 @@ from urllib.parse import quote, unquote import aiohttp from RTN import parse, title_match +from comet.core.execution import get_executor from comet.core.logger import logger from comet.core.models import settings +from comet.debrid.exceptions import DebridAuthError from comet.services.debrid_cache import cache_availability from comet.services.torrent_manager import torrent_update_queue from comet.utils.parsing import is_video +def batch_parse(filenames): + return [parse(f) for f in filenames] + + class StremThru: def __init__( self, @@ -43,17 +49,29 @@ class StremThru: async def check_premium(self): try: - user = await self.session.get( + response = await self.session.get( f"{self.base_url}/user?client_ip={self.client_ip}" ) - user = await user.json() - return user["data"]["subscription_status"] == "premium" - except Exception as e: - logger.warning( - f"Exception while checking premium status on {self.name}: {e}" - ) + user = await response.json() - return False + if "data" not in user: + raise DebridAuthError( + self.name, + f"{self.name}: Invalid API key.\nPlease check your configuration.", + ) + + if user["data"]["subscription_status"] != "premium": + raise DebridAuthError( + self.name, + f"{self.name}: No active subscription.\nPlease renew your debrid account.", + ) + except DebridAuthError: + raise + except Exception as e: + raise DebridAuthError( + self.name, + f"{self.name}: Failed to check account status.\n{e}", + ) async def get_instant(self, magnets: list): try: @@ -72,8 +90,7 @@ class StremThru: tracker_map: dict, sources_map: dict, ): - if not await self.check_premium(): - return [] + await self.check_premium() chunk_size = 50 chunks = [ @@ -95,6 +112,26 @@ class StremThru: is_offcloud = self.real_debrid_name == "offcloud" + filenames_to_parse = [] + if not is_offcloud: + for result in availability: + for torrent in result: + if torrent["status"] != "cached": + continue + for file in torrent["files"]: + filename = file["name"].split("/")[-1] + if not is_video(filename) or "sample" in filename.lower(): + continue + filenames_to_parse.append(filename) + + parsed_iter = iter([]) + if filenames_to_parse: + loop = asyncio.get_running_loop() + parsed_results = await loop.run_in_executor( + get_executor(), batch_parse, filenames_to_parse + ) + parsed_iter = iter(parsed_results) + files = [] cached_count = 0 for result in availability: @@ -127,7 +164,7 @@ class StremThru: if not is_video(filename) or "sample" in filename.lower(): continue - filename_parsed = parse(filename) + filename_parsed = next(parsed_iter) season = ( filename_parsed.seasons[0] diff --git a/comet/main.py b/comet/main.py index bac6f53..260e718 100644 --- a/comet/main.py +++ b/comet/main.py @@ -44,20 +44,18 @@ def run_with_uvicorn(): workers=settings.FASTAPI_WORKERS, log_config=None, ) - server = Server(config=config) + server = uvicorn.Server(config=config) - with server.run_in_thread(): - log_startup_info(settings) - try: - while True: - time.sleep(1) # Keep the main thread alive - except KeyboardInterrupt: - logger.log("COMET", "Server stopped by user") - except Exception as e: - logger.error(f"Unexpected error: {e}") - logger.exception(traceback.format_exc()) - finally: - logger.log("COMET", "Server Shutdown") + log_startup_info(settings) + try: + server.run() + except KeyboardInterrupt: + logger.log("COMET", "Server stopped by user") + except Exception as e: + logger.error(f"Unexpected error: {e}") + logger.exception(traceback.format_exc()) + finally: + logger.log("COMET", "Server Shutdown") def run_with_gunicorn(): diff --git a/comet/metadata/filter.py b/comet/metadata/filter.py new file mode 100644 index 0000000..5c1be70 --- /dev/null +++ b/comet/metadata/filter.py @@ -0,0 +1,100 @@ +import time +from datetime import datetime + +from comet.core.logger import logger +from comet.core.models import database, settings +from comet.metadata.tmdb import TMDBApi + + +class DigitalReleaseFilter: + async def check_is_released( + self, + session, + media_type: str, + media_id: str, + season: int = None, + episode: int = None, + ): + try: + cached_date = await database.fetch_val( + """ + SELECT release_date FROM digital_release_cache + WHERE media_id = :media_id + AND timestamp + :cache_ttl >= :current_time + """, + { + "media_id": media_id, + "cache_ttl": settings.METADATA_CACHE_TTL, + "current_time": time.time(), + }, + ) + + if cached_date is not None: + return self._is_released(cached_date) + + tmdb_id = None + + tmdb = TMDBApi(session) + + if media_id.startswith("tt"): + imdb_id = media_id.split(":")[0] + tmdb_id = await tmdb.get_tmdb_id_from_imdb(imdb_id) + else: + # Other formats (e.g. kitsu) are not supported + return True + + if not tmdb_id: + logger.warning( + f"DigitalReleaseFilter: Could not resolve {media_id} to TMDB ID. Allowing search." + ) + return True + + release_date_str = None + if media_type == "movie": + release_date_str = await tmdb.get_upcoming_movie_release_date(tmdb_id) + elif media_type == "series": + release_date_str = await tmdb.get_episode_air_date( + tmdb_id, season, episode + ) + + cache_timestamp = int(time.time()) + if release_date_str is None: + # Not found, treat as released in far future to block + release_date_timestamp = 253402300799 # 9999-12-31 + # Cache for only 1 day (86400s) to recheck later + if settings.METADATA_CACHE_TTL > 86400: + cache_timestamp = int( + time.time() - settings.METADATA_CACHE_TTL + 86400 + ) + else: + release_date_timestamp = int( + datetime.strptime(release_date_str, "%Y-%m-%d").timestamp() + ) + + await database.execute( + """ + INSERT INTO digital_release_cache (media_id, release_date, timestamp) + VALUES (:media_id, :release_date, :timestamp) + ON CONFLICT (media_id) DO UPDATE SET release_date = :release_date, timestamp = :timestamp + """, + { + "media_id": media_id, + "release_date": release_date_timestamp, + "timestamp": cache_timestamp, + }, + ) + + return self._is_released(release_date_timestamp) + except Exception as e: + logger.error( + f"DigitalReleaseFilter: Error checking release status for {media_id}: {e}" + ) + return True + + def _is_released(self, release_timestamp: float): + if release_timestamp is None: + return True + return release_timestamp <= time.time() + + +release_filter = DigitalReleaseFilter() diff --git a/comet/metadata/tmdb.py b/comet/metadata/tmdb.py new file mode 100644 index 0000000..ba79ad8 --- /dev/null +++ b/comet/metadata/tmdb.py @@ -0,0 +1,79 @@ +import aiohttp + +from comet.core.logger import logger +from comet.core.models import settings + +DEFAULT_TMDB_READ_ACCESS_TOKEN = "eyJhbGciOiJIUzI1NiJ9.eyJhdWQiOiJlNTkxMmVmOWFhM2IxNzg2Zjk3ZTE1NWY1YmQ3ZjY1MSIsInN1YiI6IjY1M2NjNWUyZTg5NGE2MDBmZjE2N2FmYyIsInNjb3BlcyI6WyJhcGlfcmVhZCJdLCJ2ZXJzaW9uIjoxfQ.xrIXsMFJpI1o1j5g2QpQcFP1X3AfRjFA5FlBFO5Naw8" + + +class TMDBApi: + def __init__(self, session: aiohttp.ClientSession): + self.session = session + self.base_url = "https://api.themoviedb.org/3" + self.headers = { + "Authorization": f"Bearer {settings.TMDB_READ_ACCESS_TOKEN if settings.TMDB_READ_ACCESS_TOKEN else DEFAULT_TMDB_READ_ACCESS_TOKEN}", + "Content-Type": "application/json", + } + + async def get_upcoming_movie_release_date(self, tmdb_id: str): + try: + url = f"{self.base_url}/movie/{tmdb_id}/release_dates" + async with self.session.get(url, headers=self.headers) as response: + if response.status != 200: + return None + + data = await response.json() + + release_dates = [] + for result in data.get("results", []): + for release in result.get("release_dates", []): + if release.get("type") in [4, 5]: # Digital or Physical + date_str = release.get("release_date", "").split("T")[0] + if date_str: + release_dates.append(date_str) + + if release_dates: + return min(release_dates) + + return None + except Exception as e: + logger.error(f"TMDB: Error getting movie release date for {tmdb_id}: {e}") + return None + + async def get_episode_air_date(self, tmdb_id: str, season: int, episode: int): + try: + url = f"{self.base_url}/tv/{tmdb_id}/season/{season}/episode/{episode}" + async with self.session.get(url, headers=self.headers) as response: + if response.status != 200: + return None + + data = await response.json() + return data.get("air_date") + except Exception as e: + logger.error( + f"TMDB: Error getting episode air date for {tmdb_id} S{season}E{episode}: {e}" + ) + return None + + async def get_tmdb_id_from_imdb(self, imdb_id: str): + try: + url = f"{self.base_url}/find/{imdb_id}?external_source=imdb_id" + async with self.session.get(url, headers=self.headers) as response: + if response.status != 200: + text = await response.text() + logger.error( + f"TMDB: Failed to get TMDB ID from IMDB ID {imdb_id}: {text}" + ) + return None + + data = await response.json() + + if data.get("movie_results"): + return str(data["movie_results"][0]["id"]) + if data.get("tv_results"): + return str(data["tv_results"][0]["id"]) + + return None + except Exception as e: + logger.error(f"TMDB: Error converting IMDB ID {imdb_id}: {e}") + return None diff --git a/comet/services/orchestration.py b/comet/services/orchestration.py index ba57094..9edb1fe 100644 --- a/comet/services/orchestration.py +++ b/comet/services/orchestration.py @@ -5,12 +5,13 @@ import aiohttp import orjson from RTN import DefaultRanking, ParsedData +from comet.core.execution import get_executor from comet.core.logger import logger from comet.core.models import CometSettingsModel, database, settings from comet.scrapers.manager import scraper_manager from comet.services.filtering import filter_worker from comet.services.ranking import rank_worker -from comet.utils.parsing import default_dump +from comet.services.torrent_manager import torrent_update_queue class TorrentManager: @@ -72,7 +73,7 @@ class TorrentManager: async for scraper_name, results in scraper_manager.scrape_all(request, session): await self.filter_manager(scraper_name, results) - await self.cache_torrents() + asyncio.create_task(self.cache_torrents()) for torrent in self.ready_to_cache: season = torrent["parsed"].seasons[0] if torrent["parsed"].seasons else None @@ -128,41 +129,24 @@ class TorrentManager: } async def cache_torrents(self): - current_time = time.time() - values = [ - { - "media_id": self.media_only_id, + for torrent in self.ready_to_cache: + file_info = { "info_hash": torrent["infoHash"], - "file_index": int(torrent["fileIndex"]) - if torrent["fileIndex"] is not None - else None, + "index": torrent["fileIndex"], + "title": torrent["title"], + "size": torrent["size"], "season": torrent["parsed"].seasons[0] if torrent["parsed"].seasons else self.season, "episode": torrent["parsed"].episodes[0] if torrent["parsed"].episodes else None, - "title": torrent["title"], - "seeders": int(torrent["seeders"]) - if torrent["seeders"] is not None - else None, - "size": int(torrent["size"]) if torrent["size"] is not None else None, + "parsed": torrent["parsed"], + "seeders": torrent["seeders"], "tracker": torrent["tracker"], - "sources": orjson.dumps(torrent["sources"]).decode("utf-8"), - "parsed": orjson.dumps(torrent["parsed"], default_dump).decode("utf-8"), - "timestamp": current_time, + "sources": torrent["sources"], } - for torrent in self.ready_to_cache - ] - - query = f""" - INSERT {"OR REPLACE " if settings.DATABASE_TYPE == "sqlite" else ""} - INTO torrents - VALUES (:media_id, :info_hash, :file_index, :season, :episode, :title, :seeders, :size, :tracker, :sources, :parsed, :timestamp) - {" ON CONFLICT DO NOTHING" if settings.DATABASE_TYPE == "postgresql" else ""} - """ - - await database.execute_many(query, values) + await torrent_update_queue.add_torrent_info(file_info, self.media_only_id) async def filter_manager(self, scraper_name: str, torrents: list): if len(torrents) == 0: @@ -188,10 +172,10 @@ class TorrentManager: return loop = asyncio.get_running_loop() - chunk_size = 50 + chunk_size = 20 tasks = [ loop.run_in_executor( - None, + get_executor(), filter_worker, new_torrents[i : i + chunk_size], self.title, @@ -217,7 +201,7 @@ class TorrentManager: ): loop = asyncio.get_running_loop() self.ranked_torrents = await loop.run_in_executor( - None, + get_executor(), rank_worker, self.torrents, self.debrid_service, diff --git a/comet/services/torrent_manager.py b/comet/services/torrent_manager.py index c180973..25f4e3f 100644 --- a/comet/services/torrent_manager.py +++ b/comet/services/torrent_manager.py @@ -1,11 +1,10 @@ import asyncio import base64 import hashlib -import html import re import time from collections import defaultdict -from urllib.parse import parse_qs, urlparse +from urllib.parse import unquote import aiohttp import anyio @@ -20,15 +19,14 @@ from comet.core.logger import logger from comet.core.models import database, settings from comet.utils.parsing import default_dump, is_video +TRACKER_PATTERN = re.compile(r"[&?]tr=([^&]+)") INFO_HASH_PATTERN = re.compile(r"btih:([a-fA-F0-9]{40}|[a-zA-Z0-9]{32})") def extract_trackers_from_magnet(magnet_uri: str): try: - decoded_uri = html.unescape(magnet_uri) - parsed = urlparse(decoded_uri) - params = parse_qs(parsed.query) - return params.get("tr", []) + trackers = TRACKER_PATTERN.findall(magnet_uri) + return [unquote(tracker) for tracker in trackers] except Exception as e: logger.warning(f"Failed to extract trackers from magnet URI: {e}") return [] @@ -247,7 +245,7 @@ add_torrent_queue = AddTorrentQueue() class TorrentUpdateQueue: - def __init__(self, batch_size: int = 100, flush_interval: float = 5.0): + def __init__(self, batch_size: int = 1000, flush_interval: float = 5.0): self.queue = asyncio.Queue() self.batch_size = batch_size self.flush_interval = flush_interval