feat: A/V health monitor, buffer bar and stream health dashboard
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 <noreply@anthropic.com>
This commit is contained in:
+83
-12
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user