From 6d5a35dd714d476aa7de230df70ed93f8b120ba7 Mon Sep 17 00:00:00 2001 From: mhdzumair Date: Thu, 1 Feb 2024 19:07:46 +0530 Subject: [PATCH] enhance prowlarr scrapping for movies and series and refactoring helper functions --- db/config.py | 3 +- db/crud.py | 101 ++++----- scrappers/helpers.py | 103 +++++++-- scrappers/prowlarr.py | 465 ++++++++++++++++++++++++++++++++--------- scrappers/torrentio.py | 106 ++++------ utils/parser.py | 6 +- utils/torrent.py | 2 + 7 files changed, 541 insertions(+), 245 deletions(-) diff --git a/db/config.py b/db/config.py index 63669a2..2f2b608 100644 --- a/db/config.py +++ b/db/config.py @@ -16,7 +16,8 @@ class Settings(BaseSettings): premiumize_oauth_client_id: str | None = None premiumize_oauth_client_secret: str | None = None prowlarr_url: str = "http://prowlarr-service:9696" - prowlarr_api_key: str + prowlarr_api_key: str | None = None + prowlarr_search_interval_hour: int = 6 class Config: env_file = ".env" diff --git a/db/crud.py b/db/crud.py index a14f087..1042ff1 100644 --- a/db/crud.py +++ b/db/crud.py @@ -21,17 +21,15 @@ from db.models import ( ) from db.schemas import Stream, MetaIdProjection from scrappers import tamilmv -from scrappers.prowlarr import scrap_streams_from_prowlarr -from scrappers.torrentio import ( - scrap_streams_from_torrentio, - scrap_series_streams_from_torrentio, -) +from scrappers.prowlarr import get_streams_from_prowlarr +from scrappers.torrentio import get_streams_from_torrentio from utils.parser import ( parse_stream_data, get_catalogs, search_imdb, parse_tv_stream_data, fetch_downloaded_info_hashes, + get_imdb_data, ) @@ -131,37 +129,19 @@ async def get_movie_streams( .to_list() ) - if settings.is_scrap_from_torrentio: - last_torrentio_stream = next( - (stream for stream in streams if stream.source == "Torrentio"), None + streams = await get_streams_from_torrentio( + user_data, streams, video_id, "movie", background_tasks + ) + if settings.prowlarr_api_key: + movie_metadata = await get_movie_data_by_id(video_id, False) + if movie_metadata: + title, year = movie_metadata.title, movie_metadata.year + else: + title, year = get_imdb_data(video_id) + streams = await get_streams_from_prowlarr( + user_data, streams, video_id, "movie", title, year, background_tasks ) - - if ( - video_id.startswith("tt") - and "torrentio_streams" in user_data.selected_catalogs - ): - if ( - last_torrentio_stream is None - or last_torrentio_stream.updated_at < datetime.now() - timedelta(days=3) - ): - streams.extend( - await scrap_streams_from_torrentio( - video_id, "movie", background_tasks - ) - ) - - last_prowlarr_stream = next( - (stream for stream in streams if stream.source == "Prowlarr"), None - ) - if "prowlarr_streams" in user_data.selected_catalogs: - if ( - last_prowlarr_stream is None - or last_prowlarr_stream.updated_at < datetime.now() - timedelta(hours=6) - ): - streams.extend( - await scrap_streams_from_prowlarr(video_id, "movie", background_tasks) - ) return await parse_stream_data(streams, user_data, secret_str) @@ -184,32 +164,30 @@ async def get_series_streams( ) if settings.is_scrap_from_torrentio: - last_torrentio_stream = next( - ( - stream - for stream in streams - if stream.source == "Torrentio" and stream.get_episode(season, episode) - ), - None, + streams = await get_streams_from_torrentio( + user_data, streams, video_id, "series", background_tasks, season, episode + ) + if settings.prowlarr_api_key: + series_metadata = await get_series_data_by_id(video_id, False) + if series_metadata: + title, year = series_metadata.title, series_metadata.year + else: + title, year = get_imdb_data(video_id) + streams = await get_streams_from_prowlarr( + user_data, + streams, + video_id, + "series", + title, + year, + background_tasks, + season, + episode, ) - if ( - video_id.startswith("tt") - and "torrentio_streams" in user_data.selected_catalogs - ): - if ( - last_torrentio_stream is None - or last_torrentio_stream.updated_at < datetime.now() - timedelta(days=3) - ): - streams.extend( - await scrap_series_streams_from_torrentio( - video_id, "series", season, episode, background_tasks - ) - ) - - matched_episode_streams = [ - stream for stream in streams if stream.get_episode(season, episode) - ] + matched_episode_streams = filter( + lambda stream: stream.get_episode(season, episode), streams + ) return await parse_stream_data( matched_episode_streams, user_data, secret_str, season, episode @@ -334,9 +312,6 @@ async def save_movie_metadata(metadata: dict): background = existing_movie.background meta_id = existing_movie.id - # Determine file index for the main movie file (largest file) - largest_file = max(metadata["file_data"], key=lambda x: x["size"]) - if "language" in metadata: languages = ( [metadata["language"]] @@ -352,8 +327,8 @@ async def save_movie_metadata(metadata: dict): torrent_name=metadata["torrent_name"], announce_list=metadata["announce_list"], size=metadata["total_size"], - filename=largest_file["filename"], - file_index=largest_file["index"], + filename=metadata["largest_file"]["filename"], + file_index=metadata["largest_file"]["index"], languages=languages, resolution=metadata.get("resolution"), codec=metadata.get("codec"), diff --git a/scrappers/helpers.py b/scrappers/helpers.py index 5af281d..873a974 100644 --- a/scrappers/helpers.py +++ b/scrappers/helpers.py @@ -1,5 +1,6 @@ import json import logging +from datetime import datetime import cloudscraper import httpx @@ -8,8 +9,8 @@ from requests.adapters import HTTPAdapter from urllib3.util.retry import Retry from db.config import settings -from utils.torrent import extract_torrent_metadata - +from db.models import TorrentStreams, Episode, Season +from utils.torrent import extract_torrent_metadata, info_hashes_to_torrent_metadata UA_HEADER = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/118.0.0.0 Safari/537.36" @@ -53,19 +54,7 @@ async def get_page_content(page, url): return await page.content() -async def download_and_save_torrent( - torrent_element, - metadata: dict, - media_type: str, - page_link: str, - scraper=None, - page=None, -): - from db import crud # Avoid circular import - - torrent_link = torrent_element.get("href") - logging.info(f"Downloading torrent: {torrent_link}") - +async def download_torrent(torrent_link, scraper=None, page=None): if scraper: response = scraper.get(torrent_link) torrent_metadata = extract_torrent_metadata(response.content) @@ -82,18 +71,24 @@ async def download_and_save_torrent( if not torrent_metadata: logging.error(f"Info hash not found for {torrent_link}") - return False + return None + return torrent_metadata + + +async def save_torrent(torrent_metadata: dict, metadata: dict, media_type: str): metadata.update(torrent_metadata) if not metadata.get("year"): - logging.error(f"Year not found for {page_link}") + logging.error("Year not found") return False + from db import crud # Avoid circular import + # Saving the metadata if media_type == "series": if not metadata.get("season"): - logging.error(f"Season not found for {page_link}") + logging.error("Season not found") return False await crud.save_series_metadata(metadata) else: @@ -105,6 +100,26 @@ async def download_and_save_torrent( return True +async def download_and_save_torrent( + torrent_element, + metadata: dict, + media_type: str, + page_link: str, + scraper=None, + page=None, +): + torrent_link = torrent_element.get("href") + logging.info(f"Downloading torrent: {torrent_link}") + torrent_metadata = await download_torrent(torrent_link, scraper, page) + if torrent_metadata is None: + return False + + result = await save_torrent(torrent_metadata, metadata, media_type) + if not result: + logging.error(f"Failed to save torrent for {page_link}") + return result + + def get_scrapper_config(site_name: str, get_key: str) -> dict: with open("resources/json/scrapper_config.json") as file: config = json.load(file) @@ -123,3 +138,55 @@ async def add_to_bitsearch(magnet_link: str): logging.error( f"Failed to add magnet link {magnet_link}: {response.status_code}" ) + + +async def update_torrent_movie_streams_metadata(info_hashes: list[str]): + """Update torrent streams metadata.""" + streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, []) + + for stream_metadata in streams_metadata: + if not stream_metadata: + continue + + torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"]) + if torrent_stream: + torrent_stream.torrent_name = stream_metadata["torrent_name"] + torrent_stream.size = stream_metadata["total_size"] + + torrent_stream.filename = stream_metadata["largest_file"]["filename"] + torrent_stream.file_index = stream_metadata["largest_file"]["index"] + torrent_stream.updated_at = datetime.now() + await torrent_stream.save() + + +async def update_torrent_series_streams_metadata(info_hashes: list[str]): + """Update torrent streams metadata.""" + streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, []) + + for stream_metadata in streams_metadata: + if not stream_metadata: + continue + + torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"]) + if torrent_stream: + episodes = [ + Episode( + episode_number=file["episode"], + filename=file["filename"], + size=file["size"], + file_index=file["index"], + ) + for file in stream_metadata["file_data"] + if file["episode"] + ] + torrent_stream.season = Season( + season_number=stream_metadata["season"], + episodes=episodes, + ) + + torrent_stream.torrent_name = stream_metadata["torrent_name"] + torrent_stream.size = stream_metadata["total_size"] + + torrent_stream.updated_at = datetime.now() + + await torrent_stream.save() diff --git a/scrappers/prowlarr.py b/scrappers/prowlarr.py index 16d544a..6fd9b7b 100644 --- a/scrappers/prowlarr.py +++ b/scrappers/prowlarr.py @@ -1,137 +1,406 @@ -import os -import re -from datetime import datetime +import asyncio +import logging +from datetime import datetime, timedelta import PTN import httpx from fastapi import BackgroundTasks +from torf import Magnet -from db import crud from db.config import settings -from db.models import TorrentStreams -from utils.parser import convert_size_to_bytes, get_imdb_title -from utils.torrent import info_hashes_to_torrent_metadata +from db.models import TorrentStreams, Season, Episode +from db.schemas import UserData +from scrappers.helpers import ( + update_torrent_series_streams_metadata, + update_torrent_movie_streams_metadata, +) +from utils.torrent import extract_torrent_metadata -async def fetch_stream_data(url: str, params: dict, headers: dict) -> dict: +async def get_streams_from_prowlarr( + user_data: UserData, + streams: list[TorrentStreams], + video_id: str, + catalog_type: str, + title: str, + year: str, + background_tasks: BackgroundTasks, + season: int = None, + episode: int = None, +): + last_stream = next( + (stream for stream in streams if "prowlarr_streams" in stream.catalog), None + ) + if video_id.startswith("tt") and "prowlarr_streams" in user_data.selected_catalogs: + if last_stream is None or last_stream.updated_at < datetime.now() - timedelta( + hours=settings.prowlarr_search_interval_hour + ): + if catalog_type == "movie": + streams.extend( + await scrap_movies_streams_from_prowlarr( + video_id, title, year, background_tasks + ) + ) + elif catalog_type == "series": + streams.extend( + await scrap_series_streams_from_prowlarr( + video_id, title, background_tasks, season, episode + ) + ) + return streams + + +async def fetch_stream_data(url: str, params: dict) -> dict: """Fetch stream data asynchronously.""" + headers = { + "accept": "application/json", + "X-Api-Key": settings.prowlarr_api_key, + } async with httpx.AsyncClient() as client: - response = await client.get(url, timeout=100, params=params, headers=headers) + response = await client.get(url, timeout=10, params=params, headers=headers) response.raise_for_status() # Will raise an exception for 4xx/5xx responses return response.json() -async def scrap_streams_from_prowlarr( - video_id: str, catalog_type: str, background_tasks: BackgroundTasks = None, season: int = None, episode: int = None +async def scrap_movies_streams_from_prowlarr( + video_id: str, + title: str, + year: str, + background_tasks: BackgroundTasks = None, ) -> list[TorrentStreams]: """ - Get streams by IMDb ID from prowlarr. + Get movie streams by IMDb ID from prowlarr. """ url = f"{settings.prowlarr_url}/api/v1/search" stream_data = [] - prefixes = ["{" + f"ImdbId:{video_id}" + "} "] - if catalog_type == "tv": - prefixes = ["{" + f"Season:{season}" + "} ", - "{" + f"Season:{season}" + "} {" + f"Episode:{episode}" + "} "] - for prefix in prefixes: - headers = { - 'accept': 'application/json', - 'X-Api-Key': settings.prowlarr_api_key, - } - params = { - 'query': f"{prefix}{get_imdb_title(video_id.removeprefix('tt'))}", - 'categories': [ - '2000', # Movies - '5000', # TV - '8000', # Other - ], - 'type': "tvsearch" if catalog_type == "tv" else "moviesearch", - 'limit': '20', - 'offset': '0', - } - try: - stream_data.extend(await fetch_stream_data(url, params, headers)) - except (httpx.HTTPError, httpx.TimeoutException): - return [] # Return empty list in case of HTTP errors or timeouts + # Params for IMDb ID search + params_imdb = { + "query": f"{{ImdbId:{video_id}}}", + "categories": [2000], # Movies + "type": "movie", + "limit": 20, + "offset": 0, + } - return await store_and_parse_stream_data( - video_id, stream_data, background_tasks + # Params for title search + params_title = { + "query": title, + "categories": [2000], # Movies + "type": "search", + "limit": 20, + "offset": 0, + } + + try: + # Fetch data for both searches simultaneously + imdb_search, title_search = await asyncio.gather( + fetch_stream_data(url, params_imdb), + fetch_stream_data(url, params_title), + ) + stream_data.extend(imdb_search) + stream_data.extend(title_search) + except (httpx.HTTPError, httpx.TimeoutException): + return [] # Return an empty list in case of HTTP errors or timeouts + + return await store_and_parse_movie_stream_data( + video_id, title, year, stream_data, background_tasks ) -async def store_and_parse_stream_data( - video_id: str, stream_data: list, background_tasks: BackgroundTasks +async def scrap_series_streams_from_prowlarr( + video_id: str, + title: str, + background_tasks: BackgroundTasks = None, + season: int = None, + episode: int = None, +) -> list[TorrentStreams]: + """ + Get series streams by IMDb ID from prowlarr. + """ + url = f"{settings.prowlarr_url}/api/v1/search" + stream_data = [] + + # Params for IMDb ID, season, and episode search + params_imdb = { + "query": f"{{ImdbId:{video_id}}}{{Season:{season}}}{{Episode:{episode}}}", + "categories": [5000], # TV + "type": "tvsearch", + "limit": 20, + "offset": 0, + } + + # Params for title search + params_title = { + "query": title, + "categories": [5000], # TV + "type": "search", + "limit": 20, + "offset": 0, + } + + try: + # Fetch data for both searches simultaneously + imdb_search, title_search = await asyncio.gather( + fetch_stream_data(url, params_imdb), + fetch_stream_data(url, params_title), + ) + stream_data.extend(imdb_search) + stream_data.extend(title_search) + except (httpx.HTTPError, httpx.TimeoutException): + return [] # Return an empty list in case of HTTP errors or timeouts + + return await store_and_parse_series_stream_data( + video_id, title, season, episode, stream_data, background_tasks + ) + + +async def get_torrent_data_from_prowlarr(download_url: str) -> tuple[dict, bool]: + """Get torrent data from prowlarr.""" + if not download_url: + return {}, False + if download_url.startswith("magnet:"): + magnet = Magnet.from_string(download_url) + return {"info_hash": magnet.infohash, "announce_list": magnet.tr}, False + + async with httpx.AsyncClient() as client: + response = await client.get(download_url, follow_redirects=False) + + if response.status_code == 301: + magnet_url = response.headers.get("Location") + magnet = Magnet.from_string(magnet_url) + return {"info_hash": magnet.infohash, "announce_list": magnet.tr}, False + elif response.status_code == 200: + return extract_torrent_metadata(response.content), True + else: + return {}, False + + +async def prowlarr_data_parser(meta_data: dict) -> tuple[dict, bool]: + """Parse prowlarr data.""" + try: + torrent_data, is_torrent_downloaded = await get_torrent_data_from_prowlarr( + meta_data.get("downloadUrl") or meta_data.get("magnetUrl") + ) + except Exception as e: + logging.warning(f"Error parsing torrent data: {e} {e.__class__.__name__}") + if meta_data.get("infoHash"): + torrent_data = { + "info_hash": meta_data.get("infoHash"), + "announce_list": [], + } + is_torrent_downloaded = False + else: + return {}, False + + torrent_data.update( + { + "seeders": meta_data.get("seeders"), + "created_at": datetime.strptime( + meta_data.get("publishDate"), "%Y-%m-%dT%H:%M:%SZ" + ), + "source": meta_data.get("indexer"), + "poster_url": meta_data.get("posterUrl"), + } + ) + if is_torrent_downloaded is False: + torrent_data.update( + { + "torrent_name": meta_data.get("title"), + "total_size": meta_data.get("size"), + **PTN.parse(meta_data.get("title")), + } + ) + return torrent_data, is_torrent_downloaded + + +async def store_and_parse_movie_stream_data( + video_id: str, + title: str, + year: str, + stream_data: list, + background_tasks: BackgroundTasks, ) -> list[TorrentStreams]: streams = [] info_hashes = [] for stream in stream_data: - infohash = stream.get("infoHash", "").lower() - if infohash: - torrent_stream = None # TODO ADD ME LATER await TorrentStreams.get(infohash) - if torrent_stream: - # Update existing stream - torrent_stream.seeders = stream["seeders"] - torrent_stream.updated_at = datetime.now() - # TODO ADD ME LATER - await torrent_stream.save() - else: - title = stream["title"] - metadata = PTN.parse(title) - languages = [] - language = metadata.get("language", "") - languages.extend(language if type(language) == list else [language]) - # Create new stream - torrent_stream = TorrentStreams( - id=infohash, - torrent_name=title, - announce_list=[], - size=stream["size"], - filename=None, - file_index=stream.get("indexerId"), - languages=languages, - resolution=metadata.get("resolution"), - codec=metadata.get("codec"), - quality=metadata.get("quality"), - audio=metadata.get("audio"), - encoder=metadata.get("encoder"), - source="Prowlarr", - catalog=["prowlarr_streams"], - updated_at=datetime.now(), - seeders=stream["seeders"], - meta_id=video_id, - ) - # TODO ADD ME LATER - await torrent_stream.save() + parsed_data, _ = await prowlarr_data_parser(stream) + info_hash = parsed_data.get("info_hash", "").lower() + if not info_hash or not parsed_data.get("seeders"): + logging.warning( + f"Skipping {info_hash} due to missing info_hash or seeders: {parsed_data.get('seeders')}" + ) + continue - streams.append(torrent_stream) - if torrent_stream.filename is None: - info_hashes.append(infohash) + if not ( + parsed_data.get("title").lower() == title.lower() + and parsed_data.get("year") == year + ) and (str(stream.get("imdbId", "")) not in video_id): + logging.warning( + f"Skipping {info_hash} due to title mismatch: '{parsed_data.get('title')}' != '{title}' or year mismatch: '{parsed_data.get('year')}' != '{year}'" + ) + continue - # TODO ADD ME LATER - background_tasks.add_task(update_torrent_streams_metadata, info_hashes) + torrent_stream = await TorrentStreams.get(info_hash) + + if torrent_stream: + # Update existing stream + torrent_stream.seeders = parsed_data["seeders"] + torrent_stream.updated_at = datetime.now() + await torrent_stream.save() + else: + # Create new stream + torrent_stream = TorrentStreams( + id=info_hash, + torrent_name=parsed_data.get("torrent_name"), + announce_list=parsed_data.get("announce_list"), + size=parsed_data.get("total_size"), + filename=parsed_data.get("largest_file", {}).get("file_name"), + file_index=parsed_data.get("largest_file", {}).get("index"), + languages=parsed_data.get("languages", []), + resolution=parsed_data.get("resolution"), + codec=parsed_data.get("codec"), + quality=parsed_data.get("quality"), + audio=parsed_data.get("audio"), + encoder=parsed_data.get("encoder"), + source=parsed_data.get("source"), + catalog=[ + "prowlarr_streams", + "prowlarr_movies", + f"{parsed_data.get('source').lower()}_movies", + ], + updated_at=datetime.now(), + seeders=parsed_data.get("seeders"), + created_at=parsed_data.get("created_at"), + meta_id=video_id, + ) + await torrent_stream.save() + + streams.append(torrent_stream) + if torrent_stream.filename is None: + info_hashes.append(info_hash) + + background_tasks.add_task(update_torrent_movie_streams_metadata, info_hashes) return streams -async def update_torrent_streams_metadata(info_hashes: list[str]): - """Update torrent streams metadata.""" - streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, []) - - for stream_metadata in streams_metadata: - if not stream_metadata: +async def store_and_parse_series_stream_data( + video_id: str, + title: str, + season: int, + episode: int, + stream_data: list, + background_tasks: BackgroundTasks, +) -> list[TorrentStreams]: + streams = [] + info_hashes = [] + for stream in stream_data: + parsed_data, is_torrent_downloaded = await prowlarr_data_parser(stream) + info_hash = parsed_data.get("info_hash", "").lower() + if not info_hash or not parsed_data.get("seeders"): + logging.warning( + f"Skipping {info_hash} due to missing info_hash or seeders: {parsed_data.get('seeders')}" + ) continue - torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"]) - if torrent_stream: - torrent_stream.torrent_name = stream_metadata.get("torrent_name") - torrent_stream.size = stream_metadata.get("total_size") + if not (parsed_data.get("title").lower() == title.lower()) and ( + str(stream.get("imdbId", "")) not in video_id + ): + logging.warning( + f"Skipping {info_hash} due to title mismatch: '{parsed_data.get('title')}' != '{title}'" + ) + continue - largest_file = max(stream_metadata.get("file_data"), key=lambda x: x["size"]) - torrent_stream.filename = largest_file["filename"] - torrent_stream.file_index = largest_file["index"] + if parsed_data.get("season"): + if isinstance(parsed_data["season"], int): + season_number = parsed_data["season"] + else: + # Skip This Stream due to multiple seasons in one torrent. + # TODO: Handle this case later. + # Need to refactor DB and how streaming provider works. + continue + else: + season_number = season + + if is_torrent_downloaded is True: + episode_data = [ + Episode( + episode_number=file["episode"], + filename=file["filename"], + size=file["size"], + file_index=file["index"], + ) + for file in parsed_data["file_data"] + if file["episode"] + ] + elif parsed_data.get("episode"): + if isinstance(parsed_data["episode"], int): + episode_data = [Episode(episode_number=parsed_data["episode"])] + else: + episode_data = [ + Episode( + episode_number=episode_number, + ) + for episode_number in parsed_data["episode"] + ] + + else: + episode_data = [Episode(episode_number=episode)] + + torrent_stream = await TorrentStreams.get(info_hash) + + if torrent_stream: + torrent_stream.seeders = parsed_data["seeders"] torrent_stream.updated_at = datetime.now() - torrent_stream.resolution = stream_metadata.get("resolution") - torrent_stream.quality = stream_metadata.get("quality") - torrent_stream.codec = stream_metadata.get("codec") + episode_item = torrent_stream.get_episode(season, episode) + if episode_item is None: + if torrent_stream.season: + torrent_stream.season.episodes.extend(episode_data) + else: + torrent_stream.season = Season( + season_number=season_number, + episodes=episode_data, + ) + episode_item = torrent_stream.get_episode(season, episode) await torrent_stream.save() + else: + # Create new stream + torrent_stream = TorrentStreams( + id=info_hash, + torrent_name=parsed_data.get("torrent_name"), + announce_list=parsed_data.get("announce_list"), + size=parsed_data.get("total_size"), + filename=None, + languages=parsed_data.get("languages", []), + resolution=parsed_data.get("resolution"), + codec=parsed_data.get("codec"), + quality=parsed_data.get("quality"), + audio=parsed_data.get("audio"), + encoder=parsed_data.get("encoder"), + source=parsed_data.get("source"), + catalog=[ + "prowlarr_streams", + "prowlarr_series", + f"{parsed_data.get('source').lower()}_series", + ], + updated_at=datetime.now(), + seeders=parsed_data.get("seeders"), + created_at=parsed_data.get("created_at"), + meta_id=video_id, + season=Season( + season_number=season_number, + episodes=episode_data, + ), + ) + await torrent_stream.save() + episode_item = torrent_stream.get_episode(season, episode) + + streams.append(torrent_stream) + + if episode_item and episode_item.size is None: + info_hashes.append(info_hash) + + background_tasks.add_task(update_torrent_series_streams_metadata, info_hashes) + + return streams diff --git a/scrappers/torrentio.py b/scrappers/torrentio.py index 58b78df..b105516 100644 --- a/scrappers/torrentio.py +++ b/scrappers/torrentio.py @@ -1,6 +1,6 @@ import logging import re -from datetime import datetime +from datetime import datetime, timedelta from os import path import PTN @@ -9,12 +9,47 @@ from fastapi import BackgroundTasks from db.config import settings from db.models import TorrentStreams, Season, Episode -from scrappers.helpers import UA_HEADER +from db.schemas import UserData +from scrappers.helpers import ( + UA_HEADER, + update_torrent_series_streams_metadata, + update_torrent_movie_streams_metadata, +) from utils.parser import convert_size_to_bytes -from utils.torrent import info_hashes_to_torrent_metadata from utils.validation_helper import is_video_file +async def get_streams_from_torrentio( + user_data: UserData, + streams: list[TorrentStreams], + video_id: str, + catalog_type: str, + background_tasks: BackgroundTasks, + season: int = None, + episode: int = None, +): + last_stream = next( + (stream for stream in streams if stream.source == "Torrentio"), None + ) + if video_id.startswith("tt") and "torrentio_streams" in user_data.selected_catalogs: + if last_stream is None or last_stream.updated_at < datetime.now() - timedelta( + days=3 + ): + if catalog_type == "movie": + streams.extend( + await scrap_movie_streams_from_torrentio( + video_id, catalog_type, background_tasks + ) + ) + elif catalog_type == "series": + streams.extend( + await scrap_series_streams_from_torrentio( + video_id, catalog_type, season, episode, background_tasks + ) + ) + return streams + + async def fetch_stream_data(url: str) -> dict: """Fetch stream data asynchronously.""" async with httpx.AsyncClient( @@ -25,7 +60,7 @@ async def fetch_stream_data(url: str) -> dict: return response.json() -async def scrap_streams_from_torrentio( +async def scrap_movie_streams_from_torrentio( video_id: str, catalog_type: str, background_tasks: BackgroundTasks ) -> list[TorrentStreams]: """ @@ -34,11 +69,11 @@ async def scrap_streams_from_torrentio( url = f"{settings.torrentio_url}/stream/{catalog_type}/{video_id}.json" try: stream_data = await fetch_stream_data(url) - return await store_and_parse_stream_data( + return await store_and_parse_movie_stream_data( video_id, stream_data.get("streams", []), background_tasks ) except (httpx.HTTPError, httpx.TimeoutException): - return [] # Return empty list in case of HTTP errors or timeouts + return [] # Return an empty list in case of HTTP errors or timeouts except Exception as e: logging.error(f"Error while fetching stream data from torrentio: {e}") return [] @@ -61,7 +96,7 @@ async def scrap_series_streams_from_torrentio( video_id, season, episode, stream_data.get("streams", []), background_tasks ) except (httpx.HTTPError, httpx.TimeoutException): - return [] # Return empty list in case of HTTP errors or timeouts + return [] # Return an empty list in case of HTTP errors or timeouts except Exception as e: logging.error(f"Error while fetching stream data from torrentio: {e}") return [] @@ -82,7 +117,7 @@ def parse_stream_title(stream: dict) -> dict: } -async def store_and_parse_stream_data( +async def store_and_parse_movie_stream_data( video_id: str, stream_data: list, background_tasks: BackgroundTasks ) -> list[TorrentStreams]: streams = [] @@ -126,7 +161,7 @@ async def store_and_parse_stream_data( if torrent_stream.filename is None: info_hashes.append(stream["infoHash"]) - background_tasks.add_task(update_torrent_streams_metadata, info_hashes) + background_tasks.add_task(update_torrent_movie_streams_metadata, info_hashes) return streams @@ -268,56 +303,3 @@ def extract_size_string(details: str) -> str: """Extract the size string from the details.""" size_match = re.search(r"💾 (\d+(?:\.\d+)?\s*(GB|MB))", details, re.IGNORECASE) return size_match.group(1) if size_match else "" - - -async def update_torrent_streams_metadata(info_hashes: list[str]): - """Update torrent streams metadata.""" - streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, []) - - for stream_metadata in streams_metadata: - if not stream_metadata: - continue - - torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"]) - if torrent_stream: - torrent_stream.torrent_name = stream_metadata["torrent_name"] - torrent_stream.size = stream_metadata["total_size"] - - largest_file = max(stream_metadata["file_data"], key=lambda x: x["size"]) - torrent_stream.filename = largest_file["filename"] - torrent_stream.file_index = largest_file["index"] - torrent_stream.updated_at = datetime.now() - await torrent_stream.save() - - -async def update_torrent_series_streams_metadata(info_hashes: list[str]): - """Update torrent streams metadata.""" - streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, []) - - for stream_metadata in streams_metadata: - if not stream_metadata: - continue - - torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"]) - if torrent_stream: - episodes = [ - Episode( - episode_number=file["episode"], - filename=file["filename"], - size=file["size"], - file_index=file["index"], - ) - for file in stream_metadata["file_data"] - if file["episode"] - ] - torrent_stream.season = Season( - season_number=stream_metadata["season"], - episodes=episodes, - ) - - torrent_stream.torrent_name = stream_metadata["torrent_name"] - torrent_stream.size = stream_metadata["total_size"] - - torrent_stream.updated_at = datetime.now() - - await torrent_stream.save() diff --git a/utils/parser.py b/utils/parser.py index e4bf5b1..c8e0f68 100644 --- a/utils/parser.py +++ b/utils/parser.py @@ -221,9 +221,9 @@ def search_imdb(title: str, year: int, retry: int = 5) -> dict: return {} -def get_imdb_title(video_id:str) -> str: - movie = ia.get_movie(video_id) - return movie.get("title") +def get_imdb_data(video_id: str) -> tuple[str, str]: + movie = ia.get_movie(video_id.removeprefix("tt")) + return movie.get("title"), movie.get("year") def parse_tv_stream_data(tv_data: MediaFusionTVMetaData) -> list[Stream]: diff --git a/utils/torrent.py b/utils/torrent.py index 0c9de9b..1911b56 100644 --- a/utils/torrent.py +++ b/utils/torrent.py @@ -70,6 +70,7 @@ def extract_torrent_metadata(content: Torrent | bytes) -> dict: "episode": parsed_data.get("episode"), } ) + largest_file = max(file_data, key=lambda x: x["size"]) return { **PTN.parse(torrent.name), @@ -78,6 +79,7 @@ def extract_torrent_metadata(content: Torrent | bytes) -> dict: "total_size": total_size, "file_data": file_data, "torrent_name": torrent.name, + "largest_file": largest_file, } except Exception as e: logging.error(f"Error occurred: {e}")