From 57c8f6fb648b6ac244fe30d296cf21425b6e4916 Mon Sep 17 00:00:00 2001 From: joaquin Date: Mon, 18 May 2026 01:14:37 +0200 Subject: [PATCH] feat: A/V health monitor, buffer bar and stream health dashboard MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Backend (restream.py): - Add _health_check_loop task: checks data flow every 5 s, updates av_health (ok/warning/error), auto-corrects stuck ffmpeg via SIGTERM - Health thresholds: warning >12 s, error >25 s — normal reconnects (1-4 s) no longer trigger false warnings - Count ffmpeg stderr A/V-warning events (av_desync_count) as metric - Expose buffer_pct (server-side queue fill of slowest client), corrections (health-triggered ffmpeg restarts), av_health in stats() - Revert wallclock timestamp flag: -use_wallclock_as_timestamps 1 was generating hundreds of spurious demuxer warnings per session, causing ffmpeg to stall and producing continuous video freezes - Simplify ffmpeg cmd to minimal safe set: -reconnect 1 -reconnect_streamed 1 -reconnect_delay_max 2 Dashboard (Dashboard.tsx): - Add 6th StatCard "Salud A/V" with dynamic color and correction count - Add "Salud A/V" column per stream: badge ok/aviso/error, corrections counter (violet), desync event count (grey) - Add BufferBar: thin progress bar showing server queue fill % (green <40 %, yellow 40-70 %, red >70 %) - Alert banner in header when any stream is in error av_health state Co-Authored-By: Claude Sonnet 4.6 --- backend/core/restream.py | 95 ++++++++++++++++++--- frontend/src/pages/Dashboard.tsx | 137 +++++++++++++++++++++++++++---- 2 files changed, 204 insertions(+), 28 deletions(-) diff --git a/backend/core/restream.py b/backend/core/restream.py index f387e01..044d394 100644 --- a/backend/core/restream.py +++ b/backend/core/restream.py @@ -50,6 +50,20 @@ _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 @@ -128,6 +142,12 @@ class BroadcastGroup: 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 # ------------------------------------------------------------------ @@ -152,6 +172,9 @@ class BroadcastGroup: 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: @@ -187,8 +210,10 @@ class BroadcastGroup: async def stop(self) -> None: self._running = False self._status = "stopped" - if self._cpu_task and not self._cpu_task.done(): - self._cpu_task.cancel() + 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: @@ -233,6 +258,11 @@ class BroadcastGroup: # 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(), } # ------------------------------------------------------------------ @@ -266,6 +296,52 @@ class BroadcastGroup: 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) + + 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 _probe_info(self) -> None: from .probe import probe_stream try: @@ -345,19 +421,12 @@ class BroadcastGroup: 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). + "-reconnect_delay_max", "2", "-timeout", str(int(STALL_TIMEOUT * 1_000_000)), "-i", stream_url, - "-c", "copy", # remux only — zero transcoding, bit-for-bit quality + "-c", "copy", "-f", "mpegts", "pipe:1", ] @@ -423,12 +492,14 @@ class BroadcastGroup: pass async def _log_ffmpeg_stderr(self, proc: asyncio.subprocess.Process) -> None: - """Read ffmpeg stderr lines and forward them to the logger.""" + """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: diff --git a/frontend/src/pages/Dashboard.tsx b/frontend/src/pages/Dashboard.tsx index 177838f..3ff5b3c 100644 --- a/frontend/src/pages/Dashboard.tsx +++ b/frontend/src/pages/Dashboard.tsx @@ -1,5 +1,5 @@ import { useEffect, useRef, useState } from 'react' -import { Activity, Server, Users, Wifi, Clock, X, Film, Tv, Clapperboard, RefreshCw, AlertTriangle, Cpu } from 'lucide-react' +import { Activity, Server, Users, Wifi, Clock, X, Film, Tv, Clapperboard, RefreshCw, AlertTriangle, Cpu, HeartPulse, ShieldCheck } from 'lucide-react' import api from '../lib/api' interface StreamClient { @@ -32,14 +32,19 @@ interface StreamSession { bps_down: number bps_up: number running: boolean - // health (live streams only) + // stream health status?: 'connecting' | 'ok' | 'reconnecting' | 'stopped' reconnect_count?: number last_data_ago?: number stream_info?: StreamInfo - // ffmpeg process info (live streams only) + // ffmpeg process info ffmpeg_pid?: number | null ffmpeg_cpu_pct?: number + // A/V health (from health monitor) + av_health?: 'ok' | 'warning' | 'error' + av_desync_count?: number // cumulative ffmpeg AV-warning events + corrections?: number // auto-corrections by health monitor + buffer_pct?: number // server-side queue fill % (0 = client fast, 100 = client slow) } interface SlotInfo { @@ -74,13 +79,6 @@ function formatDuration(startTs: number) { function formatResolution(info: StreamInfo) { if (!info.width || !info.height) return null - const res = info.height >= 1060 ? '4K' : info.height >= 1060 ? '4K' - : info.height >= 2000 ? '4K' - : info.height >= 1060 ? 'FHD' - : info.height >= 1040 ? 'FHD' - : info.height >= 700 ? 'HD' - : info.height >= 460 ? 'SD' - : 'LD' const label = info.width >= 3840 ? '4K' : info.width >= 1920 ? 'FHD' : @@ -118,10 +116,9 @@ function StreamInfoBadge({ info }: { info: StreamInfo }) { ) } -function HealthBadge({ status, reconnectCount, lastDataAgo }: { +function HealthBadge({ status, reconnectCount }: { status?: string reconnectCount?: number - lastDataAgo?: number }) { if (!status || status === 'ok') { if (reconnectCount && reconnectCount > 0) { @@ -142,6 +139,70 @@ function HealthBadge({ status, reconnectCount, lastDataAgo }: { return {status} } +function AVHealthBadge({ health, desyncs, corrections }: { + health?: 'ok' | 'warning' | 'error' + desyncs?: number + corrections?: number +}) { + const hasCorrections = corrections != null && corrections > 0 + const hasDesyncs = desyncs != null && desyncs > 0 + + const badge = health === 'error' + ? + A/V error + + : health === 'warning' + ? + A/V aviso + + : + A/V ok + + + return ( +
+ {badge} + {hasCorrections && ( + + ×{corrections} + + )} + {hasDesyncs && ( + + {desyncs} ev. + + )} +
+ ) +} + +function BufferBar({ pct }: { pct?: number }) { + const v = pct ?? 0 + const color = + v >= 70 ? 'bg-red-500' : + v >= 40 ? 'bg-yellow-400' : + 'bg-emerald-500' + const label = + v >= 70 ? 'text-red-400' : + v >= 40 ? 'text-yellow-400' : + 'text-emerald-500' + + return ( +
+
+ Buffer + {v.toFixed(0)}% +
+
+
+
+
+ ) +} + const TYPE_ICONS: Record = { live: , movie: , @@ -203,6 +264,20 @@ export default function Dashboard() { const totalBpsDown = streams.reduce((a, s) => a + (s.bps_down || 0), 0) const totalCpu = streams.reduce((a, s) => a + (s.ffmpeg_cpu_pct ?? 0), 0) const anyReconnecting = streams.some(s => s.status === 'reconnecting') + const streamsOk = streams.filter(s => !s.av_health || s.av_health === 'ok').length + const streamsWarn = streams.filter(s => s.av_health === 'warning').length + const streamsErr = streams.filter(s => s.av_health === 'error').length + const totalCorrections = streams.reduce((a, s) => a + (s.corrections ?? 0), 0) + + const healthColor = + streamsErr > 0 ? 'text-red-400' : + streamsWarn > 0 ? 'text-yellow-400' : + 'text-emerald-400' + const healthValue = + streams.length === 0 ? '—' : + streamsErr > 0 ? `${streamsErr} error${streamsErr > 1 ? 'es' : ''}` : + streamsWarn > 0 ? `${streamsWarn} aviso${streamsWarn > 1 ? 's' : ''}` : + `${streamsOk} ok` return (
@@ -215,6 +290,12 @@ export default function Dashboard() { Reconectando stream(s)... )} + {streamsErr > 0 && ( + + + {streamsErr} stream{streamsErr > 1 ? 's' : ''} con error A/V + + )}
{connected ? 'Conectado en tiempo real' : 'Desconectado'} @@ -229,7 +310,7 @@ export default function Dashboard() { )} {/* Stats cards */} -
+
} label="Streams activos" value={streams.length} /> } label="Usuarios conectados" value={totalClients} /> } label="Slots ocupados" value={`${busySlots}/${slots.length}`} /> @@ -239,6 +320,12 @@ export default function Dashboard() { label="ffmpeg CPU total" value={`${totalCpu.toFixed(2)}%`} /> + } + label={`Salud A/V${totalCorrections > 0 ? ` · ${totalCorrections} corr.` : ''}`} + value={healthValue} + valueClass={healthColor} + />
{/* Provider slots */} @@ -280,6 +367,7 @@ export default function Dashboard() { Proveedor Stream Estado + Salud A/V Usuarios ↓ / ↑ Tiempo @@ -312,11 +400,12 @@ export default function Dashboard() { : — } + + {/* Stream status + ffmpeg */} {s.ffmpeg_pid && (
@@ -336,6 +425,17 @@ export default function Dashboard() {
)} + + {/* A/V health + buffer */} + + + + +
{s.clients.map(c => ( @@ -373,13 +473,18 @@ export default function Dashboard() { ) } -function StatCard({ icon, label, value }: { icon: React.ReactNode; label: string; value: string | number }) { +function StatCard({ icon, label, value, valueClass }: { + icon: React.ReactNode + label: string + value: string | number + valueClass?: string +}) { return (
{icon}

{label}

-

{value}

+

{value}

)