api / core /app.py
Codebuff
Gate streams on real bytes, not pipe readability
b4021c5
Raw History Blame Contribute Delete
6.21 kB
"""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),
}