mirror of
https://github.com/firestar5683/StarPilot.git
synced 2026-09-30 03:13:48 +08:00
Bring back temporalPose path, mapd wrapper, and offroad-only backup behavior
This commit is contained in:
+1
-1
@@ -1087,7 +1087,7 @@ struct ModelDataV2 {
|
||||
confidence @23: ConfidenceClass;
|
||||
|
||||
# Model perceived motion
|
||||
temporalPoseDEPRECATED @21 :Pose;
|
||||
temporalPose @21 :Pose;
|
||||
|
||||
# e2e lateral planner
|
||||
action @26: Action;
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
@@ -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.")
|
||||
@@ -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()
|
||||
|
||||
@@ -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),
|
||||
|
||||
Reference in New Issue
Block a user