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 */} -