From f2be00dec5bcbbcce7571144cd3e5fe8fb8ef567 Mon Sep 17 00:00:00 2001 From: Mohamed Zumair Date: Thu, 10 Oct 2024 06:52:59 +0530 Subject: [PATCH] Fix prowlarr scrape based on the magnet link & fix movie contain series (#317) --- db/config.py | 2 +- scrapers/base_scraper.py | 53 ++++++++++------- scrapers/prowlarr.py | 125 ++++++++++++++++++++++++--------------- scrapers/torrentio.py | 26 ++++++-- scrapers/utils.py | 18 +++--- scrapers/zilean.py | 5 +- 6 files changed, 144 insertions(+), 85 deletions(-) diff --git a/db/config.py b/db/config.py index d3e0f4e..d1ca4d2 100644 --- a/db/config.py +++ b/db/config.py @@ -62,7 +62,7 @@ class Settings(BaseSettings): # Content Filtering adult_content_regex_keywords: str = ( r"(^|\b|\s|$|[\[._-])" - r"(18\s*\+|adults?|porn|sex|xxx|nude|boobs?|pussy|ass|bigass|bigtits?|blowjob|hardfuck|onlyfans?|naked|hot|milf|slut|doggy|anal|threesome|foursome|erotic|sexy|18\s*plus|trailer)" + r"(18\s*\+|adults?|porn|sex|xxx|nude|boobs?|pussy|ass|bigass|bigtits?|blowjob|hardfuck|onlyfans?|naked|hot|milf|slut|doggy|anal|threesome|foursome|erotic|sexy|18\s*plus|trailer|RiffTrax)" r"(\b|\s|$|[\]._-])" ) diff --git a/scrapers/base_scraper.py b/scrapers/base_scraper.py index bc0dee3..47b4aa2 100644 --- a/scrapers/base_scraper.py +++ b/scrapers/base_scraper.py @@ -18,8 +18,8 @@ class ScraperError(Exception): class BaseScraper(abc.ABC): - def __init__(self, cache_key_prefix: str): - self.logger = logging.getLogger(self.__class__.__name__) + def __init__(self, cache_key_prefix: str, logger_name: str): + self.logger = logging.getLogger(logger_name) self.http_client = httpx.AsyncClient(timeout=30) self.cache_key_prefix = cache_key_prefix @@ -153,8 +153,7 @@ class BaseScraper(abc.ABC): def validate_title_and_year( self, - parsed_title: str, - parsed_year: int | None, + parsed_data: dict, metadata: MediaFusionMetaData, catalog_type: str, torrent_title: str, @@ -162,8 +161,7 @@ class BaseScraper(abc.ABC): ) -> bool: """ Validate the title and year of the parsed data against the metadata. - :param parsed_title: Parsed title - :param parsed_year: Parsed year + :param parsed_data: Parsed data dictionary :param metadata: MediaFusionMetaData object :param catalog_type: Catalog type (movie, series) :param torrent_title: Torrent title @@ -173,33 +171,46 @@ class BaseScraper(abc.ABC): """ # Check similarity ratios max_similarity_ratio = calculate_max_similarity_ratio( - parsed_title, metadata.title, metadata.aka_titles + parsed_data["title"], metadata.title, metadata.aka_titles ) # Log and return False if similarity ratios is below the expected threshold if max_similarity_ratio < expected_ratio: self.logger.debug( - f"Title mismatch: '{parsed_title}' vs. '{metadata.title}'. Torrent title: '{torrent_title}'" + f"Title mismatch: '{parsed_data['title']}' vs. '{metadata.title}'. Torrent title: '{torrent_title}'" ) return False # Validate year based on a catalog type - if catalog_type == "movie" and parsed_year != metadata.year: - self.logger.debug( - f"Year mismatch for movie: {parsed_title} ({parsed_year}) vs. {metadata.title} ({metadata.year}). Torrent title: '{torrent_title}'" - ) - return False - - if catalog_type == "series" and parsed_year: - # If end_year exists, check parsed_year is within the range; otherwise, check parsed_year >= metadata.year - if ( - metadata.end_year - and not (metadata.year <= parsed_year <= metadata.end_year) - ) or (not metadata.end_year and parsed_year < metadata.year): + if catalog_type == "movie": + if parsed_data.get("year") != metadata.year: self.logger.debug( - f"Year mismatch for series: {parsed_title} ({parsed_year}) vs. {metadata.title} ({metadata.year} - {metadata.end_year}). Torrent title: '{torrent_title}'" + f"Year mismatch for movie: {parsed_data['title']} ({parsed_data.get('year')}) vs. {metadata.title} ({metadata.year}). Torrent title: '{torrent_title}'" ) return False + if parsed_data.get("season"): + self.logger.debug( + f"Season found for movie: {parsed_data['title']} ({parsed_data.get('season')}). Torrent title: '{torrent_title}'" + ) + return False + + if ( + catalog_type == "series" + and parsed_data.get("year") + and ( + ( + metadata.end_year + and not ( + metadata.year <= parsed_data.get("year") <= metadata.end_year + ) + ) + or (not metadata.end_year and parsed_data.get("year") < metadata.year) + ) + ): + self.logger.debug( + f"Year mismatch for series: {parsed_data['title']} ({parsed_data.get('year')}) vs. {metadata.title} ({metadata.year} - {metadata.end_year}). Torrent title: '{torrent_title}'" + ) + return False return True diff --git a/scrapers/prowlarr.py b/scrapers/prowlarr.py index 62cd436..da80f95 100644 --- a/scrapers/prowlarr.py +++ b/scrapers/prowlarr.py @@ -1,6 +1,6 @@ import asyncio from datetime import timedelta, datetime -from typing import List, Dict, Any, AsyncGenerator, Literal +from typing import List, Dict, Any, AsyncGenerator, Literal, AsyncIterable import PTT import dramatiq @@ -21,6 +21,7 @@ from scrapers.base_scraper import BaseScraper from scrapers.imdb_data import get_episode_by_date, get_season_episodes from utils.network import CircuitBreaker, batch_process_with_circuit_breaker from utils.parser import is_contain_18_plus_keywords +from utils.runtime_const import REDIS_ASYNC_CLIENT from utils.torrent import extract_torrent_metadata from utils.wrappers import minimum_run_interval @@ -61,7 +62,9 @@ class ProwlarrScraper(BaseScraper): ] def __init__(self): - super().__init__(cache_key_prefix="prowlarr") + super().__init__( + cache_key_prefix="prowlarr", logger_name=self.__class__.__name__ + ) self.base_url = f"{settings.prowlarr_url}/api/v1/search" @BaseScraper.cache( @@ -87,6 +90,12 @@ class ProwlarrScraper(BaseScraper): episode, ): results.append(stream) + except httpx.ReadTimeout: + self.logger.warning("Timeout while fetching search results") + except httpx.HTTPStatusError as e: + self.logger.error( + f"Error fetching search results: {e.response.text}, status code: {e.response.status_code}" + ) except Exception as e: self.logger.exception(f"An error occurred during scraping: {str(e)}") @@ -213,7 +222,7 @@ class ProwlarrScraper(BaseScraper): active_generators = len(stream_generators) processed_info_hashes = set() - async def producer(gen: AsyncGenerator, generator_id: int): + async def producer(gen: AsyncIterable, generator_id: int): try: async for stream_item in gen: await queue.put((stream_item, generator_id)) @@ -457,8 +466,7 @@ class ProwlarrScraper(BaseScraper): parsed_data = self.parse_title_data(stream_data.get("title")) if not self.validate_title_and_year( - parsed_data["title"], - parsed_data.get("year"), + parsed_data, metadata, catalog_type, stream_data.get("title"), @@ -562,22 +570,26 @@ class ProwlarrScraper(BaseScraper): async def parse_prowlarr_data( self, prowlarr_data: dict, catalog_type: str, parsed_data: dict ) -> dict | None: - download_url = prowlarr_data.get("downloadUrl") or prowlarr_data.get( - "magnetUrl" - ) - + download_url = await self.get_download_url(prowlarr_data) if not download_url: - download_url = await self.get_download_url(prowlarr_data) + return None try: torrent_data, is_torrent_downloaded = await self.get_torrent_data( download_url, prowlarr_data.get("indexer") ) - except httpx.TimeoutException: - self.logger.warning("Timeout while getting torrent data") + except httpx.HTTPStatusError as error: + if error.response.status_code in [429, 500]: + raise error + self.logger.error( + f"HTTP Error getting torrent data: {error.response.text}, status code: {error.response.status_code}" + ) return None + except httpx.TimeoutException as error: + self.logger.warning("Timeout while getting torrent data") + raise error except Exception as e: - self.logger.error(f"Error getting torrent data: {e}") + self.logger.exception(f"Error getting torrent data: {e}") return None info_hash = torrent_data.get("info_hash", "").lower() @@ -604,30 +616,24 @@ class ProwlarrScraper(BaseScraper): @staticmethod async def get_download_url(prowlarr_data: dict) -> str: - if prowlarr_data.get("indexer") in [ - "Torlock", - "YourBittorrent", - "The Pirate Bay", - "RuTracker.RU", - "BitSearch", - "BitRu", - "iDope", - "RuTor", - "Internet Archive", - "52BT", - ]: - return prowlarr_data.get("guid") - else: - if not prowlarr_data.get("magnetUrl") and not prowlarr_data.get( - "downloadUrl", "" - ).startswith("magnet:"): - torrent_info_data = await torrent_info.get_torrent_info( - prowlarr_data.get("infoUrl"), prowlarr_data.get("indexer") - ) - return torrent_info_data.get("magnetUrl") or torrent_info_data.get( - "downloadUrl" - ) - return prowlarr_data.get("magnetUrl") or prowlarr_data.get("downloadUrl") + guid = prowlarr_data.get("guid") or "" + magnet_url = prowlarr_data.get("magnetUrl") or "" + download_url = prowlarr_data.get("downloadUrl") or "" + + if guid and guid.startswith("magnet:"): + return guid + + if not magnet_url.startswith("magnet:") and not download_url.startswith( + "magnet:" + ): + torrent_info_data = await torrent_info.get_torrent_info( + prowlarr_data["infoUrl"], prowlarr_data["indexer"] + ) + return torrent_info_data.get("magnetUrl") or torrent_info_data.get( + "downloadUrl" + ) + + return magnet_url or download_url async def get_torrent_data( self, download_url: str, indexer: str @@ -695,9 +701,7 @@ class ProwlarrScraper(BaseScraper): @minimum_run_interval(hours=settings.prowlarr_search_interval_hour) @dramatiq.actor( - time_limit=10 * 60 * 1000, # 10 minutes - min_backoff=2 * 60 * 1000, # 2 minutes - max_backoff=10 * 60 * 1000, # 10 minutes + time_limit=30 * 60 * 1000, # 30 minutes priority=100, ) async def background_movie_title_search( @@ -718,8 +722,22 @@ async def background_movie_title_search( for query in scraper.MOVIE_SEARCH_QUERY_TEMPLATES ] - async for stream in scraper.process_streams(*title_streams_generators): - await scraper.store_streams([stream]) + try: + async for stream in scraper.process_streams(*title_streams_generators): + await scraper.store_streams([stream]) + except httpx.ReadTimeout: + scraper.logger.warning( + f"Timeout while fetching search results for movie {metadata.title} ({metadata.year}), retrying later" + ) + task_cache_key = f"background_tasks:background_movie_title_search:{metadata_id}" + await REDIS_ASYNC_CLIENT.delete(task_cache_key) + background_movie_title_search.send_with_options( + kwargs={"metadata_id": metadata_id}, delay=timedelta(minutes=5) + ) + except httpx.HTTPStatusError as e: + scraper.logger.error( + f"Error fetching search results: {e.response.text}, status code: {e.response.status_code}" + ) scraper.logger.info( f"Background title search completed for {metadata.title} ({metadata.year})" @@ -728,9 +746,7 @@ async def background_movie_title_search( @minimum_run_interval(hours=settings.prowlarr_search_interval_hour) @dramatiq.actor( - time_limit=10 * 60 * 1000, # 10 minutes - min_backoff=2 * 60 * 1000, # 2 minutes - max_backoff=10 * 60 * 1000, # 10 minutes + time_limit=30 * 60 * 1000, # 30 minutes priority=100, ) async def background_series_title_search( @@ -759,8 +775,23 @@ async def background_series_title_search( for query in scraper.SERIES_SEARCH_QUERY_TEMPLATES ] - async for stream in scraper.process_streams(*title_streams_generators): - await scraper.store_streams([stream]) + try: + async for stream in scraper.process_streams(*title_streams_generators): + await scraper.store_streams([stream]) + except httpx.ReadTimeout: + scraper.logger.warning( + f"Timeout while fetching search results for {metadata.title} S{season}E{episode}, retrying later" + ) + task_cache_key = f"background_tasks:background_series_title_search:metadata_id={metadata_id}_season={season}_episode={episode}" + await REDIS_ASYNC_CLIENT.delete(task_cache_key) + background_series_title_search.send_with_options( + kwargs={"metadata_id": metadata_id, "season": season, "episode": episode}, + delay=timedelta(minutes=5), + ) + except httpx.HTTPStatusError as e: + scraper.logger.error( + f"Error fetching search results: {e.response.text}, status code: {e.response.status_code}" + ) scraper.logger.info( f"Background title search completed for {metadata.title} S{season}E{episode}" diff --git a/scrapers/torrentio.py b/scrapers/torrentio.py index 556a774..446497d 100644 --- a/scrapers/torrentio.py +++ b/scrapers/torrentio.py @@ -16,10 +16,14 @@ from utils.parser import ( ) from utils.validation_helper import is_video_file +from scrapeops_python_requests.scrapeops_requests import ScrapeOpsRequests + class TorrentioScraper(BaseScraper): def __init__(self): - super().__init__(cache_key_prefix="torrentio") + super().__init__( + cache_key_prefix="torrentio", logger_name=self.__class__.__name__ + ) self.base_url = settings.torrentio_url self.semaphore = asyncio.Semaphore(10) @@ -35,8 +39,16 @@ class TorrentioScraper(BaseScraper): episode: 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}" + + scrapeops_logger = ScrapeOpsRequests( + scrapeops_api_key=settings.scrapeops_api_key, + spider_name="Torrentio Scraper", + job_name=job_name, + ) try: response = await self.make_request(url) @@ -46,14 +58,21 @@ class TorrentioScraper(BaseScraper): self.logger.warning(f"Invalid response received for {url}") return [] - return await self.parse_response( + stream_data = await self.parse_response( data, metadata, catalog_type, season, episode ) + for stream in stream_data: + scrapeops_logger.item_scraped( + item=stream.model_dump(include={"id", "meta_id"}), response=response + ) + return stream_data except (ScraperError, RetryError): return [] except Exception as e: self.logger.exception(f"Error occurred while fetching {url}: {e}") return [] + finally: + scrapeops_logger.logger.close_sdk() def validate_response(self, response: Dict[str, Any]) -> bool: return "streams" in response and isinstance(response["streams"], list) @@ -98,8 +117,7 @@ class TorrentioScraper(BaseScraper): if source != "Torrentio": if not self.validate_title_and_year( - parsed_data.get("title"), - parsed_data.get("metadata").get("year"), + parsed_data, metadata, catalog_type, parsed_data.get("torrent_name"), diff --git a/scrapers/utils.py b/scrapers/utils.py index f74e870..d0fe5e5 100644 --- a/scrapers/utils.py +++ b/scrapers/utils.py @@ -17,15 +17,6 @@ async def run_scrapers( ) -> set[TorrentStreams]: scraper_tasks = [] - if ( - settings.is_scrap_from_torrentio - and "torrentio_streams" in user_data.selected_catalogs - ): - torrentio_scraper = TorrentioScraper() - scraper_tasks.append( - torrentio_scraper.scrape_and_parse(metadata, catalog_type, season, episode) - ) - if settings.prowlarr_api_key and "prowlarr_streams" in user_data.selected_catalogs: prowlarr_scraper = ProwlarrScraper() scraper_tasks.append( @@ -41,6 +32,15 @@ async def run_scrapers( zilean_scraper.scrape_and_parse(metadata, catalog_type, season, episode) ) + if ( + settings.is_scrap_from_torrentio + and "torrentio_streams" in user_data.selected_catalogs + ): + torrentio_scraper = TorrentioScraper() + scraper_tasks.append( + torrentio_scraper.scrape_and_parse(metadata, catalog_type, season, episode) + ) + scraped_streams = await asyncio.gather(*scraper_tasks) scraped_streams = [stream for sublist in scraped_streams for stream in sublist] unique_streams = set(scraped_streams) diff --git a/scrapers/zilean.py b/scrapers/zilean.py index b0ceed9..f8a7357 100644 --- a/scrapers/zilean.py +++ b/scrapers/zilean.py @@ -15,7 +15,7 @@ from utils.parser import ( class ZileanScraper(BaseScraper): def __init__(self): - super().__init__(cache_key_prefix="zilean") + super().__init__(cache_key_prefix="zilean", logger_name=self.__class__.__name__) self.base_url = f"{settings.zilean_url}/dmm/search" self.semaphore = asyncio.Semaphore(10) @@ -86,8 +86,7 @@ class ZileanScraper(BaseScraper): torrent_data = PTT.parse_title(stream["raw_title"], True) if not self.validate_title_and_year( - torrent_data.get("title"), - torrent_data.get("year"), + torrent_data, metadata, catalog_type, stream["raw_title"],