diff --git a/backend/core/restream.py b/backend/core/restream.py index be38bb8..08b38c8 100644 --- a/backend/core/restream.py +++ b/backend/core/restream.py @@ -72,6 +72,34 @@ _HEALTH_WARN_SECS = 12.0 # > 12 s without data → warning (normal reconnects _HEALTH_ERROR_SECS = 25.0 # > 25 s → error, trigger correction +def _is_provider_session_end(last_stderr: str) -> bool: + """True when ffmpeg's last stderr line indicates the provider closed the TCP connection. + + "Error during demuxing: Connection timed out" is the provider's session boundary + signal regardless of how long that session lasted. It is NOT an error — the + provider's CDN simply ended the HTTP/TCP session (common every 30-300 s). + """ + low = last_stderr.lower() + return "connection timed out" in low + + +def _is_permanent_domain_failure(exc: Exception) -> bool: + """True only for failures that won't self-heal in a few seconds. + + DNS NXDOMAIN and connection-refused are structural — the domain is unreachable. + HTTP 4xx (458, 401, etc.) and EOF are ephemeral: often a provider token-refresh + race condition that resolves on the next attempt or within seconds. + """ + msg = str(exc).lower() + return ( + "name or service not known" in msg + or "could not resolve" in msg + or "getaddrinfo" in msg + or "connection refused" in msg + or "errno 111" in msg + ) + + def _classify_error(exc: Exception, pump_duration: float) -> str: """Return a human-readable Spanish cause for a stream failure.""" msg = str(exc).lower() @@ -87,10 +115,12 @@ def _classify_error(exc: Exception, pump_duration: float) -> str: return "Canal no encontrado en el proveedor (404)" if "401" in msg: return "Autenticación rechazada por el proveedor (401)" + if "458" in msg: + return "Token de sesión no disponible en el proveedor (HTTP 458)" + if "end of file" in msg: + return "Proveedor cerró la conexión (EOF)" 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]}" @@ -179,6 +209,9 @@ class BroadcastGroup: self._corrections: int = 0 # how many times health monitor killed stuck ffmpeg self._health_task: asyncio.Task | None = None + # Last stderr line from ffmpeg — used to classify pump exit reason + self._last_ffmpeg_stderr: str = "" + # ------------------------------------------------------------------ # Public API # ------------------------------------------------------------------ @@ -471,6 +504,7 @@ class BroadcastGroup: while self._running: current_url = self._stream_urls[url_idx % n] self._active_url = current_url + self._last_ffmpeg_stderr = "" # reset so exit classification is fresh pump_start = time.time() try: await self._pump_once(current_url) @@ -484,10 +518,15 @@ class BroadcastGroup: pump_duration = time.time() - pump_start self._reconnect_count += 1 - is_real_failure = pump_duration < PROVIDER_SESSION_MIN + # Primary classification: if ffmpeg's last stderr says + # "Connection timed out" it is always a provider session boundary — + # the duration doesn't matter (sessions can be 20 s or 200 s). + # Only fall back to the duration threshold for other error types. + is_session_end = _is_provider_session_end(self._last_ffmpeg_stderr) + is_real_failure = (not is_session_end) and (pump_duration < PROVIDER_SESSION_MIN) if is_real_failure: - # Quick exit → real error (DNS, refused, HTTP 4xx, stall) + # Quick exit with a non-session-timeout error self._failure_count += 1 failure_attempts += 1 cause = _classify_error(e, pump_duration) @@ -505,8 +544,11 @@ class BroadcastGroup: f"{cause} (pump={pump_duration:.1f}s)" ) - # Mark domain as error immediately so other streams skip it - asyncio.create_task(self._mark_domain_down(current_url)) + # Only mark domain as permanently down for DNS/refused failures. + # HTTP 458 / EOF are ephemeral (provider token refresh race) — + # marking them as error causes unnecessary domain rotation. + if _is_permanent_domain_failure(e): + asyncio.create_task(self._mark_domain_down(current_url)) # Rotate to next domain url_idx = (url_idx + 1) % n @@ -523,10 +565,10 @@ class BroadcastGroup: # else: try next domain immediately else: - # Normal provider session timeout (>25 s) — reconnect to same domain + # Provider session boundary (TCP closed by provider, any duration) self._session_count += 1 - preferred_idx = url_idx # remember: this domain works - url_idx = preferred_idx # stay on it + preferred_idx = url_idx + url_idx = preferred_idx # stay on same domain if in_failure_mode: outage_secs = int(time.time() - (failure_start or time.time())) @@ -550,7 +592,9 @@ class BroadcastGroup: ) self._status = "reconnecting" - # Reconnect immediately — no backoff for normal session rotations + # Brief pause before reconnect — gives provider token-refresh + # system time to update before we open a new session. + await asyncio.sleep(0.4) finally: self._status = "stopped" @@ -654,6 +698,7 @@ class BroadcastGroup: async for line in proc.stderr: text = line.decode(errors="replace").rstrip() if text: + self._last_ffmpeg_stderr = text # track for exit classification logger.warning(f"[ch{self.channel_id}] ffmpeg: {text}") if any(p in text for p in _AV_WARN_PATTERNS): self._av_desync_count += 1