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
1 change: 1 addition & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ ts6-stream-bot/
│ │ ├── audio.py <- PulseAudio sink helpers (introspection)
│ │ ├── audio_capture.py <- PulseAudio -> Opus -> TS3 voice frames
│ │ ├── video_capture.py <- x11grab + Pulse via aiortc MediaPlayers
│ │ ├── video_broadcaster.py <- Single libvpx encoder, per-viewer av.Packet fan-out
│ │ ├── stream_signaling.py <- TS6 stream signaling (setupstream etc.)
│ │ └── stream_publisher.py <- Per-viewer aiortc RTCPeerConnection
│ │
Expand Down
24 changes: 23 additions & 1 deletion src/ts6_stream_bot/pipeline/controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,10 @@
from ts6_stream_bot.pipeline.browser import BrowserManager
from ts6_stream_bot.pipeline.stream_publisher import StreamPublisher
from ts6_stream_bot.pipeline.stream_signaling import StreamSignaling
from ts6_stream_bot.pipeline.video_broadcaster import (
VideoBroadcaster,
VideoBroadcasterConfig,
)
from ts6_stream_bot.pipeline.video_capture import VideoCapture, VideoCaptureConfig
from ts6_stream_bot.sources import StreamSource, resolve_source
from ts6_stream_bot.ts3lib.client import Ts3Client, Ts3ClientOptions
Expand Down Expand Up @@ -254,7 +258,25 @@ async def _allocate_stream(self) -> None:
pulse_source=f"{settings.PULSE_SINK}.monitor",
)
capture = VideoCapture(capture_config)
publisher = StreamPublisher(client=self._ts3_client, signaling=signaling, capture=capture)
# Single libvpx instance shared across viewers. STREAM_BITRATE is
# in kbps (mirrors the TS6 setupstream parameter); the codec
# context wants bps. The factory closure defers reading
# capture.video_track until publisher.start() has booted capture.
broadcaster = VideoBroadcaster(
source_track_factory=lambda: capture.video_track,
config=VideoBroadcasterConfig(
bitrate=settings.STREAM_BITRATE * 1000,
width=settings.SCREEN_WIDTH,
height=settings.SCREEN_HEIGHT,
framerate=settings.SCREEN_FPS,
),
)
publisher = StreamPublisher(
client=self._ts3_client,
signaling=signaling,
capture=capture,
video_broadcaster=broadcaster,
)

log.info(
"controller.stream_setup",
Expand Down
34 changes: 24 additions & 10 deletions src/ts6_stream_bot/pipeline/stream_publisher.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
SignalingType,
StreamSignaling,
)
from ts6_stream_bot.pipeline.video_broadcaster import VideoBroadcaster
from ts6_stream_bot.pipeline.video_capture import VideoCapture
from ts6_stream_bot.ts3lib.client import Ts3Client

Expand Down Expand Up @@ -78,10 +79,12 @@ def __init__(
client: Ts3Client,
signaling: StreamSignaling,
capture: VideoCapture,
video_broadcaster: VideoBroadcaster,
) -> None:
self._client = client
self._signaling = signaling
self._capture = capture
self._video_broadcaster = video_broadcaster

self._stream_id: str | None = None
self._stream_started_event = asyncio.Event()
Expand All @@ -90,9 +93,11 @@ def __init__(
self._lock = asyncio.Lock()
# Track in-flight tasks so the GC doesn't clean them up mid-flight.
self._tasks: set[asyncio.Task[None]] = set()
# MediaRelay fans out one source track to many per-viewer subscribers
# so two PeerConnections don't both call recv() on the same underlying
# track at the same time (would race on the parec / x11grab subprocess).
# MediaRelay still handles audio fan-out: one parec subprocess
# delivers raw PCM and we want each viewer's RTCRtpSender to read
# its own subscriber queue rather than race on the source. Video
# bypasses the relay entirely - it goes through the broadcaster
# which encodes once and ships pre-encoded packets per viewer.
self._relay = MediaRelay()

# Wire up signaling callbacks. We chain so existing handlers stay alive.
Expand Down Expand Up @@ -144,8 +149,10 @@ async def start(
return self._stream_id

# Make sure the underlying ffmpeg pipelines are running before we
# accept any join requests.
# accept any join requests. Capture has to come up first so the
# broadcaster can resolve the source track factory.
await self._capture.start()
await self._video_broadcaster.start()

self._stream_started_event.clear()
self._signaling.send_setup_stream(
Expand All @@ -160,6 +167,7 @@ async def start(
try:
await asyncio.wait_for(self._stream_started_event.wait(), timeout=timeout)
except TimeoutError:
await self._video_broadcaster.stop()
await self._capture.stop()
raise

Expand All @@ -170,6 +178,7 @@ async def start(
async def stop(self) -> None:
"""Kick all viewers, stopstream, and tear down capture."""
if self._stream_id is None:
await self._video_broadcaster.stop()
await self._capture.stop()
return

Expand All @@ -193,6 +202,7 @@ async def stop(self) -> None:
# we tear down ffmpeg out from under the encoder.
await asyncio.sleep(0.5)

await self._video_broadcaster.stop()
await self._capture.stop()
self._stream_id = None
log.info("stream_publisher.stopped", stream_id=stream_id)
Expand Down Expand Up @@ -314,12 +324,16 @@ async def _on_ice_state_change() -> None:
ice=pc.iceConnectionState,
)

# Each viewer needs its OWN track that pulls from the shared source
# via MediaRelay. Sharing the source track directly across PCs makes
# both senders call recv() concurrently and crashes parec's
# readexactly() with a "another coroutine is already waiting" error.
if self._capture.video_track is not None:
pc.addTrack(self._relay.subscribe(self._capture.video_track))
# Video: subscribe to the broadcaster's pre-encoded fan-out. The
# returned track yields ``av.Packet`` objects, which RTCRtpSender
# routes through ``encoder.pack()`` (RTP packetization only, no
# libvpx call) - so each viewer's encoder spend is microseconds
# rather than the ~750 MB per-viewer libvpx setup we used to pay
# with MediaRelay-on-raw-frames.
pc.addTrack(self._video_broadcaster.subscribe())
# Audio still goes through MediaRelay: parec produces raw PCM and
# each viewer's RTCRtpSender encodes Opus on its own. Opus is
# cheap enough that the broadcaster pattern doesn't pay back here.
if self._capture.audio_track is not None:
pc.addTrack(self._relay.subscribe(self._capture.audio_track))

Expand Down
Loading
Loading