From d9051b3bc935baeb5a2135f69d5ff19618f2cf14 Mon Sep 17 00:00:00 2001 From: mhdzumair Date: Mon, 1 Jul 2024 12:19:16 +0530 Subject: [PATCH] refactoring: separate it out scrapy pipelines --- mediafusion_scrapy/pipelines/__init__.py | 18 + .../pipelines/duplicates_pipeline.py | 15 + .../formula_parser_pipeline.py} | 573 +----------------- .../live_stream_resolver_pipeline.py | 60 ++ .../pipelines/moto_gp_parser_pipeline.py | 135 +++++ .../pipelines/redis_cache_pipeline.py | 26 + .../pipelines/sport_video_parser_pipeline.py | 86 +++ .../pipelines/store_pipelines.py | 216 +++++++ .../pipelines/torrent_parser_pipeline.py | 73 +++ 9 files changed, 630 insertions(+), 572 deletions(-) create mode 100644 mediafusion_scrapy/pipelines/__init__.py create mode 100644 mediafusion_scrapy/pipelines/duplicates_pipeline.py rename mediafusion_scrapy/{pipelines.py => pipelines/formula_parser_pipeline.py} (50%) create mode 100644 mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py create mode 100644 mediafusion_scrapy/pipelines/moto_gp_parser_pipeline.py create mode 100644 mediafusion_scrapy/pipelines/redis_cache_pipeline.py create mode 100644 mediafusion_scrapy/pipelines/sport_video_parser_pipeline.py create mode 100644 mediafusion_scrapy/pipelines/store_pipelines.py create mode 100644 mediafusion_scrapy/pipelines/torrent_parser_pipeline.py diff --git a/mediafusion_scrapy/pipelines/__init__.py b/mediafusion_scrapy/pipelines/__init__.py new file mode 100644 index 0000000..d4be0d5 --- /dev/null +++ b/mediafusion_scrapy/pipelines/__init__.py @@ -0,0 +1,18 @@ +from .duplicates_pipeline import TorrentDuplicatesPipeline +from .formula_parser_pipeline import FormulaParserPipeline +from .moto_gp_parser_pipeline import MotoGPParserPipeline +from .store_pipelines import ( + QueueBasedPipeline, + EventSeriesStorePipeline, + TVStorePipeline, + MovieStorePipeline, + SeriesStorePipeline, + LiveEventStorePipeline, +) +from .redis_cache_pipeline import RedisCacheURLPipeline +from .live_stream_resolver_pipeline import LiveStreamResolverPipeline +from .torrent_parser_pipeline import ( + TorrentDownloadAndParsePipeline, + MagnetDownloadAndParsePipeline, +) +from .sport_video_parser_pipeline import SportVideoParserPipeline diff --git a/mediafusion_scrapy/pipelines/duplicates_pipeline.py b/mediafusion_scrapy/pipelines/duplicates_pipeline.py new file mode 100644 index 0000000..33ea5f3 --- /dev/null +++ b/mediafusion_scrapy/pipelines/duplicates_pipeline.py @@ -0,0 +1,15 @@ +from itemadapter import ItemAdapter +from scrapy.exceptions import DropItem + + +class TorrentDuplicatesPipeline: + def __init__(self): + self.info_hashes_seen = set() + + def process_item(self, item, spider): + adapter = ItemAdapter(item) + if adapter["info_hash"] in self.info_hashes_seen: + raise DropItem(f"Duplicate item found: {adapter['info_hash']}") + else: + self.info_hashes_seen.add(adapter["info_hash"]) + return item diff --git a/mediafusion_scrapy/pipelines.py b/mediafusion_scrapy/pipelines/formula_parser_pipeline.py similarity index 50% rename from mediafusion_scrapy/pipelines.py rename to mediafusion_scrapy/pipelines/formula_parser_pipeline.py index cd28703..74b060e 100644 --- a/mediafusion_scrapy/pipelines.py +++ b/mediafusion_scrapy/pipelines/formula_parser_pipeline.py @@ -1,45 +1,14 @@ -import asyncio import logging -import random +import logging import re from datetime import datetime -from uuid import uuid4 -import redis.asyncio as redis_async -import scrapy -from beanie import WriteRules -from itemadapter import ItemAdapter -from scrapy import signals from scrapy.exceptions import DropItem -from scrapy.http.request import NO_CALLBACK -from scrapy.utils.defer import maybe_deferred_to_future -from db import crud -from db.config import settings -from db.crud import is_torrent_stream_exists from db.models import ( - TorrentStreams, - Season, - MediaFusionSeriesMetaData, Episode, ) -from db.schemas import TVMetaData -from utils import torrent, const from utils.parser import convert_size_to_bytes -from utils.runtime_const import SPORTS_ARTIFACTS - - -class TorrentDuplicatesPipeline: - def __init__(self): - self.info_hashes_seen = set() - - def process_item(self, item, spider): - adapter = ItemAdapter(item) - if adapter["info_hash"] in self.info_hashes_seen: - raise DropItem(f"Duplicate item found: {adapter['info_hash']}") - else: - self.info_hashes_seen.add(adapter["info_hash"]) - return item class FormulaParserPipeline: @@ -471,543 +440,3 @@ class FormulaParserPipeline: ) torrent_data["episodes"] = episodes - - -class MotoGPParserPipeline: - def __init__(self): - self.name_parser_patterns = { - "smcgill1969": [ - re.compile( - r"MotoGP" # Series name with flexible space or dot - r"\.(?P\d{4})" # Year - r"x(?P\d{2})" # Round - r"\.(?P.+?)" # Event and Episode name - r"\.(?PBTSportHD|TNTSportsHD)" # Broadcaster - r"\.(?P4K|SD|1080p)" # Resolution and Quality - ), - re.compile( - r"MotoGP" # Series name with flexible space or dot - r"\.(?P\d{4})" # Year - r"(?:-\d{2}-\d{2})?" # ignore Date part starting with a hyphen - r"\.(?P.+?)" # Event and Episode name - r"(?:\.(?PWEB-DL|HDTV))?" # Optional Quality - r"(?:\.(?PBTSportHD|TNTSportsHD))?" # Optional Broadcaster - r"\.(?P4K|SD|1080p)" # Resolution - ), - ], - } - self.default_poster = random.choice(SPORTS_ARTIFACTS["MotoGP"]["poster"]) - - self.smcgill1969_resolutions = { - "4K": "4K", - "SD": "576p", - "1080p": "1080p", - } - self.known_countries_first_words = ["San", "Great"] - - self.title_parser_functions = { - "smcgill1969": self.parse_smcgill1969_title, - } - self.description_parser_functions = { - "smcgill1969": self.parse_smcgill1969_description, - } - - def process_item(self, item, spider): - uploader = item["uploader"] - title = re.sub(r"\.\.+", ".", item["torrent_name"]) - self.title_parser_functions[uploader](title, item) - if not item.get("title"): - raise DropItem(f"Title not parsed: {title}") - self.description_parser_functions[uploader](item) - if not item.get("episodes"): - raise DropItem(f"Episodes not parsed: {item!r}") - return item - - def event_and_episode_name_parser(self, event_and_episode_name): - parts = event_and_episode_name.split(".") - if parts[0] in self.known_countries_first_words and len(parts) > 1: - event = f"{parts[0]} {parts[1]}" - episode_name = " ".join(parts[2:]) - else: - event = parts[0] - episode_name = " ".join(parts[1:]) - return event, episode_name - - def parse_smcgill1969_title(self, title, torrent_data: dict): - for pattern in self.name_parser_patterns.get("smcgill1969"): - match = pattern.match(title) - if match: - data = match.groupdict() - series = "MotoGP" - event, episode_name = self.event_and_episode_name_parser( - data["EventAndEpisodeName"] - ) - - torrent_data.update( - { - "title": " ".join( - filter( - None, - [ - series, - data.get("Year"), - event, - "(smcgill1969)", # add the uploader to the title for uniqueness - ], - ) - ), - "year": int(data["Year"]), - "resolution": self.smcgill1969_resolutions.get( - data["Resolution"] - ), - "event": event, - "episode_name": episode_name, - } - ) - return - logging.warning(f"Failed to parse title: {title}") - - def parse_smcgill1969_description(self, torrent_data: dict): - torrent_description = torrent_data.get("description") - - codec_matches = re.findall(r"Codec ID\s*:\s*(\S+)", torrent_description) - if codec_matches: - torrent_data["codec"] = codec_matches[0] - torrent_data["audio"] = codec_matches[1] if len(codec_matches) > 1 else None - - torrent_data["poster"] = self.default_poster - torrent_data["is_add_title_to_poster"] = True - - episodes = [] - for index, file_detail in enumerate(torrent_data.get("file_details", [])): - file_name = file_detail.get("file_name") - file_size = file_detail.get("file_size") - - episodes.append( - Episode( - episode_number=index + 1, - filename=file_name, - size=convert_size_to_bytes(file_size), - file_index=index, - title=" ".join(file_name.split(".")[1:-1]), - released=torrent_data.get("created_at"), - ) - ) - - torrent_data["episodes"] = episodes - - -class QueueBasedPipeline: - def __init__(self): - self.queue = asyncio.Queue() - self.processing_task = None - - async def init(self): - self.processing_task = asyncio.create_task(self.process_queue()) - - async def close(self): - await self.queue.join() - self.processing_task.cancel() - - @classmethod - def from_crawler(cls, crawler): - p = cls() - crawler.signals.connect(p.init, signal=signals.spider_opened) - crawler.signals.connect(p.close, signal=signals.spider_closed) - return p - - async def process_item(self, item, spider): - # Instead of processing the item directly, add it to the queue - await self.queue.put((item, spider)) - return item - - async def process_queue(self): - logging.info("Starting processing queue") - while True: - item, spider = await self.queue.get() - try: - await self.parse_item(item, spider) - except Exception as e: - logging.error(f"Error processing item: {e}", exc_info=True) - finally: - self.queue.task_done() - - async def parse_item(self, item, spider): - raise NotImplementedError - - -class EventSeriesStorePipeline(QueueBasedPipeline): - def __init__(self): - super().__init__() - self.redis = redis_async.Redis.from_url(settings.redis_url) - - async def close(self): - await super().close() - await self.redis.aclose() - - async def parse_item(self, item, spider): - if "title" not in item: - logging.warning(f"title not found in item: {item}") - raise DropItem(f"title not found in item: {item}") - - series = await MediaFusionSeriesMetaData.find_one( - {"title": item["title"]}, fetch_links=True - ) - - if not series: - meta_id = f"mf{uuid4().fields[-1]}" - poster = item.get("poster") - background = item.get("background") - - # Create an initial entry for the series - series = MediaFusionSeriesMetaData( - id=meta_id, - title=item["title"], - year=item["year"], - poster=poster, - background=background, - streams=[], - is_poster_working=bool(poster), - is_add_title_to_poster=item.get("is_add_title_to_poster", False), - ) - await series.insert() - logging.info("Added series %s", series.title) - - meta_id = series.id - - stream = next((s for s in series.streams if s.id == item["info_hash"]), None) - if stream is None: - # Create the stream - stream = TorrentStreams( - id=item["info_hash"], - torrent_name=item["torrent_name"], - announce_list=item["announce_list"], - size=item["total_size"], - languages=item["languages"], - resolution=item.get("resolution"), - codec=item.get("codec"), - quality=item.get("quality"), - audio=item.get("audio"), - encoder=item.get("encoder"), - source=item["source"], - catalog=item["catalog"], - created_at=item["created_at"], - season=Season(season_number=1, episodes=item["episodes"]), - meta_id=meta_id, - seeders=item["seeders"], - ) - # Add the stream to the series - series.streams.append(stream) - logging.info( - "Added stream %s to series %s", stream.torrent_name, series.title - ) - - self.organize_episodes(series) - - await series.save(link_rule=WriteRules.WRITE) - logging.info("Updated series %s", series.title) - await self.redis.sadd(item["scraped_info_hash_key"], item["info_hash"]) - - return item - - def organize_episodes(self, series): - # Flatten all episodes from all streams and sort by release date - all_episodes = sorted( - ( - episode - for stream in series.streams - for episode in stream.season.episodes - ), - key=lambda e: ( - e.released.date(), - e.filename, - ), # Sort primarily by released date, then by filename - ) - - # Assign episode numbers, ensuring the same title across different qualities gets the same number - episode_number = 1 - last_title = None - for episode in all_episodes: - if episode.title != last_title: - if last_title: - episode_number = episode_number + 1 - last_title = episode.title - - episode.episode_number = episode_number - - # Now distribute episodes back to their respective streams, ensuring they are in the correct order - for stream in series.streams: - stream.season.episodes.sort( - key=lambda e: e.episode_number - ) # Ensure episodes are ordered by episode number - - -class TVStorePipeline(QueueBasedPipeline): - async def parse_item(self, item, spider): - if "title" not in item: - logging.warning(f"title not found in item: {item}") - raise DropItem(f"title not found in item: {item}") - - tv_metadata = TVMetaData.model_validate(item) - await crud.save_tv_channel_metadata(tv_metadata) - return item - - -class TorrentDownloadAndParsePipeline: - async def process_item(self, item, spider): - adapter = ItemAdapter(item) - torrent_link = adapter.get("torrent_link") - - if not torrent_link: - raise DropItem(f"No torrent link found in item: {item}") - - response = await maybe_deferred_to_future( - spider.crawler.engine.download( - scrapy.Request(torrent_link, callback=NO_CALLBACK) - ) - ) - - if response.status != 200: - spider.logger.error( - f"Failed to download torrent file: {response.url} with status {response.status}" - ) - raise DropItem(f"Failed to download torrent file: {response.url}") - - # Validate the content-type of the response - if "application/x-bittorrent" not in response.headers.get( - "Content-Type", b"" - ).decode("utf-8", "ignore"): - spider.logger.error( - f"Unexpected Content-Type for {response.url}: {response.headers.get('Content-Type')}" - ) - raise DropItem(f"Unexpected Content-Type for {response.url}") - - torrent_metadata = torrent.extract_torrent_metadata( - response.body, item.get("is_parse_ptn", True) - ) - - if not torrent_metadata: - raise DropItem(f"Failed to extract torrent metadata: {item}") - - item.update(torrent_metadata) - return item - - -class MagnetDownloadAndParsePipeline: - async def process_item(self, item, spider): - adapter = ItemAdapter(item) - magnet_link = adapter.get("magnet_link") - - if not magnet_link: - raise DropItem(f"No magnet link found in item: {item}") - - info_hash, trackers = torrent.parse_magnet(magnet_link) - if not info_hash: - raise DropItem(f"Failed to parse info_hash from magnet link: {magnet_link}") - if await is_torrent_stream_exists(info_hash): - raise DropItem(f"Torrent stream already exists: {info_hash}") - - torrent_metadata = await torrent.info_hashes_to_torrent_metadata( - [info_hash], trackers - ) - - if not torrent_metadata: - raise DropItem(f"Failed to extract torrent metadata: {item}") - - item.update(torrent_metadata[0]) - return item - - -class SportVideoParserPipeline: - RESOLUTIONS = { - "3840x2160": "4K", - "2560x1440": "1440p", - "1920x1080": "1080p", - "1280x720": "720p", - "854x480": "480p", - "640x360": "360p", - "426x240": "240p", - } - - def __init__(self): - self.title_regex = re.compile(r"^.*?\s(\d{2}\.\d{2}\.\d{4})") - - def process_item(self, item, spider): - adapter = ItemAdapter(item) - if "title" not in adapter: - raise DropItem(f"title not found in item: {item}") - - match = self.title_regex.search(adapter["title"]) or self.title_regex.search( - adapter["torrent_name"] - ) - if match: - date_str = match.group(1) - date_obj = datetime.strptime(date_str, "%d.%m.%Y") - else: - date_obj = datetime.now() - - # Try to match or infer resolution - raw_resolution = adapter.get("aspect_ratio", "").replace(" ", "") - resolution = self.infer_resolution(raw_resolution) - - if "video_codec" in adapter: - codec = re.sub(r"\s+", "", adapter["video_codec"]).replace(",", " - ") - item["codec"] = codec - - if "length" in adapter or "lenght" in adapter: - item["runtime"] = re.sub( - r"\s+", "", adapter.get("length", adapter.get("lenght", "")) - ) - - item.update( - { - "created_at": date_obj.strftime("%Y-%m-%d"), - "year": date_obj.year, - "languages": [re.sub(r"[ .]", "", item["language"])], - "is_imdb": False, - "resolution": resolution, - } - ) - return item - - def infer_resolution(self, aspect_ratio): - # Normalize the aspect_ratio string by replacing common variants of "x" with the standard one - normalized_aspect_ratio = aspect_ratio.replace("×", "x").replace("х", "x") - - # Direct match in RESOLUTIONS dictionary - resolution = self.RESOLUTIONS.get(normalized_aspect_ratio) - if resolution: - return resolution - - # Attempt to infer based on height - height = ( - normalized_aspect_ratio.split("x")[-1] - if "x" in normalized_aspect_ratio - else None - ) - if height: - for res, label in self.RESOLUTIONS.items(): - if res.endswith(height): - return label - - # Fallback to checking against common labels - for label in ["4K", "2160p", "1440p", "1080p", "720p", "480p", "360p", "240p"]: - if label in normalized_aspect_ratio: - return label - - # If no match found - return None - - -class MovieStorePipeline(QueueBasedPipeline): - async def parse_item(self, item, spider): - if "title" not in item: - raise DropItem(f"title not found in item: {item}") - - if item.get("type") != "movie": - return item - - await crud.save_movie_metadata(item, item.get("is_imdb", True)) - return item - - -class SeriesStorePipeline(QueueBasedPipeline): - async def parse_item(self, item, spider): - if "title" not in item: - raise DropItem(f"title not found in item: {item}") - - if item.get("type") != "series": - return item - await crud.save_series_metadata(item) - return item - - -class LiveEventStorePipeline(QueueBasedPipeline): - def __init__(self): - super().__init__() - self.redis = redis_async.Redis.from_url(settings.redis_url) - - async def close(self): - await super().close() - await self.redis.aclose() - - async def parse_item(self, item, spider): - if "title" not in item: - raise DropItem(f"name not found in item: {item}") - - await crud.save_events_data(self.redis, item) - return item - - -class RedisCacheURLPipeline: - def __init__(self): - self.redis = redis_async.Redis.from_url(settings.redis_url) - - async def close(self): - await self.redis.aclose() - - @classmethod - def from_crawler(cls, crawler): - p = cls() - crawler.signals.connect(p.close, signal=signals.spider_closed) - return p - - async def process_item(self, item, spider): - if "webpage_url" not in item: - raise DropItem(f"webpage_url not found in item: {item}") - - await self.redis.sadd(item["scraped_url_key"], item["webpage_url"]) - return item - - -class LiveStreamResolverPipeline: - async def process_item(self, item, spider): - adapter = ItemAdapter(item) - stream_url = adapter.get("stream_url") - stream_headers = adapter.get("stream_headers") - if not stream_headers: - referer = adapter.get("referer") - stream_headers = {"Referer": referer} if referer else {} - - if not stream_url: - raise DropItem(f"No stream URL found in item: {item}") - - response = await maybe_deferred_to_future( - spider.crawler.engine.download( - scrapy.Request( - stream_url, - callback=NO_CALLBACK, - headers=stream_headers, - method="HEAD", - dont_filter=True, - ) - ) - ) - content_type = response.headers.get("Content-Type", b"").decode().lower() - - if response.status == 200 and content_type in const.M3U8_VALID_CONTENT_TYPES: - stream_headers.update( - { - "User-Agent": response.request.headers.get("User-Agent").decode(), - "Referer": response.request.headers.get("Referer").decode(), - } - ) - - item["streams"].append( - { - "name": adapter["stream_name"], - "url": adapter["stream_url"], - "source": adapter["stream_source"], - "behaviorHints": { - "notWebReady": True, - "proxyHeaders": { - "request": stream_headers, - }, - }, - } - ) - return item - else: - raise DropItem( - f"Invalid M3U8 URL: {stream_url} with Content-Type: {content_type} response: {response.status}" - ) diff --git a/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py b/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py new file mode 100644 index 0000000..afe5a93 --- /dev/null +++ b/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py @@ -0,0 +1,60 @@ +import scrapy +from itemadapter import ItemAdapter +from scrapy.exceptions import DropItem +from scrapy.http.request import NO_CALLBACK +from scrapy.utils.defer import maybe_deferred_to_future + +from utils import const + + +class LiveStreamResolverPipeline: + async def process_item(self, item, spider): + adapter = ItemAdapter(item) + stream_url = adapter.get("stream_url") + stream_headers = adapter.get("stream_headers") + if not stream_headers: + referer = adapter.get("referer") + stream_headers = {"Referer": referer} if referer else {} + + if not stream_url: + raise DropItem(f"No stream URL found in item: {item}") + + response = await maybe_deferred_to_future( + spider.crawler.engine.download( + scrapy.Request( + stream_url, + callback=NO_CALLBACK, + headers=stream_headers, + method="HEAD", + dont_filter=True, + ) + ) + ) + content_type = response.headers.get("Content-Type", b"").decode().lower() + + if response.status == 200 and content_type in const.M3U8_VALID_CONTENT_TYPES: + stream_headers.update( + { + "User-Agent": response.request.headers.get("User-Agent").decode(), + "Referer": response.request.headers.get("Referer").decode(), + } + ) + + item["streams"].append( + { + "name": adapter["stream_name"], + "url": adapter["stream_url"], + "source": adapter["stream_source"], + "behaviorHints": { + "notWebReady": True, + "proxyHeaders": { + "request": stream_headers, + }, + }, + } + ) + return item + else: + raise DropItem( + f"Invalid M3U8 URL: {stream_url} with Content-Type: {content_type} response: {response.status}" + ) diff --git a/mediafusion_scrapy/pipelines/moto_gp_parser_pipeline.py b/mediafusion_scrapy/pipelines/moto_gp_parser_pipeline.py new file mode 100644 index 0000000..85f8d63 --- /dev/null +++ b/mediafusion_scrapy/pipelines/moto_gp_parser_pipeline.py @@ -0,0 +1,135 @@ +import logging +import random +import re + +from scrapy.exceptions import DropItem + +from db.models import ( + Episode, +) +from utils.parser import convert_size_to_bytes +from utils.runtime_const import SPORTS_ARTIFACTS + + +class MotoGPParserPipeline: + def __init__(self): + self.name_parser_patterns = { + "smcgill1969": [ + re.compile( + r"MotoGP" # Series name with flexible space or dot + r"\.(?P\d{4})" # Year + r"x(?P\d{2})" # Round + r"\.(?P.+?)" # Event and Episode name + r"\.(?PBTSportHD|TNTSportsHD)" # Broadcaster + r"\.(?P4K|SD|1080p)" # Resolution and Quality + ), + re.compile( + r"MotoGP" # Series name with flexible space or dot + r"\.(?P\d{4})" # Year + r"(?:-\d{2}-\d{2})?" # ignore Date part starting with a hyphen + r"\.(?P.+?)" # Event and Episode name + r"(?:\.(?PWEB-DL|HDTV))?" # Optional Quality + r"(?:\.(?PBTSportHD|TNTSportsHD))?" # Optional Broadcaster + r"\.(?P4K|SD|1080p)" # Resolution + ), + ], + } + self.default_poster = random.choice(SPORTS_ARTIFACTS["MotoGP"]["poster"]) + + self.smcgill1969_resolutions = { + "4K": "4K", + "SD": "576p", + "1080p": "1080p", + } + self.known_countries_first_words = ["San", "Great"] + + self.title_parser_functions = { + "smcgill1969": self.parse_smcgill1969_title, + } + self.description_parser_functions = { + "smcgill1969": self.parse_smcgill1969_description, + } + + def process_item(self, item, spider): + uploader = item["uploader"] + title = re.sub(r"\.\.+", ".", item["torrent_name"]) + self.title_parser_functions[uploader](title, item) + if not item.get("title"): + raise DropItem(f"Title not parsed: {title}") + self.description_parser_functions[uploader](item) + if not item.get("episodes"): + raise DropItem(f"Episodes not parsed: {item!r}") + return item + + def event_and_episode_name_parser(self, event_and_episode_name): + parts = event_and_episode_name.split(".") + if parts[0] in self.known_countries_first_words and len(parts) > 1: + event = f"{parts[0]} {parts[1]}" + episode_name = " ".join(parts[2:]) + else: + event = parts[0] + episode_name = " ".join(parts[1:]) + return event, episode_name + + def parse_smcgill1969_title(self, title, torrent_data: dict): + for pattern in self.name_parser_patterns.get("smcgill1969"): + match = pattern.match(title) + if match: + data = match.groupdict() + series = "MotoGP" + event, episode_name = self.event_and_episode_name_parser( + data["EventAndEpisodeName"] + ) + + torrent_data.update( + { + "title": " ".join( + filter( + None, + [ + series, + data.get("Year"), + event, + "(smcgill1969)", # add the uploader to the title for uniqueness + ], + ) + ), + "year": int(data["Year"]), + "resolution": self.smcgill1969_resolutions.get( + data["Resolution"] + ), + "event": event, + "episode_name": episode_name, + } + ) + return + logging.warning(f"Failed to parse title: {title}") + + def parse_smcgill1969_description(self, torrent_data: dict): + torrent_description = torrent_data.get("description") + + codec_matches = re.findall(r"Codec ID\s*:\s*(\S+)", torrent_description) + if codec_matches: + torrent_data["codec"] = codec_matches[0] + torrent_data["audio"] = codec_matches[1] if len(codec_matches) > 1 else None + + torrent_data["poster"] = self.default_poster + torrent_data["is_add_title_to_poster"] = True + + episodes = [] + for index, file_detail in enumerate(torrent_data.get("file_details", [])): + file_name = file_detail.get("file_name") + file_size = file_detail.get("file_size") + + episodes.append( + Episode( + episode_number=index + 1, + filename=file_name, + size=convert_size_to_bytes(file_size), + file_index=index, + title=" ".join(file_name.split(".")[1:-1]), + released=torrent_data.get("created_at"), + ) + ) + + torrent_data["episodes"] = episodes diff --git a/mediafusion_scrapy/pipelines/redis_cache_pipeline.py b/mediafusion_scrapy/pipelines/redis_cache_pipeline.py new file mode 100644 index 0000000..8d59c9c --- /dev/null +++ b/mediafusion_scrapy/pipelines/redis_cache_pipeline.py @@ -0,0 +1,26 @@ +import redis.asyncio as redis_async +from scrapy import signals +from scrapy.exceptions import DropItem + +from db.config import settings + + +class RedisCacheURLPipeline: + def __init__(self): + self.redis = redis_async.Redis.from_url(settings.redis_url) + + async def close(self): + await self.redis.aclose() + + @classmethod + def from_crawler(cls, crawler): + p = cls() + crawler.signals.connect(p.close, signal=signals.spider_closed) + return p + + async def process_item(self, item, spider): + if "webpage_url" not in item: + raise DropItem(f"webpage_url not found in item: {item}") + + await self.redis.sadd(item["scraped_url_key"], item["webpage_url"]) + return item diff --git a/mediafusion_scrapy/pipelines/sport_video_parser_pipeline.py b/mediafusion_scrapy/pipelines/sport_video_parser_pipeline.py new file mode 100644 index 0000000..62eb29f --- /dev/null +++ b/mediafusion_scrapy/pipelines/sport_video_parser_pipeline.py @@ -0,0 +1,86 @@ +import re +from datetime import datetime + +from itemadapter import ItemAdapter +from scrapy.exceptions import DropItem + + +class SportVideoParserPipeline: + RESOLUTIONS = { + "3840x2160": "4K", + "2560x1440": "1440p", + "1920x1080": "1080p", + "1280x720": "720p", + "854x480": "480p", + "640x360": "360p", + "426x240": "240p", + } + + def __init__(self): + self.title_regex = re.compile(r"^.*?\s(\d{2}\.\d{2}\.\d{4})") + + def process_item(self, item, spider): + adapter = ItemAdapter(item) + if "title" not in adapter: + raise DropItem(f"title not found in item: {item}") + + match = self.title_regex.search(adapter["title"]) or self.title_regex.search( + adapter["torrent_name"] + ) + if match: + date_str = match.group(1) + date_obj = datetime.strptime(date_str, "%d.%m.%Y") + else: + date_obj = datetime.now() + + # Try to match or infer resolution + raw_resolution = adapter.get("aspect_ratio", "").replace(" ", "") + resolution = self.infer_resolution(raw_resolution) + + if "video_codec" in adapter: + codec = re.sub(r"\s+", "", adapter["video_codec"]).replace(",", " - ") + item["codec"] = codec + + if "length" in adapter or "lenght" in adapter: + item["runtime"] = re.sub( + r"\s+", "", adapter.get("length", adapter.get("lenght", "")) + ) + + item.update( + { + "created_at": date_obj.strftime("%Y-%m-%d"), + "year": date_obj.year, + "languages": [re.sub(r"[ .]", "", item["language"])], + "is_imdb": False, + "resolution": resolution, + } + ) + return item + + def infer_resolution(self, aspect_ratio): + # Normalize the aspect_ratio string by replacing common variants of "x" with the standard one + normalized_aspect_ratio = aspect_ratio.replace("×", "x").replace("х", "x") + + # Direct match in RESOLUTIONS dictionary + resolution = self.RESOLUTIONS.get(normalized_aspect_ratio) + if resolution: + return resolution + + # Attempt to infer based on height + height = ( + normalized_aspect_ratio.split("x")[-1] + if "x" in normalized_aspect_ratio + else None + ) + if height: + for res, label in self.RESOLUTIONS.items(): + if res.endswith(height): + return label + + # Fallback to checking against common labels + for label in ["4K", "2160p", "1440p", "1080p", "720p", "480p", "360p", "240p"]: + if label in normalized_aspect_ratio: + return label + + # If no match found + return None diff --git a/mediafusion_scrapy/pipelines/store_pipelines.py b/mediafusion_scrapy/pipelines/store_pipelines.py new file mode 100644 index 0000000..c86c2b4 --- /dev/null +++ b/mediafusion_scrapy/pipelines/store_pipelines.py @@ -0,0 +1,216 @@ +import asyncio +import logging +from uuid import uuid4 + +import logging +from uuid import uuid4 + +import redis.asyncio as redis_async +from beanie import WriteRules +from scrapy import signals +from scrapy.exceptions import DropItem + +from db import crud +from db.config import settings +from db.models import ( + TorrentStreams, + Season, + MediaFusionSeriesMetaData, +) +from db.schemas import TVMetaData + + +class QueueBasedPipeline: + def __init__(self): + self.queue = asyncio.Queue() + self.processing_task = None + + async def init(self): + self.processing_task = asyncio.create_task(self.process_queue()) + + async def close(self): + await self.queue.join() + self.processing_task.cancel() + + @classmethod + def from_crawler(cls, crawler): + p = cls() + crawler.signals.connect(p.init, signal=signals.spider_opened) + crawler.signals.connect(p.close, signal=signals.spider_closed) + return p + + async def process_item(self, item, spider): + # Instead of processing the item directly, add it to the queue + await self.queue.put((item, spider)) + return item + + async def process_queue(self): + logging.info("Starting processing queue") + while True: + item, spider = await self.queue.get() + try: + await self.parse_item(item, spider) + except Exception as e: + logging.error(f"Error processing item: {e}", exc_info=True) + finally: + self.queue.task_done() + + async def parse_item(self, item, spider): + raise NotImplementedError + + +class EventSeriesStorePipeline(QueueBasedPipeline): + def __init__(self): + super().__init__() + self.redis = redis_async.Redis.from_url(settings.redis_url) + + async def close(self): + await super().close() + await self.redis.aclose() + + async def parse_item(self, item, spider): + if "title" not in item: + logging.warning(f"title not found in item: {item}") + raise DropItem(f"title not found in item: {item}") + + series = await MediaFusionSeriesMetaData.find_one( + {"title": item["title"]}, fetch_links=True + ) + + if not series: + meta_id = f"mf{uuid4().fields[-1]}" + poster = item.get("poster") + background = item.get("background") + + # Create an initial entry for the series + series = MediaFusionSeriesMetaData( + id=meta_id, + title=item["title"], + year=item["year"], + poster=poster, + background=background, + streams=[], + is_poster_working=bool(poster), + is_add_title_to_poster=item.get("is_add_title_to_poster", False), + ) + await series.insert() + logging.info("Added series %s", series.title) + + meta_id = series.id + + stream = next((s for s in series.streams if s.id == item["info_hash"]), None) + if stream is None: + # Create the stream + stream = TorrentStreams( + id=item["info_hash"], + torrent_name=item["torrent_name"], + announce_list=item["announce_list"], + size=item["total_size"], + languages=item["languages"], + resolution=item.get("resolution"), + codec=item.get("codec"), + quality=item.get("quality"), + audio=item.get("audio"), + encoder=item.get("encoder"), + source=item["source"], + catalog=item["catalog"], + created_at=item["created_at"], + season=Season(season_number=1, episodes=item["episodes"]), + meta_id=meta_id, + seeders=item["seeders"], + ) + # Add the stream to the series + series.streams.append(stream) + logging.info( + "Added stream %s to series %s", stream.torrent_name, series.title + ) + + self.organize_episodes(series) + + await series.save(link_rule=WriteRules.WRITE) + logging.info("Updated series %s", series.title) + await self.redis.sadd(item["scraped_info_hash_key"], item["info_hash"]) + + return item + + def organize_episodes(self, series): + # Flatten all episodes from all streams and sort by release date + all_episodes = sorted( + ( + episode + for stream in series.streams + for episode in stream.season.episodes + ), + key=lambda e: ( + e.released.date(), + e.filename, + ), # Sort primarily by released date, then by filename + ) + + # Assign episode numbers, ensuring the same title across different qualities gets the same number + episode_number = 1 + last_title = None + for episode in all_episodes: + if episode.title != last_title: + if last_title: + episode_number = episode_number + 1 + last_title = episode.title + + episode.episode_number = episode_number + + # Now distribute episodes back to their respective streams, ensuring they are in the correct order + for stream in series.streams: + stream.season.episodes.sort( + key=lambda e: e.episode_number + ) # Ensure episodes are ordered by episode number + + +class TVStorePipeline(QueueBasedPipeline): + async def parse_item(self, item, spider): + if "title" not in item: + logging.warning(f"title not found in item: {item}") + raise DropItem(f"title not found in item: {item}") + + tv_metadata = TVMetaData.model_validate(item) + await crud.save_tv_channel_metadata(tv_metadata) + return item + + +class MovieStorePipeline(QueueBasedPipeline): + async def parse_item(self, item, spider): + if "title" not in item: + raise DropItem(f"title not found in item: {item}") + + if item.get("type") != "movie": + return item + + await crud.save_movie_metadata(item, item.get("is_imdb", True)) + return item + + +class SeriesStorePipeline(QueueBasedPipeline): + async def parse_item(self, item, spider): + if "title" not in item: + raise DropItem(f"title not found in item: {item}") + + if item.get("type") != "series": + return item + await crud.save_series_metadata(item) + return item + + +class LiveEventStorePipeline(QueueBasedPipeline): + def __init__(self): + super().__init__() + self.redis = redis_async.Redis.from_url(settings.redis_url) + + async def close(self): + await super().close() + await self.redis.aclose() + + async def parse_item(self, item, spider): + if "title" not in item: + raise DropItem(f"name not found in item: {item}") + + await crud.save_events_data(self.redis, item) + return item diff --git a/mediafusion_scrapy/pipelines/torrent_parser_pipeline.py b/mediafusion_scrapy/pipelines/torrent_parser_pipeline.py new file mode 100644 index 0000000..43b06b3 --- /dev/null +++ b/mediafusion_scrapy/pipelines/torrent_parser_pipeline.py @@ -0,0 +1,73 @@ +import scrapy +from itemadapter import ItemAdapter +from scrapy.exceptions import DropItem +from scrapy.http.request import NO_CALLBACK +from scrapy.utils.defer import maybe_deferred_to_future + +from db.crud import is_torrent_stream_exists +from utils import torrent + + +class TorrentDownloadAndParsePipeline: + async def process_item(self, item, spider): + adapter = ItemAdapter(item) + torrent_link = adapter.get("torrent_link") + + if not torrent_link: + raise DropItem(f"No torrent link found in item: {item}") + + response = await maybe_deferred_to_future( + spider.crawler.engine.download( + scrapy.Request(torrent_link, callback=NO_CALLBACK) + ) + ) + + if response.status != 200: + spider.logger.error( + f"Failed to download torrent file: {response.url} with status {response.status}" + ) + raise DropItem(f"Failed to download torrent file: {response.url}") + + # Validate the content-type of the response + if "application/x-bittorrent" not in response.headers.get( + "Content-Type", b"" + ).decode("utf-8", "ignore"): + spider.logger.error( + f"Unexpected Content-Type for {response.url}: {response.headers.get('Content-Type')}" + ) + raise DropItem(f"Unexpected Content-Type for {response.url}") + + torrent_metadata = torrent.extract_torrent_metadata( + response.body, item.get("is_parse_ptn", True) + ) + + if not torrent_metadata: + raise DropItem(f"Failed to extract torrent metadata: {item}") + + item.update(torrent_metadata) + return item + + +class MagnetDownloadAndParsePipeline: + async def process_item(self, item, spider): + adapter = ItemAdapter(item) + magnet_link = adapter.get("magnet_link") + + if not magnet_link: + raise DropItem(f"No magnet link found in item: {item}") + + info_hash, trackers = torrent.parse_magnet(magnet_link) + if not info_hash: + raise DropItem(f"Failed to parse info_hash from magnet link: {magnet_link}") + if await is_torrent_stream_exists(info_hash): + raise DropItem(f"Torrent stream already exists: {info_hash}") + + torrent_metadata = await torrent.info_hashes_to_torrent_metadata( + [info_hash], trackers + ) + + if not torrent_metadata: + raise DropItem(f"Failed to extract torrent metadata: {item}") + + item.update(torrent_metadata[0]) + return item