Add support for scraping debrid cache status from Torrentio & Upstream MediaFusion

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