feat: show ffmpeg PID and CPU% per stream in dashboard
Backend samples /proc/{pid}/stat every 2s to calculate CPU usage of
each ffmpeg process and exposes ffmpeg_pid + ffmpeg_cpu_pct in stats().
Dashboard shows a cyan ffmpeg badge with live CPU% per stream row,
plus a new StatCard with total ffmpeg CPU across all active streams.
This commit is contained in:
@@ -29,6 +29,7 @@ is re-encoded — quality is bit-for-bit identical to the source.
|
||||
"""
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import time
|
||||
import uuid
|
||||
@@ -103,6 +104,11 @@ class BroadcastGroup:
|
||||
self._status: str = "connecting"
|
||||
self._last_data_at: float = time.time()
|
||||
|
||||
# ffmpeg process tracking
|
||||
self._ffmpeg_pid: int | None = None
|
||||
self._ffmpeg_cpu: float = 0.0
|
||||
self._cpu_task: asyncio.Task | None = None
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Public API
|
||||
# ------------------------------------------------------------------
|
||||
@@ -124,6 +130,9 @@ class BroadcastGroup:
|
||||
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}"
|
||||
)
|
||||
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:
|
||||
@@ -159,6 +168,8 @@ 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()
|
||||
if self._task and not self._task.done():
|
||||
self._task.cancel()
|
||||
try:
|
||||
@@ -197,12 +208,43 @@ class BroadcastGroup:
|
||||
"reconnect_count": self._reconnect_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, 1),
|
||||
}
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Internal helpers
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def _monitor_cpu(self) -> None:
|
||||
"""Sample /proc/{pid}/stat every 2 s and update self._ffmpeg_cpu (%)."""
|
||||
clk_tck = float(os.sysconf(os.sysconf_names.get("SC_CLK_TCK", 2)) or 100)
|
||||
prev_ticks: int | None = None
|
||||
prev_wall: float | None = None
|
||||
|
||||
while self._running:
|
||||
await asyncio.sleep(2)
|
||||
pid = self._ffmpeg_pid
|
||||
if pid is None:
|
||||
self._ffmpeg_cpu = 0.0
|
||||
prev_ticks = prev_wall = None
|
||||
continue
|
||||
try:
|
||||
with open(f"/proc/{pid}/stat") as f:
|
||||
parts = f.read().split()
|
||||
ticks = int(parts[13]) + int(parts[14]) # utime + stime
|
||||
wall = time.time()
|
||||
if prev_ticks is not None and prev_wall is not None:
|
||||
dt = wall - prev_wall
|
||||
if dt > 0:
|
||||
self._ffmpeg_cpu = (ticks - prev_ticks) / clk_tck / dt * 100
|
||||
prev_ticks = ticks
|
||||
prev_wall = wall
|
||||
except (FileNotFoundError, IndexError, ValueError, OSError):
|
||||
self._ffmpeg_cpu = 0.0
|
||||
prev_ticks = prev_wall = None
|
||||
|
||||
async def _probe_info(self) -> None:
|
||||
from .probe import probe_stream
|
||||
try:
|
||||
@@ -288,6 +330,7 @@ class BroadcastGroup:
|
||||
|
||||
self._status = "ok"
|
||||
self._last_data_at = time.time()
|
||||
self._ffmpeg_pid = proc.pid
|
||||
logger.info(
|
||||
f"[ch{self.channel_id}] ffmpeg started pid={proc.pid} "
|
||||
f"(reconnects so far: {self._reconnect_count})"
|
||||
@@ -325,6 +368,8 @@ class BroadcastGroup:
|
||||
await self._dispatch(chunk, now)
|
||||
|
||||
finally:
|
||||
self._ffmpeg_pid = None
|
||||
self._ffmpeg_cpu = 0.0
|
||||
stderr_task.cancel()
|
||||
try:
|
||||
proc.kill()
|
||||
|
||||
Reference in New Issue
Block a user