diff --git a/openpilot/system/athena/athenad.py b/openpilot/system/athena/athenad.py index 828ae6b84..1351f45c2 100755 --- a/openpilot/system/athena/athenad.py +++ b/openpilot/system/athena/athenad.py @@ -795,7 +795,7 @@ def startStream(sdp: str, enabled: bool) -> dict: bridge_services_in = [] # stale car params case taken care of by webrtcd being shut off on ignition - cp_bytes = Params().get("CarParamsPersistent") + cp_bytes = params.get("CarParamsPersistent") if cp_bytes is not None: with car.CarParams.from_bytes(cp_bytes) as CP: if CP.notCar: @@ -808,7 +808,11 @@ def startStream(sdp: str, enabled: bool) -> dict: # webrtcd clears IsLiveStreaming when the session ends params.put_bool("IsLiveStreaming", True) # wait for webrtcd end points to wake up - wait_for_webrtcd() + try: + wait_for_webrtcd() + except TimeoutError: + cloudlog.event("athena.startStream.webrtcd_offroad_start_timeout", error=True) + raise return post_stream_request(StreamRequestBody(sdp, ["wideRoad"], enabled, bridge_services_in, ["carState", "deviceState"])) diff --git a/openpilot/system/webrtc/webrtcd.py b/openpilot/system/webrtc/webrtcd.py index 3c0c3037f..56b4ac43c 100755 --- a/openpilot/system/webrtc/webrtcd.py +++ b/openpilot/system/webrtc/webrtcd.py @@ -21,10 +21,16 @@ from typing import Any from openpilot.system.webrtc.helpers import StreamRequestBody from openpilot.system.webrtc.schema import generate_field from openpilot.common.params import Params +from openpilot.common.swaglog import cloudlog from openpilot.cereal import messaging, log SESSION_TIMEOUT_SECONDS = 300 + +# ice candidate parser for logging +def _ice_candidates(sdp: str) -> list[str]: + return [line.removeprefix("a=") for line in sdp.splitlines() if line.startswith("a=candidate:")] + # socket trick: route lookup for 8.8.8.8 (nothing is sent or actually connected to) # return the source interfaces IP which is the default interface of the device def _default_route_ip() -> str | None: @@ -253,7 +259,7 @@ class StreamSession: self._cleanup_lock = asyncio.Lock() self._cleanup_done = False self.logger = logging.getLogger("webrtcd") - self.logger.info( + cloudlog.warning( "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, ) @@ -341,14 +347,18 @@ class StreamSession: if self.bitrate_controller is not None: self.bitrate_controller.start() - self.logger.info("Stream session (%s) connected", self.identifier) + with cloudlog.ctx(session_id=self.identifier): + cloudlog.warning("webrtcd.session.connected") if self.is_body: await self.run_body_session() else: await self.run_normal_session() - self.logger.info("Stream session (%s) ended", self.identifier) + with cloudlog.ctx(session_id=self.identifier): + cloudlog.warning("webrtcd.session.ended") except Exception: self.logger.exception("Stream session failure") + with cloudlog.ctx(session_id=self.identifier): + cloudlog.exception("webrtcd.session.exception") finally: await self.post_run_cleanup() @@ -422,15 +432,25 @@ async def handle_get_stream(state: ServerState, raw_body: bytes, content_type: s stream_dict[session.identifier] = session try: answer = await asyncio.wait_for(session.get_answer(), timeout=30) + cloudlog.event( + "webrtcd.session.ice_candidates", + session_id=session.identifier, + offer_candidates=_ice_candidates(body.sdp), + answer_candidates=_ice_candidates(answer.sdp), + ) except TimeoutError: await session.stop() stream_dict.pop(session.identifier, None) logging.getLogger("webrtcd").exception("Timed out creating stream answer") + with cloudlog.ctx(session_id=session.identifier): + cloudlog.warning("webrtcd.session.answer_timeout") raise except Exception: await session.stop() stream_dict.pop(session.identifier, None) logging.getLogger("webrtcd").exception("Failed to create stream answer") + with cloudlog.ctx(session_id=session.identifier): + cloudlog.exception("webrtcd.session.answer_exception") raise session.start()