ab038030eb
- Connection timed out en stderr de ffmpeg = siempre fin de sesión del proveedor, independientemente de la duración (sesiones de 20s o 200s, mismo mecanismo) - HTTP 458 / EOF son fallos efímeros (race condition de token refresh) — no se marcan como dominios permanentemente caídos, solo DNS fail y connection refused lo son - _mark_domain_down solo se llama para fallos permanentes (DNS NXDOMAIN, refused) - Pausa de 0.4s antes de reconectar tras sesión normal: da tiempo al proveedor a refrescar el token antes de que abramos una nueva sesión - _last_ffmpeg_stderr: nueva variable de estado para clasificar el motivo de salida Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
747 lines
31 KiB
Python
747 lines
31 KiB
Python
"""BroadcastGroup: fan-out a single provider stream to N simultaneous clients.
|
||
|
||
Recovery model
|
||
--------------
|
||
The pump loop runs indefinitely while self._running is True. Each iteration
|
||
spawns an ffmpeg process in copy (remux) mode. ffmpeg owns the upstream
|
||
HTTP connection. When the provider closes the session (typically every
|
||
~2 minutes) ffmpeg exits and the outer loop spawns a fresh process — the
|
||
provider delivers from the current live edge so timestamps are always
|
||
forward-moving and the player never sees a backward seek (rewind).
|
||
|
||
If ffmpeg itself exits (unrecoverable error or stall beyond its own timeout)
|
||
the outer loop drains all client queues, backs off briefly, and spawns a
|
||
fresh ffmpeg process — again without disconnecting clients.
|
||
|
||
SENTINEL is only pushed when stop() is called or the pump truly gives up.
|
||
|
||
Stall detection
|
||
---------------
|
||
ffmpeg's -timeout option (socket-level read timeout) handles provider stalls.
|
||
A safety-net asyncio.wait_for wraps each stdout read in case ffmpeg hangs
|
||
without producing output or exiting.
|
||
|
||
A/V sync
|
||
--------
|
||
ffmpeg re-muxes the incoming MPEG-TS into a fresh output stream. On every
|
||
reconnect the output PTS/DTS sequence is continuous regardless of provider-
|
||
side timestamp resets, eliminating audio desync entirely. No video or audio
|
||
is re-encoded — quality is bit-for-bit identical to the source.
|
||
"""
|
||
import asyncio
|
||
import logging
|
||
import os
|
||
import shutil
|
||
import time
|
||
import uuid
|
||
from dataclasses import dataclass, field
|
||
from typing import AsyncIterator
|
||
|
||
from ..config import settings
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
CHUNK_SIZE = settings.stream_chunk_size
|
||
QUEUE_MAX = settings.stream_queue_maxsize
|
||
STALL_TIMEOUT = 10.0 # seconds — passed to ffmpeg -timeout (in µs below)
|
||
MAX_RETRIES = 999 # effectively infinite; pool stops the group when 0 clients
|
||
_BACKOFF_CAP = 5.0 # max seconds between outer retries
|
||
|
||
_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.
|
||
_AV_WARN_PATTERNS = (
|
||
"timestamp discontinuity",
|
||
"Packet corrupt",
|
||
"PES packet size mismatch",
|
||
"DTS, out of order",
|
||
)
|
||
|
||
# Health check: data-age thresholds
|
||
_HEALTH_WARN_SECS = 12.0 # > 12 s without data → warning (normal reconnects are 1-4 s)
|
||
_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()
|
||
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 "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"
|
||
return f"Error en {int(pump_duration)}s: {str(exc)[:200]}"
|
||
|
||
|
||
def _url_netloc(url: str) -> str:
|
||
from urllib.parse import urlparse
|
||
try:
|
||
return urlparse(url).netloc
|
||
except Exception:
|
||
return url
|
||
|
||
|
||
@dataclass
|
||
class ClientHandle:
|
||
client_id: str
|
||
queue: asyncio.Queue = field(default_factory=lambda: asyncio.Queue(maxsize=QUEUE_MAX))
|
||
connected_at: float = field(default_factory=time.time)
|
||
is_priority: bool = False
|
||
username: str = ""
|
||
_overflow_count: int = 0
|
||
_overflow_since: float | None = None
|
||
|
||
async def read(self) -> AsyncIterator[bytes]:
|
||
while True:
|
||
chunk = await self.queue.get()
|
||
if chunk is _SENTINEL:
|
||
break
|
||
yield chunk
|
||
|
||
|
||
class BroadcastGroup:
|
||
"""Streams one provider URL and distributes chunks to all registered clients."""
|
||
|
||
def __init__(
|
||
self,
|
||
channel_id: int,
|
||
stream_urls: list[str],
|
||
provider_account_id: int,
|
||
channel_name: str = "",
|
||
provider_name: str = "",
|
||
# kept for backward compat in callers that still pass stream_url as kwarg
|
||
stream_url: str | None = None,
|
||
):
|
||
self.channel_id = channel_id
|
||
self.channel_name = channel_name
|
||
# Support old callers that pass a single stream_url
|
||
if stream_urls:
|
||
self._stream_urls = stream_urls
|
||
elif stream_url:
|
||
self._stream_urls = [stream_url]
|
||
else:
|
||
self._stream_urls = []
|
||
self.stream_url = self._stream_urls[0] if self._stream_urls else ""
|
||
self._active_url: str = self.stream_url
|
||
self.provider_account_id = provider_account_id
|
||
self.provider_name = provider_name
|
||
self.started_at = time.time()
|
||
|
||
self._clients: dict[str, ClientHandle] = {}
|
||
self._lock = asyncio.Lock()
|
||
self._task: asyncio.Task | None = None
|
||
self._running = False
|
||
|
||
# Throughput tracking
|
||
self._bytes_pumped: int = 0
|
||
self._bytes_prev: int = 0
|
||
self._bps_time: float = time.time()
|
||
self._bps_down: float = 0.0
|
||
|
||
# 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()
|
||
|
||
# ffmpeg process tracking
|
||
self._ffmpeg_pid: int | None = None
|
||
self._ffmpeg_start_time: float | None = None # wall time when ffmpeg started
|
||
self._ffmpeg_cpu: float = 0.0
|
||
self._cpu_task: asyncio.Task | None = None
|
||
|
||
# A/V health tracking
|
||
self._av_health: str = "ok" # ok / warning / error
|
||
self._av_desync_count: int = 0 # cumulative stderr A/V-warning events
|
||
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
|
||
# ------------------------------------------------------------------
|
||
|
||
@property
|
||
def client_count(self) -> int:
|
||
return len(self._clients)
|
||
|
||
@property
|
||
def client_ids(self) -> list[str]:
|
||
return list(self._clients.keys())
|
||
|
||
@property
|
||
def all_non_priority(self) -> bool:
|
||
return all(not h.is_priority for h in self._clients.values())
|
||
|
||
async def start(self) -> None:
|
||
self._running = True
|
||
self._task = asyncio.create_task(
|
||
self._pump_with_retry(), name=f"pump-ch{self.channel_id}"
|
||
)
|
||
self._cpu_task = asyncio.create_task(
|
||
self._monitor_cpu(), name=f"cpu-ch{self.channel_id}"
|
||
)
|
||
self._health_task = asyncio.create_task(
|
||
self._health_check_loop(), name=f"health-ch{self.channel_id}"
|
||
)
|
||
asyncio.create_task(self._probe_info(), name=f"probe-ch{self.channel_id}")
|
||
|
||
async def add_client(self, user_id: str, username: str = "", is_priority: bool = False) -> ClientHandle:
|
||
client_id = f"{user_id}-{uuid.uuid4().hex[:8]}"
|
||
handle = ClientHandle(client_id=client_id, is_priority=is_priority, username=username)
|
||
async with self._lock:
|
||
self._clients[client_id] = handle
|
||
logger.info(
|
||
f"Client {client_id} ({username}) joined channel {self.channel_id} "
|
||
f"(total={len(self._clients)})"
|
||
)
|
||
return handle
|
||
|
||
async def force_disconnect_client(self, client_id: str) -> bool:
|
||
async with self._lock:
|
||
handle = self._clients.get(client_id)
|
||
if handle is None:
|
||
return False
|
||
_drain_queue(handle.queue)
|
||
try:
|
||
handle.queue.put_nowait(_SENTINEL)
|
||
except asyncio.QueueFull:
|
||
pass
|
||
return True
|
||
|
||
async def remove_client(self, client_id: str) -> int:
|
||
async with self._lock:
|
||
self._clients.pop(client_id, None)
|
||
remaining = len(self._clients)
|
||
logger.info(f"Client {client_id} left channel {self.channel_id} (remaining={remaining})")
|
||
return remaining
|
||
|
||
async def stop(self) -> None:
|
||
self._running = False
|
||
self._status = "stopped"
|
||
for task in (self._cpu_task, self._health_task):
|
||
if task and not task.done():
|
||
task.cancel()
|
||
self._cpu_task = self._health_task = None
|
||
if self._task and not self._task.done():
|
||
self._task.cancel()
|
||
try:
|
||
await self._task
|
||
except asyncio.CancelledError:
|
||
pass
|
||
async with self._lock:
|
||
for handle in self._clients.values():
|
||
try:
|
||
handle.queue.put_nowait(_SENTINEL)
|
||
except asyncio.QueueFull:
|
||
pass
|
||
|
||
def stats(self) -> dict:
|
||
n = self.client_count or 1
|
||
return {
|
||
"stream_type": "live",
|
||
"channel_id": self.channel_id,
|
||
"channel_name": self.channel_name or f"Canal #{self.channel_id}",
|
||
"provider_account_id": self.provider_account_id,
|
||
"provider_name": self.provider_name,
|
||
"client_count": self.client_count,
|
||
"clients": [
|
||
{
|
||
"client_id": h.client_id,
|
||
"username": h.username or h.client_id.rsplit("-", 1)[0],
|
||
}
|
||
for h in self._clients.values()
|
||
],
|
||
"started_at": self.started_at,
|
||
"bytes_pumped": self._bytes_pumped,
|
||
"bps_down": self._bps_down,
|
||
"bps_up": self._bps_down * n,
|
||
"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
|
||
"ffmpeg_pid": self._ffmpeg_pid,
|
||
"ffmpeg_cpu_pct": round(self._ffmpeg_cpu, 2),
|
||
# multi-domain failover info
|
||
"active_url_domain": _url_netloc(self._active_url),
|
||
"url_count": len(self._stream_urls),
|
||
# A/V health
|
||
"av_health": self._av_health,
|
||
"av_desync_count": self._av_desync_count,
|
||
"corrections": self._corrections,
|
||
"buffer_pct": self._buffer_pct(),
|
||
"buffer_server_bytes": self._buffer_server_bytes(),
|
||
"buffer_server_secs": self._buffer_server_secs(),
|
||
"buffer_capacity": QUEUE_MAX * CHUNK_SIZE,
|
||
}
|
||
|
||
# ------------------------------------------------------------------
|
||
# Internal helpers
|
||
# ------------------------------------------------------------------
|
||
|
||
async def _monitor_cpu(self) -> None:
|
||
"""Update self._ffmpeg_cpu every 3 s using cumulative average since start.
|
||
|
||
Delta-based sampling fails for low-CPU processes (copy-mode ffmpeg uses
|
||
<1% CPU) because the tick delta in 2 s is 0-1 integer ticks, making the
|
||
result always 0.0%. Cumulative average ((total_cpu_ticks / clk_tck) /
|
||
elapsed_wall) gives a stable, accurate reading — the same method ps uses.
|
||
"""
|
||
clk_tck = float(os.sysconf(os.sysconf_names.get("SC_CLK_TCK", 2)) or 100)
|
||
|
||
while self._running:
|
||
await asyncio.sleep(3)
|
||
pid = self._ffmpeg_pid
|
||
start_time = self._ffmpeg_start_time
|
||
if pid is None or start_time is None:
|
||
self._ffmpeg_cpu = 0.0
|
||
continue
|
||
try:
|
||
with open(f"/proc/{pid}/stat") as f:
|
||
parts = f.read().split()
|
||
cpu_ticks = int(parts[13]) + int(parts[14]) # utime + stime
|
||
elapsed = time.time() - start_time
|
||
if elapsed > 0:
|
||
self._ffmpeg_cpu = cpu_ticks / clk_tck / elapsed * 100
|
||
except (FileNotFoundError, IndexError, ValueError, OSError):
|
||
self._ffmpeg_cpu = 0.0
|
||
|
||
def _buffer_pct(self) -> float:
|
||
"""Return queue fill % of the slowest (most backlogged) client (0–100)."""
|
||
if not self._clients:
|
||
return 0.0
|
||
max_fill = max(h.queue.qsize() for h in self._clients.values())
|
||
return round(max_fill / QUEUE_MAX * 100, 1)
|
||
|
||
def _buffer_server_bytes(self) -> int:
|
||
"""Bytes buffered for the least-buffered client (first to run dry on provider drop)."""
|
||
if not self._clients:
|
||
return 0
|
||
return min(h.queue.qsize() for h in self._clients.values()) * CHUNK_SIZE
|
||
|
||
def _buffer_server_secs(self) -> float:
|
||
"""Estimated seconds of coverage at current bitrate."""
|
||
if not self._clients or self._bps_down <= 0:
|
||
return 0.0
|
||
return round(self._buffer_server_bytes() / self._bps_down, 1)
|
||
|
||
async def _health_check_loop(self) -> None:
|
||
"""Update _av_health every 5 s based on data-flow age.
|
||
|
||
If ffmpeg is running but producing no output for _HEALTH_ERROR_SECS,
|
||
kill it so _pump_with_retry can restart it immediately — this is the
|
||
auto-correction the caller sees as a 'corrections' counter tick.
|
||
"""
|
||
consecutive_error = 0
|
||
while self._running:
|
||
await asyncio.sleep(5)
|
||
if not self._running:
|
||
break
|
||
age = time.time() - self._last_data_at
|
||
if age < _HEALTH_WARN_SECS:
|
||
self._av_health = "ok"
|
||
consecutive_error = 0
|
||
elif age < _HEALTH_ERROR_SECS:
|
||
self._av_health = "warning"
|
||
consecutive_error = 0
|
||
else:
|
||
self._av_health = "error"
|
||
consecutive_error += 1
|
||
# After 2 consecutive error checks (≈10 s) kill a stuck ffmpeg
|
||
# so _pump_with_retry can restart it without waiting for the
|
||
# 25-second safety-net timeout in _pump_once.
|
||
if consecutive_error >= 2 and self._ffmpeg_pid is not None:
|
||
import signal as _sig
|
||
try:
|
||
os.kill(self._ffmpeg_pid, _sig.SIGTERM)
|
||
self._corrections += 1
|
||
logger.warning(
|
||
f"[ch{self.channel_id}] Health correction #{self._corrections}: "
|
||
f"killed stuck ffmpeg pid={self._ffmpeg_pid} "
|
||
f"(no data for {age:.1f}s)"
|
||
)
|
||
except (ProcessLookupError, PermissionError):
|
||
pass
|
||
consecutive_error = 0 # reset so we don't spam kills
|
||
|
||
async def _log_stream_event(
|
||
self,
|
||
event_type: str,
|
||
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:
|
||
from ..database import AsyncSessionLocal
|
||
from ..models.stream_event import StreamEvent
|
||
async with AsyncSessionLocal() as db:
|
||
ev = StreamEvent(
|
||
channel_id=self.channel_id,
|
||
channel_name=self.channel_name or f"Canal #{self.channel_id}",
|
||
provider_name=self.provider_name,
|
||
event_type=event_type,
|
||
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()
|
||
except Exception as exc:
|
||
logger.debug(f"[ch{self.channel_id}] Failed to persist stream event: {exc}")
|
||
|
||
async def _probe_info(self) -> None:
|
||
from .probe import probe_stream
|
||
try:
|
||
info = await probe_stream(self.stream_url)
|
||
self.stream_info.update(info)
|
||
except Exception as e:
|
||
logger.debug(f"probe_info failed for ch{self.channel_id}: {e}")
|
||
|
||
async def _mark_domain_down(self, url: str) -> None:
|
||
"""Update ProviderUrl.status='error' immediately so other streams skip this domain."""
|
||
from urllib.parse import urlparse
|
||
from datetime import datetime, timezone
|
||
from ..database import AsyncSessionLocal
|
||
from ..models.provider import ProviderUrl
|
||
from sqlalchemy import select
|
||
netloc = urlparse(url).netloc
|
||
if not netloc:
|
||
return
|
||
try:
|
||
async with AsyncSessionLocal() as db:
|
||
result = await db.execute(
|
||
select(ProviderUrl).where(
|
||
ProviderUrl.provider_account_id == self.provider_account_id,
|
||
ProviderUrl.url.ilike(f"%{netloc}%"),
|
||
)
|
||
)
|
||
pu = result.scalar_one_or_none()
|
||
if pu and pu.status != "error":
|
||
pu.status = "error"
|
||
pu.last_checked_at = datetime.now(timezone.utc)
|
||
await db.commit()
|
||
logger.info(f"[ch{self.channel_id}] Dominio {netloc} marcado como error (fallo rápido)")
|
||
except Exception as exc:
|
||
logger.debug(f"[ch{self.channel_id}] _mark_domain_down failed: {exc}")
|
||
|
||
async def _pump_with_retry(self) -> None:
|
||
# Domain stickiness: on normal provider session timeouts (>PROVIDER_SESSION_MIN)
|
||
# we reconnect to the same domain — providers rotate DNS/CDN themselves, so
|
||
# staying on the same host avoids unnecessary domain cycling.
|
||
# Only on real quick failures (<PROVIDER_SESSION_MIN) do we advance to the
|
||
# next domain and immediately mark the failed one as 'error' in the DB.
|
||
in_failure_mode = False
|
||
failure_start: float | None = None
|
||
failure_attempts = 0
|
||
preferred_idx = 0 # index of last domain that served a full session
|
||
base_wait = 0.5
|
||
url_idx = 0
|
||
n = len(self._stream_urls) or 1
|
||
|
||
try:
|
||
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)
|
||
url_idx = 0
|
||
break
|
||
except asyncio.CancelledError:
|
||
raise
|
||
except Exception as e:
|
||
if not self._running:
|
||
break
|
||
pump_duration = time.time() - pump_start
|
||
self._reconnect_count += 1
|
||
|
||
# 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 with a non-session-timeout error
|
||
self._failure_count += 1
|
||
failure_attempts += 1
|
||
cause = _classify_error(e, pump_duration)
|
||
|
||
if not in_failure_mode:
|
||
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} (pump={pump_duration:.1f}s)")
|
||
else:
|
||
logger.warning(
|
||
f"[ch{self.channel_id}] Fallo continuado intento {failure_attempts} — "
|
||
f"{cause} (pump={pump_duration:.1f}s)"
|
||
)
|
||
|
||
# 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
|
||
|
||
cycle_done = (failure_attempts % n == 0)
|
||
self._status = "reconnecting"
|
||
if cycle_done:
|
||
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)
|
||
# else: try next domain immediately
|
||
|
||
else:
|
||
# Provider session boundary (TCP closed by provider, any duration)
|
||
self._session_count += 1
|
||
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()))
|
||
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, dominio estable: {_url_netloc(current_url)}"
|
||
)
|
||
in_failure_mode = False
|
||
failure_start = None
|
||
failure_attempts = 0
|
||
else:
|
||
logger.info(
|
||
f"[ch{self.channel_id}] Sesión del proveedor cerrada tras "
|
||
f"{pump_duration:.1f}s (normal, sesión #{self._session_count}, "
|
||
f"dominio: {_url_netloc(current_url)})"
|
||
)
|
||
|
||
self._status = "reconnecting"
|
||
# 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"
|
||
self._running = False
|
||
async with self._lock:
|
||
for handle in self._clients.values():
|
||
_drain_queue(handle.queue)
|
||
try:
|
||
handle.queue.put_nowait(_SENTINEL)
|
||
except asyncio.QueueFull:
|
||
pass
|
||
|
||
async def _pump_once(self, url: str | None = None) -> None:
|
||
"""
|
||
Spawn ffmpeg in copy (remux) mode and pump its stdout to all client
|
||
queues until ffmpeg exits or self._running goes False.
|
||
|
||
Reconnection is intentionally NOT delegated to ffmpeg (-reconnect_streamed
|
||
is omitted). When the provider closes the TCP session (typically every
|
||
~2 minutes), ffmpeg exits cleanly and _pump_with_retry immediately spawns
|
||
a fresh process that opens a brand-new HTTP connection. The provider then
|
||
delivers the stream from the current live edge — timestamps are always
|
||
forward-moving. If ffmpeg were to reconnect internally it would often
|
||
receive a few seconds of already-seen data (the provider rewinds to the
|
||
last keyframe boundary), causing a visible backward seek in the player.
|
||
"""
|
||
stream_url = url if url is not None else self.stream_url
|
||
cmd = [
|
||
_FFMPEG,
|
||
"-hide_banner", "-loglevel", "warning",
|
||
"-timeout", str(int(STALL_TIMEOUT * 1_000_000)),
|
||
"-i", stream_url,
|
||
"-c", "copy",
|
||
"-f", "mpegts",
|
||
"pipe:1",
|
||
]
|
||
|
||
proc = await asyncio.create_subprocess_exec(
|
||
*cmd,
|
||
stdout=asyncio.subprocess.PIPE,
|
||
stderr=asyncio.subprocess.PIPE,
|
||
)
|
||
|
||
self._status = "ok"
|
||
self._last_data_at = time.time()
|
||
self._ffmpeg_pid = proc.pid
|
||
self._ffmpeg_start_time = time.time()
|
||
logger.info(
|
||
f"[ch{self.channel_id}] ffmpeg started pid={proc.pid} "
|
||
f"(reconnects so far: {self._reconnect_count})"
|
||
)
|
||
|
||
# Drain stderr in background so the pipe never fills and blocks ffmpeg.
|
||
stderr_task = asyncio.create_task(self._log_ffmpeg_stderr(proc))
|
||
|
||
try:
|
||
while self._running:
|
||
try:
|
||
chunk = await asyncio.wait_for(
|
||
proc.stdout.read(CHUNK_SIZE),
|
||
timeout=STALL_TIMEOUT + 15, # safety net beyond ffmpeg's own timeout
|
||
)
|
||
except asyncio.TimeoutError:
|
||
raise RuntimeError(
|
||
f"ffmpeg produced no output for {STALL_TIMEOUT + 15:.0f}s"
|
||
)
|
||
|
||
if not chunk:
|
||
# ffmpeg exited cleanly or with error
|
||
rc = proc.returncode
|
||
raise RuntimeError(f"ffmpeg exited (rc={rc})")
|
||
|
||
# Throughput accounting
|
||
self._bytes_pumped += len(chunk)
|
||
self._last_data_at = time.time()
|
||
now = self._last_data_at
|
||
if now - self._bps_time >= 2.0:
|
||
self._bps_down = (self._bytes_pumped - self._bytes_prev) / (now - self._bps_time)
|
||
self._bytes_prev = self._bytes_pumped
|
||
self._bps_time = now
|
||
|
||
await self._dispatch(chunk, now)
|
||
|
||
finally:
|
||
self._ffmpeg_pid = None
|
||
self._ffmpeg_start_time = None
|
||
self._ffmpeg_cpu = 0.0
|
||
stderr_task.cancel()
|
||
try:
|
||
proc.kill()
|
||
except ProcessLookupError:
|
||
pass
|
||
try:
|
||
await asyncio.wait_for(proc.wait(), timeout=5)
|
||
except asyncio.TimeoutError:
|
||
pass
|
||
|
||
async def _log_ffmpeg_stderr(self, proc: asyncio.subprocess.Process) -> None:
|
||
"""Read ffmpeg stderr lines, log them and count A/V-quality events."""
|
||
try:
|
||
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
|
||
except asyncio.CancelledError:
|
||
pass
|
||
except Exception:
|
||
pass
|
||
|
||
async def _dispatch(self, chunk: bytes, now: float) -> None:
|
||
"""Push chunk to every client queue; evict clients that have been full for >30 s."""
|
||
async with self._lock:
|
||
dead: list[str] = []
|
||
for cid, handle in self._clients.items():
|
||
try:
|
||
handle.queue.put_nowait(chunk)
|
||
handle._overflow_count = 0
|
||
handle._overflow_since = None
|
||
except asyncio.QueueFull:
|
||
handle._overflow_count += 1
|
||
if handle._overflow_since is None:
|
||
handle._overflow_since = now
|
||
elif now - handle._overflow_since > 30:
|
||
dead.append(cid)
|
||
|
||
for cid in dead:
|
||
logger.warning(
|
||
f"[ch{self.channel_id}] Evicting unresponsive client {cid} "
|
||
f"after 30 s of sustained queue overflow"
|
||
)
|
||
dropped = self._clients.pop(cid, None)
|
||
if dropped is not None:
|
||
_drain_queue(dropped.queue)
|
||
try:
|
||
dropped.queue.put_nowait(_SENTINEL)
|
||
except asyncio.QueueFull:
|
||
pass
|
||
|
||
|
||
def _drain_queue(q: asyncio.Queue) -> None:
|
||
"""Empty a queue non-blockingly."""
|
||
while not q.empty():
|
||
try:
|
||
q.get_nowait()
|
||
except asyncio.QueueEmpty:
|
||
break
|