From 5179dbfa6d9f75bbfd07d4d1761418ca30ef13e0 Mon Sep 17 00:00:00 2001 From: joaquin Date: Sun, 17 May 2026 23:45:11 +0200 Subject: [PATCH] feat: replace aiohttp pump with ffmpeg copy-mode remux Uses ffmpeg -c copy (no transcoding) with -reconnect_streamed to handle provider session-token refreshes transparently. ffmpeg normalises PTS/DTS across reconnects, eliminating A/V desync permanently. Quality is bit-for- bit identical to the source stream. --- backend/core/restream.py | 208 +++++++++++++++++++++++---------------- 1 file changed, 122 insertions(+), 86 deletions(-) diff --git a/backend/core/restream.py b/backend/core/restream.py index 8922881..d129b86 100644 --- a/backend/core/restream.py +++ b/backend/core/restream.py @@ -3,44 +3,51 @@ Recovery model -------------- The pump loop runs indefinitely while self._running is True. Each iteration -opens a fresh upstream connection. If the connection fails or stalls the loop -sends SENTINEL to every client (forcing the player to reconnect from scratch), -sleeps briefly, then retries the upstream connection. +spawns an ffmpeg process in copy (remux) mode. ffmpeg owns the upstream +HTTP connection and uses its built-in -reconnect_streamed flag to handle +short provider interruptions (e.g. expiring session tokens) transparently, +normalising PTS/DTS in the output so clients never see an A/V discontinuity. -Sending SENTINEL on reconnect is intentional: keeping the TCP connection alive -while splicing in chunks from a new provider connection causes a PTS/DTS jump -that most players cannot reconcile silently — the result is permanent audio -desync. A clean player reconnect gives the player a fresh timestamp timeline -and guarantees A/V sync from the first frame of the new connection. +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 --------------- -Each chunk read is wrapped in asyncio.wait_for(STALL_TIMEOUT). If the -upstream stops sending bytes for STALL_TIMEOUT seconds we treat it as a dead -stream and reconnect immediately rather than waiting for a TCP timeout that -could take minutes. +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 shutil import time import uuid from dataclasses import dataclass, field from typing import AsyncIterator -import aiohttp - 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 of silence → reconnect (fast stall detection) -MAX_RETRIES = 999 # effectively infinite — pool stops the group when 0 clients remain -# Back-off: 0.5 s, 1 s, 2 s, 4 s … capped at 5 s -_BACKOFF_CAP = 5.0 +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 +_SENTINEL = object() # signals end-of-stream to client queues + +_FFMPEG = shutil.which("ffmpeg") or "ffmpeg" @dataclass @@ -91,9 +98,9 @@ class BroadcastGroup: self._bps_down: float = 0.0 # Health / metadata - self.stream_info: dict = {} # populated by background ffprobe + self.stream_info: dict = {} self._reconnect_count: int = 0 - self._status: str = "connecting" # connecting | ok | reconnecting | stopped + self._status: str = "connecting" self._last_data_at: float = time.time() # ------------------------------------------------------------------ @@ -117,7 +124,6 @@ class BroadcastGroup: self._task = asyncio.create_task( self._pump_with_retry(), name=f"pump-ch{self.channel_id}" ) - # Probe stream info in background — does NOT delay the first chunk 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: @@ -187,11 +193,9 @@ class BroadcastGroup: "bps_down": self._bps_down, "bps_up": self._bps_down * n, "running": self._running, - # Health "status": self._status, "reconnect_count": self._reconnect_count, "last_data_ago": round(time.time() - self._last_data_at, 1), - # Stream metadata from ffprobe "stream_info": self.stream_info, } @@ -208,19 +212,13 @@ class BroadcastGroup: logger.debug(f"probe_info failed for ch{self.channel_id}: {e}") async def _pump_with_retry(self) -> None: - """ - Outer retry loop. Each iteration calls _pump_once() which runs until - the upstream connection dies or stalls. On failure we back off and - retry without touching client connections. - """ consecutive = 0 - base_wait = 1.0 + base_wait = 0.5 try: while self._running: try: await self._pump_once() - # _pump_once returned cleanly (self._running went False) break except asyncio.CancelledError: raise @@ -232,28 +230,18 @@ class BroadcastGroup: wait = min(base_wait * (2 ** (consecutive - 1)), _BACKOFF_CAP) self._status = "reconnecting" logger.warning( - f"[ch{self.channel_id}] Upstream error (attempt {consecutive}), " - f"reconnecting in {wait:.1f}s: {e}" + f"[ch{self.channel_id}] ffmpeg exited (attempt {consecutive}), " + f"restarting in {wait:.1f}s: {e}" ) - # Send SENTINEL to every client so the player drops its - # own buffers and reconnects from scratch. This is the - # only reliable way to prevent A/V desync: if we keep the - # TCP connection alive but splice in chunks from a new - # provider connection, the player sees a PTS/DTS jump it - # cannot reconcile silently. A clean player reconnect - # gives it a fresh timeline and guarantees sync. + # Drain queues so stale buffered chunks are not delivered + # after the fresh ffmpeg process starts. async with self._lock: for handle in self._clients.values(): _drain_queue(handle.queue) - try: - handle.queue.put_nowait(_SENTINEL) - except asyncio.QueueFull: - pass await asyncio.sleep(wait) finally: self._status = "stopped" self._running = False - # Signal all waiting clients that the stream is over async with self._lock: for handle in self._clients.values(): _drain_queue(handle.queue) @@ -264,50 +252,100 @@ class BroadcastGroup: async def _pump_once(self) -> None: """ - Open one upstream connection and pump until EOF, stall, or error. - Resets consecutive-failure counter on first successful chunk. + Spawn ffmpeg in copy (remux) mode and pump its stdout to all client + queues until ffmpeg exits or self._running goes False. + + ffmpeg handles HTTP reconnects internally (-reconnect_streamed), which + covers provider session-token refreshes without any visible interruption + to clients. PTS/DTS values are normalised by ffmpeg across reconnects, + eliminating A/V desync entirely. """ - timeout = aiohttp.ClientTimeout( - connect=settings.stream_connect_timeout, - sock_read=None, # we control read timeouts via wait_for ourselves + cmd = [ + _FFMPEG, + "-hide_banner", "-loglevel", "warning", + # Follow HTTP redirects and reconnect when the connection drops. + # -reconnect_streamed is the key flag: it retries even on non- + # seekable (live) streams, re-issuing the original GET so that + # providers using short-lived 302 tokens get a fresh token. + "-reconnect", "1", + "-reconnect_streamed", "1", + "-reconnect_delay_max", "4", + # Socket-level read timeout in microseconds. If the provider + # stops sending bytes for STALL_TIMEOUT seconds ffmpeg closes + # the connection and tries to reconnect (or exits if it gives up). + "-timeout", str(int(STALL_TIMEOUT * 1_000_000)), + "-i", self.stream_url, + "-c", "copy", # remux only — zero transcoding, bit-for-bit quality + "-f", "mpegts", + "pipe:1", + ] + + proc = await asyncio.create_subprocess_exec( + *cmd, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, ) - async with aiohttp.ClientSession(timeout=timeout) as session: - async with session.get(self.stream_url, ssl=False) as resp: - if resp.status not in (200, 206): - raise RuntimeError(f"Provider returned HTTP {resp.status}") - self._status = "ok" - self._last_data_at = time.time() - logger.info( - f"[ch{self.channel_id}] Stream connected " - f"(reconnects so far: {self._reconnect_count})" - ) + self._status = "ok" + self._last_data_at = time.time() + logger.info( + f"[ch{self.channel_id}] ffmpeg started pid={proc.pid} " + f"(reconnects so far: {self._reconnect_count})" + ) - while self._running: - # --- read one chunk with stall timeout --- - try: - chunk = await asyncio.wait_for( - resp.content.readany(), timeout=STALL_TIMEOUT - ) - except asyncio.TimeoutError: - raise RuntimeError( - f"Stream stalled: no data for {STALL_TIMEOUT:.0f}s" - ) + # Drain stderr in background so the pipe never fills and blocks ffmpeg. + stderr_task = asyncio.create_task(self._log_ffmpeg_stderr(proc)) - if not chunk: - raise RuntimeError("Provider closed the stream (EOF)") + 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" + ) - # --- 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 + if not chunk: + # ffmpeg exited cleanly or with error + rc = proc.returncode + raise RuntimeError(f"ffmpeg exited (rc={rc})") - # --- distribute to clients --- - await self._dispatch(chunk, now) + # 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: + 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 and forward them to the logger.""" + try: + async for line in proc.stderr: + text = line.decode(errors="replace").rstrip() + if text: + logger.warning(f"[ch{self.channel_id}] ffmpeg: {text}") + 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.""" @@ -319,8 +357,6 @@ class BroadcastGroup: handle._overflow_count = 0 handle._overflow_since = None except asyncio.QueueFull: - # Drop-newest: discard this chunk for the slow client. - # Preserves the already-buffered sequence → no A/V desync from gaps. handle._overflow_count += 1 if handle._overflow_since is None: handle._overflow_since = now