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
{label}
-{value}
+{value}