diff --git a/db/config.py b/db/config.py index 90d463b..bad8ce8 100644 --- a/db/config.py +++ b/db/config.py @@ -16,6 +16,7 @@ class Settings(BaseSettings): prowlarr_url: str = "http://prowlarr-service:9696" zilean_url: str = "http://zilean.zilean:8181" playwright_cdp_url: str = "ws://browserless:3000?blockAds=true&stealth=true" + flaresolverr_url: str = "http://flaresolverr:8191/v1" # External API keys and secrets secret_key: str = Field(..., max_length=32, min_length=32) @@ -65,7 +66,7 @@ class Settings(BaseSettings): disable_validate_tv_streams_in_db: bool = False sport_video_scheduler_crontab: str = "*/20 * * * *" disable_sport_video_scheduler: bool = False - streamed_scheduler_crontab: str = "*/15 * * * *" + streamed_scheduler_crontab: str = "*/30 * * * *" disable_streamed_scheduler: bool = False mrgamingstreams_scheduler_crontab: str = "*/15 * * * *" disable_mrgamingstreams_scheduler: bool = True # Disabled due to site being down. diff --git a/db/crud.py b/db/crud.py index 0d05167..32d0d40 100644 --- a/db/crud.py +++ b/db/crud.py @@ -116,7 +116,7 @@ async def get_tv_meta_list( # Define query filters for TVStreams query_filters = { "is_working": True, - "namespace": {"$in": [namespace, "mediafusion", None]}, + "namespaces": {"$in": [namespace, "mediafusion", None]}, } poster_path = f"{settings.poster_host_url}/poster/tv/" @@ -428,7 +428,7 @@ async def get_tv_streams(video_id: str, namespace: str, user_data) -> list[Strea { "meta_id": video_id, "is_working": True, - "namespace": {"$in": [namespace, "mediafusion", None]}, + "namespaces": {"$in": [namespace, "mediafusion", None]}, }, ).to_list() @@ -848,7 +848,7 @@ async def process_tv_search_query(search_query: str, namespace: str) -> dict: { "$match": { "is_working": True, - "namespace": {"$in": [namespace, "mediafusion", None]}, + "namespaces": {"$in": [namespace, "mediafusion", None]}, } }, {"$count": "num_working_streams"}, @@ -907,11 +907,25 @@ async def save_tv_channel_metadata(tv_metadata: schemas.TVMetaData) -> str: # Ensure the channel document is upserted try: - await MediaFusionTVMetaData.find_one( - MediaFusionTVMetaData.id == channel_id - ).upsert( - Set({}), # Update operation is a no-op for now - on_insert=MediaFusionTVMetaData( + channel_data = await MediaFusionTVMetaData.get(channel_id) + if channel_data: + if channel_data.is_poster_working is False: + background = ( + {"background": tv_metadata.background} + if tv_metadata.background + else {} + ) + await channel_data.update( + Set( + { + "poster": tv_metadata.poster, + "is_poster_working": True, + **background, + } + ) + ) + else: + channel_data = MediaFusionTVMetaData( id=channel_id, title=tv_metadata.title, poster=tv_metadata.poster, @@ -922,13 +936,13 @@ async def save_tv_channel_metadata(tv_metadata: schemas.TVMetaData) -> str: genres=genres, type="tv", streams=[], - ), - ) + ) + await channel_data.create() except DuplicateKeyError: pass # Stream processing - stream_ids = [] + bulk_writer = BulkWriter() for stream in tv_metadata.streams: # Define stream document with meta_id stream_doc = TVStreams( @@ -943,7 +957,9 @@ async def save_tv_channel_metadata(tv_metadata: schemas.TVMetaData) -> str: source=stream.source, country=stream.country, meta_id=channel_id, - namespace=tv_metadata.namespace, + namespaces=[tv_metadata.namespace], + drm_key_id=stream.drm_key_id, + drm_key=stream.drm_key, ) # Check if the stream exists (by URL or ytId) and upsert accordingly @@ -951,29 +967,29 @@ async def save_tv_channel_metadata(tv_metadata: schemas.TVMetaData) -> str: TVStreams.url == stream.url, TVStreams.ytId == stream.ytId, ) - if existing_stream: - stream_ids.append(existing_stream.id) - else: - inserted_stream = await stream_doc.insert() - stream_ids.append(inserted_stream.id) - - # Update the TV channel with new stream links, if there are any new streams - if stream_ids: - await MediaFusionTVMetaData.find_one( - MediaFusionTVMetaData.id == channel_id - ).update( - { - "$addToSet": { - "streams": { - "$each": [ - TVStreams.link_from_id(stream_id) - for stream_id in stream_ids - ] - } + if existing_stream == stream_doc: + update_data = {} + if ( + stream.drm_key_id != existing_stream.drm_key_id + or stream.drm_key != existing_stream.drm_key + ) and tv_metadata.namespace in existing_stream.namespaces: + update_data = { + "drm_key_id": stream.drm_key_id, + "drm_key": stream.drm_key, } - } - ) + if tv_metadata.namespace not in existing_stream.namespaces: + existing_stream.namespaces.append(tv_metadata.namespace) + update_data.update({"namespaces": existing_stream.namespaces}) + if update_data: + await existing_stream.update( + Set(update_data), + bulk_writer=bulk_writer, + ) + else: + await TVStreams.insert_one(stream_doc, bulk_writer=bulk_writer) + + await bulk_writer.commit() logging.info(f"Processed TV channel {tv_metadata.title}") return channel_id @@ -991,24 +1007,17 @@ async def save_events_data(metadata: dict) -> str: existing_event_data = MediaFusionEventsMetaData.model_validate_json( existing_event_json ) - # Use a dictionary keyed by 'url' to ensure uniqueness of streams - existing_streams = { - stream.url or stream.ytId or stream.externalUrl: stream - for stream in existing_event_data.streams - } + existing_streams = set(existing_event_data.streams) else: - existing_streams = {} + existing_streams = set() # Update or add streams based on the uniqueness of 'url' for stream in metadata["streams"]: # Create a TVStreams instance for each stream stream_instance = TVStreams(meta_id=meta_id, **stream) - existing_streams[ - stream_instance.url or stream_instance.ytId or stream_instance.externalUrl - ] = stream_instance.model_dump() + existing_streams.add(stream_instance) - # Update the event metadata with the updated list of streams - streams = list(existing_streams.values()) + streams = list(existing_streams) event_start_timestamp = metadata.get("event_start_timestamp", 0) diff --git a/db/models.py b/db/models.py index 1d7311d..06c2e8b 100644 --- a/db/models.py +++ b/db/models.py @@ -104,7 +104,38 @@ class TVStreams(Document): country: str | None = None is_working: Optional[bool] = True test_failure_count: int = 0 - namespace: str = "mediafusion" + namespaces: list[str] = Field(default_factory=lambda: ["mediafusion"]) + drm_key_id: str | None = None + drm_key: str | None = None + + def __eq__(self, other): + if not isinstance(other, TVStreams): + return False + return ( + self.url == other.url + and self.ytId == other.ytId + and self.externalUrl == other.externalUrl + and self.drm_key_id == other.drm_key_id + and self.drm_key == other.drm_key + ) + + def __hash__(self): + return hash( + (self.url, self.ytId, self.externalUrl, self.drm_key_id, self.drm_key) + ) + + class Settings: + indexes = [ + IndexModel( + [ + ("meta_id", ASCENDING), + ("created_at", DESCENDING), + ("namespaces", ASCENDING), + ("is_working", ASCENDING), + ] + ), + IndexModel([("url", pymongo.TEXT), ("ytId", pymongo.TEXT)]), + ] class MediaFusionMetaData(Document): diff --git a/db/schemas.py b/db/schemas.py index 7392b44..1132851 100644 --- a/db/schemas.py +++ b/db/schemas.py @@ -258,6 +258,11 @@ class MetaIdProjection(BaseModel): id: str = Field(alias="_id") +class TVMetaProjection(BaseModel): + id: str = Field(alias="_id") + title: str + + class TVStreams(BaseModel): name: str url: str | None = None @@ -265,6 +270,8 @@ class TVStreams(BaseModel): source: str country: str | None = None behaviorHints: StreamBehaviorHints | None = None + drm_key_id: str | None = None + drm_key: str | None = None @model_validator(mode="after") def validate_url_or_yt_id(self) -> "TVStreams": @@ -282,7 +289,7 @@ class TVMetaData(BaseModel): logo: Optional[str] = None genres: list[str] = [] streams: list[TVStreams] - namespace: str = "mediafusion" + namespace: str = Field(default="mediafusion") class TorrentStreamsList(BaseModel): diff --git a/mediafusion_scrapy/middlewares.py b/mediafusion_scrapy/middlewares.py index daec773..8a6b964 100644 --- a/mediafusion_scrapy/middlewares.py +++ b/mediafusion_scrapy/middlewares.py @@ -2,12 +2,14 @@ # # See documentation in: # https://docs.scrapy.org/en/latest/topics/spider-middleware.html +import asyncio import random +from urllib.parse import urlparse -# useful for handling different item types with a single interface - -from scrapy import signals +import httpx +from scrapy import signals, Request from scrapy.downloadermiddlewares.retry import RetryMiddleware +from scrapy.exceptions import IgnoreRequest from scrapy.utils.response import response_status_message from twisted.internet import reactor, defer @@ -128,7 +130,7 @@ class TooManyRequestsRetryMiddleware(RetryMiddleware): """ DEFAULT_DELAY = 30 # Default initial delay in seconds. - MAX_DELAY = 300 # Max delay between retries. + MAX_DELAY = 180 # Max delay between retries. BACKOFF_FACTOR = 3 # Exponential backoff factor. def __init__(self, settings): @@ -162,14 +164,14 @@ class TooManyRequestsRetryMiddleware(RetryMiddleware): # Calculate the delay with exponential backoff retry_after = response.headers.get("retry-after") try: - retry_after = int(retry_after) + retry_after = int(retry_after) + random.randint(1, 10) except (ValueError, TypeError): delay = min( self.MAX_DELAY, self.DEFAULT_DELAY * (self.BACKOFF_FACTOR**retries) ) else: delay = min( - self.MAX_DELAY, retry_after + random.randint(0, 10) + self.MAX_DELAY, retry_after + random.randint(5, 120) ) # Add some jitter spider.logger.info( @@ -191,3 +193,102 @@ class TooManyRequestsRetryMiddleware(RetryMiddleware): return deferred return response + + +class FlaresolverrMiddleware: + + def __init__(self, flaresolverr_url, cache_duration, max_timeout, max_attempts): + self.flaresolverr_url = flaresolverr_url + self.cache_duration = cache_duration + self.max_timeout = max_timeout + self.max_attempts = max_attempts + self.solved_domains = {} + self.client = httpx.AsyncClient() + + @classmethod + def from_crawler(cls, crawler): + return cls( + flaresolverr_url=crawler.settings.get( + "FLARESOLVERR_URL", "http://localhost:8191/v1" + ), + cache_duration=crawler.settings.get("FLARESOLVERR_CACHE_DURATION", 3600), + max_timeout=crawler.settings.get("FLARESOLVERR_MAX_TIMEOUT", 60000), + max_attempts=crawler.settings.get("FLARESOLVERR_MAX_ATTEMPTS", 3), + ) + + async def process_request(self, request, spider): + return None + + async def process_response(self, request, response, spider): + if not hasattr(spider, "use_flaresolverr") or not spider.use_flaresolverr: + return response + + if response.status == 403 or ( + response.status == 503 and "cloudflare" in response.text.lower() + ): + return await self._handle_cloudflare(request, spider) + + return response + + async def _handle_cloudflare(self, request, spider): + domain = urlparse(request.url).netloc + current_time = reactor.seconds() + + if domain in self.solved_domains: + last_solved_time, solution = self.solved_domains[domain] + if current_time - last_solved_time < self.cache_duration: + return self._apply_solution(request, solution) + + for attempt in range(self.max_attempts): + try: + timeout = min( + self.max_timeout, 30000 * (2**attempt) + ) # Exponential backoff + flaresolverr_response = await self.client.post( + self.flaresolverr_url, + headers={"Content-Type": "application/json"}, + json={ + "cmd": "request.get", + "url": request.url, + "maxTimeout": timeout, + }, + timeout=timeout / 1000 + 5, + ) + + if flaresolverr_response.status_code == 200: + solution = flaresolverr_response.json() + if solution.get("status") == "ok": + self.solved_domains[domain] = (current_time, solution) + return self._apply_solution(request, solution) + + spider.logger.error( + f"FlareSolverr attempt {attempt + 1} failed: {flaresolverr_response.text}" + ) + await asyncio.sleep(2**attempt) # Wait before next attempt + + except httpx.RequestError as e: + spider.logger.error( + f"FlareSolverr request error on attempt {attempt + 1}: {e}" + ) + await asyncio.sleep(2**attempt) + + spider.logger.error( + f"Failed to solve Cloudflare challenge for {request.url} after {self.max_attempts} attempts" + ) + raise IgnoreRequest() + + def _apply_solution(self, original_request, solution): + solution_response = solution.get("solution", {}).get("response", {}) + return Request( + url=original_request.url, + headers=solution_response.get("headers", {}), + cookies={ + cookie["name"]: cookie["value"] + for cookie in solution_response.get("cookies", []) + }, + dont_filter=True, + meta={"flaresolverr_solved": True, **original_request.meta}, + ) + + async def close_spider(self, spider): + await self.client.aclose() diff --git a/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py b/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py index afe5a93..8a97772 100644 --- a/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py +++ b/mediafusion_scrapy/pipelines/live_stream_resolver_pipeline.py @@ -32,7 +32,7 @@ class LiveStreamResolverPipeline: ) content_type = response.headers.get("Content-Type", b"").decode().lower() - if response.status == 200 and content_type in const.M3U8_VALID_CONTENT_TYPES: + if response.status == 200 and content_type in const.IPTV_VALID_CONTENT_TYPES: stream_headers.update( { "User-Agent": response.request.headers.get("User-Agent").decode(), diff --git a/mediafusion_scrapy/settings.py b/mediafusion_scrapy/settings.py index 3c4ef67..7513d32 100644 --- a/mediafusion_scrapy/settings.py +++ b/mediafusion_scrapy/settings.py @@ -6,6 +6,7 @@ # https://docs.scrapy.org/en/latest/topics/settings.html # https://docs.scrapy.org/en/latest/topics/downloader-middleware.html # https://docs.scrapy.org/en/latest/topics/spider-middleware.html +from db.config import settings BOT_NAME = "mediafusion_scrapy" @@ -51,6 +52,7 @@ SPIDER_MIDDLEWARES = { # Enable or disable downloader middlewares # See https://docs.scrapy.org/en/latest/topics/downloader-middleware.html DOWNLOADER_MIDDLEWARES = { + "mediafusion_scrapy.middlewares.FlaresolverrMiddleware": 542, "mediafusion_scrapy.middlewares.TooManyRequestsRetryMiddleware": 543, } @@ -107,5 +109,8 @@ RETRY_HTTP_CODES = [ 522, 524, 408, + 403, ] # 429 is handled by the middleware RETRY_TIMES = 5 + +FLARESOLVERR_URL = settings.flaresolverr_url diff --git a/mediafusion_scrapy/spiders/live_tv.py b/mediafusion_scrapy/spiders/live_tv.py index a5b9780..0581411 100644 --- a/mediafusion_scrapy/spiders/live_tv.py +++ b/mediafusion_scrapy/spiders/live_tv.py @@ -1,5 +1,5 @@ import re -from urllib.parse import urljoin, urlparse +from urllib.parse import urljoin, urlparse, parse_qs import scrapy @@ -9,7 +9,13 @@ from utils import const class LiveTVSpider(scrapy.Spider): fallback_pattern = re.compile( - r"source: ['\"](.*?)['\"],\s*[\s\S]*?mimeType: ['\"]application/x-mpegURL['\"]" + r"source: ['\"](.*?)['\"],\s*[\s\S]*?mimeType: ['\"]application/(x-mpegURL|vnd\.apple\.mpegURL|dash\+xml)['\"]", + re.IGNORECASE, + ) + + any_m3u8_pattern = re.compile( + r'(?:"|\')?(https?://.*?\.m3u8(?:\?[^"\']*)?)(?:"|\')?', + re.IGNORECASE, ) # this site sometimes returns html instead of image @@ -92,7 +98,7 @@ class LiveTVSpider(scrapy.Spider): def parse_channel_page(self, response): channel_data = self.extract_channel_data(response) - player_api_base = self.extract_player_api_base(response) + player_api_base, player_api_method = self.extract_player_api_base(response) if not player_api_base: self.logger.error(f"Player API base URL not found for {response.url}") return @@ -106,6 +112,7 @@ class LiveTVSpider(scrapy.Spider): meta={ "channel_data": channel_data, "player_api_base": player_api_base, + "player_api_method": player_api_method, "response": response, # Pass the original response to access player options later }, dont_filter=True, @@ -116,6 +123,7 @@ class LiveTVSpider(scrapy.Spider): original_response = meta["response"] channel_data = meta["channel_data"] player_api_base = meta["player_api_base"] + player_api_method = meta["player_api_method"] content_type = response.headers.get("Content-Type", b"").decode().lower() is_image = "image" in content_type @@ -125,7 +133,7 @@ class LiveTVSpider(scrapy.Spider): if is_image or is_allowed_url: yield from self.process_player_options( - original_response, channel_data, player_api_base + original_response, channel_data, player_api_base, player_api_method ) else: self.logger.error(f"Invalid poster URL: {response.url}") @@ -146,12 +154,26 @@ class LiveTVSpider(scrapy.Spider): return channel_data def extract_player_api_base(self, response): - """Extracts the player API base URL.""" - return ( - response.xpath("//script[contains(text(), 'player_api')]/text()") - .re_first(r'"player_api":"([^"]+)"', default="") - .replace("\\/", "/") - ) + """Extracts the admin-ajax URL or player API base for GET or POST requests.""" + # Directly extract the URL used for admin-ajax POST requests + admin_ajax = response.xpath("//script[contains(text(), 'player_api')]/text()") + admin_ajax_method = admin_ajax.re_first(r'"play_method"\s*:\s*"([^"]+)"') + if admin_ajax_method == "wp_json": + admin_ajax_url = admin_ajax.re_first(r'"player_api"\s*:\s*"([^"]+)"') + elif admin_ajax_method == "admin_ajax": + admin_ajax_url = admin_ajax.re_first(r'"url"\s*:\s*"([^"]+)"') + else: + admin_ajax_url = None + + if admin_ajax_url: + # Correctly format and return the full URL + admin_ajax_full_url = urljoin( + response.url, admin_ajax_url.replace("\\/", "/") + ) + return admin_ajax_full_url, admin_ajax_method + else: + self.logger.error("Admin AJAX URL not found for TamilUltra.") + return None, None def extract_stream_details(self, element): """Extracts stream title and country name from an element.""" @@ -163,13 +185,17 @@ class LiveTVSpider(scrapy.Spider): country_name = get_country_name(country_code) return stream_title, country_name - def process_player_options(self, response, channel_data, player_api_base): + def process_player_options( + self, response, channel_data, player_api_base, player_api_method + ): for element in response.css("#playeroptionsul > li.dooplay_player_option"): yield from self.process_player_option( - element, channel_data, player_api_base + element, channel_data, player_api_base, player_api_method ) - def process_player_option(self, element, channel_data, player_api_base): + def process_player_option( + self, element, channel_data, player_api_base, player_api_method + ): """Processes each player option element to yield API request.""" stream_title, country_name = self.extract_stream_details(element) data_post, data_nume, data_type = ( @@ -179,16 +205,34 @@ class LiveTVSpider(scrapy.Spider): ) if all([data_post, data_nume, data_type]): - api_url = f"{player_api_base}{data_post}/{data_type}/{data_nume}" - yield scrapy.Request( - url=api_url, - callback=self.parse_api_response, - meta={ - "channel_data": channel_data, - "stream_title": stream_title, - "country_name": country_name, - }, - ) + if player_api_method == "wp_json": + api_url = f"{player_api_base}{data_post}/{data_type}/{data_nume}" + yield scrapy.Request( + url=api_url, + callback=self.parse_api_response, + meta={ + "channel_data": channel_data, + "stream_title": stream_title, + "country_name": country_name, + }, + ) + else: + form_data = { + "action": "doo_player_ajax", + "post": data_post, + "nume": data_nume, + "type": data_type, + } + yield scrapy.FormRequest( + url=player_api_base, + formdata=form_data, + callback=self.parse_api_response, + meta={ + "channel_data": channel_data, + "stream_title": stream_title, + "country_name": country_name, + }, + ) def parse_api_response(self, response): channel_data = response.meta.get("channel_data") @@ -213,17 +257,95 @@ class LiveTVSpider(scrapy.Spider): }, ) - def extract_m3u8_urls(self, response): - """Extracts M3U8 URLs using direct and fallback regex patterns.""" + def extract_m3u8_or_mpd_urls(self, response) -> tuple[list[dict], dict]: + """Extracts M3U8 and MPD URLs with appropriate metadata.""" parsed_url = urlparse(response.url) - channel_id = parsed_url.query.split("=")[-1] + parsed_query = parse_qs(parsed_url.query) - m3u8_urls = re.findall( - rf"{re.escape(channel_id)}['\"]:\s*{{\s*url:\s*['\"](.*?\.m3u8.*?)['\"]", - response.text, - ) - if not m3u8_urls: - m3u8_urls = self.fallback_pattern.findall(response.text) + if ( + response.headers.get("Content-Type", b"").decode().lower() + in const.IPTV_VALID_CONTENT_TYPES[:2] + ): + # If the content type is M3U8, return the URL directly + return [{"url": response.url, "type": "m3u8"}], {} + + stream_data = [] + + if "source" in parsed_query: + stream_data.append( + { + "url": urljoin(response.url, parsed_query["source"][0]), + "type": "unknown", # We'll determine the type later + } + ) + elif "zy" in parsed_query and ".mpd``" in parsed_query["zy"][0]: + data = parsed_query["zy"][0] + url, key_data = data.split("``") + drm_key_id, drm_key = key_data.split(":") + stream_data.append( + { + "url": url, + "type": "mpd", + "drm_key_id": drm_key_id, + "drm_key": drm_key, + } + ) + elif "tamilultra" in response.url: + query_string = parsed_url.query + stream_data.append( + { + "url": urljoin(response.url, query_string), + "type": "m3u8", + } + ) + + else: + channel_id = parsed_query.get("id", [""])[0] + + # Pattern to match both M3U8 and MPD URLs + pattern = rf"{re.escape(channel_id)}['\"]:\s*{{\s*['\"]?url['\"]?\s*:\s*['\"](.*?)['\"]" + matches = re.findall(pattern, response.text, re.DOTALL) + + if not matches: + matches = self.fallback_pattern.findall(response.text) + + if not matches: + matches = self.any_m3u8_pattern.findall(response.text) + + for match in matches: + url = match[0] if isinstance(match, tuple) else match + stream_info = { + "url": url, + "type": "unknown", # We'll determine the type later + } + + # Extract clearkeys if it's an MPD stream + if url.endswith(".mpd"): + drm_data = self.extract_drm_keys(response.text, channel_id) + if not drm_data: + self.logger.error( + f"Clearkeys not found for MPD URL: {url} on channel page: {response.url}" + ) + continue + stream_info.update(drm_data) + + stream_data.append(stream_info) + + if not stream_data: + self.logger.error( + "No M3U8 or MPD URLs found for channel page: %s", + response.meta["channel_data"]["channel_page_url"], + ) + return [], {} + + # Determine stream types + for stream in stream_data: + if stream["url"].endswith(".mpd"): + stream["type"] = "mpd" + elif any(ext in stream["url"] for ext in [".m3u8", ".m3u"]): + stream["type"] = "m3u8" + else: + stream["type"] = "m3u8" user_agent = response.request.headers.get("User-Agent").decode() parsed_url = urlparse(response.url) @@ -239,34 +361,37 @@ class LiveTVSpider(scrapy.Spider): }, } - return m3u8_urls, behavior_hints + return stream_data, behavior_hints def request_and_extract_video_url(self, response): channel_data = response.meta.get("channel_data") stream_title = response.meta.get("stream_title") country_name = response.meta.get("country_name") - # Extract M3U8 URLs - m3u8_urls, behavior_hints = self.extract_m3u8_urls(response) - if not m3u8_urls: + # Extract M3U8 & MPD URLs and behavior hints + streams_data, behavior_hints = self.extract_m3u8_or_mpd_urls(response) + if not streams_data: self.logger.error( - "No M3U8 URLs found for channel url: %s, stream title: %s", + "No M3U8 or MPD URLs found for channel url: %s, stream title: %s", channel_data["channel_page_url"], stream_title, ) return - for index, url in enumerate(m3u8_urls, 1): + for index, stream_data in enumerate(streams_data, 1): + url = stream_data["url"] full_url = urljoin(response.url, url) # Instead of appending to streams_info, initiate validation request yield scrapy.Request( url=full_url, - headers=behavior_hints["proxyHeaders"]["request"], - callback=self.validate_m3u8_url, - errback=self.handle_m3u8_failure, + headers=behavior_hints.get("proxyHeaders", {}).get("request", {}), + callback=self.validate_m3u8_or_mpd_url, + errback=self.handle_m3u8_or_mpd_failure, meta={ - "index": index if len(m3u8_urls) > 1 else None, + "index": index if len(streams_data) > 1 else None, "stream_title": stream_title, + "drm_key_id": stream_data.get("drm_key_id"), + "drm_key": stream_data.get("drm_key"), "full_url": full_url, "country_name": country_name, "channel_data": channel_data, @@ -275,17 +400,21 @@ class LiveTVSpider(scrapy.Spider): dont_filter=True, ) - def validate_m3u8_url(self, response): + def validate_m3u8_or_mpd_url(self, response): meta = response.meta content_type = response.headers.get("Content-Type", b"").decode().lower() - if response.status == 200 and content_type in const.M3U8_VALID_CONTENT_TYPES: + if response.status == 200 and content_type in const.IPTV_VALID_CONTENT_TYPES: # Content type is valid, proceed with adding the stream stream_info = { - "name": f"{meta['stream_title']} - {meta['index']}" - if meta["index"] - else meta["stream_title"], + "name": ( + f"{meta['stream_title']} - {meta['index']}" + if meta["index"] + else meta["stream_title"] + ), "url": meta["full_url"], + "drm_key_id": meta["drm_key_id"], + "drm_key": meta["drm_key"], "country": meta["country_name"], "behaviorHints": meta["behavior_hints"], "source": meta["channel_data"]["source"], @@ -300,9 +429,9 @@ class LiveTVSpider(scrapy.Spider): f"Invalid M3U8 URL: {meta['full_url']} with Content-Type: {content_type}" ) - def handle_m3u8_failure(self, failure): + def handle_m3u8_or_mpd_failure(self, failure): self.logger.error( - "Failed to get m3u8 URL from channel page: %s stream title: %s", + "Failed to get m3u8 or MPD URL from channel page: %s stream title: %s", failure.request.meta["channel_data"]["channel_page_url"], failure.request.meta["stream_title"], ) @@ -317,3 +446,35 @@ class LiveTVSpider(scrapy.Spider): channel_data_copy["country"] = country_name return channel_data_copy + + @staticmethod + def extract_drm_keys(response_text: str, channel_id: str) -> dict: + stream_info = {} + + # Pattern for channel entry + channel_pattern = rf'"{re.escape(channel_id)}"\s*:\s*{{[^}}]+}}' + channel_match = re.search(channel_pattern, response_text, re.DOTALL) + + if channel_match: + channel_data = channel_match.group(0) + + # Pattern for clearkeys + clearkey_pattern = ( + r'["\']?clearkeys["\']?\s*:\s*{\s*["\'](.+?)["\']\s*:\s*["\'](.+?)["\']' + ) + clearkey_match = re.search(clearkey_pattern, channel_data, re.DOTALL) + + # Pattern for k1 and k2 + k1k2_pattern = r'["\']?k1["\']?\s*:\s*["\'](.+?)["\'],\s*["\']?k2["\']?\s*:\s*["\'](.+?)["\']' + k1k2_match = re.search(k1k2_pattern, channel_data) + + if clearkey_match: + key_id, key = clearkey_match.groups() + stream_info["drm_key_id"] = key_id + stream_info["drm_key"] = key + elif k1k2_match: + key_id, key = k1k2_match.groups() + stream_info["drm_key_id"] = key_id + stream_info["drm_key"] = key + + return stream_info diff --git a/mediafusion_scrapy/spiders/streamed.py b/mediafusion_scrapy/spiders/streamed.py index 9eaaf22..10b6125 100644 --- a/mediafusion_scrapy/spiders/streamed.py +++ b/mediafusion_scrapy/spiders/streamed.py @@ -1,34 +1,48 @@ +import json import random import re - -import scrapy from datetime import datetime +from urllib.parse import urljoin + +import pytz +import scrapy from utils.runtime_const import SPORTS_ARTIFACTS class StreamedSpider(scrapy.Spider): name = "streamed" - allowed_domains = ["streamed.su"] - categories = { - "American Football": "https://streamed.su/category/american-football", - "Basketball": "https://streamed.su/category/basketball", - "Baseball": "https://streamed.su/category/baseball", - "Cricket": "https://streamed.su/category/cricket", - "Football": "https://streamed.su/category/football", - "Fighting": "https://streamed.su/category/fight", - "Hockey": "https://streamed.su/category/hockey", - "Tennis": "https://streamed.su/category/tennis", - "Rugby": "https://streamed.su/category/rugby", - "Golf": "https://streamed.su/category/golf", - "Dart": "https://streamed.su/category/darts", - "Afl": "https://streamed.su/category/afl", - "Motor Sport": "https://streamed.su/category/motor-sports", - "Other Sports": "https://streamed.su/category/other", + + api_base_url = "https://streamed.su/api" + live_matches_url = f"{api_base_url}/matches/live" + stream_url_template = f"{api_base_url}/stream/{{source}}/{{id}}" + image_url_template = ( + f"{api_base_url}/images/poster/{{batch_id_1}}/{{batch_id_2}}.webp" + ) + + m3u8_base_url = "https://{domain}/{source}/js/{id}/{stream_no}/playlist.m3u8" + mediafusion_referer = "https://mediafusion.addon/" + + category_mapping = { + "afl": "Afl", + "american-football": "American Football", + "baseball": "Baseball", + "basketball": "Basketball", + "billiards": "Billiards", + "cricket": "Cricket", + "darts": "Dart", + "fight": "Fighting", + "football": "Football", + "golf": "Golf", + "hockey": "Hockey", + "motor-sports": "Motor Sport", + "other": "Other Sports", + "rugby": "Rugby", + "tennis": "Tennis", } - m3u8_base_url = "https://rr.vipstreams.in/alpha/js" - mediafusion_referer = "https://mediafusion.addon/" + domains = None + domain_host = None custom_settings = { "ITEM_PIPELINES": { @@ -39,85 +53,122 @@ class StreamedSpider(scrapy.Spider): "DEFAULT_REQUEST_HEADERS": {"Referer": mediafusion_referer}, } - def __init__(self, *args, **kwargs): - super(StreamedSpider, self).__init__(*args, **kwargs) - def start_requests(self): - for category, url in self.categories.items(): - yield scrapy.Request(url, self.parse, meta={"category": category}) + yield scrapy.Request(self.live_matches_url, self.parse_live_matches) - def parse(self, response, **kwargs): - category = response.meta["category"] - events = response.xpath('//a[contains(@href,"/watch/")]') + def parse_live_matches(self, response): + matches = response.json() + for match in matches: + if not match.get("sources"): + self.logger.info(f"No sources available for match: {match['title']}") + continue - for event in events: - event_name = event.xpath(".//h1/text()").get().strip() - event_url = event.xpath(".//@href").get() + category = self.category_mapping.get(match["category"], "Other") item = { "stream_source": "Streamed (streamed.su)", "genres": [category], - "poster": random.choice(SPORTS_ARTIFACTS[category]["poster"]), - "background": random.choice(SPORTS_ARTIFACTS[category]["background"]), - "logo": random.choice(SPORTS_ARTIFACTS[category]["logo"]), - "is_add_title_to_poster": True, - "title": event_name, - "url": response.urljoin(event_url), + "title": match["title"], + "event_start_timestamp": match["date"] / 1000, # Convert to seconds "streams": [], } - yield response.follow( - event_url, - self.parse_event, - meta={"item": item}, - headers={"Referer": self.mediafusion_referer}, - ) + # Handle poster image + if "poster" in match and match["poster"]: + item["poster"] = urljoin(self.api_base_url, match["poster"]) + elif ( + match.get("teams") + and match["teams"].get("home", {}).get("badge") + and match["teams"].get("away", {}).get("badge") + ): + item["poster"] = self.image_url_template.format( + batch_id_1=match["teams"]["home"]["badge"], + batch_id_2=match["teams"]["away"]["badge"], + ) + else: + item["poster"] = random.choice(SPORTS_ARTIFACTS[category]["poster"]) - def parse_event(self, response): - script_text = response.xpath( - '//script[contains(text(), "const data =")]/text()' - ).get() + item["background"] = random.choice(SPORTS_ARTIFACTS[category]["background"]) + item["logo"] = random.choice(SPORTS_ARTIFACTS[category]["logo"]) + item["is_add_title_to_poster"] = True - event_start_timestamp = 0 - if script_text: - # Use a regular expression to find the timestamp - timestamp_match = re.search(r"date:(\d+)", script_text) + for source in match["sources"]: + yield scrapy.Request( + self.stream_url_template.format( + source=source["source"], id=source["id"] + ), + self.parse_stream, + meta={"item": item}, + ) - if timestamp_match: - # Extract the timestamp and convert to UTC datetime - event_timestamp_ms = int(timestamp_match.group(1)) - event_start_timestamp = event_timestamp_ms / 1000 + def parse_stream(self, response): + item = response.meta["item"].copy() + stream_data_list = response.json() - if event_start_timestamp != 0: - event_start_time = datetime.fromtimestamp(event_start_timestamp).strftime( - "%I:%M%p GMT" - ) - description = f'{response.meta["item"]["title"]} - {event_start_time}' + for stream_data in stream_data_list: + if self.domains and self.domain_host: + yield self.create_stream_item(stream_data, item) + else: + yield scrapy.Request( + stream_data["embedUrl"], + self.parse_embed, + meta={"item": item, "stream_data": stream_data}, + headers={"Referer": self.mediafusion_referer}, + ) + + def parse_embed(self, response): + item = response.meta["item"] + stream_data = response.meta["stream_data"] + + if self.domains is None or self.domain_host is None: + self.extract_domain_info(response) + + if self.domains and self.domain_host: + yield self.create_stream_item(stream_data, item) else: - description = response.meta["item"]["title"] - - # If no timer, proceed to scrape available stream links - stream_links = response.xpath('//a[contains(@href, "/watch/")]') - if not stream_links: - # No streams available atm - self.logger.info(f"No streams available for this event yet. {response.url}") - return - - for link in stream_links: - stream_name = link.xpath(".//h1/text()").get().strip() - stream_url = link.xpath(".//@href").get() - stream_quality = link.xpath(".//h2/text()").get().strip() - language = link.xpath(".//div[last()]/text()").get().strip() - - m3u8_url = f"{self.m3u8_base_url}{stream_url.replace('/watch', '').replace('/alpha', '')}/playlist.m3u8" - item = response.meta["item"].copy() - item.update( - { - "stream_name": f"{stream_name}\nšŸ“ŗ {stream_quality} - 🌐 {language}", - "stream_url": m3u8_url, - "referer": self.mediafusion_referer, - "description": description, - "event_start_timestamp": event_start_timestamp, - } + self.logger.error( + f"Failed to extract domain information for stream: {stream_data['id']}" ) - yield item + + def extract_domain_info(self, response): + script_content = response.xpath( + '//script[contains(text(), "var k=")]/text()' + ).get() + if script_content: + vars_match = re.search( + r'var k="(\w+)",i="([^"]+)",s="(\d+)",l=(\[.+?\]),h="([^"]+)";', + script_content, + ) + if vars_match: + _, _, _, domains, domain_host = vars_match.groups() + self.domains = json.loads(domains) + self.domain_host = domain_host + else: + self.logger.warning( + "Failed to extract domain variables from script content." + ) + else: + self.logger.warning( + "Failed to find script content with domain variables in response." + ) + + def create_stream_item(self, stream_data, item): + m3u8_url = self.m3u8_base_url.format( + domain=f"{random.choice(self.domains)}.{self.domain_host}", + source=stream_data["source"], + id=stream_data["id"], + stream_no=stream_data["streamNo"], + ) + + stream_item = item.copy() + stream_item.update( + { + "stream_name": f"{'HD' if stream_data['hd'] else 'SD'} - 🌐 {stream_data['language']}\n" + f"šŸ”— {stream_data['source'].title()} Stream {stream_data['streamNo']}", + "stream_url": m3u8_url, + "referer": self.mediafusion_referer, + "description": f"{item['title']} - {datetime.fromtimestamp(item['event_start_timestamp'], tz=pytz.utc).strftime('%I:%M%p GMT')}", + } + ) + + return stream_item diff --git a/mediafusion_scrapy/spiders/tamilbulb.py b/mediafusion_scrapy/spiders/tamilbulb.py new file mode 100644 index 0000000..1eec8d9 --- /dev/null +++ b/mediafusion_scrapy/spiders/tamilbulb.py @@ -0,0 +1,7 @@ +from mediafusion_scrapy.spiders.live_tv import LiveTVSpider + + +class TamilBulbSpider(LiveTVSpider): + name = "tamilbulb" + start_urls = ["https://tamilbulb.tv/"] + use_flaresolverr = True diff --git a/mediafusion_scrapy/spiders/tamilultra.py b/mediafusion_scrapy/spiders/tamilultra.py index bf77f64..5696893 100644 --- a/mediafusion_scrapy/spiders/tamilultra.py +++ b/mediafusion_scrapy/spiders/tamilultra.py @@ -1,74 +1,6 @@ -from urllib.parse import urljoin, urlparse - -import scrapy - from mediafusion_scrapy.spiders.live_tv import LiveTVSpider class TamilUltraSpider(LiveTVSpider): name = "tamilultra" start_urls = ["https://tamilultra.tv/"] - - def extract_player_api_base(self, response): - """Extracts the admin-ajax URL for POST requests.""" - # Directly extract the URL used for admin-ajax POST requests - admin_ajax_url = response.xpath( - "//script[contains(text(), 'player_api')]/text()" - ).re_first(r'"url":"([^"]+)"') - if admin_ajax_url: - # Correctly format and return the full URL - admin_ajax_full_url = urljoin( - response.url, admin_ajax_url.replace("\\/", "/") - ) - return admin_ajax_full_url - else: - self.logger.error("Admin AJAX URL not found for TamilUltra.") - return None - - def process_player_option(self, element, channel_data, player_api_post_url): - """Processes each player option element to send a POST request.""" - stream_title, country_name = self.extract_stream_details(element) - data_post, data_nume, data_type = ( - element.attrib.get("data-post"), - element.attrib.get("data-nume"), - element.attrib.get("data-type"), - ) - - if all([data_post, data_nume, data_type]): - form_data = { - "action": "doo_player_ajax", - "post": data_post, - "nume": data_nume, - "type": data_type, - } - yield scrapy.FormRequest( - url=player_api_post_url, - formdata=form_data, - callback=self.parse_api_response, - meta={ - "channel_data": channel_data, - "stream_title": stream_title, - "country_name": country_name, - }, - ) - - def extract_m3u8_urls(self, response): - """Extracts M3U8 URLs using direct and fallback regex patterns.""" - query_string = urlparse(response.url).query - m3u8_urls = [urljoin(response.url, query_string)] - - user_agent = response.request.headers.get("User-Agent").decode() - parsed_url = urlparse(response.url) - referer = f"{parsed_url.scheme}://{parsed_url.netloc}" - - behavior_hints = { - "notWebReady": True, - "proxyHeaders": { - "request": { - "User-Agent": user_agent, - "Referer": referer, - } - }, - } - - return m3u8_urls, behavior_hints diff --git a/resources/exceptions/mediaflow_proxy_required.mp4 b/resources/exceptions/mediaflow_proxy_required.mp4 new file mode 100644 index 0000000..a2a45c8 Binary files /dev/null and b/resources/exceptions/mediaflow_proxy_required.mp4 differ diff --git a/resources/js/scraperControl.js b/resources/js/scraperControl.js index 38bbf91..0994221 100644 --- a/resources/js/scraperControl.js +++ b/resources/js/scraperControl.js @@ -57,7 +57,7 @@ function addStreamInput() {
- +
@@ -76,6 +76,14 @@ function addStreamInput() {
+
+ + +
+
+ + +
@@ -266,6 +274,8 @@ function constructTvMetadata() { ytId: document.getElementById(`streamYtId-${index}`).value.trim(), source: document.getElementById(`streamSource-${index}`).value.trim(), country: document.getElementById(`streamCountry-${index}`).value.trim(), + drm_key_id: document.getElementById(`streamDrmKeyId-${index}`).value.trim(), + drm_key: document.getElementById(`streamDrmKey-${index}`).value.trim(), behaviorHints: { proxyHeaders: document.getElementById(`streamProxyHeaders-${index}`) ? JSON.parse(document.getElementById(`streamProxyHeaders-${index}`).value || "{}") : {}, notWebReady: true diff --git a/resources/json/sports_artifacts.json b/resources/json/sports_artifacts.json index d252454..6376789 100644 --- a/resources/json/sports_artifacts.json +++ b/resources/json/sports_artifacts.json @@ -88,6 +88,19 @@ "https://external-content.duckduckgo.com/iu/?u=http%3A//cliparts.co/cliparts/Aib/Kr4/AibKr4gGT.png&f=1&nofb=1" ] }, + "Billiards": { + "background": [ + "https://external-content.duckduckgo.com/iu/?u=http%3A%2F%2Feskipaper.com%2Fimages%2Fbilliards-15.jpg&f=1&nofb=1&ipt=8b9307b8db760ec15d07a6ed66717495b73db55a0dfcc4bd688a901dcff463ed&ipo=images" + ], + "poster": [ + "https://external-content.duckduckgo.com/iu/?u=https%3A%2F%2Fak.picdn.net%2Fshutterstock%2Fvideos%2F9725927%2Fthumb%2F1.jpg%3Fip%3Dx480&f=1&nofb=1&ipt=36b55f7a9456378df6799f8f9edab09a45cb9c8ad7670353c30c3c10049b01c1&ipo=images", + "https://external-content.duckduckgo.com/iu/?u=https%3A%2F%2Fwallpapercave.com%2Fwp%2Fwp3177240.jpg&f=1&nofb=1&ipt=5dc29a21348a01f83b30f5b67ad204fbd1371b3109cf69eaa90e0803ba7c4f6d&ipo=images" + ], + "logo": [ + "https://external-content.duckduckgo.com/iu/?u=https%3A%2F%2Fstatic.vecteezy.com%2Fsystem%2Fresources%2Fthumbnails%2F009%2F730%2F725%2Fsmall%2Fbilliards-championship-sports-badge-design-logo-and-simple-text-billiard-room-or-pool-club-and-team-billiard-ball-icon-symbol-template-vector.jpg&f=1&nofb=1&ipt=e7e888ffd1287f424e7cd3d05bd0bce1a8dbd2f4af4cfd2d82c4a51e8875d519&ipo=images", + "https://external-content.duckduckgo.com/iu/?u=https%3A%2F%2Ftse1.mm.bing.net%2Fth%3Fid%3DOIP.D6r9x55zsgJtybnoheuHZAAAAA%26pid%3DApi&f=1&ipt=9ab2e7fd4ce8760177a69b6d493646f9142b3657b899685fed8d36f8b230b616&ipo=images" + ] + }, "Bowling": { "background": [ "https://i.pinimg.com/736x/fa/2d/21/fa2d21f56d7c44e03760a3db74a306e6.jpg" diff --git a/resources/manifest.json b/resources/manifest.json index 3c04768..1efa49c 100644 --- a/resources/manifest.json +++ b/resources/manifest.json @@ -1403,20 +1403,21 @@ "name": "genre", "isRequired": false, "options": [ - "American Football", - "Basketball", - "Baseball", - "Cricket", - "Football", - "Fighting", - "Hockey", - "Tennis", - "Rugby", - "Golf", - "Dart", "Afl", + "American Football", + "Baseball", + "Basketball", + "Billiards", + "Cricket", + "Dart", + "Fighting", + "Football", + "Golf", + "Hockey", "Motor Sport", - "Other Sports" + "Other Sports", + "Rugby", + "Tennis" ] } ] diff --git a/scrapers/tv.py b/scrapers/tv.py index fd014c5..51eefac 100644 --- a/scrapers/tv.py +++ b/scrapers/tv.py @@ -1,18 +1,21 @@ import asyncio +import difflib import logging import re import dramatiq +import httpx from beanie import BulkWriter +from beanie.odm.operators.update.general import Set from ipytv import playlist from ipytv.channel import IPTVAttr from db import schemas, crud -from db.models import TVStreams +from db.models import TVStreams, MediaFusionTVMetaData from utils import validation_helper from utils.parser import is_contain_18_plus_keywords from utils.runtime_const import REDIS_ASYNC_CLIENT -from utils.validation_helper import validate_m3u8_url +from utils.validation_helper import validate_live_stream_url async def add_tv_metadata(batch, namespace: str): @@ -133,23 +136,25 @@ async def validate_tv_streams_in_db(page=0, page_size=25, *args, **kwargs): logging.info(f"No TV streams to validate on page {page}") return - async def validate_and_update_tv_stream(stream, bulk_writer): - is_valid = await validate_m3u8_url(stream.url, stream.behaviorHints or {}) + async def validate_and_update_tv_stream(stream, bw): + is_valid = await validate_live_stream_url( + stream.url, stream.behaviorHints or {} + ) logging.info(f"Stream: {stream.name}, Status: {is_valid}") if is_valid: stream.is_working = is_valid stream.test_failure_count = 0 - await stream.replace(bulk_writer=bulk_writer) + await stream.replace(bulk_writer=bw) return stream.test_failure_count += 1 if stream.test_failure_count >= 3: - await stream.delete(bulk_writer=bulk_writer) + await stream.delete(bulk_writer=bw) logging.error(f"{stream.name} has failed 3 times and deleting it.") return stream.is_working = is_valid - await stream.replace(bulk_writer=bulk_writer) + await stream.replace(bulk_writer=bw) logging.error(f"Stream failed validation: {stream.name}") bulk_writer = BulkWriter() @@ -166,3 +171,60 @@ async def validate_tv_streams_in_db(page=0, page_size=25, *args, **kwargs): validate_tv_streams_in_db.send_with_options( args=(page + 1, page_size), delay=2 * 6000 ) + + +@dramatiq.actor(time_limit=30 * 60 * 1000, priority=5) +async def update_tv_posters_in_db(*args, **kwargs): + """Validate TV posters in the database.""" + not_working_posters = await MediaFusionTVMetaData.find( + MediaFusionTVMetaData.is_poster_working == False + ).to_list() + logging.info(f"Found {len(not_working_posters)} TV posters to update.") + + if not not_working_posters: + logging.info("No TV posters to update.") + return + + async with httpx.AsyncClient() as client: + response = await client.get( + "https://iptv-org.github.io/api/channels.json", timeout=30 + ) + iptv_channels = response.json() + + iptv_org_channel_data = {} + for channel in iptv_channels: + iptv_org_channel_data[channel["name"].casefold()] = channel + if channel.get("alt_names"): + for alt_name in channel["alt_names"]: + iptv_org_channel_data[alt_name.casefold()] = channel + iptv_org_channel_names = iptv_org_channel_data.keys() + + def get_similar_channel_name(name, cutoff=0.8): + name = name.casefold() + name = re.sub(r"\s*\[.*?]|\s*\(.*?\)", "", name) + name = name.split(" – ")[0] + matches = difflib.get_close_matches( + name, iptv_org_channel_names, n=1, cutoff=cutoff + ) + return matches[0] if matches else None + + bulk_writer = BulkWriter() + for stream in not_working_posters: + iptv_org_channel_name = get_similar_channel_name(stream.title, cutoff=0.8) + if not iptv_org_channel_name: + logging.error(f"Channel not found in iptv-org: {stream.title}") + continue + + iptv_org_channel = iptv_org_channel_data[iptv_org_channel_name] + poster = iptv_org_channel.get("logo") + if not poster: + logging.error(f"Poster not found for channel: {stream.title}") + continue + + await stream.update( + Set({"poster": poster, "is_poster_working": True}), bulk_writer=bulk_writer + ) + + logging.info(f"Committing {len(bulk_writer.operations)} updates to the database") + await bulk_writer.commit() + logging.info("Updated TV posters in the database.") diff --git a/utils/const.py b/utils/const.py index 03c5eec..5a63504 100644 --- a/utils/const.py +++ b/utils/const.py @@ -147,11 +147,12 @@ UA_HEADER = { } -M3U8_VALID_CONTENT_TYPES = [ +IPTV_VALID_CONTENT_TYPES = [ "application/vnd.apple.mpegurl", "application/x-mpegurl", "video/mp2t", "application/octet-stream", + "application/dash+xml", ] SCRAPY_SPIDERS = { diff --git a/utils/crypto.py b/utils/crypto.py index 8b1e00b..1b57ac8 100644 --- a/utils/crypto.py +++ b/utils/crypto.py @@ -1,9 +1,14 @@ +import base64 import hashlib +import json +import time + import zlib from base64 import urlsafe_b64encode, urlsafe_b64decode from Crypto.Cipher import AES from Crypto.Random import get_random_bytes +from Crypto.Util.Padding import pad from db.schemas import UserData from utils.runtime_const import SECRET_KEY @@ -45,3 +50,17 @@ def decrypt_user_data(secret_str: str | None = None) -> UserData: def get_text_hash(text: str, full_hash: bool = False) -> str: hash_str = hashlib.sha256(text.encode()).hexdigest() return hash_str if full_hash else hash_str[:10] + + +def encrypt_data( + secret_key: str, data: dict, expiration: int = None, ip: str = None +) -> str: + if expiration: + data["exp"] = int(time.time()) + expiration + if ip: + data["ip"] = ip + json_data = json.dumps(data).encode("utf-8") + iv = get_random_bytes(16) + cipher = AES.new(secret_key.encode("utf-8").ljust(32)[:32], AES.MODE_CBC, iv) + encrypted_data = cipher.encrypt(pad(json_data, AES.block_size)) + return base64.urlsafe_b64encode(iv + encrypted_data).decode("utf-8") diff --git a/utils/network.py b/utils/network.py index 643bd2c..70f3640 100644 --- a/utils/network.py +++ b/utils/network.py @@ -2,12 +2,14 @@ import asyncio import logging from typing import Callable, AsyncGenerator, Any, Tuple from urllib import parse +from urllib.parse import urlencode import httpx from fastapi.requests import Request from db.schemas import UserData from utils import crypto +from utils.crypto import encrypt_data from utils.runtime_const import PRIVATE_CIDR, REDIS_ASYNC_CLIENT @@ -253,6 +255,9 @@ def encode_mediaflow_proxy_url( query_params: dict | None = None, request_headers: dict | None = None, response_headers: dict | None = None, + encryption_api_password: str = None, + expiration: int = None, + ip: str = None, ) -> str: query_params = query_params or {} if destination_url is not None: @@ -267,8 +272,16 @@ def encode_mediaflow_proxy_url( query_params.update( {f"r_{key}": value for key, value in response_headers.items()} ) - # Encode the query parameters - encoded_params = parse.urlencode(query_params, quote_via=parse.quote) + + if encryption_api_password: + if "api_password" not in query_params: + query_params["api_password"] = encryption_api_password + encrypted_token = encrypt_data( + encryption_api_password, query_params, expiration, ip + ) + encoded_params = urlencode({"token": encrypted_token}) + else: + encoded_params = urlencode(query_params) # Construct the full URL base_url = parse.urljoin(mediaflow_proxy_url, endpoint) diff --git a/utils/parser.py b/utils/parser.py index 58097e9..f44403c 100644 --- a/utils/parser.py +++ b/utils/parser.py @@ -1,6 +1,8 @@ import asyncio import functools import logging +from typing import Optional, List + import math import re @@ -14,7 +16,7 @@ from utils import const from utils.const import STREAMING_PROVIDERS_SHORT_NAMES from utils.network import encode_mediaflow_proxy_url from utils.runtime_const import ADULT_CONTENT_KEYWORDS, TRACKERS -from utils.validation_helper import validate_m3u8_url_with_cache +from utils.validation_helper import validate_m3u8_or_mpd_url_with_cache async def filter_and_sort_streams( @@ -317,60 +319,124 @@ def convert_size_to_bytes(size_str: str) -> int: async def parse_tv_stream_data( - tv_streams: list[TVStreams], user_data: UserData -) -> list[Stream]: - stream_list = [] + tv_streams: List[TVStreams], user_data: UserData +) -> List[Stream]: is_mediaflow_proxy_enabled = ( user_data.mediaflow_config and user_data.mediaflow_config.proxy_live_streams ) addon_name = ( f"{settings.addon_name} {'šŸ•µšŸ¼ā€ā™‚ļø' if is_mediaflow_proxy_enabled else 'šŸ“”'}" ) - for stream in tv_streams[::-1]: - if settings.validate_m3u8_urls_liveness: - is_working = await validate_m3u8_url_with_cache( - stream.url, stream.behaviorHints or {} - ) - if not is_working: - continue - if is_mediaflow_proxy_enabled: - stream.url = encode_mediaflow_proxy_url( - user_data.mediaflow_config.proxy_url, - "/proxy/hls", - stream.url, - request_headers=stream.behaviorHints.get("proxyHeaders", {}).get( - "request", {} - ), - query_params={"api_password": user_data.mediaflow_config.api_password}, - ) - stream.behaviorHints.update({"proxyHeaders": {}}) + stream_processor = functools.partial( + process_stream, + is_mediaflow_proxy_enabled=is_mediaflow_proxy_enabled, + mediaflow_config=user_data.mediaflow_config, + addon_name=addon_name, + ) - country_info = f"\n🌐 {stream.country}" if stream.country else "" + processed_streams = await asyncio.gather( + *[stream_processor(stream) for stream in reversed(tv_streams)] + ) - stream_list.append( - Stream( - name=addon_name, - description=f"šŸ“ŗ {stream.name}{country_info}\nšŸ”— {stream.source}", - url=stream.url, - ytId=stream.ytId, - behaviorHints=stream.behaviorHints, - ) - ) + stream_list = [] + is_mediaflow_needed = False + + for result in processed_streams: + if result: + if isinstance(result, Stream): + stream_list.append(result) + elif result == "MEDIAFLOW_NEEDED": + is_mediaflow_needed = True if not stream_list: - stream_list.append( - Stream( - name=settings.addon_name, - description="🚫 No streams are live at the moment.", - url=f"{settings.host_url}/static/exceptions/no_streams_live.mp4", - behaviorHints={"notWebReady": True}, + if is_mediaflow_needed: + stream_list.append( + create_exception_stream( + addon_name, + "🚫 MediaFlow Proxy is required to watch this stream.", + "mediaflow_proxy_required.mp4", + ) + ) + else: + stream_list.append( + create_exception_stream( + addon_name, + "🚫 No streams are live at the moment.", + "no_streams_live.mp4", + ) ) - ) return stream_list +async def process_stream( + stream: TVStreams, + is_mediaflow_proxy_enabled: bool, + mediaflow_config, + addon_name: str, +) -> Optional[Stream | str]: + if settings.validate_m3u8_urls_liveness: + is_working = await validate_m3u8_or_mpd_url_with_cache( + stream.url, stream.behaviorHints or {} + ) + if not is_working: + return None + + stream_url, behavior_hints = stream.url, stream.behaviorHints + behavior_hints = behavior_hints if behavior_hints else {} + + if stream.drm_key or "dlhd" in stream.source: + if not is_mediaflow_proxy_enabled: + return "MEDIAFLOW_NEEDED" + stream_url = get_proxy_url(stream, mediaflow_config) + behavior_hints["proxyHeaders"] = None + elif is_mediaflow_proxy_enabled: + stream_url = get_proxy_url(stream, mediaflow_config) + behavior_hints["proxyHeaders"] = None + + country_info = f"\n🌐 {stream.country}" if stream.country else "" + + return Stream( + name=addon_name, + description=f"šŸ“ŗ {stream.name}{country_info}\nšŸ”— {stream.source}", + url=stream_url, + ytId=stream.ytId, + behaviorHints=behavior_hints, + ) + + +def get_proxy_url(stream: TVStreams, mediaflow_config) -> str: + endpoint = ( + "/proxy/mpd/manifest.m3u8" if stream.drm_key else "/proxy/hls/manifest.m3u8" + ) + query_params = {} + if stream.drm_key: + query_params = {"key_id": stream.drm_key_id, "key": stream.drm_key} + elif "dlhd" in stream.source: + query_params = {"use_request_proxy": False} + + return encode_mediaflow_proxy_url( + mediaflow_config.proxy_url, + endpoint, + stream.url, + query_params=query_params, + request_headers=stream.behaviorHints.get("proxyHeaders", {}).get("request", {}), + encryption_api_password=mediaflow_config.api_password, + ) + + +def create_exception_stream( + addon_name: str, description: str, exc_file_name: str +) -> Stream: + return Stream( + name=addon_name, + description=description, + url=f"{settings.host_url}/static/exceptions/{exc_file_name}", + behaviorHints={"notWebReady": True}, + ) + + async def fetch_downloaded_info_hashes( user_data: UserData, user_ip: str | None ) -> list[str]: diff --git a/utils/validation_helper.py b/utils/validation_helper.py index cb470c1..f113aa1 100644 --- a/utils/validation_helper.py +++ b/utils/validation_helper.py @@ -34,7 +34,7 @@ async def validate_image_url(url: str) -> bool: return is_valid_url(url) and await does_url_exist(url) -async def validate_m3u8_url( +async def validate_live_stream_url( url: str, behaviour_hint: dict, validate_url: bool = False ) -> bool: if validate_url and not is_valid_url(url): @@ -47,26 +47,29 @@ async def validate_m3u8_url( url, allow_redirects=True, headers=headers, - timeout=aiohttp.ClientTimeout(total=30), + timeout=aiohttp.ClientTimeout(total=20), ) as response: - content_type = response.headers.get("Content-Type", "").lower() - - is_valid = content_type in const.M3U8_VALID_CONTENT_TYPES + content_type = response.headers.get("content-type", "").lower() + is_valid = content_type in const.IPTV_VALID_CONTENT_TYPES return is_valid except (ClientError, asyncio.TimeoutError) as err: logging.error(err) return False -async def validate_m3u8_url_with_cache(url: str, behaviour_hint: dict): - cache_key = f"m3u8_url:{parse.urlparse(url).netloc}" - cache_data = await REDIS_ASYNC_CLIENT.get(cache_key) - if cache_data: - return json.loads(cache_data) +async def validate_m3u8_or_mpd_url_with_cache(url: str, behaviour_hint: dict): + try: + cache_key = f"m3u8_url:{parse.urlparse(url).netloc}" + cache_data = await REDIS_ASYNC_CLIENT.get(cache_key) + if cache_data: + return json.loads(cache_data) - is_valid = await validate_m3u8_url(url, behaviour_hint) - await REDIS_ASYNC_CLIENT.set(cache_key, json.dumps(is_valid), ex=180) - return is_valid + is_valid = await validate_live_stream_url(url, behaviour_hint) + await REDIS_ASYNC_CLIENT.set(cache_key, json.dumps(is_valid), ex=180) + return is_valid + except Exception as e: + logging.exception(e) + return False class ValidationError(Exception): @@ -91,7 +94,7 @@ async def validate_tv_metadata(metadata: schemas.TVMetaData) -> list[schemas.TVS for stream in metadata.streams: if stream.url: stream_validation_tasks.append( - validate_m3u8_url( + validate_live_stream_url( stream.url, ( stream.behaviorHints.model_dump(exclude_none=True)