fix: distinguir sesiones normales del proveedor de fallos reales en logs y dashboard

- Solo registra "caída" cuando ffmpeg sale en <25s (error real: DNS, 403, auth, stall)
- Timeouts de sesión del proveedor (~2 min) NO se registran como fallos — son normales
- "Recuperado" se registra cuando el stream vuelve a funcionar tras un fallo real
- Añade clasificación automática de causa: conexión rechazada, DNS, 403, timeout...
- Dashboard: "reinicios" ahora distingue sesiones del proveedor (gris, normal) vs fallos (rojo)
- Registro: nueva columna "Duración fallo", proveedor+dominio juntos, filas con color de fondo
- StreamEvent: campo pump_duration_secs para persistir duración del fallo o interrupción

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
KiraStream
2026-05-19 15:12:04 +00:00
parent f965ccafa8
commit 4bfc7592a6
12 changed files with 474 additions and 362 deletions
+111 -45
View File
@@ -51,6 +51,12 @@ _SENTINEL = object() # signals end-of-stream to client queues
_FFMPEG = shutil.which("ffmpeg") or "ffmpeg"
# Pumps that run for less than this many seconds before exiting are classified
# as real failures (connection refused, DNS error, auth expired, etc.).
# Pumps that run longer are normal provider session timeouts (~2 min) and are
# NOT logged as failures — they are transparent to the user thanks to the buffer.
PROVIDER_SESSION_MIN = 25.0
# ffmpeg stderr patterns that indicate input-side A/V issues (demuxer warnings).
# These don't affect output quality with -use_wallclock_as_timestamps but are
# counted as a health metric to surface in the dashboard.
@@ -66,6 +72,28 @@ _HEALTH_WARN_SECS = 12.0 # > 12 s without data → warning (normal reconnects
_HEALTH_ERROR_SECS = 25.0 # > 25 s → error, trigger correction
def _classify_error(exc: Exception, pump_duration: float) -> str:
"""Return a human-readable Spanish cause for a stream failure."""
msg = str(exc).lower()
if "connection refused" in msg or "errno 111" in msg:
return "Conexión rechazada por el servidor del proveedor"
if "name or service not known" in msg or "could not resolve" in msg or "getaddrinfo" in msg:
return "No se pudo resolver el dominio (DNS)"
if "no output" in msg or ("timeout" in msg and "no" in msg):
return f"Sin datos del proveedor durante {int(pump_duration)}s (stall)"
if "403" in msg:
return "Acceso denegado por el proveedor (403 Forbidden)"
if "404" in msg:
return "Canal no encontrado en el proveedor (404)"
if "401" in msg:
return "Autenticación rechazada por el proveedor (401)"
if "rc=1" in msg or "rc=-1" in msg:
return f"Error del proveedor (ffmpeg rc≠0) tras {int(pump_duration)}s"
if "rc=0" in msg:
return f"Sesión del proveedor cerrada tras {int(pump_duration)}s"
return f"Error en {int(pump_duration)}s: {str(exc)[:200]}"
def _url_netloc(url: str) -> str:
from urllib.parse import urlparse
try:
@@ -134,6 +162,8 @@ class BroadcastGroup:
# Health / metadata
self.stream_info: dict = {}
self._reconnect_count: int = 0
self._session_count: int = 0 # normal provider session timeouts (>PROVIDER_SESSION_MIN)
self._failure_count: int = 0 # real quick failures (<PROVIDER_SESSION_MIN)
self._status: str = "connecting"
self._last_data_at: float = time.time()
@@ -251,6 +281,8 @@ class BroadcastGroup:
"running": self._running,
"status": self._status,
"reconnect_count": self._reconnect_count,
"session_count": self._session_count,
"failure_count": self._failure_count,
"last_data_ago": round(time.time() - self._last_data_at, 1),
"stream_info": self.stream_info,
# ffmpeg process info
@@ -364,6 +396,7 @@ class BroadcastGroup:
domain: str,
error_message: str,
attempt: int,
pump_duration_secs: int | None = None,
) -> None:
"""Persist a stream failure/recovery event to the database (fire-and-forget)."""
try:
@@ -378,6 +411,7 @@ class BroadcastGroup:
domain=domain,
error_message=error_message[:500] if error_message else None,
attempt_number=attempt,
pump_duration_secs=pump_duration_secs,
)
db.add(ev)
await db.commit()
@@ -393,10 +427,15 @@ class BroadcastGroup:
logger.debug(f"probe_info failed for ch{self.channel_id}: {e}")
async def _pump_with_retry(self) -> None:
consecutive = 0
base_wait = 0.5
url_idx = 0
n = len(self._stream_urls) or 1
# Failure-mode tracking: a "real failure" is any pump that exits in
# less than PROVIDER_SESSION_MIN seconds. Longer pumps are normal
# provider session timeouts (~2 min) and are invisible to the user.
in_failure_mode = False # True once at least one real failure has occurred
failure_start: float | None = None
failure_attempts = 0 # consecutive real failures in this failure episode
base_wait = 0.5
url_idx = 0
n = len(self._stream_urls) or 1
try:
while self._running:
@@ -413,54 +452,81 @@ class BroadcastGroup:
if not self._running:
break
pump_duration = time.time() - pump_start
# Recovery detection: had failures before but this attempt ran stably
# for >30s — the stream was healthy and dropped again (e.g. provider
# session timeout). Log recovery, then treat next failure as fresh drop.
if consecutive > 1 and pump_duration > 30:
asyncio.create_task(self._log_stream_event(
"recuperado",
_url_netloc(current_url),
f"Estable {int(pump_duration)}s antes de caer de nuevo",
consecutive,
))
consecutive = 0 # reset so next failure is logged as a fresh "caida"
consecutive += 1
self._reconnect_count += 1
next_idx = (url_idx + 1) % n
cycle_done = (consecutive % n == 0)
# Log the first drop of each failure sequence
if consecutive == 1:
asyncio.create_task(self._log_stream_event(
"caida",
_url_netloc(current_url),
str(e),
1,
))
is_real_failure = pump_duration < PROVIDER_SESSION_MIN
if n > 1:
logger.warning(
f"[ch{self.channel_id}] Domain {_url_netloc(current_url)} failed "
f"(attempt {consecutive}), trying {_url_netloc(self._stream_urls[next_idx])}: {e}"
)
if is_real_failure:
# Quick exit → real provider/network error
self._failure_count += 1
failure_attempts += 1
cause = _classify_error(e, pump_duration)
if not in_failure_mode:
# First real failure of this episode → log caída
in_failure_mode = True
failure_start = time.time()
asyncio.create_task(self._log_stream_event(
"caida",
_url_netloc(current_url),
cause,
1,
int(pump_duration),
))
logger.warning(
f"[ch{self.channel_id}] CAÍDA — {cause} "
f"(pump={pump_duration:.1f}s)"
)
else:
# Still failing — don't spam the log table
logger.warning(
f"[ch{self.channel_id}] Fallo continuado intento {failure_attempts} — "
f"{cause} (pump={pump_duration:.1f}s)"
)
else:
logger.warning(
f"[ch{self.channel_id}] Upstream error (attempt {consecutive}): {e}"
)
# Long pump → normal provider session timeout
self._session_count += 1
url_idx = next_idx
if in_failure_mode:
# Was in failure mode, now ran long → stream recovered
outage_secs = int(time.time() - (failure_start or time.time()))
asyncio.create_task(self._log_stream_event(
"recuperado",
_url_netloc(current_url),
f"Restablecido tras {failure_attempts} intento(s) y {outage_secs}s de interrupción",
failure_attempts,
outage_secs,
))
logger.info(
f"[ch{self.channel_id}] RECUPERADO — tras {failure_attempts} "
f"intento(s) en {outage_secs}s"
)
in_failure_mode = False
failure_start = None
failure_attempts = 0
else:
# Healthy session timeout — normal, no event logged
logger.info(
f"[ch{self.channel_id}] Sesión del proveedor cerrada tras "
f"{pump_duration:.1f}s (normal, sesión #{self._session_count})"
)
next_idx = (url_idx + 1) % n
url_idx = next_idx
self._status = "reconnecting"
if cycle_done:
# All domains tried — back off before next cycle.
# Do NOT drain client queues: the buffered data (up to 32 MB) covers
# the reconnect window so players never see a black screen.
wait = min(base_wait * (2 ** (consecutive // n - 1)), _BACKOFF_CAP)
logger.warning(f"[ch{self.channel_id}] All {n} domain(s) failed, retrying in {wait:.1f}s")
await asyncio.sleep(wait)
# else: try next domain immediately (no sleep)
if is_real_failure:
cycle_done = (failure_attempts % n == 0)
if cycle_done:
# All domains tried — back off before next cycle
wait = min(base_wait * (2 ** (failure_attempts // n - 1)), _BACKOFF_CAP)
logger.warning(
f"[ch{self.channel_id}] All {n} domain(s) failed "
f"(attempt {failure_attempts}), retrying in {wait:.1f}s"
)
await asyncio.sleep(wait)
# Normal session timeout: reconnect immediately (no sleep needed)
finally:
self._status = "stopped"
self._running = False