From 9844075bb27fee142a8e123a9e03bdf1afd1e008 Mon Sep 17 00:00:00 2001 From: stef <19478336+stefpi@users.noreply.github.com> Date: Fri, 22 May 2026 19:20:22 -0700 Subject: [PATCH] remote teleop multi-stream on 1 video track (#38013) * athenad and webrtcd updates * remove feature stream services from webrtcd split * stream encoder thread * reduce diff * wire webrtc to livestream camera encoder * request livestream camera switch service * remove camera list in favour of init camera field * remove cors * clean * remove unused * remove extra try except * add back exception trace * add stream road camera info to stream cameras * fix * clean diff * clean diff * add testJoystick only on body * fix camera list * remove reference to future service * encode all cameras and swap in video track in webrtc * clean * explicitly gate bridge send * clean leftover * make local bodyteleop work still --- system/athena/athenad.py | 2 +- system/loggerd/loggerd.h | 4 ++-- system/webrtc/device/video.py | 8 +++++++- system/webrtc/webrtcd.py | 24 ++++++++++++++---------- tools/bodyteleop/web.py | 2 +- 5 files changed, 25 insertions(+), 15 deletions(-) diff --git a/system/athena/athenad.py b/system/athena/athenad.py index 36f1d848c..6bccaae63 100755 --- a/system/athena/athenad.py +++ b/system/athena/athenad.py @@ -580,7 +580,7 @@ def startStream(sdp: str) -> dict: if CP.notCar: bridge_services_in.append("testJoystick") - body = StreamRequestBody(sdp, ["driver"], bridge_services_in, ["carState"]) + body = StreamRequestBody(sdp, "wideRoad", bridge_services_in, ["carState"]) try: resp = requests.post(f"http://localhost:{WEBRTCD_PORT}/stream", json=asdict(body), timeout=10) diff --git a/system/loggerd/loggerd.h b/system/loggerd/loggerd.h index 6aa0c8be4..01bce2c9e 100644 --- a/system/loggerd/loggerd.h +++ b/system/loggerd/loggerd.h @@ -47,8 +47,8 @@ struct EncoderSettings { } static EncoderSettings StreamEncoderSettings() { - int _stream_bitrate = getenv("STREAM_BITRATE") ? atoi(getenv("STREAM_BITRATE")) : 1'000'000; - return EncoderSettings{.encode_type = cereal::EncodeIndex::Type::QCAMERA_H264, .bitrate = _stream_bitrate , .gop_size = 15}; + int _stream_bitrate = getenv("STREAM_BITRATE") ? atoi(getenv("STREAM_BITRATE")) : 4'000'000; + return EncoderSettings{.encode_type = cereal::EncodeIndex::Type::QCAMERA_H264, .bitrate = _stream_bitrate , .gop_size = 5}; } }; diff --git a/system/webrtc/device/video.py b/system/webrtc/device/video.py index 50feab4f4..b46d02220 100644 --- a/system/webrtc/device/video.py +++ b/system/webrtc/device/video.py @@ -19,10 +19,16 @@ class LiveStreamVideoStreamTrack(TiciVideoStreamTrack): dt = DT_DMON if camera_type == "driver" else DT_MDL super().__init__(camera_type, dt) - self._sock = messaging.sub_sock(self.camera_to_sock_mapping[camera_type], conflate=True) + self._sock = self._make_sock(camera_type) self._pts = 0 self._t0_ns = time.monotonic_ns() + def _make_sock(self, camera_type: str) -> messaging.SubSocket: + return messaging.sub_sock(self.camera_to_sock_mapping[camera_type], conflate=True) + + def switch_camera(self, camera_type: str) -> None: + self._sock = self._make_sock(camera_type) + async def recv(self): while True: msg = messaging.recv_one_or_none(self._sock) diff --git a/system/webrtc/webrtcd.py b/system/webrtc/webrtcd.py index 7a51d4503..1239a1759 100755 --- a/system/webrtc/webrtcd.py +++ b/system/webrtc/webrtcd.py @@ -124,18 +124,15 @@ class DynamicPubMaster(messaging.PubMaster): class StreamSession: shared_pub_master = DynamicPubMaster([]) - def __init__(self, sdp: str, cameras: list[str], incoming_services: list[str], outgoing_services: list[str], debug_mode: bool = False): + def __init__(self, sdp: str, init_camera: str, incoming_services: list[str], outgoing_services: list[str], debug_mode: bool = False): from aiortc.mediastreams import VideoStreamTrack from openpilot.system.webrtc.device.video import LiveStreamVideoStreamTrack from teleoprtc import WebRTCAnswerBuilder - from teleoprtc.info import parse_info_from_offer - config = parse_info_from_offer(sdp) builder = WebRTCAnswerBuilder(sdp) - assert len(cameras) == config.n_expected_camera_tracks, "Incoming stream has misconfigured number of video tracks" - for cam in cameras: - builder.add_video_stream(cam, LiveStreamVideoStreamTrack(cam) if not debug_mode else VideoStreamTrack()) + self.video_track = LiveStreamVideoStreamTrack(init_camera) if not debug_mode else VideoStreamTrack() + builder.add_video_stream(init_camera, self.video_track) self.stream = builder.stream() self.identifier = str(uuid.uuid4()) @@ -155,8 +152,8 @@ class StreamSession: self._cleanup_done = False self.logger = logging.getLogger("webrtcd") self.logger.info( - "New stream session (%s), cameras %s, incoming services %s, outgoing services %s", - self.identifier, cameras, incoming_services, outgoing_services, + "New stream session (%s), init camera %s, incoming services %s, outgoing services %s", + self.identifier, init_camera, incoming_services, outgoing_services, ) def start(self): @@ -176,6 +173,13 @@ class StreamSession: def message_handler(self, message: bytes): assert self.incoming_bridge is not None try: + msg_json = json.loads(message) + if msg_json.get("type") == "livestreamCameraSwitch" and hasattr(self.video_track, "switch_camera"): + self.video_track.switch_camera(msg_json["data"]["camera"]) + return + + if msg_json.get("type") not in self.incoming_bridge_services: + return self.incoming_bridge.send(message) except Exception: self.logger.exception("Cereal incoming proxy failure") @@ -214,7 +218,7 @@ class StreamSession: @dataclass class StreamRequestBody: sdp: str - cameras: list[str] + initCamera: str bridge_services_in: list[str] = field(default_factory=list) bridge_services_out: list[str] = field(default_factory=list) @@ -236,7 +240,7 @@ async def get_stream(request: 'web.Request'): await s.stop() del stream_dict[sid] - session = StreamSession(body.sdp, body.cameras, body.bridge_services_in, body.bridge_services_out, debug_mode) + session = StreamSession(body.sdp, body.initCamera, body.bridge_services_in, body.bridge_services_out, debug_mode) try: answer = await session.get_answer() except ValueError as e: diff --git a/tools/bodyteleop/web.py b/tools/bodyteleop/web.py index 29d23184b..d357561b2 100644 --- a/tools/bodyteleop/web.py +++ b/tools/bodyteleop/web.py @@ -56,7 +56,7 @@ async def ping(request: 'web.Request'): async def offer(request: 'web.Request'): params = await request.json() - body = StreamRequestBody(params["sdp"], ["driver"], ["testJoystick"], ["carState"]) + body = StreamRequestBody(params["sdp"], "driver", ["testJoystick"], ["carState"]) body_json = json.dumps(dataclasses.asdict(body)) logger.info("Sending offer to webrtcd...")