diff --git a/cereal/log.capnp b/cereal/log.capnp index b58f1ff297..6cd6599a6a 100644 --- a/cereal/log.capnp +++ b/cereal/log.capnp @@ -1087,7 +1087,7 @@ struct ModelDataV2 { confidence @23: ConfidenceClass; # Model perceived motion - temporalPoseDEPRECATED @21 :Pose; + temporalPose @21 :Pose; # e2e lateral planner action @26: Action; diff --git a/frogpilot/frogpilot_process.py b/frogpilot/frogpilot_process.py index c32fe0074f..8cc2db85d8 100644 --- a/frogpilot/frogpilot_process.py +++ b/frogpilot/frogpilot_process.py @@ -13,7 +13,6 @@ from openpilot.system.athena.registration import UNREGISTERED_DONGLE_ID from openpilot.frogpilot.assets.model_manager import MODEL_DOWNLOAD_ALL_PARAM, MODEL_DOWNLOAD_PARAM, ModelManager from openpilot.frogpilot.assets.theme_manager import THEME_COMPONENT_PARAMS, ThemeManager -from openpilot.frogpilot.common.frogpilot_backups import backup_toggles from openpilot.frogpilot.common.frogpilot_functions import capture_report, update_maps, update_openpilot from openpilot.frogpilot.common.frogpilot_utilities import ThreadManager, flash_panda, is_url_pingable, lock_doors, use_konik_server from openpilot.frogpilot.common.frogpilot_variables import ERROR_LOGS_PATH, FrogPilotVariables @@ -129,9 +128,6 @@ def update_toggles(frogpilot_variables, started, theme_manager, thread_manager, theme_manager.theme_updated = False theme_manager.update_active_theme(time_validated, frogpilot_toggles, randomize_theme=randomize_theme) - if time_validated: - thread_manager.run_with_lock(backup_toggles, (params)) - return frogpilot_toggles def frogpilot_thread(): @@ -230,8 +226,8 @@ def frogpilot_thread(): theme_manager.update_active_theme(time_validated, frogpilot_toggles) - thread_manager.run_with_lock(backup_toggles, (params, True)) - thread_manager.run_with_lock(send_stats) + if not started: + thread_manager.run_with_lock(send_stats) thread_manager.run_with_lock(update_checks, (now, model_manager, theme_manager, thread_manager, params, params_memory, frogpilot_toggles, True)) rate_keeper.keep_time() diff --git a/frogpilot/navigation/mapd_wrapper.py b/frogpilot/navigation/mapd_wrapper.py new file mode 100644 index 0000000000..f0c24fc338 --- /dev/null +++ b/frogpilot/navigation/mapd_wrapper.py @@ -0,0 +1,164 @@ +#!/usr/bin/env python3 +import json +import os +import signal +import subprocess +import sys +import time + +from collections import defaultdict, deque +from pathlib import Path + +from openpilot.common.basedir import BASEDIR +from openpilot.common.swaglog import cloudlog + +MAPD_DIR = Path(BASEDIR) / "frogpilot/navigation" +MAPD_BIN = MAPD_DIR / "mapd" +OFFLINE_ROOT = Path("/data/media/0/osm/offline") +RESTART_DELAY_S = 0.25 +MISSING_TILE_BACKOFF_S = 30.0 +FAILURE_WINDOW_S = 3.0 +FAILURE_THRESHOLD = 3 + + +def extract_bounds_filename(line: str) -> str | None: + try: + payload = json.loads(line) + except json.JSONDecodeError: + return None + + if payload.get("msg") != "Loading bounds file": + return None + + filename = payload.get("filename") + return filename if isinstance(filename, str) else None + + +def is_offline_read_error(line: str) -> bool: + try: + payload = json.loads(line) + except json.JSONDecodeError: + return False + + return payload.get("msg") == "could not unmarshal offline data" + + +class CorruptTileMonitor: + def __init__(self, threshold: int = FAILURE_THRESHOLD, window_s: float = FAILURE_WINDOW_S): + self.threshold = threshold + self.window_s = window_s + self.current_filename: str | None = None + self.failures: dict[str, deque[float]] = defaultdict(deque) + + def observe(self, line: str, now: float | None = None) -> str | None: + filename = extract_bounds_filename(line) + if filename is not None: + self.current_filename = filename + return None + + if not is_offline_read_error(line) or self.current_filename is None: + return None + + ts = time.monotonic() if now is None else now + failures = self.failures[self.current_filename] + failures.append(ts) + + cutoff = ts - self.window_s + while failures and failures[0] < cutoff: + failures.popleft() + + if len(failures) >= self.threshold: + return self.current_filename + return None + + +def quarantine_offline_tile(filename: str) -> Path | None: + tile_path = Path(filename) + try: + tile_path.relative_to(OFFLINE_ROOT) + except ValueError: + cloudlog.warning(f"mapd_wrapper refusing to quarantine unexpected path: {filename}") + return None + + if not tile_path.exists(): + return None + + quarantined = tile_path.with_name(f"{tile_path.name}.corrupt.{int(time.time())}") + tile_path.rename(quarantined) + return quarantined + + +def terminate_child(proc: subprocess.Popen[str]) -> None: + if proc.poll() is not None: + return + + proc.terminate() + try: + proc.wait(timeout=2) + except subprocess.TimeoutExpired: + proc.kill() + proc.wait(timeout=2) + + +def run_mapd_once() -> int: + proc = subprocess.Popen( + [MAPD_BIN.as_posix()], + cwd=MAPD_DIR, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + bufsize=1, + ) + assert proc.stdout is not None + + def _handle_signal(signum, _frame): + terminate_child(proc) + raise SystemExit(128 + signum) + + signal.signal(signal.SIGTERM, _handle_signal) + signal.signal(signal.SIGINT, _handle_signal) + + monitor = CorruptTileMonitor() + + for line in proc.stdout: + print(line, end="") + bad_tile = monitor.observe(line) + if bad_tile is None: + continue + + quarantined = quarantine_offline_tile(bad_tile) + if quarantined is None: + if not OFFLINE_ROOT.exists(): + cloudlog.warning( + f"mapd_wrapper detected repeated offline read failures for {bad_tile}, " + f"but {OFFLINE_ROOT} does not exist; backing off mapd restarts" + ) + terminate_child(proc) + return 2 + + cloudlog.warning(f"mapd_wrapper detected repeated offline read failures for {bad_tile}, but could not quarantine it") + else: + message = f"mapd_wrapper quarantined corrupt offline tile: {bad_tile} -> {quarantined}" + print(message, flush=True) + cloudlog.warning(message) + + terminate_child(proc) + return 1 if quarantined is not None else 2 + + return proc.wait() + + +def main() -> None: + while True: + exit_code = run_mapd_once() + if exit_code == 1: + time.sleep(RESTART_DELAY_S) + continue + if exit_code == 2: + time.sleep(MISSING_TILE_BACKOFF_S) + continue + raise SystemExit(exit_code) + + +if __name__ == "__main__": + main() diff --git a/frogpilot/navigation/test_mapd_wrapper.py b/frogpilot/navigation/test_mapd_wrapper.py new file mode 100644 index 0000000000..e1c591c666 --- /dev/null +++ b/frogpilot/navigation/test_mapd_wrapper.py @@ -0,0 +1,42 @@ +#!/usr/bin/env python3 +import json + +from pathlib import Path + +from openpilot.frogpilot.navigation.mapd_wrapper import CorruptTileMonitor, quarantine_offline_tile + + +def _loading_line(filename: str) -> str: + return json.dumps({"msg": "Loading bounds file", "filename": filename}) + + +def _error_line() -> str: + return json.dumps({"msg": "could not unmarshal offline data", "error": "EOF"}) + + +def test_corrupt_tile_monitor_triggers_after_repeated_failures(): + filename = "/data/media/0/osm/offline/36/-98/37.500000_-98.000000_37.750000_-97.750000" + monitor = CorruptTileMonitor(threshold=3, window_s=3.0) + + assert monitor.observe(_loading_line(filename), now=0.0) is None + assert monitor.observe(_error_line(), now=0.1) is None + assert monitor.observe(_loading_line(filename), now=0.2) is None + assert monitor.observe(_error_line(), now=0.3) is None + assert monitor.observe(_loading_line(filename), now=0.4) is None + assert monitor.observe(_error_line(), now=0.5) == filename + + +def test_quarantine_offline_tile_renames_file(tmp_path, monkeypatch): + offline_root = tmp_path / "offline" + tile = offline_root / "36/-98/37.500000_-98.000000_37.750000_-97.750000" + tile.parent.mkdir(parents=True) + tile.write_text("bad") + + monkeypatch.setattr("openpilot.frogpilot.navigation.mapd_wrapper.OFFLINE_ROOT", offline_root) + + quarantined = quarantine_offline_tile(tile.as_posix()) + + assert quarantined is not None + assert not tile.exists() + assert Path(quarantined).exists() + assert Path(quarantined).name.startswith(f"{tile.name}.corrupt.") diff --git a/selfdrive/modeld/fill_model_msg.py b/selfdrive/modeld/fill_model_msg.py index 82c4c92b1d..6d113a93f2 100644 --- a/selfdrive/modeld/fill_model_msg.py +++ b/selfdrive/modeld/fill_model_msg.py @@ -123,6 +123,13 @@ def fill_model_msg(base_msg: capnp._DynamicStructBuilder, extended_msg: capnp._D lead.prob = net_output_data['lead_prob'][0,i].tolist() lead.probTime = ModelConstants.LEAD_T_OFFSETS[i] + # temporal pose + temporal_pose = modelV2.temporalPose + temporal_pose.trans = net_output_data['plan'][0,0,Plan.VELOCITY].tolist() + temporal_pose.transStd = net_output_data['plan_stds'][0,0,Plan.VELOCITY].tolist() + temporal_pose.rot = net_output_data['plan'][0,0,Plan.ORIENTATION_RATE].tolist() + temporal_pose.rotStd = net_output_data['plan_stds'][0,0,Plan.ORIENTATION_RATE].tolist() + # meta meta = modelV2.meta meta.desireState = net_output_data['desire_state'][0].reshape(-1).tolist() diff --git a/system/manager/process_config.py b/system/manager/process_config.py index f2b1c10aa5..f203cbd567 100644 --- a/system/manager/process_config.py +++ b/system/manager/process_config.py @@ -136,7 +136,7 @@ else: procs += [ PythonProcess("device_syncd", "frogpilot.system.device_syncd", always_run), PythonProcess("frogpilot_process", "frogpilot.frogpilot_process", always_run), - NativeProcess("mapd", "frogpilot/navigation", ["./mapd"], always_run), + PythonProcess("mapd", "frogpilot.navigation.mapd_wrapper", always_run), PythonProcess("the_pond", "frogpilot.system.the_pond.the_pond", always_run, nice=19), PythonProcess("galaxy", "frogpilot.system.galaxy.galaxy", always_run, nice=19), PythonProcess("speed_limit_filler", "frogpilot.system.speed_limit_filler", run_speed_limit_filler),