"""FFprobe-based stream metadata — resolution, FPS, codecs, bitrate.""" import asyncio import json import logging import time logger = logging.getLogger(__name__) # Simple in-process cache: url → (info_dict, fetched_at) _cache: dict[str, tuple[dict, float]] = {} _CACHE_TTL = 3600 # 1 hour — re-probe if URL expires or stream restarts # Locks per URL to avoid parallel probes for the same stream _locks: dict[str, asyncio.Lock] = {} async def probe_stream(url: str) -> dict: """ Run ffprobe against *url* and return a dict with video/audio metadata. Returns {} on failure (non-fatal — the stream still works). Results are cached for CACHE_TTL seconds. """ now = time.time() cached = _cache.get(url) if cached and now - cached[1] < _CACHE_TTL: return cached[0] if url not in _locks: _locks[url] = asyncio.Lock() async with _locks[url]: # Re-check after acquiring lock (another task may have already probed) cached = _cache.get(url) if cached and now - cached[1] < _CACHE_TTL: return cached[0] info = await _run_ffprobe(url) _cache[url] = (info, time.time()) return info async def _run_ffprobe(url: str) -> dict: try: proc = await asyncio.create_subprocess_exec( "ffprobe", "-v", "quiet", "-print_format", "json", "-show_streams", "-show_format", "-analyzeduration", "3000000", # 3 s — fast enough for live streams "-probesize", "1000000", # 1 MB url, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL, ) try: stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=20) except asyncio.TimeoutError: try: proc.kill() except ProcessLookupError: pass logger.debug(f"ffprobe timed out for {url}") return {} data = json.loads(stdout) except Exception as e: logger.debug(f"ffprobe failed for {url}: {e}") return {} info: dict = {} for stream in data.get("streams", []): ctype = stream.get("codec_type", "") if ctype == "video" and "video_codec" not in info: info["video_codec"] = stream.get("codec_name", "") info["video_profile"] = stream.get("profile", "") info["width"] = stream.get("width", 0) info["height"] = stream.get("height", 0) # FPS: prefer avg_frame_rate, fallback r_frame_rate for fps_field in ("avg_frame_rate", "r_frame_rate"): fps_str = stream.get(fps_field, "") if fps_str and "/" in fps_str: try: num, den = fps_str.split("/") fps = round(int(num) / int(den), 3) if int(den) else 0 if 1 < fps < 200: # sanity check info["fps"] = round(fps, 2) break except (ValueError, ZeroDivisionError): pass elif ctype == "audio" and "audio_codec" not in info: info["audio_codec"] = stream.get("codec_name", "") info["audio_channels"] = stream.get("channels", 0) info["audio_sample_rate"] = int(stream.get("sample_rate", 0) or 0) fmt = data.get("format", {}) try: br = int(fmt.get("bit_rate", 0) or 0) if br > 0: info["bitrate_bps"] = br except (ValueError, TypeError): pass logger.debug(f"Probed {url}: {info}") return info def invalidate(url: str) -> None: """Evict a URL from the cache (call when stream URL changes).""" _cache.pop(url, None)