"""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" # 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 _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._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 # ------------------------------------------------------------------ # 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, "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, ) -> 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, ) 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 _pump_with_retry(self) -> None: consecutive = 0 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 pump_start = time.time() try: await self._pump_once(current_url) url_idx = 0 # reset to primary after successful pump break except asyncio.CancelledError: raise except Exception as e: 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, )) 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}" ) else: logger.warning( f"[ch{self.channel_id}] Upstream error (attempt {consecutive}): {e}" ) 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) 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: 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