Files
KiraStream b969b7e5af Mejoras de estabilidad de stream IPTV y nuevo registro de fallos
- restream.py: eliminar -reconnect_streamed de ffmpeg para evitar rebobinados
  al reconectar (el proveedor re-enviaba desde keyframe anterior causando PTS
  backward). El loop de Python gestiona las reconexiones con conexión fresca.
- restream.py: eliminar _drain_queue del loop de reconexión para que el buffer
  de la cola cubra el tiempo de reconexión sin pantalla negra en el player.
- restream.py: añadir buffer_server_bytes/secs/capacity a stats() y método
  _log_stream_event() para persistir eventos de caída y recuperación en BD.
- models/stream_event.py: nuevo modelo StreamEvent para registro de fallos.
- database.py: registrar StreamEvent en init_db.
- api/admin/logs.py: nuevos endpoints GET/DELETE /logs/events.
- Dashboard.tsx: mostrar reserva de buffer del servidor (capacidad + estado).
- Logs.tsx: añadir pestaña "Fallos de stream" con tabla de eventos persistidos.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-19 14:43:31 +00:00

113 lines
3.7 KiB
Python

"""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)