From a48b988703c2f0e11d02549a52b45e72389fa317 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 1 May 2026 23:15:55 +0000 Subject: [PATCH] Bound the source-side frame backlog: encode in executor + frame skip Live deploy on the upgraded host (6 vCPU, 8 GB) leaked ~850 MB of RSS in 16 seconds the moment a viewer tried to connect: 23:10:51 rss_mb=1163 23:11:07 rss_mb=2012 (+849 MB, +53 MB/s) The leak isn't ours - it's the unbounded asyncio.Queue inside aiortc's MediaPlayer track (contrib/media.py:229, ``self._queue: asyncio.Queue = asyncio.Queue()`` with no maxsize). ffmpeg's worker thread pumps decoded VideoFrames into that queue; our broadcaster pump consumes one per loop iteration. At 1080p the raw frames are ~8 MB each (bgr0 from x11grab); even slightly slower-than-realtime consumption piles up megabytes per second. Two fixes that together cap the backlog at one frame: 1. ``codec.encode()`` now runs in ``loop.run_in_executor`` rather than blocking the event loop. While libvpx works in a worker thread, the loop can keep draining the source queue. This is what aiortc itself does for per-sender encoders in ``RTCRtpSender._next_encoded_frame``. 2. ``_drain_to_latest`` peeks at the MediaPlayer's private ``_queue`` (stable in aiortc 1.x) and pulls every already-buffered frame, returning only the newest. Older frames in the backlog are counted into a new ``stale_frames_dropped`` metric and silently discarded - far better than the alternative of a slowly-growing queue that ends in OOM. The counter shows up in the ``video_broadcaster.stopped`` log so an operator can tell at a glance whether their host is keeping up with the configured resolution / framerate. Tests: * New ``test_drain_to_latest_skips_backlog_keeps_newest`` seeds the source's queue with two stale frames + one current and asserts only the latest is returned with the correct drop count. * New ``test_drain_to_latest_no_backlog_returns_current`` confirms the empty-queue fast path is a no-op. --- .../pipeline/video_broadcaster.py | 59 ++++++++++++++++++- tests/test_video_broadcaster.py | 53 +++++++++++++++++ 2 files changed, 110 insertions(+), 2 deletions(-) diff --git a/src/ts6_stream_bot/pipeline/video_broadcaster.py b/src/ts6_stream_bot/pipeline/video_broadcaster.py index 6904e9d..a56159f 100644 --- a/src/ts6_stream_bot/pipeline/video_broadcaster.py +++ b/src/ts6_stream_bot/pipeline/video_broadcaster.py @@ -47,6 +47,7 @@ import av import structlog from aiortc.mediastreams import MediaStreamError, MediaStreamTrack +from av.frame import Frame from av.packet import Packet from av.video.codeccontext import VideoCodecContext from av.video.frame import PictureType, VideoFrame @@ -138,6 +139,11 @@ def __init__( self._force_keyframe = False self._subscribers: list[_Subscriber] = [] self._encoded_frames = 0 + # How many frames we discarded because the source queue had + # piled up while the encoder was busy. Logged on stop() so an + # operator can tell whether their host is keeping up with the + # configured resolution / framerate. + self._stale_frames_dropped = 0 @property def is_alive(self) -> bool: @@ -184,7 +190,11 @@ async def stop(self) -> None: for _ in self._codec.encode(None): pass self._codec = None - log.info("video_broadcaster.stopped", encoded_frames=self._encoded_frames) + log.info( + "video_broadcaster.stopped", + encoded_frames=self._encoded_frames, + stale_frames_dropped=self._stale_frames_dropped, + ) # --- subscription ----------------------------------------------------- @@ -265,6 +275,7 @@ async def _pump_loop(self) -> None: async def _pump(self) -> None: assert self._source is not None + loop = asyncio.get_running_loop() while not self._stopped: try: frame = await self._source.recv() @@ -272,6 +283,21 @@ async def _pump(self) -> None: log.info("video_broadcaster.source_ended") return + # Skip backlog: aiortc's MediaPlayer keeps an unbounded + # ``asyncio.Queue`` between the ffmpeg worker thread and our + # consumer. If the encoder is even slightly slower than the + # source frame rate, raw frames pile up there - at 1080p + # that's ~8 MB each (bgr0 from x11grab) and the live deploy + # leaked ~850 MB in 16 s before OOM. By draining whatever + # extra frames are already queued and keeping only the most + # recent one, we cap the backlog to a single frame and drop + # stale ones rather than the whole pipeline drowning. The + # subscriber-side encoder backlog is independently bounded + # by ``queue_size`` per subscriber. + frame, dropped = self._drain_to_latest(frame) + if dropped: + self._stale_frames_dropped += dropped + if not isinstance(frame, VideoFrame): # Source drift - shouldn't happen with x11grab, log and skip. log.warning("video_broadcaster.unexpected_frame_type", got=type(frame).__name__) @@ -287,8 +313,14 @@ async def _pump(self) -> None: frame.pict_type = PictureType.I self._force_keyframe = False + # libvpx's encode is the expensive call. Running it in the + # default thread pool lets the event loop keep draining the + # source queue (via the next ``self._source.recv()``) while + # the encode runs in a worker thread - which is what aiortc + # itself does for per-sender encoders in + # rtcrtpsender._next_encoded_frame. try: - packets = list(self._codec.encode(frame)) + packets = await loop.run_in_executor(None, self._encode_frame, frame) except av.error.FFmpegError as exc: log.warning("video_broadcaster.encode_failed", error=str(exc)) continue @@ -297,6 +329,29 @@ async def _pump(self) -> None: self._encoded_frames += 1 self._fanout(packet) + def _drain_to_latest(self, current: Frame | Packet) -> tuple[Frame | Packet, int]: + """Reach into the source's internal asyncio queue and pull every + already-buffered frame. Return the newest one (the rest are + discarded). The MediaPlayer track aiortc gives us has a private + ``_queue`` attribute - we rely on it being a plain + ``asyncio.Queue`` because it has been since aiortc 1.0.x and + nothing else exposes a non-blocking peek.""" + latest = current + dropped = 0 + queue = getattr(self._source, "_queue", None) + if queue is None: + return latest, dropped + while True: + try: + latest = queue.get_nowait() + except asyncio.QueueEmpty: + return latest, dropped + dropped += 1 + + def _encode_frame(self, frame: VideoFrame) -> list[Packet]: + assert self._codec is not None + return list(self._codec.encode(frame)) + def _fanout(self, packet: Packet) -> None: """Push a packet into every subscriber queue. On QueueFull, drop the oldest packet for that subscriber rather than block - one diff --git a/tests/test_video_broadcaster.py b/tests/test_video_broadcaster.py index df21d20..be117bd 100644 --- a/tests/test_video_broadcaster.py +++ b/tests/test_video_broadcaster.py @@ -235,6 +235,59 @@ async def test_is_alive_true_after_start() -> None: await bc.stop() +async def test_drain_to_latest_skips_backlog_keeps_newest() -> None: + """The MediaPlayer queue grows unbounded if we don't keep up. We + pull whatever's already buffered and keep only the newest frame + so the source queue never accumulates beyond a single frame. + Live regression: 850 MB RAM growth in 16 s under 1080p load + came from this exact backlog.""" + bc = VideoBroadcaster(lambda: _FrameSource(0), _make_config()) + + # Fake source with a queue carrying older frames already buffered. + class _SourceWithQueue: + kind = "video" + + def __init__(self) -> None: + self._queue: asyncio.Queue[Any] = asyncio.Queue() + + src = _SourceWithQueue() + bc._source = src # type: ignore[assignment] + + older = MagicMock(name="older-frame") + middle = MagicMock(name="middle-frame") + latest = MagicMock(name="latest-frame") + src._queue.put_nowait(middle) + src._queue.put_nowait(latest) + + # ``current`` simulates the frame already taken via recv() before + # we noticed the queue was full of stale ones. + chosen, dropped = bc._drain_to_latest(older) + + assert chosen is latest + assert dropped == 2 # middle + latest were both pulled; older was discarded + + +async def test_drain_to_latest_no_backlog_returns_current() -> None: + """When the source queue is empty, the freshly received frame + passes through unchanged with no extra drops.""" + bc = VideoBroadcaster(lambda: _FrameSource(0), _make_config()) + + class _SourceWithQueue: + kind = "video" + + def __init__(self) -> None: + self._queue: asyncio.Queue[Any] = asyncio.Queue() + + src = _SourceWithQueue() + bc._source = src # type: ignore[assignment] + + only = MagicMock(name="only-frame") + chosen, dropped = bc._drain_to_latest(only) + + assert chosen is only + assert dropped == 0 + + async def test_is_alive_false_after_source_ends() -> None: """When the source raises MediaStreamError, the publisher needs to know so it can refuse new joins instead of attaching them to a