diff --git a/db/config.py b/db/config.py index f66d523..8a7ef9e 100644 --- a/db/config.py +++ b/db/config.py @@ -78,6 +78,7 @@ class Settings(BaseSettings): is_scrap_from_mediafusion: bool = False mediafusion_search_interval_days: int = 3 mediafusion_url: str = "https://mediafusion.elfhosted.com" + mediafusion_api_password: str | None = None sync_debrid_cache_streams: bool = False # Zilean Settings diff --git a/db/crud.py b/db/crud.py index 873f2d2..e18007e 100644 --- a/db/crud.py +++ b/db/crud.py @@ -562,11 +562,13 @@ async def get_streams_base( # Handle live search and stream updates if live_search_streams: - scraper_args = [metadata, content_type] - if content_type == "series": - scraper_args.extend([season, episode]) - - new_streams = await run_scrapers(*scraper_args) + new_streams = await run_scrapers( + user_data=user_data, + metadata=metadata, + catalog_type=content_type, + season=season, + episode=episode, + ) all_streams = list(set(cached_streams).union(new_streams)) if new_streams: diff --git a/scrapers/base_scraper.py b/scrapers/base_scraper.py index 5225275..adbb3c5 100644 --- a/scrapers/base_scraper.py +++ b/scrapers/base_scraper.py @@ -29,6 +29,7 @@ from db.models import ( MediaFusionMetaData, ) from db.redis_database import REDIS_ASYNC_CLIENT +from db.schemas import UserData from scrapers import torrent_info from scrapers.imdb_data import get_episode_by_date from utils.network import batch_process_with_circuit_breaker, CircuitBreaker @@ -267,6 +268,7 @@ class BaseScraper(abc.ABC): async def scrape_and_parse( self, + user_data: UserData, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, @@ -281,7 +283,7 @@ class BaseScraper(abc.ABC): self.metrics.episode = episode try: result = await self._scrape_and_parse( - metadata, catalog_type, season, episode + user_data, metadata, catalog_type, season, episode ) if isinstance(result, list): return result @@ -486,6 +488,7 @@ class BaseScraper(abc.ABC): async def parse_response( self, response: Dict[str, Any], + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, @@ -494,6 +497,7 @@ class BaseScraper(abc.ABC): """ Parse the response into TorrentStreams objects. :param response: Response dictionary + :param user_data: UserData object :param metadata: MediaFusionMetaData object :param catalog_type: Catalog type (movie, series) :param season: Season number (for series) @@ -504,6 +508,7 @@ class BaseScraper(abc.ABC): def get_cache_key( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: str = None, @@ -761,6 +766,7 @@ class IndexerBaseScraper(BaseScraper, abc.ABC): async def _scrape_and_parse( self, + user_data: UserData, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, diff --git a/scrapers/bt4g.py b/scrapers/bt4g.py index cb14a87..26163ea 100644 --- a/scrapers/bt4g.py +++ b/scrapers/bt4g.py @@ -45,6 +45,7 @@ class BT4GScraper(BaseScraper): @BaseScraper.rate_limit(calls=2, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: Optional[int] = None, diff --git a/scrapers/jackett.py b/scrapers/jackett.py index 13670fc..dbd7a78 100644 --- a/scrapers/jackett.py +++ b/scrapers/jackett.py @@ -82,13 +82,16 @@ class JackettScraper(IndexerBaseScraper): @IndexerBaseScraper.rate_limit(calls=5, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, episode: int = None, ) -> List[TorrentStreams]: """Scrape and parse Jackett indexers for torrent streams""" - return await super()._scrape_and_parse(metadata, catalog_type, season, episode) + return await super()._scrape_and_parse( + user_data, metadata, catalog_type, season, episode + ) async def get_healthy_indexers(self) -> List[dict]: """Fetch and return list of healthy Jackett indexers with their capabilities""" diff --git a/scrapers/mediafusion.py b/scrapers/mediafusion.py index 793dcf7..da32261 100644 --- a/scrapers/mediafusion.py +++ b/scrapers/mediafusion.py @@ -6,7 +6,9 @@ import httpx from db.config import settings from db.models import MediaFusionMetaData, TorrentStreams +from db.schemas import UserData from scrapers.stremio_addons import StremioScraper +from utils.crypto import crypto_utils from utils.parser import convert_size_to_bytes from utils.runtime_const import MEDIAFUSION_SEARCH_TTL @@ -24,16 +26,41 @@ class MediafusionScraper(StremioScraper): timeout=30, proxy=settings.requests_proxy_url ) + def _generate_url( + self, + user_data: UserData, + metadata: MediaFusionMetaData, + catalog_type: str, + season: Optional[int] = None, + episode: Optional[int] = None, + ) -> str: + url = f"{self.base_url}/stream/{catalog_type}/{metadata.id}.json" + if catalog_type == "series": + url = f"{self.base_url}/stream/{catalog_type}/{metadata.id}:{season}:{episode}.json" + upstream_user_data = UserData( + api_password=settings.mediafusion_api_password, + streaming_provider=user_data.streaming_provider, + nudity_filter=user_data.nudity_filter, + certification_filter=user_data.certification_filter, + ) + self.http_client.headers.update( + {"encoded_user_data": crypto_utils.encode_user_data(upstream_user_data)} + ) + return url + @StremioScraper.cache(ttl=MEDIAFUSION_SEARCH_TTL) @StremioScraper.rate_limit(calls=5, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: Optional[int] = None, episode: Optional[int] = None, ) -> List[TorrentStreams]: - return await super()._scrape_and_parse(metadata, catalog_type, season, episode) + return await super()._scrape_and_parse( + user_data, metadata, catalog_type, season, episode + ) def get_adult_content_field(self, stream_data: Dict[str, Any]) -> str: return stream_data["description"] @@ -41,21 +68,30 @@ class MediafusionScraper(StremioScraper): def get_scraper_name(self) -> str: return "Mediafusion" - def parse_stream_title(self, stream: dict) -> dict: + def parse_stream_title(self, stream: dict) -> tuple[dict, bool]: description = stream["description"].splitlines() torrent_name = description[0].removeprefix("📂 ").split(" ┈➤ ")[0] metadata = PTT.parse_title(torrent_name, True) source = stream["name"].split()[0].title() + info_hash = stream.get("infoHash") + if not info_hash: + url = stream["url"] + if "/stream/" not in url: + return {}, False + info_hash = url.split("/")[6] + is_cached = "⚡️" in stream["name"] - return { - "torrent_name": torrent_name, - "title": metadata.get("title"), - "size": convert_size_to_bytes( - self.extract_size_string(stream["description"]) - ), - "seeders": self.extract_seeders(stream["description"]), - "languages": metadata.get("languages", []), - "metadata": metadata, - "filename": (stream.get("behaviorHints", {}).get("filename")), - "source": source, - } + metadata.update( + { + "source": source, + "info_hash": info_hash, + "size": convert_size_to_bytes( + self.extract_size_string(stream["description"]) + ), + "torrent_name": torrent_name, + "seeders": self.extract_seeders(stream["description"]), + "filename": stream.get("behaviorHints", {}).get("filename"), + } + ) + + return metadata, is_cached diff --git a/scrapers/prowlarr.py b/scrapers/prowlarr.py index 73b3d02..7d42ad5 100644 --- a/scrapers/prowlarr.py +++ b/scrapers/prowlarr.py @@ -27,13 +27,16 @@ class ProwlarrScraper(IndexerBaseScraper): @IndexerBaseScraper.rate_limit(calls=5, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, episode: int = None, ) -> List[TorrentStreams]: """Scrape and parse Prowlarr indexers for torrent streams""" - return await super()._scrape_and_parse(metadata, catalog_type, season, episode) + return await super()._scrape_and_parse( + user_data, metadata, catalog_type, season, episode + ) @property def live_title_search_enabled(self) -> bool: diff --git a/scrapers/scraper_tasks.py b/scrapers/scraper_tasks.py index 72a328e..41aa06a 100644 --- a/scrapers/scraper_tasks.py +++ b/scrapers/scraper_tasks.py @@ -10,6 +10,7 @@ import dramatiq from db.config import settings from db.models import TorrentStreams, MediaFusionMetaData +from db.schemas import UserData from scrapers.base_scraper import BaseScraper from scrapers.bt4g import BT4GScraper from scrapers.imdb_data import get_imdb_title_data, search_imdb, search_multiple_imdb @@ -52,6 +53,7 @@ CACHED_DATA = [ async def run_scrapers( + user_data: UserData, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, @@ -65,7 +67,9 @@ async def run_scrapers( # Create tasks for enabled scrapers tasks = [ tg.create_task( - scraper_cls().scrape_and_parse(metadata, catalog_type, season, episode), + scraper_cls().scrape_and_parse( + user_data, metadata, catalog_type, season, episode + ), name=f"{scraper_cls.__name__}", ) for is_enabled, scraper_cls in SCRAPERS diff --git a/scrapers/stremio_addons.py b/scrapers/stremio_addons.py index 733f50a..55344e5 100644 --- a/scrapers/stremio_addons.py +++ b/scrapers/stremio_addons.py @@ -1,4 +1,5 @@ import asyncio +import logging import re from abc import abstractmethod from typing import List, Dict, Any, Optional @@ -6,7 +7,9 @@ from typing import List, Dict, Any, Optional from tenacity import RetryError from db.models import TorrentStreams, MediaFusionMetaData, EpisodeFile +from db.schemas import UserData from scrapers.base_scraper import BaseScraper, ScraperError +from streaming_providers.cache_helpers import store_cached_info_hashes from utils.parser import is_contain_18_plus_keywords @@ -16,19 +19,25 @@ class StremioScraper(BaseScraper): self.base_url = base_url self.semaphore = asyncio.Semaphore(10) + def _generate_url( + self, + user_data, + metadata: MediaFusionMetaData, + catalog_type: str, + season: Optional[int] = None, + episode: Optional[int] = None, + ) -> str: + pass + async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: Optional[int] = None, episode: Optional[int] = None, ) -> List[TorrentStreams]: - url = f"{self.base_url}/stream/{catalog_type}/{metadata.id}.json" - job_name = f"{metadata.title}:{metadata.id}" - if catalog_type == "series": - url = f"{self.base_url}/stream/{catalog_type}/{metadata.id}:{season}:{episode}.json" - job_name += f":{season}:{episode}" - + url = self._generate_url(user_data, metadata, catalog_type, season, episode) try: response = await self.make_request(url) response.raise_for_status() @@ -41,7 +50,7 @@ class StremioScraper(BaseScraper): self.metrics.record_found_items(len(data.get("streams", []))) return await self.parse_response( - data, metadata, catalog_type, season, episode + data, user_data, metadata, catalog_type, season, episode ) except (ScraperError, RetryError): self.metrics.record_error("request_failed") @@ -57,6 +66,7 @@ class StremioScraper(BaseScraper): async def parse_response( self, response: Dict[str, Any], + user_data: UserData, metadata: MediaFusionMetaData, catalog_type: str, season: Optional[int] = None, @@ -67,7 +77,19 @@ class StremioScraper(BaseScraper): for stream_data in response.get("streams", []) ] results = await asyncio.gather(*tasks) - return [stream for stream in results if stream is not None] + streams = [] + cached_info_hashes = [] + for stream, is_cached in results: + if stream: + streams.append(stream) + if is_cached: + cached_info_hashes.append(stream.id) + + logging.info( + f"Found {len(streams)} streams for {metadata.title} on {self.get_scraper_name()} with {len(cached_info_hashes)} cached streams for {user_data.streaming_provider.service}" + ) + await store_cached_info_hashes(user_data.streaming_provider, cached_info_hashes) + return streams async def process_stream( self, @@ -76,7 +98,7 @@ class StremioScraper(BaseScraper): catalog_type: str, season: Optional[int] = None, episode: Optional[int] = None, - ) -> Optional[TorrentStreams]: + ) -> tuple[Optional[TorrentStreams], bool]: async with self.semaphore: try: adult_content_field = self.get_adult_content_field(stream_data) @@ -85,9 +107,9 @@ class StremioScraper(BaseScraper): self.logger.warning( f"Stream contains 18+ keywords: {adult_content_field}" ) - return None + return None, False - parsed_data = self.parse_stream_title(stream_data) + parsed_data, is_cached = self.parse_stream_title(stream_data) source = parsed_data["source"] if self.get_scraper_name() not in source: @@ -97,7 +119,7 @@ class StremioScraper(BaseScraper): catalog_type, parsed_data.get("torrent_name"), ): - return None + return None, False stream = self.create_torrent_stream(stream_data, parsed_data, metadata) @@ -105,18 +127,18 @@ class StremioScraper(BaseScraper): if not self.process_series_data( stream, parsed_data, season, episode, stream_data ): - return None + return None, False # Record metrics for successful processing self.metrics.record_processed_item() self.metrics.record_quality(stream.quality) self.metrics.record_source(source) - return stream + return stream, is_cached except Exception as e: self.metrics.record_error("stream_processing_error") self.logger.exception(f"Error processing stream: {e}") - return None + return None, False def create_torrent_stream( self, @@ -125,18 +147,18 @@ class StremioScraper(BaseScraper): metadata: MediaFusionMetaData, ) -> TorrentStreams: return TorrentStreams( - id=stream_data["infoHash"], + id=parsed_data["info_hash"], meta_id=metadata.id, torrent_name=parsed_data["torrent_name"], size=parsed_data["size"], filename=parsed_data["filename"], file_index=stream_data.get("fileIdx"), languages=parsed_data["languages"], - resolution=parsed_data["metadata"].get("resolution"), - codec=parsed_data["metadata"].get("codec"), - quality=parsed_data["metadata"].get("quality"), - audio=parsed_data["metadata"].get("audio"), - hdr=parsed_data["metadata"].get("hdr"), + resolution=parsed_data.get("resolution"), + codec=parsed_data.get("codec"), + quality=parsed_data.get("quality"), + audio=parsed_data.get("audio"), + hdr=parsed_data.get("hdr"), source=parsed_data["source"], uploader=parsed_data.get("uploader"), catalog=[f"{self.cache_key_prefix}_streams"], @@ -158,7 +180,7 @@ class StremioScraper(BaseScraper): ) -> bool: season_number = season - if parsed_data["metadata"].get("episodes"): + if parsed_data.get("episodes"): episode_data = [ EpisodeFile( season_number=season_number, @@ -169,7 +191,7 @@ class StremioScraper(BaseScraper): else None ), ) - for episode_number in parsed_data["metadata"]["episodes"] + for episode_number in parsed_data["episodes"] ] else: episode_data = [ @@ -199,7 +221,7 @@ class StremioScraper(BaseScraper): raise NotImplementedError @abstractmethod - def parse_stream_title(self, stream: Dict[str, Any]) -> Dict[str, Any]: + def parse_stream_title(self, stream: Dict[str, Any]) -> tuple[Dict[str, Any], bool]: raise NotImplementedError @abstractmethod diff --git a/scrapers/torrentio.py b/scrapers/torrentio.py index 204dcb2..3cf51de 100644 --- a/scrapers/torrentio.py +++ b/scrapers/torrentio.py @@ -1,17 +1,27 @@ from datetime import timedelta -from typing import Dict, Any, List +from typing import Dict, Any, List, Optional import PTT import httpx from db.config import settings from db.models import MediaFusionMetaData, TorrentStreams +from db.schemas import UserData from scrapers.stremio_addons import StremioScraper from utils.parser import ( convert_size_to_bytes, ) from utils.runtime_const import TORRENTIO_SEARCH_TTL +SUPPORTED_DEBRID_SERVICE = { + "realdebrid", + "premiumize", + "alldebrid", + "debridlink", + "offcloud", + "torbox", +} + class TorrentioScraper(StremioScraper): cache_key_prefix = "torrentio" @@ -26,16 +36,37 @@ class TorrentioScraper(StremioScraper): timeout=30, proxy=settings.requests_proxy_url ) + def _generate_url( + self, + user_data: UserData, + metadata: MediaFusionMetaData, + catalog_type: str, + season: Optional[int] = None, + episode: Optional[int] = None, + ) -> str: + user_data_str = ( + f"/{user_data.streaming_provider.service}={user_data.streaming_provider.token.rstrip('=')}" + if user_data.streaming_provider.service in SUPPORTED_DEBRID_SERVICE + else "" + ) + url = f"{self.base_url}{user_data_str}/stream/{catalog_type}/{metadata.id}.json" + if catalog_type == "series": + url = f"{self.base_url}{user_data_str}/stream/{catalog_type}/{metadata.id}:{season}:{episode}.json" + return url + @StremioScraper.cache(ttl=TORRENTIO_SEARCH_TTL) @StremioScraper.rate_limit(calls=5, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, episode: int = None, ) -> List[TorrentStreams]: - return await super()._scrape_and_parse(metadata, catalog_type, season, episode) + return await super()._scrape_and_parse( + user_data, metadata, catalog_type, season, episode + ) def get_adult_content_field(self, stream_data: Dict[str, Any]) -> str: return stream_data["title"] @@ -43,23 +74,31 @@ class TorrentioScraper(StremioScraper): def get_scraper_name(self) -> str: return "Torrentio" - def parse_stream_title(self, stream: dict) -> dict: + def parse_stream_title(self, stream: dict) -> tuple[dict, bool]: try: descriptions = stream.get("title") torrent_name = descriptions.splitlines()[0] metadata = PTT.parse_title(torrent_name, True) - source = stream["name"].split()[0].title() + source = stream["name"].splitlines()[0].split()[-1] + info_hash = stream.get("infoHash") + if not info_hash: + url = stream["url"] + info_hash = url.split("/")[5] - return { - "torrent_name": torrent_name, - "title": metadata.get("title"), - "size": convert_size_to_bytes(self.extract_size_string(descriptions)), - "seeders": self.extract_seeders(descriptions), - "languages": metadata["languages"], - "metadata": metadata, - "filename": stream.get("behaviorHints", {}).get("filename"), - "source": source, - } + is_cached = "+" in stream["name"] + metadata.update( + { + "source": source, + "info_hash": info_hash, + "size": convert_size_to_bytes( + self.extract_size_string(descriptions) + ), + "torrent_name": torrent_name, + "seeders": self.extract_seeders(descriptions), + "filename": stream.get("behaviorHints", {}).get("filename"), + } + ) + return metadata, is_cached except Exception as e: self.metrics.record_error("title_parsing_error") raise e diff --git a/scrapers/yts.py b/scrapers/yts.py index 706af50..8624648 100644 --- a/scrapers/yts.py +++ b/scrapers/yts.py @@ -18,6 +18,7 @@ class YTSScraper(BaseScraper): @BaseScraper.rate_limit(calls=2, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: Optional[int] = None, diff --git a/scrapers/zilean.py b/scrapers/zilean.py index 40daa9e..19b7956 100644 --- a/scrapers/zilean.py +++ b/scrapers/zilean.py @@ -26,6 +26,7 @@ class ZileanScraper(BaseScraper): @BaseScraper.rate_limit(calls=5, period=timedelta(seconds=1)) async def _scrape_and_parse( self, + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: int = None, @@ -90,7 +91,7 @@ class ZileanScraper(BaseScraper): try: streams = await self.parse_response( - stream_data, metadata, catalog_type, season, episode + stream_data, user_data, metadata, catalog_type, season, episode ) return streams except (ScraperError, RetryError): @@ -105,6 +106,7 @@ class ZileanScraper(BaseScraper): async def parse_response( self, response: List[Dict[str, Any]], + user_data, metadata: MediaFusionMetaData, catalog_type: str, season: int = None,