From e502d0f10e34682b558beb0dea5fd6208d4d1da5 Mon Sep 17 00:00:00 2001 From: Mohamed Zumair Date: Thu, 29 Aug 2024 10:02:01 +0530 Subject: [PATCH] Add support for MediaFlow Proxy for Debrid stream & Live stream (#271) * Add support for MediaFlow Proxy for Debrid stream & Live stream * Add emoji for mediaflow proxy stream and add error logging --- api/main.py | 14 ++-- db/crud.py | 12 ++-- db/models.py | 9 ++- db/schemas.py | 15 ++++- resources/html/configure.html | 73 +++++++++++++++++---- resources/js/config_script.js | 29 +++++++-- streaming_providers/proxy_handlers.py | 92 -------------------------- streaming_providers/routes.py | 93 ++++++++++----------------- utils/network.py | 82 ++++++++++++++++++++++- utils/parser.py | 37 ++++++++--- utils/runtime_const.py | 6 ++ 11 files changed, 266 insertions(+), 196 deletions(-) delete mode 100644 streaming_providers/proxy_handlers.py diff --git a/api/main.py b/api/main.py index 465489c..a330740 100644 --- a/api/main.py +++ b/api/main.py @@ -32,7 +32,7 @@ from utils.lock import ( maintain_heartbeat, release_scheduler_lock, ) -from utils.network import get_request_namespace, get_user_public_ip +from utils.network import get_request_namespace, get_user_public_ip, get_user_data from utils.parser import generate_manifest from utils.runtime_const import DELETE_ALL_META, DELETE_ALL_META_ITEM, TEMPLATES @@ -73,10 +73,6 @@ app.add_middleware(middleware.UserDataMiddleware) app.mount("/static", StaticFiles(directory="resources"), name="static") -def get_user_data(request: Request) -> schemas.UserData: - return request.user - - @app.on_event("startup") async def init_server(): await database.init() @@ -301,7 +297,7 @@ async def get_catalog( await crud.get_events_meta_list(request.app.state.redis, genre, skip) ) else: - user_ip = get_user_public_ip(request) + user_ip = await get_user_public_ip(request, user_data) metas.metas.extend( await crud.get_meta_list( user_data, @@ -471,7 +467,7 @@ async def get_streams( episode: int = None, user_data: schemas.UserData = Depends(get_user_data), ): - user_ip = get_user_public_ip(request) + user_ip = await get_user_public_ip(request, user_data) user_feeds = [] if season is None or episode is None: season = episode = 1 @@ -527,13 +523,13 @@ async def get_streams( fetched_streams.extend(user_feeds) elif catalog_type == "events": fetched_streams = await crud.get_event_streams( - request.app.state.redis, video_id + request.app.state.redis, video_id, user_data ) response.headers.update(const.NO_CACHE_HEADERS) else: response.headers.update(const.NO_CACHE_HEADERS) fetched_streams = await crud.get_tv_streams( - request.app.state.redis, video_id, namespace=get_request_namespace(request) + request.app.state.redis, video_id, get_request_namespace(request), user_data ) return {"streams": fetched_streams} diff --git a/db/crud.py b/db/crud.py index 00e9d20..c525642 100644 --- a/db/crud.py +++ b/db/crud.py @@ -425,7 +425,9 @@ async def get_series_streams( ) -async def get_tv_streams(redis: Redis, video_id: str, namespace: str) -> list[Stream]: +async def get_tv_streams( + redis: Redis, video_id: str, namespace: str, user_data +) -> list[Stream]: tv_streams = await TVStreams.find( { "meta_id": video_id, @@ -434,7 +436,7 @@ async def get_tv_streams(redis: Redis, video_id: str, namespace: str) -> list[St }, ).to_list() - return await parse_tv_stream_data(tv_streams, redis) + return await parse_tv_stream_data(tv_streams, redis, user_data) async def get_movie_meta(meta_id: str, redis: Redis, user_data: schemas.UserData): @@ -1102,14 +1104,14 @@ async def get_event_data_by_id(redis, meta_id: str) -> MediaFusionEventsMetaData return MediaFusionEventsMetaData.model_validate_json(events_json) -async def get_event_streams(redis, meta_id: str) -> list[Stream]: +async def get_event_streams(redis, meta_id: str, user_data) -> list[Stream]: event_key = f"event:{meta_id}" event_json = await redis.get(event_key) if not event_json: - return await parse_tv_stream_data([], redis) + return await parse_tv_stream_data([], redis, user_data) event_data = MediaFusionEventsMetaData.model_validate_json(event_json) - return await parse_tv_stream_data(event_data.streams, redis) + return await parse_tv_stream_data(event_data.streams, redis, user_data) async def get_genres(catalog_type: str, redis: Redis) -> list[str]: diff --git a/db/models.py b/db/models.py index b06d794..3d6583a 100644 --- a/db/models.py +++ b/db/models.py @@ -3,7 +3,7 @@ from typing import Optional, Any import pymongo from beanie import Document, Link -from pydantic import BaseModel, Field, ConfigDict +from pydantic import BaseModel, Field, ConfigDict, field_validator from pymongo import IndexModel, ASCENDING, DESCENDING @@ -44,6 +44,13 @@ class TorrentStreams(Document): seeders: Optional[int] = None cached: Optional[bool] = Field(default=False, exclude=True) + @field_validator("audio", mode="before") + def validate_audio(cls, v): + # Ensure audio is a string + if v and isinstance(v, list): + return v[0] + return v + class Settings: indexes = [ IndexModel( diff --git a/db/schemas.py b/db/schemas.py index 7310e27..05c6356 100644 --- a/db/schemas.py +++ b/db/schemas.py @@ -106,6 +106,17 @@ class QBittorrentConfig(BaseModel): populate_by_name = True +class MediaFlowConfig(BaseModel): + proxy_url: str | None = Field(alias="pu") + api_password: str | None = Field(alias="ap") + proxy_live_streams: bool = Field(default=False, alias="pls") + proxy_debrid_streams: bool = Field(default=False, alias="pds") + + class Config: + extra = "ignore" + populate_by_name = True + + class StreamingProvider(BaseModel): service: Literal[ "realdebrid", @@ -168,13 +179,13 @@ class UserData(BaseModel): ] ] = Field(default=["Adults"], alias="cf") api_password: str | None = Field(default=None, alias="ap") - proxy_debrid_stream: bool = Field(default=False, alias="pds") language_sorting: list[str | None] = Field( - default=const.SUPPORTED_LANGUAGES, alias="ls" + default=list(const.SUPPORTED_LANGUAGES), alias="ls" ) quality_filter: list[str] = Field( default=list(const.QUALITY_GROUPS.keys()), alias="qf" ) + mediaflow_config: MediaFlowConfig | None = Field(default=None, alias="mfc") @field_validator("selected_resolutions", mode="after") def validate_selected_resolutions(cls, v): diff --git a/resources/html/configure.html b/resources/html/configure.html index 4eeaf2d..ea1dcc5 100644 --- a/resources/html/configure.html +++ b/resources/html/configure.html @@ -231,19 +231,6 @@ - - -
-
Proxy debrid stream:
-
- - -
-
@@ -474,6 +461,66 @@ +
+

External Services Configuration

+
+ +
+
MediaFlow Configuration
+
+ + +
+
+
+ + +
+ Please enter a valid MediaFlow Proxy URL. +
+
+
+ +
+ + +
+
+ Please enter the MediaFlow API password. +
+
+
+ + +
+
+ + +
+
+
+
+ {% if authentication_required %}
diff --git a/resources/js/config_script.js b/resources/js/config_script.js index 0bbc656..eae87dd 100644 --- a/resources/js/config_script.js +++ b/resources/js/config_script.js @@ -124,6 +124,12 @@ function setElementDisplay(elementId, displayStatus) { document.getElementById(elementId).style.display = displayStatus; } +function validateUrl(url) { + // This regex supports domain names, IPv4, and IPv6 addresses + const urlPattern = /^(https?:\/\/)?(([a-z0-9]([a-z0-9-]*[a-z0-9])?\.)+[a-z]{2,}|localhost|\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}|\[(?:[0-9a-fA-F]{1,4}:){7}[0-9a-fA-F]{1,4}\])(:?\d+)?(\/[-a-z\d%_.~+]*)*(\?[;&a-z\d%_.~+=-]*)?(\#[-a-z\d_]*)?$/i; + return urlPattern.test(url); +} + // Function to format bytes into a human-readable format function formatBytes(bytes, decimals = 2) { if (bytes === "0") return '0 Bytes'; @@ -189,13 +195,11 @@ function updateProviderFields(isChangeEvent = false) { setElementDisplay('qbittorrent_config', 'none'); } setElementDisplay('watchlist_section', 'block'); - setElementDisplay('proxy_debrid_section', 'block'); watchlistLabel.textContent = `Enable ${provider.charAt(0).toUpperCase() + provider.slice(1)} Watchlist`; } else { setElementDisplay('credentials', 'none'); setElementDisplay('token_input', 'none'); setElementDisplay('watchlist_section', 'none'); - setElementDisplay('proxy_debrid_section', 'none'); setElementDisplay('qbittorrent_config', 'none'); } @@ -296,11 +300,23 @@ function getUserData() { } streamingProviderData.service = provider; streamingProviderData.enable_watchlist_catalogs = document.getElementById('enable_watchlist').checked; - streamingProviderData.proxy_debrid_stream = document.getElementById('proxy_debrid_stream').checked; } else { streamingProviderData = null; } + const mediaflowEnabled = document.getElementById('enable_mediaflow').checked + let mediaflowConfig = null; + if (mediaflowEnabled) { + mediaflowConfig = { + proxy_url: document.getElementById('mediaflow_proxy_url').value, + api_password: document.getElementById('mediaflow_api_password').value, + proxy_live_streams: document.getElementById('proxy_live_streams').checked, + proxy_debrid_streams: document.getElementById('proxy_debrid_streams').checked + }; + validateInput('mediaflow_proxy_url', validateUrl(mediaflowConfig.proxy_url)); + validateInput('mediaflow_api_password', mediaflowConfig.api_password.trim() !== ''); + } + // Collect and validate other user data const maxSizeSlider = document.getElementById('max_size_slider'); const maxSizeValue = maxSizeSlider.value; @@ -341,7 +357,6 @@ function getUserData() { selected_catalogs: Array.from(document.querySelectorAll('input[name="selected_catalogs"]:checked')).map(el => el.value), selected_resolutions: Array.from(document.querySelectorAll('input[name="selected_resolutions"]:checked')).map(el => el.value || null), enable_catalogs: document.getElementById('enable_catalogs').checked, - proxy_debrid_stream: document.getElementById('proxy_debrid_stream').checked, max_size: maxSizeBytes, max_streams_per_resolution: maxStreamsPerResolution, torrent_sorting_priority: selectedSortingOptions, @@ -351,6 +366,7 @@ function getUserData() { language_sorting: languageSorting, quality_filter: selectedQualityFilters, api_password: apiPassword, + mediaflow_config: mediaflowConfig, }; } @@ -383,6 +399,10 @@ document.getElementById('provider_token').addEventListener('input', function () adjustOAuthSectionDisplay(); }); +document.getElementById('enable_mediaflow').addEventListener('change', function () { + setElementDisplay('mediaflow_config', this.checked ? 'block' : 'none'); +}); + // Event listener for the slider document.getElementById('max_size_slider').addEventListener('input', updateSizeOutput); @@ -472,6 +492,7 @@ document.addEventListener('DOMContentLoaded', function () { } setupPasswordToggle('qbittorrent_password', 'toggleQbittorrentPassword', 'toggleQbittorrentPasswordIcon'); setupPasswordToggle('webdav_password', 'toggleWebdavPassword', 'toggleWebdavPasswordIcon'); + setupPasswordToggle('mediaflow_api_password', 'toggleMediaFlowPassword', 'toggleMediaFlowPasswordIcon'); }); diff --git a/streaming_providers/proxy_handlers.py b/streaming_providers/proxy_handlers.py deleted file mode 100644 index 4d84f86..0000000 --- a/streaming_providers/proxy_handlers.py +++ /dev/null @@ -1,92 +0,0 @@ -import logging - -import httpx -from fastapi import Response -from fastapi.responses import StreamingResponse -from starlette.background import BackgroundTask - -logger = logging.getLogger(__name__) - - -async def handle_head_request(video_url: str) -> Response: - async with httpx.AsyncClient() as client: - try: - head_response = await client.head(video_url) - head_response.raise_for_status() - return Response( - headers={ - "Content-Length": head_response.headers.get("content-length", ""), - "Accept-Ranges": head_response.headers.get( - "accept-ranges", "bytes" - ), - }, - status_code=head_response.status_code, - ) - except httpx.HTTPStatusError as e: - logger.error(f"Upstream service error while handling HEAD request: {e}") - return Response(status_code=502, content=f"Upstream service error: {e}") - except Exception as e: - logger.error(f"Internal server error while handling HEAD request: {e}") - return Response(status_code=500, content=f"Internal server error: {e}") - - -class Streamer: - def __init__(self, client): - self.client = client - self.response = None - - async def stream_content( - self, video_url: str, headers: dict, chunk_size: int = 65536 - ): - try: - async with self.client.stream( - "GET", video_url, headers=headers - ) as self.response: - self.response.raise_for_status() - async for chunk in self.response.aiter_bytes(chunk_size): - yield chunk - except httpx.HTTPStatusError as e: - logger.error(f"HTTP error while streaming content: {e}") - raise e - except Exception as e: - logger.error(f"Unexpected error while streaming content: {e}") - raise e - - async def close(self): - try: - if self.response is not None: - await self.response.aclose() - if self.client is not None: - await self.client.aclose() - except Exception as e: - logger.error(f"Error while closing Streamer: {e}") - - -async def handle_get_request( - video_url: str, range_header: str -) -> StreamingResponse | Response: - client = httpx.AsyncClient() - try: - head_response = await client.head(video_url, headers={"Range": range_header}) - head_response.raise_for_status() - if head_response.status_code == 206: - streamer = Streamer(client) - return StreamingResponse( - streamer.stream_content(video_url, {"Range": range_header}), - status_code=206, - headers={ - "Content-Range": head_response.headers["content-range"], - "Content-Length": head_response.headers["content-length"], - "Accept-Ranges": "bytes", - }, - background=BackgroundTask(streamer.close), - ) - except httpx.HTTPStatusError as e: - logger.error(f"Upstream service error while handling GET request: {e}") - return Response(status_code=502, content=f"Upstream service error: {e}") - except Exception as e: - logger.error(f"Internal server error while handling GET request: {e}") - return Response(status_code=500, content=f"Internal server error: {e}") - - logger.warning(f"Resource not found for video URL: {video_url}") - return Response(status_code=404) diff --git a/streaming_providers/routes.py b/streaming_providers/routes.py index d8ef651..164f4e1 100644 --- a/streaming_providers/routes.py +++ b/streaming_providers/routes.py @@ -6,13 +6,14 @@ from fastapi import ( Response, HTTPException, APIRouter, + Depends, ) from fastapi.responses import RedirectResponse from redis.asyncio import Redis -from db import crud +from db import crud, schemas from db.config import settings -from streaming_providers import mapper, proxy_handlers +from streaming_providers import mapper from streaming_providers.debridlink.api import router as debridlink_router from streaming_providers.exceptions import ProviderException from streaming_providers.premiumize.api import router as premiumize_router @@ -20,7 +21,7 @@ from streaming_providers.realdebrid.api import router as realdebrid_router from streaming_providers.seedr.api import router as seedr_router from utils import crypto, torrent, wrappers, const from utils.lock import acquire_redis_lock, release_redis_lock -from utils.network import get_user_public_ip +from utils.network import get_user_public_ip, get_user_data, encode_mediaflow_proxy_url # Seconds until when the Video URLs are cached URL_CACHE_EXP = 3600 @@ -46,14 +47,14 @@ async def streaming_provider_endpoint( request: Request, season: int = None, episode: int = None, + user_data: schemas.UserData = Depends(get_user_data), ): response.headers.update(const.NO_CACHE_HEADERS) - user_data = request.scope.get("user", crypto.decrypt_user_data(secret_str)) if not user_data.streaming_provider: raise HTTPException(status_code=400, detail="No streaming provider set.") - user_ip = get_user_public_ip(request) + user_ip = await get_user_public_ip(request, user_data) redirect_status_code = 302 cached_stream_url_key = "streaming_provider_" + crypto.get_text_hash( f"{user_ip}_{secret_str}_{info_hash}_{season}_{episode}", @@ -64,7 +65,16 @@ async def streaming_provider_endpoint( if cached_stream_url := await get_cached_stream_url( request.app.state.redis, cached_stream_url_key ): - logging.info("Redirecting to cached URL...") + if ( + user_data.mediaflow_config + and user_data.mediaflow_config.proxy_debrid_streams + ): + cached_stream_url = encode_mediaflow_proxy_url( + user_data.mediaflow_config.proxy_url, + "/proxy/stream", + cached_stream_url, + query_params={"api_password": user_data.mediaflow_config.api_password}, + ) return RedirectResponse( url=cached_stream_url, headers=response.headers, @@ -121,6 +131,17 @@ async def streaming_provider_endpoint( await request.app.state.redis.set( cached_stream_url_key, video_url.encode("utf-8"), ex=URL_CACHE_EXP ) + if ( + user_data.mediaflow_config + and user_data.mediaflow_config.proxy_debrid_streams + ): + video_url = encode_mediaflow_proxy_url( + user_data.mediaflow_config.proxy_url, + "/proxy/stream", + video_url, + query_params={"api_password": user_data.mediaflow_config.api_password}, + ) + except ProviderException as error: logging.error( "Exception occurred for %s: %s", @@ -143,65 +164,17 @@ async def streaming_provider_endpoint( ) -# This endpoint is to proxy the debrid provider via this instance. -@router.head("/{secret_str}/proxy_stream", tags=["streaming_provider"]) -@router.get("/{secret_str}/proxy_stream", tags=["streaming_provider"]) -@wrappers.exclude_rate_limit -@wrappers.auth_required -async def proxy_streaming_provider_endpoint( - secret_str: str, - info_hash: str, - response: Response, - request: Request, - season: int = None, - episode: int = None, -): - response.headers.update(const.NO_CACHE_HEADERS) - - user_data = request.scope.get("user", crypto.decrypt_user_data(secret_str)) - user_ip = get_user_public_ip(request) - cached_stream_url_key = "streaming_provider_" + crypto.get_text_hash( - f"{user_ip}_{secret_str}_{info_hash}_{season}_{episode}", - full_hash=True, - ) - - # Check if the URL is already cached - if cached_stream_url := await get_cached_stream_url( - request.app.state.redis, cached_stream_url_key - ): - video_url = cached_stream_url - else: - streaming_provider_response = await streaming_provider_endpoint( - secret_str, info_hash, response, request, season, episode - ) - if streaming_provider_response.status_code != 302: - return streaming_provider_response - - video_url = streaming_provider_response.headers.get("location", "") - - # Validate if proxying is allowed - if ( - settings.is_public_instance is False - and user_data.proxy_debrid_stream is True - and user_data.streaming_provider - ): - if request.method == "HEAD": - return await proxy_handlers.handle_head_request(video_url) - else: - range_header = request.headers.get("range", "bytes=0-") - return await proxy_handlers.handle_get_request(video_url, range_header) - - return RedirectResponse(url=video_url, headers=response.headers, status_code=302) - - @router.get("/{secret_str}/delete_all_watchlist", tags=["streaming_provider"]) @wrappers.exclude_rate_limit @wrappers.auth_required -async def delete_all_watchlist(request: Request, response: Response, secret_str: str): +async def delete_all_watchlist( + request: Request, + response: Response, + user_data: schemas.UserData = Depends(get_user_data), +): response.headers.update(const.NO_CACHE_HEADERS) - user_data = request.scope.get("user", crypto.decrypt_user_data(secret_str)) - user_ip = get_user_public_ip(request) + user_ip = get_user_public_ip(request, user_data) if not user_data.streaming_provider: raise HTTPException(status_code=400, detail="No streaming provider set.") diff --git a/utils/network.py b/utils/network.py index 5f1049c..b44cf00 100644 --- a/utils/network.py +++ b/utils/network.py @@ -1,11 +1,14 @@ import asyncio import logging from typing import Callable +from urllib import parse import httpx from fastapi.requests import Request -from utils.runtime_const import PRIVATE_CIDR +from db.schemas import UserData +from utils import crypto +from utils.runtime_const import PRIVATE_CIDR, REDIS_CLIENT class CircuitBreakerOpenException(Exception): @@ -157,7 +160,54 @@ def get_client_ip(request: Request) -> str | None: return request.client.host if request.client else "127.0.0.1" -def get_user_public_ip(request: Request): +async def get_mediaflow_proxy_public_ip( + mediaflow_proxy_url: str, api_password +) -> str | None: + """ + Get the public IP address of the MediaFlow proxy server. + """ + cache_key = crypto.get_text_hash( + f"{mediaflow_proxy_url}:{api_password}", full_hash=True + ) + if public_ip := await REDIS_CLIENT.getex(cache_key, ex=300): + return public_ip + + try: + async with httpx.AsyncClient() as client: + response = await client.get( + mediaflow_proxy_url + "/proxy/ip", + params={"api_password": api_password}, + timeout=10, + ) + response.raise_for_status() + public_ip = response.json().get("ip") + if public_ip: + await REDIS_CLIENT.set(cache_key, public_ip, ex=300) + return public_ip + except httpx.HTTPStatusError as e: + logging.error(f"HTTP error occurred: {e}") + except httpx.RequestError as e: + logging.error(f"Request error occurred: {e}") + except Exception as e: + logging.error(f"An unexpected error occurred: {e}") + return None + + +async def get_user_public_ip( + request: Request, user_data: UserData | None = None +) -> str: + # Check if user has mediaflow config + if ( + user_data + and user_data.mediaflow_config + and user_data.mediaflow_config.proxy_debrid_streams + ): + public_ip = await get_mediaflow_proxy_public_ip( + user_data.mediaflow_config.proxy_url, + user_data.mediaflow_config.api_password, + ) + if public_ip: + return public_ip # Get the user's public IP address user_ip = get_client_ip(request) # check if the user's IP address is a private IP address @@ -183,3 +233,31 @@ def get_request_namespace(request: Request) -> str: namespace = f"tenant-{parts[0]}" return namespace + + +def get_user_data(request: Request) -> UserData: + return request.user + + +def encode_mediaflow_proxy_url( + mediaflow_proxy_url: str, + endpoint: str, + destination_url: str | None = None, + query_params: dict | None = None, + request_headers: dict | None = None, +) -> str: + query_params = query_params or {} + if destination_url is not None: + query_params["d"] = destination_url + + # Add headers if provided + if request_headers: + query_params.update( + {f"h_{key}": value for key, value in request_headers.items()} + ) + # Encode the query parameters + encoded_params = parse.urlencode(query_params, quote_via=parse.quote) + + # Construct the full URL + base_url = parse.urljoin(mediaflow_proxy_url, endpoint) + return f"{base_url}?{encoded_params}" diff --git a/utils/parser.py b/utils/parser.py index b29102d..954091c 100644 --- a/utils/parser.py +++ b/utils/parser.py @@ -13,6 +13,7 @@ from db.schemas import Stream, UserData from streaming_providers import mapper 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 @@ -154,15 +155,15 @@ async def parse_stream_data( base_proxy_url_template = "" if has_streaming_provider: - stream_path = "stream" if ( - settings.is_public_instance is False - and user_data.proxy_debrid_stream is True + user_data.mediaflow_config + and user_data.mediaflow_config.proxy_debrid_streams ): streaming_provider_name += " šŸ•µšŸ¼ā€ā™‚ļø" - stream_path = "proxy_stream" - base_proxy_url_template = f"{settings.host_url}/streaming_provider/{secret_str}/{stream_path}?info_hash={{}}" + base_proxy_url_template = ( + f"{settings.host_url}/streaming_provider/{secret_str}/stream?info_hash={{}}" + ) stream_list = [] for stream_data in streams: @@ -211,7 +212,9 @@ async def parse_stream_data( file_size = stream_data.size size_info = convert_bytes_to_readable(file_size) - languages = f"🌐 {' + '.join(stream_data.languages)}" + languages = ( + f"🌐 {' + '.join(stream_data.languages)}" if stream_data.languages else None + ) source_info = f"šŸ”— {stream_data.source}" description = "\n".join( @@ -288,9 +291,15 @@ def convert_size_to_bytes(size_str: str) -> int: async def parse_tv_stream_data( - tv_streams: list[TVStreams], redis: Redis + tv_streams: list[TVStreams], redis: Redis, user_data: UserData ) -> list[Stream]: stream_list = [] + 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( @@ -299,11 +308,23 @@ async def parse_tv_stream_data( 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": {}}) + country_info = f"\n🌐 {stream.country}" if stream.country else "" stream_list.append( Stream( - name=settings.addon_name, + name=addon_name, description=f"šŸ“ŗ {stream.name}{country_info}\nšŸ”— {stream.source}", url=stream.url, ytId=stream.ytId, diff --git a/utils/runtime_const.py b/utils/runtime_const.py index 8876861..e8fd9e5 100644 --- a/utils/runtime_const.py +++ b/utils/runtime_const.py @@ -1,6 +1,7 @@ import re from fastapi.templating import Jinja2Templates +import redis.asyncio as redis from db import schemas from db.config import settings @@ -31,3 +32,8 @@ DELETE_ALL_META_ITEM = { TRACKERS = get_json_data("resources/json/trackers.json") SECRET_KEY = settings.secret_key.encode("utf-8") + + +REDIS_CLIENT = redis.Redis( + connection_pool=redis.ConnectionPool.from_url(settings.redis_url) +)