Merge branch 'master' into spatial-feat

This commit is contained in:
James Vecellio-Grant
2026-08-23 15:25:34 -07:00
committed by GitHub
7 changed files with 51 additions and 29 deletions
+3 -1
View File
@@ -24,7 +24,9 @@ function agnos_init {
if $AGNOS_PY --verify $MANIFEST; then
sudo reboot
fi
$DIR/openpilot/common/hardware/comma/updater $AGNOS_PY $MANIFEST
while true; do
$DIR/openpilot/common/hardware/comma/updater $AGNOS_PY $MANIFEST
done
fi
}
+1 -1
View File
@@ -59,7 +59,7 @@ inline static std::unordered_map<std::string, ParamKeyAttributes> keys = {
{"IsDriverViewEnabled", {CLEAR_ON_MANAGER_START, BOOL}},
{"IsEngaged", {PERSISTENT, BOOL}},
{"IsLdwEnabled", {PERSISTENT | BACKUP, BOOL}},
{"IsLiveStreaming", {CLEAR_ON_MANAGER_START, BOOL}},
{"IsLiveStreaming", {CLEAR_ON_MANAGER_START | CLEAR_ON_IGNITION_ON, BOOL}},
{"IsMetric", {PERSISTENT | BACKUP, BOOL}},
{"IsOffroad", {CLEAR_ON_MANAGER_START, BOOL}},
{"IsRhdDetected", {PERSISTENT, BOOL}},
+9 -11
View File
@@ -149,11 +149,15 @@ class BigButton(Widget):
def set_touch_valid_callback(self, touch_callback: Callable[[], bool]) -> None:
super().set_touch_valid_callback(lambda: touch_callback() and self._grow_animation_until is None)
def _width_hint(self) -> int:
# A value moves the title to the top, where it shares space with the icon.
def _title_width_hint(self) -> int:
# A value moves the title to the top, where it shares space with the icon
icon_size = self._txt_icon.width if self._txt_icon and self.value else 0
return int(self._rect.width - self.LABEL_HORIZONTAL_PADDING * 2 - icon_size)
def _subtitle_width_hint(self) -> int:
# Bottom aligned, so it sits below the icon
return int(self._rect.width - self.LABEL_HORIZONTAL_PADDING * 2)
def _get_label_font_size(self):
if len(self.text) <= 18:
return 48
@@ -228,14 +232,14 @@ class BigButton(Widget):
label_color = LABEL_COLOR if self.enabled else rl.Color(255, 255, 255, int(255 * 0.35))
self._label.set_color(label_color)
label_rect = rl.Rectangle(label_x, btn_y + self.LABEL_VERTICAL_PADDING, self._width_hint(),
label_rect = rl.Rectangle(label_x, btn_y + self.LABEL_VERTICAL_PADDING, self._title_width_hint(),
self._rect.height - self.LABEL_VERTICAL_PADDING * 2)
self._label.render(label_rect)
if self.value:
label_y = btn_y + self.LABEL_VERTICAL_PADDING + self._label.get_content_height(self._width_hint())
label_y = label_rect.y + self._label.get_content_height(int(label_rect.width))
sub_label_height = btn_y + self._rect.height - self.LABEL_VERTICAL_PADDING - label_y
sub_label_rect = rl.Rectangle(label_x, label_y, self._width_hint(), sub_label_height)
sub_label_rect = rl.Rectangle(label_x, label_y, self._subtitle_width_hint(), sub_label_height)
self._sub_label.render(sub_label_rect)
# ICON -------------------------------------------------------------------
@@ -312,9 +316,6 @@ class BigMultiToggle(BigToggle):
self.set_value(self._options[0])
def _width_hint(self) -> int:
return int(self._rect.width - self.LABEL_HORIZONTAL_PADDING * 2 - self._txt_enabled_toggle.width)
def _handle_mouse_release(self, mouse_pos: MousePos):
super()._handle_mouse_release(mouse_pos)
cur_idx = self._options.index(self.value)
@@ -363,9 +364,6 @@ class GreyBigButton(BigButton):
def LABEL_VERTICAL_PADDING(self):
return BigButton.LABEL_VERTICAL_PADDING if self._label.text else 18
def _width_hint(self) -> int:
return int(self._rect.width - self.LABEL_HORIZONTAL_PADDING * 2)
def _get_label_font_size(self):
return 36
+6 -4
View File
@@ -828,20 +828,22 @@ 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:
bridge_services_in.append("testJoystick")
else:
raise Exception("failed to get CarParamsPersistent")
if params.get_bool("IsOffroad"):
# manager owns camerad/stream_encoderd/webrtcd; flip the param and let it bring them up.
# 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"]))
+2 -5
View File
@@ -106,15 +106,12 @@ def or_(*fns):
def and_(*fns):
return lambda *args: operator.and_(*(fn(*args) for fn in fns))
def not_(*fns):
return lambda *args: operator.not_(*(fn(*args) for fn in fns))
procs = [
DaemonProcess("manage_athenad", "openpilot.system.athena.manage_athenad", "AthenadPid"),
NativeProcess("loggerd", "openpilot/system/loggerd", ["./loggerd"], logging),
NativeProcess("encoderd", "openpilot/system/loggerd", ["./encoderd"], only_onroad),
NativeProcess("stream_encoderd", "openpilot/system/loggerd", ["./encoderd", "--stream"], or_(and_(livestream, not_(iscar)), notcar)),
NativeProcess("stream_encoderd", "openpilot/system/loggerd", ["./encoderd", "--stream"], or_(livestream, notcar)),
PythonProcess("logmessaged", "openpilot.system.logmessaged", always_run),
NativeProcess("camerad", "openpilot/system/camerad", ["./camerad"], or_(driverview, livestream), enabled=not WEBCAM),
@@ -159,7 +156,7 @@ procs = [
# debug procs
NativeProcess("bridge", "openpilot/cereal/messaging", ["./bridge"], notcar),
PythonProcess("webrtcd", "openpilot.system.webrtc.webrtcd", or_(and_(livestream, not_(iscar)), notcar)),
PythonProcess("webrtcd", "openpilot.system.webrtc.webrtcd", or_(livestream, notcar)),
PythonProcess("joystick", "openpilot.tools.joystick.joystick_control", and_(joystick, iscar)),
# sunnylink <3
+3 -3
View File
@@ -23,9 +23,9 @@ def post_stream_request(body: StreamRequestBody) -> dict:
ret["time"] = (t_end - t_start) * 1000
return ret
except requests.ConnectTimeout as e:
raise Exception("webrtc took too long to respond.") from e
raise Exception("device took too long to respond.") from e
except requests.ConnectionError as e:
raise Exception("webrtc server on device is not running.") from e
raise Exception("turn car ignition off to use livestreaming.") from e
def wait_for_webrtcd(max_retries: float = 10) -> None:
@@ -37,4 +37,4 @@ def wait_for_webrtcd(max_retries: float = 10) -> None:
except requests.ConnectionError:
attempts += 1
time.sleep(0.5)
raise TimeoutError("webrtcd did not initialize in time.")
raise TimeoutError("livestreaming service did not initialize in time.")
+27 -4
View File
@@ -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,
)
@@ -329,9 +335,12 @@ class StreamSession:
async def run(self):
try:
self.params.put("LivestreamRequestKeyframe", True)
# avoid datachannel race by adding messange_handler immediately
self.stream.set_message_handler(self.message_handler)
await asyncio.wait_for(self.stream.wait_for_connection(), timeout=15)
if self.stream.has_messaging_channel():
self.stream.set_message_handler(self.message_handler)
if self.incoming_bridge is not None:
await self.shared_pub_master.add_services_if_needed(self.incoming_bridge_services)
if self.outgoing_bridge is not None:
@@ -341,14 +350,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 +435,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()