From 1e66d34eb2a04257355d0087e9ec70ae26f2dbfc Mon Sep 17 00:00:00 2001 From: Jason Wen Date: Fri, 26 Jun 2026 21:13:26 -0400 Subject: [PATCH] Revert "remote body teleop from connect support (#38011)" This reverts commit aff9f9ffae3f80a9d162a6a0b462e621de86a100. --- system/athena/athenad.py | 42 +-------------------- system/webrtc/webrtcd.py | 81 ++++++++++------------------------------ 2 files changed, 21 insertions(+), 102 deletions(-) diff --git a/system/athena/athenad.py b/system/athena/athenad.py index 52d009c341..52fb3f8c2b 100755 --- a/system/athena/athenad.py +++ b/system/athena/athenad.py @@ -29,7 +29,7 @@ from websocket import (ABNF, WebSocket, WebSocketException, WebSocketTimeoutExce create_connection) import cereal.messaging as messaging -from cereal import car, log +from cereal import log from cereal.services import SERVICE_LIST from openpilot.common.api import Api, get_key_pair from openpilot.common.utils import CallbackReader, get_upload_stream @@ -45,7 +45,6 @@ from openpilot.system.hardware.hw import Paths ATHENA_HOST = os.getenv('ATHENA_HOST', 'wss://athena.comma.ai') HANDLER_THREADS = int(os.getenv('HANDLER_THREADS', "4")) LOCAL_PORT_WHITELIST = {22, } # SSH -WEBRTCD_PORT = 5001 LOG_ATTR_NAME = 'user.upload' LOG_ATTR_VALUE_MAX_UNIX_TIME = int.to_bytes(2147483647, 4, sys.byteorder) @@ -568,16 +567,6 @@ def getSshAuthorizedKeys() -> str: def getGithubUsername() -> str: return cast(str, Params().get("GithubUsername") or "") - -@dispatcher.add_method -def getNotCar() -> bool: - cp_bytes = Params().get("CarParamsPersistent") - if cp_bytes is not None: - with car.CarParams.from_bytes(cp_bytes) as CP: - return CP.notCar - return False - - @dispatcher.add_method def getSimInfo(): return HARDWARE.get_sim_info() @@ -599,35 +588,6 @@ def getNetworks(): return HARDWARE.get_networks() -@dispatcher.add_method -def startStream(sdp: str) -> dict: - from openpilot.system.webrtc.webrtcd import StreamRequestBody - bridge_services_in = [] - - # get live car params to avoid stale notCar edge case - cp_bytes = Params().get("CarParams") - if cp_bytes is not None: - with car.CarParams.from_bytes(cp_bytes) as CP: - if CP.notCar: - bridge_services_in.append("testJoystick") - - body = StreamRequestBody(sdp, ["driver"], bridge_services_in, ["carState"]) - try: - resp = requests.post(f"http://localhost:{WEBRTCD_PORT}/stream", - json=asdict(body), timeout=10) - if not resp.ok: - try: - error_body = resp.json() - raise Exception(error_body.get("message", f"webrtcd returned {resp.status_code}")) - except ValueError: - resp.raise_for_status() - return resp.json() - except requests.ConnectTimeout as e: - raise Exception("webrtc took too long to respond. is it on?") from e - except requests.ConnectionError as e: - raise Exception("webrtc is not running. turn on comma body ignition.") from e - - @dispatcher.add_method def takeSnapshot() -> str | dict[str, str] | None: from openpilot.system.camerad.snapshot import jpeg_write, snapshot diff --git a/system/webrtc/webrtcd.py b/system/webrtc/webrtcd.py index 7a51d45034..d2c90cafb5 100755 --- a/system/webrtc/webrtcd.py +++ b/system/webrtc/webrtcd.py @@ -2,7 +2,6 @@ import argparse import asyncio -import contextlib import json import uuid import logging @@ -84,16 +83,11 @@ class CerealProxyRunner: assert self.task is None self.task = asyncio.create_task(self.run()) - async def stop(self): - if self.task is None: + def stop(self): + if self.task is None or self.task.done(): return - task = self.task + self.task.cancel() self.task = None - if task.done(): - return - task.cancel() - with contextlib.suppress(asyncio.CancelledError): - await task async def run(self): from aiortc.exceptions import InvalidStateError @@ -151,29 +145,24 @@ class StreamSession: self.outgoing_bridge_runner = CerealProxyRunner(self.outgoing_bridge) self.run_task: asyncio.Task | None = None - self._cleanup_lock = asyncio.Lock() - 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, - ) + self.logger.info("New stream session (%s), cameras %s, incoming services %s, outgoing services %s", + self.identifier, cameras, incoming_services, outgoing_services) def start(self): self.run_task = asyncio.create_task(self.run()) - async def stop(self): - if self.run_task is not None and not self.run_task.done() and self.run_task is not asyncio.current_task(): - self.run_task.cancel() - with contextlib.suppress(asyncio.CancelledError): - await self.run_task + def stop(self): + if self.run_task.done(): + return + self.run_task.cancel() self.run_task = None - await self.post_run_cleanup() + asyncio.run(self.post_run_cleanup()) async def get_answer(self): return await self.stream.start() - def message_handler(self, message: bytes): + async def message_handler(self, message: bytes): assert self.incoming_bridge is not None try: self.incoming_bridge.send(message) @@ -194,21 +183,16 @@ class StreamSession: self.logger.info("Stream session (%s) connected", self.identifier) await self.stream.wait_for_disconnection() + await self.post_run_cleanup() self.logger.info("Stream session (%s) ended", self.identifier) except Exception: self.logger.exception("Stream session failure") - finally: - await self.post_run_cleanup() async def post_run_cleanup(self): - async with self._cleanup_lock: - if self._cleanup_done: - return - self._cleanup_done = True - if self.outgoing_bridge_runner is not None: - await self.outgoing_bridge_runner.stop() - await self.stream.stop() + await self.stream.stop() + if self.outgoing_bridge is not None: + self.outgoing_bridge_runner.stop() @dataclass @@ -224,33 +208,11 @@ async def get_stream(request: 'web.Request'): raw_body = await request.json() body = StreamRequestBody(**raw_body) - async with request.app['stream_lock']: - # Fully disconnect any other active stream before starting the replacement. - for sid, s in list(stream_dict.items()): - if s.run_task and not s.run_task.done(): - try: - ch = s.stream.get_messaging_channel() - ch.send(json.dumps({"type": "connectionReplaced", "data": "Another device has connected, closing this session."})) - except Exception: - pass - await s.stop() - del stream_dict[sid] + session = StreamSession(body.sdp, body.cameras, body.bridge_services_in, body.bridge_services_out, debug_mode) + answer = await session.get_answer() + session.start() - session = StreamSession(body.sdp, body.cameras, body.bridge_services_in, body.bridge_services_out, debug_mode) - try: - answer = await session.get_answer() - except ValueError as e: - await session.stop() - raise web.HTTPBadRequest( - text=json.dumps({"error": "invalid_sdp", "message": str(e)}), - content_type="application/json", - ) from e - except Exception: - await session.stop() - raise - session.start() - - stream_dict[session.identifier] = session + stream_dict[session.identifier] = session return web.json_response({"sdp": answer.sdp, "type": answer.type}) @@ -262,7 +224,6 @@ async def get_schema(request: 'web.Request'): schema_dict = {s: generate_field(log.Event.schema.fields[s]) for s in services} return web.json_response(schema_dict) - async def post_notify(request: 'web.Request'): try: payload = await request.json() @@ -278,10 +239,9 @@ async def post_notify(request: 'web.Request'): return web.Response(status=200, text="OK") - async def on_shutdown(app: 'web.Application'): for session in app['streams'].values(): - await session.stop() + session.stop() del app['streams'] @@ -294,7 +254,6 @@ def webrtcd_thread(host: str, port: int, debug: bool): app = web.Application() app['streams'] = dict() - app['stream_lock'] = asyncio.Lock() app['debug'] = debug app.on_shutdown.append(on_shutdown) app.router.add_post("/stream", get_stream)