From 29348f5fa5b032c8f32cb7dc0caa2d444d21bac0 Mon Sep 17 00:00:00 2001 From: mhdzumair Date: Sun, 10 Nov 2024 21:40:54 +0530 Subject: [PATCH] Refactor extractors and routes for improved structure Reorganized extractor routes into dedicated modules for better maintainability. Introduced a base class for extractors, enabling consistent request handling across different services. Additionally, updated configurations and error handling in extractors, enhancing code readability and robustness. --- mediaflow_proxy/configs.py | 4 ++ mediaflow_proxy/extractors/__init__.py | 0 mediaflow_proxy/extractors/base.py | 42 ++++++++++++++ mediaflow_proxy/extractors/doodstream.py | 51 ++++++++++------- mediaflow_proxy/extractors/factory.py | 24 ++++++++ mediaflow_proxy/extractors/mixdrop.py | 44 ++++++++------- mediaflow_proxy/extractors/uqload.py | 22 +++++--- mediaflow_proxy/extractors_routes.py | 55 ------------------- mediaflow_proxy/main.py | 14 ++--- mediaflow_proxy/routes/__init__.py | 2 + mediaflow_proxy/routes/extractor.py | 32 +++++++++++ .../{routes.py => routes/proxy.py} | 19 ++++++- mediaflow_proxy/schemas.py | 8 +++ mediaflow_proxy/utils/http_utils.py | 4 ++ 14 files changed, 205 insertions(+), 116 deletions(-) create mode 100644 mediaflow_proxy/extractors/__init__.py create mode 100644 mediaflow_proxy/extractors/base.py create mode 100644 mediaflow_proxy/extractors/factory.py delete mode 100644 mediaflow_proxy/extractors_routes.py create mode 100644 mediaflow_proxy/routes/__init__.py create mode 100644 mediaflow_proxy/routes/extractor.py rename mediaflow_proxy/{routes.py => routes/proxy.py} (91%) diff --git a/mediaflow_proxy/configs.py b/mediaflow_proxy/configs.py index 4d3bb7f..5293a19 100644 --- a/mediaflow_proxy/configs.py +++ b/mediaflow_proxy/configs.py @@ -6,6 +6,10 @@ class Settings(BaseSettings): proxy_url: str | None = None # The URL of the proxy server to route requests through. enable_streaming_progress: bool = False # Whether to enable streaming progress tracking. + user_agent: str = ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/58.0.3029.110 Safari/537.3" # The user agent to use for HTTP requests. + ) + class Config: env_file = ".env" extra = "ignore" diff --git a/mediaflow_proxy/extractors/__init__.py b/mediaflow_proxy/extractors/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/mediaflow_proxy/extractors/base.py b/mediaflow_proxy/extractors/base.py new file mode 100644 index 0000000..f7e6609 --- /dev/null +++ b/mediaflow_proxy/extractors/base.py @@ -0,0 +1,42 @@ +from abc import ABC, abstractmethod +from typing import Dict, Tuple, Optional + +import httpx + +from mediaflow_proxy.configs import settings + + +class BaseExtractor(ABC): + """Base class for all URL extractors.""" + + def __init__(self, proxy_enabled: bool = False): + self.proxy_url = settings.proxy_url if proxy_enabled else None + self.base_headers = { + "User-Agent": settings.user_agent, + "Accept-Language": "en-US,en;q=0.5", + } + + async def _make_request( + self, url: str, headers: Optional[Dict] = None, follow_redirects: bool = True, **kwargs + ) -> httpx.Response: + """Make HTTP request with error handling.""" + try: + async with httpx.AsyncClient(proxy=self.proxy_url) as client: + response = await client.get( + url, + headers={**self.base_headers, **(headers or {})}, + follow_redirects=follow_redirects, + timeout=30, + **kwargs, + ) + response.raise_for_status() + return response + except httpx.HTTPError as e: + raise ValueError(f"HTTP request failed: {str(e)}") + except Exception as e: + raise ValueError(f"Request failed: {str(e)}") + + @abstractmethod + async def extract(self, url: str) -> Tuple[str, Dict[str, str]]: + """Extract final URL and required headers.""" + pass diff --git a/mediaflow_proxy/extractors/doodstream.py b/mediaflow_proxy/extractors/doodstream.py index e4b995e..e0019ce 100644 --- a/mediaflow_proxy/extractors/doodstream.py +++ b/mediaflow_proxy/extractors/doodstream.py @@ -1,25 +1,34 @@ -import httpx -import time import re -from mediaflow_proxy.configs import settings +import time +from typing import Tuple, Dict + +from mediaflow_proxy.extractors.base import BaseExtractor -async def doodstream_url(d: str, use_request_proxy: bool): - async with httpx.AsyncClient(proxy=settings.proxy_url if use_request_proxy else None) as client: - headers = { - "Range": "bytes=0-", - "Referer": "https://d000d.com/", - } +class DoodStreamExtractor(BaseExtractor): + """DoodStream URL extractor.""" - response = await client.get(d, follow_redirects=True) - if response.status_code == 200: - # Get unique timestamp for the request - real_time = str(int(time.time())) - pattern = r"(\/pass_md5\/.*?)'.*(\?token=.*?expiry=)" - match = re.search(pattern, response.text, re.DOTALL) - if match: - url = f"https://d000d.com{match[1]}" - rebobo = await client.get(url, headers=headers, follow_redirects=True) - final_url = f"{rebobo.text}123456789{match[2]}{real_time}" - doodstream_dict = {"Referer": "https://d000d.com/"} - return final_url, doodstream_dict + def __init__(self, proxy_enabled: bool = False): + super().__init__(proxy_enabled) + self.base_url = "https://d000d.com" + + async def extract(self, url: str) -> Tuple[str, Dict[str, str]]: + """Extract DoodStream URL.""" + response = await self._make_request(url) + + # Extract URL pattern + pattern = r"(\/pass_md5\/.*?)'.*(\?token=.*?expiry=)" + match = re.search(pattern, response.text, re.DOTALL) + if not match: + raise ValueError("Failed to extract URL pattern") + + # Build final URL + pass_url = f"{self.base_url}{match[1]}" + referer = f"{self.base_url}/" + headers = {"Range": "bytes=0-", "Referer": referer} + + rebobo_response = await self._make_request(pass_url, headers=headers) + timestamp = str(int(time.time())) + final_url = f"{rebobo_response.text}123456789{match[2]}{timestamp}" + + return final_url, {"Referer": referer} diff --git a/mediaflow_proxy/extractors/factory.py b/mediaflow_proxy/extractors/factory.py new file mode 100644 index 0000000..eb7a94a --- /dev/null +++ b/mediaflow_proxy/extractors/factory.py @@ -0,0 +1,24 @@ +from typing import Dict, Type + +from mediaflow_proxy.extractors.base import BaseExtractor +from mediaflow_proxy.extractors.doodstream import DoodStreamExtractor +from mediaflow_proxy.extractors.mixdrop import MixdropExtractor +from mediaflow_proxy.extractors.uqload import UqloadExtractor + + +class ExtractorFactory: + """Factory for creating URL extractors.""" + + _extractors: Dict[str, Type[BaseExtractor]] = { + "Doodstream": DoodStreamExtractor, + "Uqload": UqloadExtractor, + "Mixdrop": MixdropExtractor, + } + + @classmethod + def get_extractor(cls, host: str, proxy_enabled: bool = False) -> BaseExtractor: + """Get appropriate extractor instance for the given host.""" + extractor_class = cls._extractors.get(host) + if not extractor_class: + raise ValueError(f"Unsupported host: {host}") + return extractor_class(proxy_enabled) diff --git a/mediaflow_proxy/extractors/mixdrop.py b/mediaflow_proxy/extractors/mixdrop.py index 796c885..6c3085b 100644 --- a/mediaflow_proxy/extractors/mixdrop.py +++ b/mediaflow_proxy/extractors/mixdrop.py @@ -1,27 +1,31 @@ -import httpx import re import string -from mediaflow_proxy.configs import settings +from typing import Dict, Tuple + +from mediaflow_proxy.extractors.base import BaseExtractor -async def mixdrop_url(d: str, use_request_proxy: bool): - headers = { - "User-Agent": "Mozilla/5.0 (Windows NT 10.10; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/58.0.3029.110 Safari/537.3", - "Accept-Language": "en-US,en;q=0.5", - } - async with httpx.AsyncClient(proxy=settings.proxy_url if use_request_proxy else None) as client: - response = await client.get(d, headers=headers, follow_redirects=True, timeout=30) - [s1, s2] = re.search(r"\}\('(.+)',.+,'(.+)'\.split", response.text).group(1, 2) +class MixdropExtractor(BaseExtractor): + """Mixdrop URL extractor.""" + + async def extract(self, url: str) -> Tuple[str, Dict[str, str]]: + """Extract Mixdrop URL.""" + response = await self._make_request(url) + + # Extract and decode URL + match = re.search(r"\}\('(.+)',.+,'(.+)'\.split", response.text) + if not match: + raise ValueError("Failed to extract URL components") + + s1, s2 = match.group(1, 2) schema = s1.split(";")[2][5:-1] terms = s2.split("|") + + # Build character mapping charset = string.digits + string.ascii_letters - d = dict() - for i in range(len(terms)): - d[charset[i]] = terms[i] or charset[i] - final_url = "https:" - for c in schema: - final_url += d[c] if c in d else c - headers_dict = { - "User-Agent": "Mozilla/5.0 (Windows NT 10.10; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/58.0.3029.110 Safari/537.3" - } - return final_url, headers_dict + char_map = {charset[i]: terms[i] or charset[i] for i in range(len(terms))} + + # Construct final URL + final_url = "https:" + "".join(char_map.get(c, c) for c in schema) + + return final_url, {"User-Agent": self.base_headers["User-Agent"]} diff --git a/mediaflow_proxy/extractors/uqload.py b/mediaflow_proxy/extractors/uqload.py index 03a5b98..d85b816 100644 --- a/mediaflow_proxy/extractors/uqload.py +++ b/mediaflow_proxy/extractors/uqload.py @@ -1,14 +1,18 @@ -import httpx import re -from mediaflow_proxy.configs import settings +from typing import Dict, Tuple + +from mediaflow_proxy.extractors.base import BaseExtractor -async def uqload_url(d: str, use_request_proxy: bool): - async with httpx.AsyncClient(proxy=settings.proxy_url if use_request_proxy else None) as client: +class UqloadExtractor(BaseExtractor): + """Uqload URL extractor.""" + + async def extract(self, url: str) -> Tuple[str, Dict[str, str]]: + """Extract Uqload URL.""" + response = await self._make_request(url) - response = await client.get(d, follow_redirects=True) video_url_match = re.search(r'sources: \["(.*?)"\]', response.text) - if video_url_match: - final_url = video_url_match.group(1) - uqload_dict = {"Referer": "https://uqload.to/"} - return final_url, uqload_dict + if not video_url_match: + raise ValueError("Failed to extract video URL") + + return video_url_match.group(1), {"Referer": "https://uqload.to/"} diff --git a/mediaflow_proxy/extractors_routes.py b/mediaflow_proxy/extractors_routes.py deleted file mode 100644 index 3b814b6..0000000 --- a/mediaflow_proxy/extractors_routes.py +++ /dev/null @@ -1,55 +0,0 @@ -from fastapi import APIRouter, Query -from fastapi.responses import JSONResponse, RedirectResponse -from .extractors.doodstream import doodstream_url -from .extractors.uqload import uqload_url -from .extractors.mixdrop import mixdrop_url -from mediaflow_proxy.configs import settings - -extractor_router = APIRouter() -host_map = {"Doodstream": doodstream_url, "Mixdrop": mixdrop_url, "Uqload": uqload_url} - - -@extractor_router.get("/extractor") -async def doodstream_extractor( - d: str = Query(..., description="Extract Clean Link from various Hosts"), - use_request_proxy: bool = Query(False, description="Whether to use the MediaFlow proxy configuration."), - host: str = Query( - ..., description='From which Host the URL comes from, here avaiable ones: "Doodstream","Mixdrop","Uqload"' - ), - redirect_stream: bool = Query( - False, - description="If enabled the response will be redirected to stream endpoint automatically and the stream will be proxied", - ), -): - """ - Extract a clean link from DoodStream,Mixdrop,Uqload - - Args: request (Request): The incoming HTTP request - - Returns: The clean link (url) and the headers needed to access the url - - N.B. You can't use a rotating proxy if type is set to "Doodstream" - """ - try: - final_url, headers_dict = await host_map[host](d, use_request_proxy) - except Exception as e: - return JSONResponse(content={"error": str(e)}) - if redirect_stream == True: - formatted_headers = format_headers(headers_dict) - redirected_stream = f"/proxy/stream?api_password={settings.api_password}&d={final_url}&{formatted_headers}" - return RedirectResponse(url=redirected_stream) - elif redirect_stream == False: - return JSONResponse(content={"url": final_url, "headers": headers_dict}) - - -def format_headers(headers): - """ - Format the headers dictionary into a query string format with 'h_' prefix. - - Args: - - headers: A dictionary of headers. - - Returns: - - A query string formatted string of headers. - """ - return "&".join(f"h_{key}={value}" for key, value in headers.items()) diff --git a/mediaflow_proxy/main.py b/mediaflow_proxy/main.py index fa0ac04..c985061 100644 --- a/mediaflow_proxy/main.py +++ b/mediaflow_proxy/main.py @@ -1,6 +1,6 @@ import logging -from importlib import resources import uuid +from importlib import resources from fastapi import FastAPI, Depends, Security, HTTPException, BackgroundTasks from fastapi.security import APIKeyQuery, APIKeyHeader @@ -9,12 +9,11 @@ from starlette.responses import RedirectResponse, JSONResponse from starlette.staticfiles import StaticFiles from mediaflow_proxy.configs import settings -from mediaflow_proxy.routes import proxy_router -from mediaflow_proxy.extractors_routes import extractor_router +from mediaflow_proxy.routes import proxy_router, extractor_router from mediaflow_proxy.schemas import GenerateUrlRequest from mediaflow_proxy.utils.crypto_utils import EncryptionHandler, EncryptionMiddleware -from mediaflow_proxy.utils.rd_speedtest import run_speedtest, prune_task, results from mediaflow_proxy.utils.http_utils import encode_mediaflow_proxy_url +from mediaflow_proxy.utils.rd_speedtest import run_speedtest, prune_task, results logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s") app = FastAPI() @@ -57,7 +56,7 @@ async def trigger_speedtest(background_tasks: BackgroundTasks, api_password: str # Generate a random UUID as task_id task_id = str(uuid.uuid4()) # Generate unique task ID background_tasks.add_task(run_speedtest, task_id) - + # Schedule the task to be pruned after 1 hour background_tasks.add_task(prune_task, task_id) @@ -97,8 +96,7 @@ async def generate_encrypted_or_encoded_url(request: GenerateUrlRequest): app.include_router(proxy_router, prefix="/proxy", tags=["proxy"], dependencies=[Depends(verify_api_key)]) -app.include_router(extractor_router, tags=["extractors"], dependencies=[Depends(verify_api_key)]) - +app.include_router(extractor_router, prefix="/extractor", tags=["extractors"], dependencies=[Depends(verify_api_key)]) static_path = resources.files("mediaflow_proxy").joinpath("static") app.mount("/", StaticFiles(directory=str(static_path), html=True), name="static") @@ -111,4 +109,4 @@ def run(): if __name__ == "__main__": - run() \ No newline at end of file + run() diff --git a/mediaflow_proxy/routes/__init__.py b/mediaflow_proxy/routes/__init__.py new file mode 100644 index 0000000..4a828f6 --- /dev/null +++ b/mediaflow_proxy/routes/__init__.py @@ -0,0 +1,2 @@ +from .proxy import proxy_router +from .extractor import extractor_router diff --git a/mediaflow_proxy/routes/extractor.py b/mediaflow_proxy/routes/extractor.py new file mode 100644 index 0000000..553976d --- /dev/null +++ b/mediaflow_proxy/routes/extractor.py @@ -0,0 +1,32 @@ +from typing import Annotated + +from fastapi import APIRouter, Query, HTTPException +from fastapi.responses import RedirectResponse + +from mediaflow_proxy.configs import settings +from mediaflow_proxy.extractors.factory import ExtractorFactory +from mediaflow_proxy.schemas import ExtractorURLParams + +extractor_router = APIRouter() + + +@extractor_router.get("/video") +async def extract_url( + extractor_params: Annotated[ExtractorURLParams, Query()], +): + """Extract clean links from various video hosting services.""" + try: + extractor = ExtractorFactory.get_extractor(extractor_params.host, extractor_params.use_request_proxy) + final_url, headers = await extractor.extract(extractor_params.destination) + + if extractor_params.redirect_stream: + formatted_headers = "&".join(f"h_{k}={v}" for k, v in headers.items()) + stream_url = f"/proxy/stream?api_password={settings.api_password}&d={final_url}&{formatted_headers}" + return RedirectResponse(url=stream_url) + + return {"url": final_url, "headers": headers} + + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) + except Exception as e: + raise HTTPException(status_code=500, detail=f"Extraction failed: {str(e)}") diff --git a/mediaflow_proxy/routes.py b/mediaflow_proxy/routes/proxy.py similarity index 91% rename from mediaflow_proxy/routes.py rename to mediaflow_proxy/routes/proxy.py index a24eb2a..4efda44 100644 --- a/mediaflow_proxy/routes.py +++ b/mediaflow_proxy/routes/proxy.py @@ -2,9 +2,22 @@ from typing import Annotated from fastapi import Request, Depends, APIRouter, Query, HTTPException -from .handlers import handle_hls_stream_proxy, proxy_stream, get_manifest, get_playlist, get_segment, get_public_ip -from .schemas import MPDSegmentParams, MPDPlaylistParams, HLSManifestParams, ProxyStreamParams, MPDManifestParams -from .utils.http_utils import get_proxy_headers, ProxyRequestHeaders +from mediaflow_proxy.handlers import ( + handle_hls_stream_proxy, + proxy_stream, + get_manifest, + get_playlist, + get_segment, + get_public_ip, +) +from mediaflow_proxy.schemas import ( + MPDSegmentParams, + MPDPlaylistParams, + HLSManifestParams, + ProxyStreamParams, + MPDManifestParams, +) +from mediaflow_proxy.utils.http_utils import get_proxy_headers, ProxyRequestHeaders proxy_router = APIRouter() diff --git a/mediaflow_proxy/schemas.py b/mediaflow_proxy/schemas.py index b584ee8..58efa88 100644 --- a/mediaflow_proxy/schemas.py +++ b/mediaflow_proxy/schemas.py @@ -1,3 +1,5 @@ +from typing import Literal + from pydantic import BaseModel, Field, IPvAnyAddress, ConfigDict @@ -55,3 +57,9 @@ class MPDSegmentParams(GenericParams): mime_type: str = Field(..., description="The MIME type of the segment.") key_id: str | None = Field(None, description="The DRM key ID (optional).") key: str | None = Field(None, description="The DRM key (optional).") + + +class ExtractorURLParams(GenericParams): + host: Literal["Doodstream", "Mixdrop", "Uqload"] = Field(..., description="The host to extract the URL from.") + destination: str = Field(..., description="The URL of the stream.", alias="d") + redirect_stream: bool = Field(False, description="Whether to redirect to the stream endpoint automatically.") diff --git a/mediaflow_proxy/utils/http_utils.py b/mediaflow_proxy/utils/http_utils.py index ffea361..09e928f 100644 --- a/mediaflow_proxy/utils/http_utils.py +++ b/mediaflow_proxy/utils/http_utils.py @@ -125,6 +125,10 @@ class Streamer: else: async for chunk in self.response.aiter_bytes(): yield chunk + self.bytes_transferred += len(chunk) + except httpx.TimeoutException: + logger.warning(f"Timeout while streaming {url}") + raise DownloadError(409, f"Timeout while streaming {url}") except GeneratorExit: logger.info("Streaming session stopped by the user") except Exception as e: