"""FastAPI surface. Scope is deliberately two endpoints: ``GET /formats`` and ``GET /stream``. ``POST /download`` and ``GET /files/{name}`` are gone, because writing a 10 hour VOD to a Space's ephemeral disk is not something this service should offer. ``/debug/*`` is temporary scaffolding used to prove the ladder against the live Space; it is removed once Kick and Twitch are verified. """ from __future__ import annotations import os import shutil import time from collections import deque from fastapi import FastAPI, HTTPException, Query from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import StreamingResponse from core.logging import log, track from core.naming import name_from_url from core.pipeline import ( FFMPEG, FFMPEG_FIRST_BYTE, Relay, iter_body, kill_proc, pump_stderr, start_ffmpeg, tail, ) from sources.base import STREAMLINK, YTDLP, StreamOpenError, list_variants, sniff ENABLE_DEBUG = os.environ.get("ENABLE_DEBUG", "1") != "0" app = FastAPI(title="VOD Downloader", version="2.0") app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"], expose_headers=["Content-Disposition", "X-Upstream"], ) def pick_source(url: str): value = (url or "").lower() if "kick.com" in value: import sources.kick as module elif "twitch.tv" in value: import sources.twitch as module elif "youtube.com" in value or "youtu.be" in value: import sources.youtube as module else: raise HTTPException(400, "unsupported host (kick.com, twitch.tv, youtube.com only)") return module def source_context(url: str): module = pick_source(url) return ( module, getattr(module, "ORDER", ("ytdlp",)), getattr(module, "RESOLVERS", None), ) @app.get("/") def index(): return { "service": "vod-downloader", "endpoints": ["/formats?url=", "/stream?url=&format_id=", "/health"], } @app.get("/health") def health(): return { "ok": True, "debug": ENABLE_DEBUG, "tools": { "ffmpeg": shutil.which(FFMPEG) or None, "streamlink": shutil.which(STREAMLINK) or None, "yt_dlp": shutil.which(YTDLP) or None, }, "cookies": os.path.isfile("cookies.txt") or os.path.isfile("/data/cookies.txt"), } @app.get("/formats") def formats(url: str = Query(..., description="Kick / Twitch / YouTube VOD url")): started = time.monotonic() try: module = pick_source(url) variants = module.formats(url) except HTTPException: raise except Exception as exc: # noqa: BLE001 - surface extractor errors to the client track("formats", url=url, status="error", detail=str(exc)[:200]) raise HTTPException(502, str(exc)[:600]) from exc track( "formats", url=url, status="ok", detail=f"{len(variants)} variants in {int((time.monotonic() - started) * 1000)}ms", ) return variants @app.get("/stream") def stream( url: str = Query(..., description="Kick / Twitch / YouTube VOD url"), format_id: str = Query("best", description="an id from /formats"), ): module = pick_source(url) started = time.monotonic() try: upstream = module.open(url, format_id) except StreamOpenError as exc: track("stream", url=url, status="extract_failed", detail=str(exc)[:200]) raise HTTPException( 502, { "error": "no rung produced data", "format_id": format_id, "tried": str(exc)[:900], }, ) from exc except HTTPException: raise except Exception as exc: # noqa: BLE001 track("stream", url=url, status="error", detail=str(exc)[:200]) raise HTTPException(502, str(exc)[:600]) from exc ff = start_ffmpeg(upstream) ff_sink: deque = deque() pump_stderr(ff, ff_sink) ff_relay = Relay(ff) # Gate on ffmpeg's first *byte*, not on pipe readability: a pipe whose # writer exited also reads as ready, which is how a dead stream used to # become a 200 with an empty body. if not ff_relay.wait_first(FFMPEG_FIRST_BYTE): kill_proc(ff) kill_proc(upstream.proc) ff_relay.close() upstream.relay.close() detail = { "error": "remux produced no output", "upstream": upstream.name, "upstream_stderr": tail(upstream.sink), "ffmpeg_stderr": tail(ff_sink), } track("stream", url=url, status="remux_failed", detail=str(detail)[:200]) raise HTTPException(502, detail) filename = f"{name_from_url(url)}.mp4" log.info( "stream ok upstream=%s format=%s in %dms", upstream.name, format_id, int((time.monotonic() - started) * 1000), ) track("stream", url=url, status="ok", detail=upstream.name) return StreamingResponse( iter_body(upstream, ff, ff_relay, ff_sink), media_type="video/mp4", headers={ "Content-Disposition": f'attachment; filename="{filename}"', "X-Upstream": upstream.name, }, ) if ENABLE_DEBUG: @app.get("/debug/probe") def debug_probe( url: str = Query(...), format_id: str = Query("best"), seconds: float = Query(12.0, ge=1, le=30), ): """Report what every listing rung and every streaming candidate really does.""" module, order, resolvers = source_context(url) listing: list[dict] = [] variants: list[dict] = [] listed = None try: listed, variants = list_variants(url, order, listing) except Exception as exc: # noqa: BLE001 listed = None listing.append({"rung": "all", "ok": False, "ms": 0, "detail": str(exc)[:300]}) return { "platform": module.__name__, "order": list(order), "listing": listing, "listed_by": listed, "variants": variants, "sniff": sniff(url, format_id, order, resolvers, seconds=seconds), }