2026-09-16 18:56:17 +00:00
|
|
|
"""
|
|
|
|
|
Encapsulates all HTTP/Socket.IO calls to the external MeTube service.
|
|
|
|
|
|
|
|
|
|
Verified against the real MeTube source (alexta69/metube, app/main.py + app/ytdl.py):
|
|
|
|
|
|
|
|
|
|
- POST /add only accepts/returns {"status": "ok"|"error", "msg": ...} — no job id.
|
|
|
|
|
The job id/url/status are only observable via Socket.IO events or GET /history.
|
|
|
|
|
- MeTube itself dedups by URL (a second /add for a queued URL is a no-op "ok").
|
|
|
|
|
- Socket.IO 'added'/'updated'/'completed' events carry DownloadInfo.to_public_dict(),
|
|
|
|
|
JSON-*string*-encoded (json.JSONEncoder().encode(...)), not a raw object — must
|
|
|
|
|
json.loads() the payload. Relevant keys: id, title, url, status, msg, percent
|
|
|
|
|
(float 0-100 or None), filename (already relative to DOWNLOAD_DIR), error.
|
|
|
|
|
'canceled'/'cleared' carry just an id (also JSON-string-encoded).
|
|
|
|
|
- MeTube status vocabulary: pending/preparing/scheduled/downloading/postprocessing/
|
|
|
|
|
finished/error — mapped to our own vocabulary in sync with download_jobs.
|
|
|
|
|
- GET /history returns {"queue": [...], "pending": [...], "done": [...]} of the
|
|
|
|
|
same to_public_dict() shape — used for reconciliation after our own restart.
|
|
|
|
|
- Downloaded files are served by aiohttp's static route at /download/<relative
|
|
|
|
|
path>, matching METUBE_CONTAINER_DOWNLOAD_DIR.
|
|
|
|
|
"""
|
|
|
|
|
|
|
|
|
|
import logging
|
|
|
|
|
from pathlib import PurePosixPath
|
|
|
|
|
from typing import Awaitable, Callable
|
|
|
|
|
from urllib.parse import quote
|
|
|
|
|
|
|
|
|
|
import httpx
|
|
|
|
|
import socketio
|
|
|
|
|
|
|
|
|
|
from app.config import settings
|
|
|
|
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
|
|
|
|
METUBE_STATUS_MAP: dict[str, str] = {
|
|
|
|
|
"pending": "queued",
|
|
|
|
|
"preparing": "queued",
|
|
|
|
|
"scheduled": "queued",
|
|
|
|
|
"downloading": "downloading",
|
|
|
|
|
"postprocessing": "postprocessing",
|
|
|
|
|
"finished": "completed",
|
|
|
|
|
"error": "failed",
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
EventHandler = Callable[[str, dict | str], Awaitable[None]]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class MeTubeClient:
|
|
|
|
|
def __init__(self) -> None:
|
|
|
|
|
self.api_base_url = settings.metube_api_base_url.rstrip("/")
|
|
|
|
|
self.public_base_url = settings.metube_public_base_url.rstrip("/")
|
|
|
|
|
self.download_dir = settings.metube_container_download_dir
|
|
|
|
|
self.timeout = settings.metube_request_timeout_seconds
|
|
|
|
|
|
|
|
|
|
def enqueue_video(self, youtube_url: str, custom_name_prefix: str) -> dict:
|
|
|
|
|
payload = {
|
|
|
|
|
"url": youtube_url,
|
|
|
|
|
"download_type": "video",
|
|
|
|
|
"codec": "auto",
|
|
|
|
|
"format": "mp4",
|
|
|
|
|
"quality": "best",
|
|
|
|
|
"auto_start": True,
|
|
|
|
|
"custom_name_prefix": custom_name_prefix,
|
|
|
|
|
}
|
|
|
|
|
response = httpx.post(f"{self.api_base_url}/add", json=payload, timeout=self.timeout)
|
|
|
|
|
response.raise_for_status()
|
|
|
|
|
return response.json()
|
|
|
|
|
|
|
|
|
|
def health(self) -> bool:
|
|
|
|
|
try:
|
|
|
|
|
response = httpx.get(self.api_base_url, timeout=5)
|
|
|
|
|
return response.status_code < 500
|
|
|
|
|
except Exception:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
def fetch_history(self) -> dict:
|
|
|
|
|
response = httpx.get(f"{self.api_base_url}/history", timeout=self.timeout)
|
|
|
|
|
response.raise_for_status()
|
|
|
|
|
return response.json()
|
|
|
|
|
|
2026-09-16 19:47:35 +00:00
|
|
|
def delete_download(self, metube_job_id: str) -> dict:
|
|
|
|
|
"""Asks MeTube to remove a finished download from its 'done' list and
|
|
|
|
|
delete the underlying file (actual file deletion additionally depends
|
|
|
|
|
on MeTube's own DELETE_FILE_ON_TRASHCAN config, which we don't
|
|
|
|
|
control). Only ever called with a metube_job_id our own app tracked
|
|
|
|
|
from a download it started -- never touches files MeTube already had
|
|
|
|
|
before we existed."""
|
|
|
|
|
response = httpx.post(
|
|
|
|
|
f"{self.api_base_url}/delete",
|
|
|
|
|
json={"ids": [metube_job_id], "where": "done"},
|
|
|
|
|
timeout=self.timeout,
|
|
|
|
|
)
|
|
|
|
|
response.raise_for_status()
|
|
|
|
|
return response.json()
|
|
|
|
|
|
2026-09-16 18:56:17 +00:00
|
|
|
def build_media_url(self, filename: str) -> str | None:
|
|
|
|
|
"""Safely turn a MeTube-reported filename into a public /download/... URL.
|
|
|
|
|
|
|
|
|
|
`filename` is expected relative to METUBE_CONTAINER_DOWNLOAD_DIR (that's
|
|
|
|
|
what MeTube itself stores), but we defensively strip an accidental
|
|
|
|
|
absolute prefix and reject any path that escapes the download dir.
|
|
|
|
|
"""
|
|
|
|
|
if not filename:
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
normalized = filename.replace("\\", "/")
|
|
|
|
|
download_dir = self.download_dir.rstrip("/")
|
|
|
|
|
if download_dir and normalized.startswith(download_dir + "/"):
|
|
|
|
|
normalized = normalized[len(download_dir) + 1 :]
|
|
|
|
|
normalized = normalized.lstrip("/")
|
|
|
|
|
|
|
|
|
|
path = PurePosixPath(normalized)
|
|
|
|
|
if ".." in path.parts or path.is_absolute():
|
|
|
|
|
logger.warning("Rejected unsafe MeTube filename: %r", filename)
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
encoded = "/".join(quote(part) for part in path.parts)
|
|
|
|
|
return f"{self.public_base_url}/download/{encoded}"
|
|
|
|
|
|
|
|
|
|
def check_media(self, media_url: str) -> bool:
|
|
|
|
|
try:
|
|
|
|
|
response = httpx.head(media_url, timeout=self.timeout, follow_redirects=True)
|
|
|
|
|
return response.status_code == 200
|
|
|
|
|
except Exception:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
async def run_event_listener(self, on_event: EventHandler) -> None:
|
|
|
|
|
"""Runs forever, reconnecting automatically, until cancelled."""
|
|
|
|
|
sio = socketio.AsyncClient(reconnection=True, reconnection_delay=5, reconnection_delay_max=30)
|
|
|
|
|
|
|
|
|
|
for event_name in ("added", "updated", "completed", "canceled", "cleared"):
|
|
|
|
|
|
|
|
|
|
async def _handler(data, _event_name=event_name):
|
|
|
|
|
await on_event(_event_name, data)
|
|
|
|
|
|
|
|
|
|
sio.on(event_name, _handler)
|
|
|
|
|
|
|
|
|
|
@sio.event
|
|
|
|
|
async def connect():
|
|
|
|
|
logger.info("Connected to MeTube Socket.IO at %s", self.api_base_url)
|
|
|
|
|
|
|
|
|
|
@sio.event
|
|
|
|
|
async def disconnect():
|
|
|
|
|
logger.warning("Disconnected from MeTube Socket.IO")
|
|
|
|
|
|
|
|
|
|
while True:
|
|
|
|
|
try:
|
|
|
|
|
await sio.connect(self.api_base_url, wait_timeout=10)
|
|
|
|
|
await sio.wait()
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.warning("MeTube Socket.IO connection failed, retrying", exc_info=True)
|
|
|
|
|
await sio.sleep(10)
|