mirror of
https://github.com/Viren070/MediaFusion.git
synced 2025-12-01 23:21:11 +01:00
Fix prowlarr scrape based on the magnet link & fix movie contain series (#317)
This commit is contained in:
+1
-1
@@ -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|$|[\]._-])"
|
||||
)
|
||||
|
||||
|
||||
+32
-21
@@ -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
|
||||
|
||||
|
||||
+78
-47
@@ -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}"
|
||||
|
||||
+22
-4
@@ -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"),
|
||||
|
||||
+9
-9
@@ -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)
|
||||
|
||||
+2
-3
@@ -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"],
|
||||
|
||||
Reference in New Issue
Block a user