From 02b6cef15cc4772401bf7226da25d6f36018f839 Mon Sep 17 00:00:00 2001 From: stef <19478336+stefpi@users.noreply.github.com> Date: Mon, 27 Jul 2026 20:47:43 -0700 Subject: [PATCH] webrtc: cleanup api (#38447) * add support for multi-track with wideRoad single track as default * remove stale debug mode * multi track support * make cameras list optional * remove multi camera support from athena, make it only local * update to teleop master --- openpilot/system/athena/athenad.py | 2 +- openpilot/system/webrtc/helpers.py | 2 +- openpilot/system/webrtc/webrtcd.py | 48 ++++++++++++++++-------------- teleoprtc_repo | 2 +- 4 files changed, 29 insertions(+), 25 deletions(-) diff --git a/openpilot/system/athena/athenad.py b/openpilot/system/athena/athenad.py index e7c318499a..49fe080bb2 100755 --- a/openpilot/system/athena/athenad.py +++ b/openpilot/system/athena/athenad.py @@ -589,7 +589,7 @@ def startStream(sdp: str, enabled: bool) -> dict: # wait for webrtcd end points to wake up wait_for_webrtcd() - return post_stream_request(StreamRequestBody(sdp, "wideRoad", enabled, bridge_services_in, ["carState", "deviceState"])) + return post_stream_request(StreamRequestBody(sdp, ["wideRoad"], enabled, bridge_services_in, ["carState", "deviceState"])) def get_logs_to_send_sorted() -> list[str]: diff --git a/openpilot/system/webrtc/helpers.py b/openpilot/system/webrtc/helpers.py index 776b64573e..de45e1c6c9 100644 --- a/openpilot/system/webrtc/helpers.py +++ b/openpilot/system/webrtc/helpers.py @@ -8,7 +8,7 @@ WEBRTCD_PORT = 5001 @dataclass class StreamRequestBody: sdp: str - init_camera: str + cameras: list[str] enabled: bool bridge_services_in: list[str] = field(default_factory=list) bridge_services_out: list[str] = field(default_factory=list) diff --git a/openpilot/system/webrtc/webrtcd.py b/openpilot/system/webrtc/webrtcd.py index 7979fbe8e3..8e022eda1c 100755 --- a/openpilot/system/webrtc/webrtcd.py +++ b/openpilot/system/webrtc/webrtcd.py @@ -221,7 +221,7 @@ class LivestreamBitrateController(AsyncTaskRunner): class StreamSession: shared_pub_master = DynamicPubMaster([]) - def __init__(self, body: StreamRequestBody, debug_mode: bool = False): + def __init__(self, body: StreamRequestBody): from openpilot.system.webrtc.device.video import LiveStreamVideoStreamTrack from teleoprtc.builder import WebRTCAnswerBuilder @@ -230,8 +230,11 @@ class StreamSession: builder = WebRTCAnswerBuilder(body.sdp, bind_address=_default_route_ip()) self.enabled = body.enabled - self.video_track = LiveStreamVideoStreamTrack(body.init_camera, self.enabled) - builder.add_video_stream(body.init_camera, self.video_track) + self.video_tracks = [] + for camera in body.cameras: + track = LiveStreamVideoStreamTrack(camera, self.enabled) + self.video_tracks.append(track) + builder.add_video_stream(camera, track) self.stream = builder.stream() self.is_body = "testJoystick" in body.bridge_services_in @@ -251,8 +254,8 @@ class StreamSession: self._cleanup_done = False self.logger = logging.getLogger("webrtcd") self.logger.info( - "New stream session (%s), init camera %s, video enabled %s, incoming services %s, outgoing services %s", - self.identifier, body.init_camera, body.enabled, body.bridge_services_in, body.bridge_services_out, + "New stream session (%s), video cameras %s, video enabled %s, incoming services %s, outgoing services %s", + self.identifier, [t.id for t in self.video_tracks], body.enabled, body.bridge_services_in, body.bridge_services_out, ) def start(self): @@ -277,14 +280,17 @@ class StreamSession: match msg_type: case "livestreamCameraSwitch": - self.video_track.switch_camera(payload["data"]["camera"]) + # only needed for 1 track stream + if len(self.video_tracks) == 1: + self.video_tracks[0].switch_camera(payload["data"]["camera"]) case "livestreamSettings": if self.bitrate_controller is not None: self.bitrate_controller.set_quality(payload["data"]["quality"]) case "livestreamVideoEnable": enabled = payload["data"]["enabled"] self.enabled = enabled - self.video_track.enable(enabled) + for track in self.video_tracks: + track.enable(enabled) if self.outgoing_bridge is not None: self.outgoing_bridge.enable(enabled) if self.bitrate_controller is not None: @@ -297,8 +303,8 @@ class StreamSession: }}) self.stream.get_messaging_channel().send(pong) case "enableTimingSei": - if hasattr(self.video_track, 'timing_sei_enabled'): - self.video_track.timing_sei_enabled = bool(payload["data"]["enabled"]) + for track in self.video_tracks: + track.timing_sei_enabled = bool(payload["data"]["enabled"]) case _: if msg_type not in self.incoming_bridge_services: return @@ -356,17 +362,16 @@ class StreamSession: await self.bitrate_controller.stop() if self.outgoing_bridge is not None: await self.outgoing_bridge.stop() - if self.video_track is not None: - self.video_track.stop() - self.video_track = None + for track in self.video_tracks: + track.stop() + self.video_tracks.clear() await self.stream.stop() class ServerState: - def __init__(self, debug: bool): + def __init__(self): self.streams: dict[str, StreamSession] = {} self.stream_lock = asyncio.Lock() - self.debug = debug self.teardown: asyncio.TimerHandle | None = None @@ -391,7 +396,7 @@ def _text_response(text: str, status: int = 200) -> tuple[int, bytes, str]: async def handle_get_stream(state: ServerState, raw_body: bytes) -> tuple[int, bytes, str]: - stream_dict, debug_mode = state.streams, state.debug + stream_dict = state.streams body = StreamRequestBody(**json.loads(raw_body)) async with state.stream_lock: @@ -410,7 +415,7 @@ async def handle_get_stream(state: ServerState, raw_body: bytes) -> tuple[int, b await s.stop() stream_dict.pop(sid, None) - session = StreamSession(body, debug_mode) + session = StreamSession(body) stream_dict[session.identifier] = session try: answer = await asyncio.wait_for(session.get_answer(), timeout=30) @@ -558,23 +563,23 @@ async def _shutdown(server: WebrtcdHTTPServer, state: ServerState, loop: asyncio loop.stop() -def prewarm_stream_session_imports(debug_mode: bool = False) -> None: +def prewarm_stream_session_imports() -> None: from openpilot.system.webrtc.device.video import LiveStreamVideoStreamTrack from teleoprtc.builder import WebRTCAnswerBuilder assert LiveStreamVideoStreamTrack assert WebRTCAnswerBuilder -def webrtcd_thread(host: str, port: int, debug: bool): +def webrtcd_thread(host: str, port: int): logging.basicConfig(level=logging.INFO, handlers=[logging.StreamHandler()]) prewarm_start = time.monotonic() - prewarm_stream_session_imports(debug) + prewarm_stream_session_imports() prewarm_end = time.monotonic() logging.getLogger("webrtcd").info(f"webrtc prewarm finished in {(prewarm_end - prewarm_start) * 1000} ms") loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) - state = ServerState(debug) + state = ServerState() server = WebrtcdHTTPServer((host, port), WebrtcdHandler) server.state = state @@ -608,10 +613,9 @@ def main(): parser = argparse.ArgumentParser(description="WebRTC daemon") parser.add_argument("--host", type=str, default="0.0.0.0", help="Host to listen on") parser.add_argument("--port", type=int, default=5001, help="Port to listen on") - parser.add_argument("--debug", action="store_true", help="Enable debug mode") args = parser.parse_args() - webrtcd_thread(args.host, args.port, args.debug) + webrtcd_thread(args.host, args.port) if __name__=="__main__": diff --git a/teleoprtc_repo b/teleoprtc_repo index 3ac144dc9a..3f0b03a1d1 160000 --- a/teleoprtc_repo +++ b/teleoprtc_repo @@ -1 +1 @@ -Subproject commit 3ac144dc9a57ace70172d338690b121c42c5add2 +Subproject commit 3f0b03a1d1c4887a0e76bc55888bc1a2eec29802