diff --git a/openpilot/selfdrive/modeld/modeld.py b/openpilot/selfdrive/modeld/modeld.py index 3c7dd935a2..68f49b4113 100755 --- a/openpilot/selfdrive/modeld/modeld.py +++ b/openpilot/selfdrive/modeld/modeld.py @@ -359,10 +359,11 @@ def main(demo=False, remote_addr: str | None = None, big_model: bool = False): fill_driving_model_data(drivingdata_send, modelv2_send) fill_pose_msg(posenet_send, model_output, meta_main.frame_id, vipc_dropped_frames, meta_main.timestamp_eof, live_calib_seen) - pm.send('modelV2', modelv2_send) - pm.send('drivingModelData', drivingdata_send) - pm.send('cameraOdometry', posenet_send) - pm.send('modelDataV2SP', mdv2sp_send) + if remote_addr is not None or not params.get_bool("WgpuEnabled"): + pm.send('modelV2', modelv2_send) + pm.send('drivingModelData', drivingdata_send) + pm.send('cameraOdometry', posenet_send) + pm.send('modelDataV2SP', mdv2sp_send) last_vipc_frame_id = meta_main.frame_id diff --git a/openpilot/system/manager/process_config.py b/openpilot/system/manager/process_config.py index 9210b99bb8..a72675385d 100644 --- a/openpilot/system/manager/process_config.py +++ b/openpilot/system/manager/process_config.py @@ -95,9 +95,6 @@ def is_stock_model(started, params, CP: car.CarParams) -> bool: """Check if the active model runner is stock.""" return bool(get_active_model_runner(params, not started) == custom.ModelManagerSP.Runner.stock) -def not_wgpu(started: bool, params: Params, CP: car.CarParams) -> bool: - return not params.get_bool("WgpuEnabled") - def mapd_ready(started: bool, params: Params, CP: car.CarParams) -> bool: return bool(os.path.exists(Paths.mapd_root())) @@ -131,7 +128,7 @@ procs = [ PythonProcess("micd", "openpilot.system.micd", iscar), PythonProcess("timed", "openpilot.system.timed", always_run, enabled=not PC), - PythonProcess("modeld", "openpilot.selfdrive.modeld.modeld", and_(and_(only_onroad, is_stock_model), not_wgpu)), + PythonProcess("modeld", "openpilot.selfdrive.modeld.modeld", and_(only_onroad, is_stock_model)), PythonProcess("dmonitoringmodeld", "openpilot.selfdrive.modeld.dmonitoringmodeld", driverview, enabled=(WEBCAM or not PC)), PythonProcess("sensord", "openpilot.system.sensord.sensord", only_onroad, enabled=not PC), diff --git a/openpilot/tools/wgpu/README.md b/openpilot/tools/wgpu/README.md index 89892913d0..25208b026b 100644 --- a/openpilot/tools/wgpu/README.md +++ b/openpilot/tools/wgpu/README.md @@ -4,9 +4,10 @@ This runs driving `modeld` on a laptop and returns its cereal outputs to a comma device over the existing Wi-Fi network. It reuses the existing HEVC camera stream, VisionIPC decoder, and cereal ZMQ bridge. -This is for offroad/bench testing only. Wi-Fi has no deterministic latency or -availability guarantee. The device restores local `modeld` when the device-side -helper exits, but that is not a seamless onroad failover. +This is for controlled bench testing only. Wi-Fi has no deterministic latency +or availability guarantee. Local `modeld` remains warm during a WGPU session; +the device-side helper switches only after receiving a fresh remote model and +restores local publication if the remote model is missing for 350 ms. ## Build @@ -38,8 +39,9 @@ the small model on that class of laptop. ## Run -Find the laptop's LAN IP address that the comma device can reach. While the -device is offroad, run: +Find the laptop's LAN IP address that the comma device can reach. The helper can +start while onroad: it forwards camera/state while local `modeld` remains active, +then switches publication after the laptop produces a fresh valid model: ```sh cd /data/openpilot @@ -56,7 +58,9 @@ python3 -m openpilot.tools.wgpu.host COMMA_IP Add `--big-model` after `COMMA_IP` to use the locally compiled big model. The first remote `carParams` packet can take up to 50 seconds. Stop either side -with Ctrl+C. Stop the device helper before changing branches or rebooting. +with Ctrl+C. A host disconnect automatically restores the warm local publisher +after a 350 ms timeout. Stop the device helper before changing branches or +rebooting. For stationary bench testing only, model lag can be changed from a blocking event to a live warning showing model age and dropped frames. After ignition is diff --git a/openpilot/tools/wgpu/device.py b/openpilot/tools/wgpu/device.py index 7be206f2e0..1d7e5820f2 100644 --- a/openpilot/tools/wgpu/device.py +++ b/openpilot/tools/wgpu/device.py @@ -5,10 +5,14 @@ import subprocess import time from pathlib import Path +import openpilot.cereal.messaging as messaging from openpilot.common.params import Params +from openpilot.common.realtime import DT_MDL +from openpilot.tools.wgpu.zmq import ZmqSubSocket MODEL_OUTPUTS = "modelV2,drivingModelData,cameraOdometry,modelDataV2SP" +REMOTE_MODEL_TIMEOUT = 0.35 ROOT = Path(__file__).resolve().parents[3] BRIDGE = ROOT / "openpilot/cereal/messaging/bridge" @@ -26,6 +30,17 @@ def handle_sigterm(*_) -> None: raise KeyboardInterrupt +def receive_fresh_model(sock: ZmqSubSocket) -> bool: + raw = sock.receive(non_blocking=True) + if raw is None: + return False + event = messaging.log_from_bytes(raw) + if event.which() != "modelV2" or not event.valid: + return False + model_age = (time.monotonic_ns() - event.modelV2.timestampEof) / 1e9 + return 0 <= model_age < REMOTE_MODEL_TIMEOUT + + def main() -> None: parser = argparse.ArgumentParser(description="Route modeld traffic between this device and a wireless host.") parser.add_argument("host", help="Laptop IP address reachable from this device") @@ -35,27 +50,45 @@ def main() -> None: raise FileNotFoundError(f"build the cereal bridge first: {BRIDGE}") params = Params() - if not params.get_bool("IsOffroad"): - raise RuntimeError("start the wgpu bridge while offroad") - - procs: list[subprocess.Popen] = [] + forward: subprocess.Popen | None = None + reverse: subprocess.Popen | None = None + wgpu_enabled = False try: + # Keep local modeld publishing while the remote model connects and warms up. + forward = subprocess.Popen([str(BRIDGE)]) + remote_model = ZmqSubSocket("modelV2", args.host, conflate=True) + print(f"forwarding camera/state to {args.host}; waiting for a fresh remote model") + while not receive_fresh_model(remote_model): + if forward.poll() is not None: + raise RuntimeError(f"forward bridge exited with status {forward.returncode}") + time.sleep(0.05) + + # Local modeld remains warm but stops publishing when this flag changes. params.put_bool("WgpuEnabled", True, block=True) - procs = [ - subprocess.Popen([str(BRIDGE)]), - subprocess.Popen([str(BRIDGE), args.host, MODEL_OUTPUTS]), - ] - print(f"wgpu enabled; forwarding camera/state to {args.host}") - print("keep this process running; Ctrl+C restores local modeld") - while all(proc.poll() is None for proc in procs): - time.sleep(0.25) - failed = next(proc for proc in procs if proc.poll() is not None) - raise RuntimeError(f"bridge exited with status {failed.returncode}") + wgpu_enabled = True + time.sleep(2 * DT_MDL) + reverse = subprocess.Popen([str(BRIDGE), args.host, MODEL_OUTPUTS]) + print("wgpu active; Ctrl+C or loss of the remote model restores local modeld") + + last_remote_model = time.monotonic() + while True: + if receive_fresh_model(remote_model): + last_remote_model = time.monotonic() + if forward.poll() is not None: + raise RuntimeError(f"forward bridge exited with status {forward.returncode}") + if reverse.poll() is not None: + raise RuntimeError(f"reverse bridge exited with status {reverse.returncode}") + if time.monotonic() - last_remote_model > REMOTE_MODEL_TIMEOUT: + raise RuntimeError("remote model timed out; restoring local modeld") + time.sleep(0.05) finally: - params.put_bool("WgpuEnabled", False, block=True) - for proc in procs: - if proc.poll() is None: - stop_process(proc) + # Stop remote publication before allowing the warm local publisher to resume. + if reverse is not None and reverse.poll() is None: + stop_process(reverse) + if wgpu_enabled: + params.put_bool("WgpuEnabled", False, block=True) + if forward is not None and forward.poll() is None: + stop_process(forward) print("wgpu disabled; local modeld restored")