From ac0669165668d87ef81b1b06b349519e215d2be6 Mon Sep 17 00:00:00 2001 From: joaquin Date: Sun, 17 May 2026 23:50:10 +0200 Subject: [PATCH] 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. --- backend/core/restream.py | 45 ++++++++++++++++++++++++++++++++ frontend/src/pages/Dashboard.tsx | 36 +++++++++++++++++++++---- 2 files changed, 76 insertions(+), 5 deletions(-) diff --git a/backend/core/restream.py b/backend/core/restream.py index d129b86..e045d12 100644 --- a/backend/core/restream.py +++ b/backend/core/restream.py @@ -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() diff --git a/frontend/src/pages/Dashboard.tsx b/frontend/src/pages/Dashboard.tsx index a0f2f43..b29076c 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 } from 'lucide-react' +import { Activity, Server, Users, Wifi, Clock, X, Film, Tv, Clapperboard, RefreshCw, AlertTriangle, Cpu } from 'lucide-react' import api from '../lib/api' interface StreamClient { @@ -37,6 +37,9 @@ interface StreamSession { reconnect_count?: number last_data_ago?: number stream_info?: StreamInfo + // ffmpeg process info (live streams only) + ffmpeg_pid?: number | null + ffmpeg_cpu_pct?: number } interface SlotInfo { @@ -195,9 +198,10 @@ export default function Dashboard() { } } - const totalClients = streams.reduce((a, s) => a + s.client_count, 0) - const busySlots = slots.filter(s => s.status === 'busy').length - const totalBpsDown = streams.reduce((a, s) => a + (s.bps_down || 0), 0) + const totalClients = streams.reduce((a, s) => a + s.client_count, 0) + const busySlots = slots.filter(s => s.status === 'busy').length + 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') return ( @@ -225,11 +229,16 @@ export default function Dashboard() { )} {/* Stats cards */} -
+
} label="Streams activos" value={streams.length} /> } label="Usuarios conectados" value={totalClients} /> } label="Slots ocupados" value={`${busySlots}/${slots.length}`} /> } label="↓ Total proveedor" value={formatBps(totalBpsDown)} /> + 80 ? 'text-red-400' : totalCpu > 40 ? 'text-yellow-400' : 'text-cyan-400'} />} + label="ffmpeg CPU total" + value={`${totalCpu.toFixed(1)}%`} + />
{/* Provider slots */} @@ -309,6 +318,23 @@ export default function Dashboard() { reconnectCount={s.reconnect_count} lastDataAgo={s.last_data_ago} /> + {s.ffmpeg_pid && ( +
+ + + ffmpeg + + {s.ffmpeg_cpu_pct !== undefined && ( + 80 ? 'text-red-400' : + s.ffmpeg_cpu_pct > 40 ? 'text-yellow-400' : + 'text-gray-400' + }`}> + {s.ffmpeg_cpu_pct.toFixed(1)}% + + )} +
+ )}