Files
KiraTV/backend/core/restream.py
T
KiraStream ab038030eb fix: clasificación correcta de fin de sesión vs fallo real de proveedor
- 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>
2026-05-19 15:57:24 +00:00

747 lines
31 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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