Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 57 additions & 2 deletions src/ts6_stream_bot/pipeline/video_broadcaster.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 -----------------------------------------------------

Expand Down Expand Up @@ -265,13 +275,29 @@ 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()
except MediaStreamError:
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__)
Expand All @@ -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
Expand All @@ -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
Expand Down
53 changes: 53 additions & 0 deletions tests/test_video_broadcaster.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading