#!/usr/bin/env python3 """sanad_api_go2 — Go2 fleet agent: TELEMETRY + MAP sync in ONE service. The single Go2 agent (one container, one systemd service) that: 1. TELEMETRY — every ~2 s POSTs the robot's live status: POST {SERVER_URL}/api/v1/fleet/ingest/telemetry { sn, name, mac, brand, type, model, battery, charging, battery_detail, motor_temp, storage, status, position, faults, map, ts } 2. MAP — a background loop (every MAP_POLL_INTERVAL, default 30 s) checks the Sanad dashboard's saved maps (web_nav3: RTAB-Map .db + places) and uploads each map ONE time (re-upload only if its content changes): POST {SERVER_URL}/api/v1/fleet/ingest/{sn}/map (multipart: db + meta) The map result is SHOWN inside every telemetry post as the "map" field: "map": { "uploaded": true|false, "state": "uploaded" | "no_map" | "failed" | "pending", "maps_found": N, "last_map": "floor-1"|null, "error": "no saved map found …"|null, "checked_ts": … } so the server always sees whether the map made it — and why not. DATA SOURCES (Unitree Go2, unitree_go DDS (⚠ UNVERIFIED on hardware)) ----------------------------------------- battery / charging : rt/lf/bmsstate (BmsState_) soc, current, voltage, temp, soh, cycles faults / motor temp: rt/lowstate (LowState_) per-motor temps + staleness position {x,y} : rt/lf/odommodestate (SportModeState_) firmware odom over DDS status : derived (charging/moving/idle/offline); optional FSM GET RPC (G1 ids: 200 walk-ready / 4 stand / 2 squat / 702 lie2stand) storage : host disk via the read-only /:/host mount (+ Sanad data dir size) maps : MAPS_DIR//*.db (+ maps_meta.json), places under DATA_DIR//places/.json — the web_nav3 stores Read-only toward the robot: never commands motion. Degrades to heartbeats if unitree_sdk2py is unavailable; --simulate fakes only the DDS side (the map scan stays real). CONFIG — environment (see .env.example). Key vars: SERVER_URL, DEVICE_TOKEN, SN (required at install), ROBOT_NAME, DDS_INTERFACE (eth0), POLL_INTERVAL (2), TELEMETRY_ENDPOINT, ROBOT (web_nav3 robot name, default sanad), MAPS_DIR, DATA_DIR, STATE_DIR, MAP_SELECT (all|active|newest), MAP_UPLOAD_MODE (multipart|base64json), MAP_ENDPOINT, MAP_POLL_INTERVAL (30), VERIFY_TLS, HTTP_TIMEOUT CLI: --simulate | --once | --dry-run | --force (map re-upload) | --list (maps) | -v """ from __future__ import annotations import argparse import base64 import datetime as _dt import hashlib import json import logging import math import os import shutil import sys import threading import time import uuid from dataclasses import dataclass, field from pathlib import Path from typing import Any, Dict, List, Optional import requests log = logging.getLogger("sanad_api_go2") # --------------------------------------------------------------------------- # # env helpers # --------------------------------------------------------------------------- # def _load_dotenv(path: str = ".env") -> None: p = Path(path) if not p.exists(): return for line in p.read_text().splitlines(): line = line.strip() if not line or line.startswith("#") or "=" not in line: continue k, _, v = line.partition("=") os.environ.setdefault(k.strip(), v.strip().strip('"').strip("'")) def _env(name: str, default: str = "") -> str: return os.environ.get(name, default).strip() def _env_bool(name: str, default: bool) -> bool: return _env(name, "1" if default else "0").lower() in ("1", "true", "yes", "on") def _now_str() -> str: """Full local date+time with UTC offset (e.g. 2026-07-13 14:20:33+04:00). TZ_OFFSET_HOURS (default +4, Dubai) keeps container clocks honest without tzdata.""" off = float(_env("TZ_OFFSET_HOURS", "4")) tz = _dt.timezone(_dt.timedelta(hours=off)) return _dt.datetime.now(tz).isoformat(sep=" ", timespec="seconds") # --------------------------------------------------------------------------- # # config # --------------------------------------------------------------------------- # @dataclass class Config: server_url: str device_token: str sn: str name: str brand: str robot_type: str model: str storage_path: str data_path: str # telemetry / DDS dds_interface: str dds_domain: int mac_interface: str read_fsm: bool position_source: str rosbridge_url: str low_soc: int motor_temp_max: float alert_log_patterns: str alert_log_cooldown: float alert_scan_interval: float alert_backfill_bytes: int poll_interval: float telemetry_endpoint: str # map sync robot: str maps_dir: Path extra_map_dirs: List[Path] # extra roots to scan for pgm+yaml sets (e.g. Nav2/Pudu maps) web_data_dir: Optional[Path] legacy_places: Optional[Path] web_nav3_url: str map_select: str map_upload_mode: str map_endpoint_tmpl: str map_poll_interval: float map_max_upload_mb: float state_dir: Path # logs + alerts alert_endpoint: str logs_endpoint: str logs_interval: float # remote dashboard (register the Sanad UI URL for the fleet to embed) remote_enable: bool remote_endpoint: str remote_kind: str remote_host: str remote_ports: str remote_url: str remote_interval: float ssh_enable: bool ssh_user: str ssh_port: int control_url: str control_enable: bool # project logs (e.g. the robot's Sanad app) shipped alongside agent logs project_log_container: str project_log_path: str project_log_label: str project_log_backfill: int project_log_exclude: str ros_distro: str # transport verify_tls: bool http_timeout: float @classmethod def from_env(cls) -> "Config": server = _env("SERVER_URL").rstrip("/") token = _env("DEVICE_TOKEN") missing = [n for n, v in (("SERVER_URL", server), ("DEVICE_TOKEN", token)) if not v] if missing: raise SystemExit(f"[config] missing required env: {', '.join(missing)}") iface = _env("DDS_INTERFACE", "eth0") data_dir = _env("DATA_DIR") legacy = _env("LEGACY_PLACES") return cls( server_url=server, device_token=token, sn=_env("SN", "go2_0000"), name=_env("ROBOT_NAME", "") or _env("SN", "go2_0000"), brand=_env("ROBOT_BRAND", "unitree"), robot_type=_env("ROBOT_TYPE", "dog"), model=_env("ROBOT_MODEL", "go2"), storage_path=_env("STORAGE_PATH", ""), data_path=_env("STORAGE_DATA_PATH", ""), dds_interface=iface, dds_domain=int(_env("DDS_DOMAIN", "0")), mac_interface=_env("MAC_INTERFACE", iface), read_fsm=_env_bool("GO2_READ_FSM", False), position_source=_env("GO2_POSITION_SOURCE", "none").lower(), rosbridge_url=_env("ROSBRIDGE_URL", "ws://127.0.0.1:9090"), low_soc=int(_env("LOW_SOC", "50")), motor_temp_max=float(_env("MOTOR_TEMP_MAX", "85")), # log-driven alerts: "CODE=regex" entries separated by ";;" (regex may # contain '|'). Scanned against the robot's project logs (sanadr1). # NOTE: matching is CASE-SENSITIVE (log levels are uppercase); use an # inline (?i) prefix for case-insensitive text (Gemini messages). alert_log_patterns=_env("ALERT_LOG_PATTERNS", "GEMINI_BILLING=(?i)prepayment credits.{0,40}deplet|Please go to AI Studio|Failed to connect to Gemini" ";;ROBOT_ERROR=\\bERROR\\b|\\bCRITICAL\\b|^Traceback"), alert_log_cooldown=float(_env("ALERT_LOG_COOLDOWN", "300")), # per-signature re-alert gap alert_scan_interval=float(_env("ALERT_SCAN_INTERVAL", "10")), alert_backfill_bytes=int(_env("ALERT_BACKFILL_BYTES", str(8 * 1024 * 1024))), poll_interval=float(_env("POLL_INTERVAL", "2")), telemetry_endpoint=_env("TELEMETRY_ENDPOINT", "/api/v1/fleet/ingest/telemetry"), robot=_env("ROBOT", "sanad"), maps_dir=Path(_env("MAPS_DIR", "/data/maps")), # colon-separated extra roots (mounted Nav2/Pudu map dirs). Any *.yaml+*.pgm # set found here is rendered to PNG and uploaded like a slam_toolbox map. extra_map_dirs=[Path(p) for p in _env("EXTRA_MAP_DIRS", "").split(":") if p.strip()], web_data_dir=Path(data_dir) if data_dir else None, legacy_places=Path(legacy) if legacy else None, web_nav3_url=_env("WEB_NAV3_URL", "").rstrip("/"), map_select=_env("MAP_SELECT", "all").lower(), map_upload_mode=_env("MAP_UPLOAD_MODE", "multipart").lower(), map_endpoint_tmpl=_env("MAP_ENDPOINT", "/api/v1/fleet/ingest/{sn}/map"), map_poll_interval=float(_env("MAP_POLL_INTERVAL", "30")), map_max_upload_mb=float(_env("MAP_MAX_UPLOAD_MB", "7")), state_dir=Path(_env("STATE_DIR", "/data/state")), alert_endpoint=_env("ALERT_ENDPOINT", "/api/v1/fleet/ingest/{sn}/alert"), logs_endpoint=_env("LOGS_ENDPOINT", "/api/v1/fleet/ingest/{sn}/logs"), logs_interval=float(_env("LOGS_INTERVAL", "60")), remote_enable=_env_bool("REMOTE_ENABLE", True), remote_endpoint=_env("REMOTE_ENDPOINT", "/api/v1/fleet/ingest/{sn}/remote"), remote_kind=_env("REMOTE_KIND", "web"), remote_host=_env("REMOTE_HOST", ""), remote_ports=_env("REMOTE_PORTS", "8001,8014,8011,8012,8013,8000,8080"), remote_url=_env("REMOTE_URL", ""), # explicit URL (e.g. a public tunnel) wins remote_interval=float(_env("REMOTE_INTERVAL", "60")), ssh_enable=_env_bool("SSH_REGISTER", True), ssh_user=_env("SSH_USER", "unitree"), ssh_port=int(_env("SSH_PORT", "22")), control_url=_env("CONTROL_STATUS_URL", ""), # Sanad /api/controller/status control_enable=_env_bool("CONTROL_ENABLE", False), # remote mode-SWITCH (motion!) — off project_log_container=_env("PROJECT_LOG_CONTAINER", "auto"), project_log_path=_env("PROJECT_LOG_PATH", ""), project_log_label=_env("PROJECT_LOG_LABEL", ""), project_log_backfill=int(_env("PROJECT_LOG_BACKFILL", "100")), # drop uvicorn/access-log noise so shipped lines match the project's # own LIVE LOGS panel exactly (app-logger lines + tracebacks only) project_log_exclude=_env("PROJECT_LOG_EXCLUDE", r"^(INFO|WARNING|ERROR|DEBUG|CRITICAL):\s"), ros_distro=_env("SOFTWARE_ROS", ""), verify_tls=_env_bool("VERIFY_TLS", True), http_timeout=float(_env("HTTP_TIMEOUT", "30")), ) def telemetry_url(self) -> str: return self.server_url + self.telemetry_endpoint def map_url(self) -> str: return self.server_url + self.map_endpoint_tmpl.format(sn=self.sn) def alert_url(self) -> str: return self.server_url + self.alert_endpoint.format(sn=self.sn) def logs_url(self) -> str: return self.server_url + self.logs_endpoint.format(sn=self.sn) def remote_url_ep(self) -> str: return self.server_url + self.remote_endpoint.format(sn=self.sn) def auth_headers(self) -> Dict[str, str]: return {"Authorization": f"Bearer {self.device_token}"} # --------------------------------------------------------------------------- # # mac + storage # --------------------------------------------------------------------------- # def read_mac(interface: str) -> str: p = Path(f"/sys/class/net/{interface}/address") try: mac = p.read_text().strip() if mac and mac != "00:00:00:00:00:00": return mac.lower() except Exception: pass n = uuid.getnode() return ":".join(f"{(n >> b) & 0xff:02x}" for b in range(40, -1, -8)) _AGENT_VERSION = "2026.07.13" _software_cache: Optional[Dict[str, Any]] = None def read_software(cfg: Config) -> Dict[str, Any]: """Robot software/OS card: ROS distro, host OS (via /host), kernel, arch. The container shares the HOST kernel; the host OS comes from /host/etc/os-release (the read-only /:/host mount).""" global _software_cache if _software_cache is not None: return _software_cache sw: Dict[str, Any] = {} # ROS distro: pinned via SOFTWARE_ROS, else detect host installs /opt/ros/* ros = cfg.ros_distro if not ros: try: for base in ("/host/opt/ros", "/opt/ros"): p = Path(base) if p.is_dir(): distros = sorted(d.name for d in p.iterdir() if d.is_dir()) if distros: ros = ",".join(distros) break except Exception: pass sw["ros"] = ros or None # host OS (firmware/OS card) os_name = os_ver = None for rel in ("/host/etc/os-release", "/etc/os-release"): try: kv = {} for line in Path(rel).read_text().splitlines(): if "=" in line: k, _, v = line.partition("=") kv[k] = v.strip().strip('"') os_name = kv.get("PRETTY_NAME") or kv.get("NAME") os_ver = kv.get("VERSION_ID") break except Exception: continue sw["os"] = os_name # e.g. "Ubuntu 20.04.6 LTS" sw["os_version"] = os_ver # e.g. "20.04" u = os.uname() sw["kernel"] = u.release # host kernel (shared with the container) sw["arch"] = u.machine # e.g. aarch64 sw["python"] = ".".join(map(str, sys.version_info[:3])) sw["agent"] = f"{log.name} {_AGENT_VERSION}" _software_cache = sw return sw _firmware_cache: Optional[Dict[str, Any]] = None def read_firmware_static() -> Dict[str, Any]: """Board-level firmware card (static, via /host): compute board model + Jetson L4T/BSP release + kernel. DDS adds live robot/bms fw versions.""" global _firmware_cache if _firmware_cache is not None: return dict(_firmware_cache) fw: Dict[str, Any] = {} # board model, e.g. "NVIDIA Orin NX Developer Kit" for p in ("/host/sys/firmware/devicetree/base/model", "/host/proc/device-tree/model", "/sys/firmware/devicetree/base/model", "/proc/device-tree/model"): try: fw["board"] = Path(p).read_bytes().decode().strip("\x00 \n") break except Exception: continue # Jetson L4T/BSP: "# R35 (release), REVISION: 3.1, ..." -> "R35.3.1" for p in ("/host/etc/nv_tegra_release", "/etc/nv_tegra_release"): try: head = Path(p).read_text().splitlines()[0] import re m = re.search(r"(R\d+).*?REVISION:\s*([\d.]+)", head) fw["l4t"] = f"{m.group(1)}.{m.group(2)}" if m else head.lstrip("# ").strip() break except Exception: continue fw["kernel"] = os.uname().release _firmware_cache = fw return dict(fw) _data_size_cache: Dict[str, Any] = {"ts": 0.0, "kb": None} def read_storage(cfg: Config) -> Optional[Dict[str, Any]]: """Disk usage of the robot's root fs + optional Sanad data-dir size. In docker, bind-mount the host root read-only at /host (the installer does) so this reports the HOST disk, not the container overlay.""" root = cfg.storage_path or ("/host" if os.path.isdir("/host") else "/") try: du = shutil.disk_usage(root) out: Dict[str, Any] = { "total_gb": round(du.total / 1e9, 2), "free_gb": round(du.free / 1e9, 2), "used_percent": round(du.used / du.total * 100, 1), } except Exception: return None if cfg.data_path and os.path.isdir(cfg.data_path): now = time.monotonic() if _data_size_cache["kb"] is None or now - _data_size_cache["ts"] > 60: try: total = 0 for r, _, files in os.walk(cfg.data_path): for f in files: try: total += os.path.getsize(os.path.join(r, f)) except OSError: pass _data_size_cache.update(ts=now, kb=round(total / 1024, 1)) except Exception: pass if _data_size_cache["kb"] is not None: out["data_kb"] = _data_size_cache["kb"] return out # --------------------------------------------------------------------------- # # DDS reader (telemetry side — degrades to heartbeats if SDK absent) # --------------------------------------------------------------------------- # class DDSReader: """Subscribes rt/lowstate (unitree_go LowState_); battery is NESTED in LowState_.bms_state (Go2 has no separate rt/lf/bmsstate). Optional position from rt/lf/sportmodestate. Passive reads only; the only RPC ever issued is GET_FSM_ID (Go2 has no loco FSM RPC, so it stays off).""" def __init__(self, cfg: Config): self.cfg = cfg self._lock = threading.Lock() self._bms: Optional[Dict[str, Any]] = None self._bms_ts = 0.0 self._low_ts = 0.0 self._temps: List[float] = [] self._max_dq = 0.0 self._xy: Optional[Dict[str, float]] = None self._fw: Dict[str, Any] = {} self._loco = None self.ok = False self._start() def _start(self) -> None: try: from unitree_sdk2py.core.channel import ( ChannelFactoryInitialize, ChannelSubscriber) from unitree_sdk2py.idl.unitree_go.msg.dds_ import LowState_ SportModeState_ = None if self.cfg.position_source in ("sportmode", "odom"): try: from unitree_sdk2py.idl.unitree_go.msg.dds_ import SportModeState_ except Exception: SportModeState_ = None except Exception as e: log.warning("unitree_sdk2py unavailable (%s) — telemetry runs in heartbeat mode", e) return try: ChannelFactoryInitialize(self.cfg.dds_domain, self.cfg.dds_interface) self._low_sub = ChannelSubscriber("rt/lowstate", LowState_) self._low_sub.Init(self._on_low, 10) if SportModeState_ is not None: self._sport_sub = ChannelSubscriber("rt/lf/sportmodestate", SportModeState_) self._sport_sub.Init(self._on_odom, 10) self.ok = True log.info("DDS up: domain=%d iface=%s (rt/lowstate; battery from bms_state%s)", self.cfg.dds_domain, self.cfg.dds_interface, " + sportmode" if SportModeState_ is not None else "") except Exception as e: log.warning("DDS init failed (%s) — heartbeat mode", e) def _init_loco(self) -> None: return # Go2 has no r1/g1-style loco FSM RPC def _on_low(self, msg) -> None: try: # battery from the nested BMS state bms = getattr(msg, "bms_state", None) or getattr(msg, "bms", None) if bms is not None: soc = int(getattr(bms, "soc", 0) or 0) cur = int(getattr(bms, "current", 0) or 0) # mA volt_v = None try: pv = getattr(msg, "power_v", None) if pv: volt_v = round(float(pv), 1) except Exception: volt_v = None temp_c = None try: ntc = [] for attr in ("bq_ntc", "mcu_ntc"): nt = getattr(bms, attr, None) if nt is not None: vals = [int(x) for x in nt] if hasattr(nt, "__iter__") else [int(nt)] ntc.extend(v for v in vals if -40 <= v <= 150) if ntc: temp_c = max(ntc) except Exception: temp_c = None try: vh, vl = getattr(bms, "version_high", None), getattr(bms, "version_low", None) if vh is not None: self._fw["bms"] = f"{int(vh)}.{int(vl or 0)}" except Exception: pass with self._lock: self._bms = { "soc": max(0, min(100, soc)), "current_a": round(cur / 1000.0, 2), "voltage_v": volt_v, "temp_c": temp_c, "soh": int(getattr(bms, "soh", 0) or 0), "cycles": int(getattr(bms, "cycle", 0) or 0), } self._bms_ts = time.monotonic() temps: List[float] = [] max_dq = 0.0 for m in (getattr(msg, "motor_state", None) or []): t = getattr(m, "temperature", None) if t is not None: try: vals = [float(x) for x in t] if hasattr(t, "__iter__") else [float(t)] # 0 = slot not reporting (unpopulated motor), not a real temp temps.extend(v for v in vals if 0 < v <= 200) except Exception: pass dq = getattr(m, "dq", None) if dq is not None: try: max_dq = max(max_dq, abs(float(dq))) except Exception: pass with self._lock: self._low_ts = time.monotonic() self._temps = temps self._max_dq = max_dq except Exception: pass def _on_odom(self, msg) -> None: try: pos = getattr(msg, "position", None) if pos is not None and len(pos) >= 2: with self._lock: self._xy = {"x": round(float(pos[0]), 3), "y": round(float(pos[1]), 3)} except Exception: pass def snapshot(self) -> Dict[str, Any]: with self._lock: now = time.monotonic() return { "bms": dict(self._bms) if self._bms else None, "low_age": (now - self._low_ts) if self._low_ts else None, "temps": list(self._temps), "max_dq": self._max_dq, "xy": dict(self._xy) if self._xy else None, "fw": dict(self._fw), } def fsm_id(self) -> Optional[int]: return None # Go2 has no loco FSM id RPC # Go2 has no r1/g1 FSM-id scheme (SportClient modes). Placeholders — UNVERIFIED. _FSM_STATUS = {} # Control-panel mode labels + the switchable set. _CONTROL_MODES = {0: "idle", 1: "stand", 2: "walk"} _CONTROL_SWITCHABLE = ["damp", "stand", "walk"] _control_cache: Dict[str, Any] = {"ts": 0.0, "data": None} def read_control(cfg: Config) -> Optional[Dict[str, Any]]: """READ-ONLY control status from the Sanad dashboard's /api/controller/status (no DDS, no motion). Reports the current loco mode + the switchable set. Remote SWITCHING is a separate, motion-capable path gated by CONTROL_ENABLE.""" now = time.monotonic() if _control_cache["data"] is not None and now - _control_cache["ts"] < 1.5: return _control_cache["data"] url = cfg.control_url if not url: port = _REMOTE_STAT.get("port") if not port: return None url = f"http://127.0.0.1:{port}/api/controller/status" try: r = requests.get(url, timeout=2) d = r.json() if r.ok else None except Exception: d = None if not isinstance(d, dict): _control_cache.update(ts=now, data=None) return None fid = d.get("fsm_id") out = { "fsm_id": fid, "mode": _CONTROL_MODES.get(fid, "unknown") if fid is not None else None, "armed": d.get("armed"), "walk_ready": d.get("walk_ready"), "teleop_active": d.get("teleop_active"), "last_velocity": d.get("last_velocity"), "sdk_available": d.get("sdk_available"), "switchable_modes": _CONTROL_SWITCHABLE, "remote_switch_enabled": bool(cfg.control_enable), # remote mode-switch armed? } _control_cache.update(ts=now, data=out) return out class RosbridgePosition: def __init__(self, cfg: Config): self.cfg = cfg self._xy: Optional[Dict[str, float]] = None self._lock = threading.Lock() self._stop = False try: import websocket # noqa: F401 except Exception as e: log.warning("websocket-client absent (%s) — rosbridge position disabled", e) self._ok = False return self._ok = True threading.Thread(target=self._run, daemon=True).start() def _run(self) -> None: import websocket sub = json.dumps({"op": "subscribe", "topic": "/odom", "type": "nav_msgs/Odometry", "throttle_rate": 500}) while not self._stop: try: ws = websocket.create_connection(self.cfg.rosbridge_url, timeout=5) ws.send(sub) while not self._stop: msg = json.loads(ws.recv()) pos = (((msg.get("msg") or {}).get("pose") or {}).get("pose") or {}).get("position") if pos: with self._lock: self._xy = {"x": round(float(pos["x"]), 3), "y": round(float(pos["y"]), 3)} except Exception as e: log.debug("rosbridge position reconnect: %s", e) time.sleep(3) def get(self) -> Optional[Dict[str, float]]: with self._lock: return dict(self._xy) if self._xy else None # --------------------------------------------------------------------------- # # map sync (web_nav3 saved maps → fleet server, uploaded ONCE per content) # --------------------------------------------------------------------------- # @dataclass class MapArtifact: path: Path # .db (rtabmap) or .yaml (slam_toolbox set) name: str stem: str size: int mtime: int fmt: str = "rtabmap_db" # rtabmap_db | slam_toolbox files: Dict[str, Path] = field(default_factory=dict) # slam set: pgm/yaml/posegraph/data description: str = "" sha256: str = "" points: List[Dict[str, Any]] = field(default_factory=list) def fingerprint(self) -> str: return f"{self.fmt}:{self.size}:{self.mtime}" def _sha256(path: Path) -> str: h = hashlib.sha256() with path.open("rb") as f: for chunk in iter(lambda: f.read(1024 * 1024), b""): h.update(chunk) return h.hexdigest() def _sha256_set(paths: List[Path]) -> str: h = hashlib.sha256() for p in sorted(paths): with p.open("rb") as f: for chunk in iter(lambda: f.read(1024 * 1024), b""): h.update(chunk) return h.hexdigest() # ---- slam_toolbox map set (office.yaml + office.pgm [+ .posegraph .data]) ---- def _parse_map_yaml(path: Path) -> Dict[str, Any]: """Tiny parser for a ROS map_server yaml (image/resolution/origin) — no pyyaml.""" out: Dict[str, Any] = {} for line in path.read_text().splitlines(): line = line.split("#", 1)[0].strip() if ":" not in line: continue k, _, v = line.partition(":") k, v = k.strip(), v.strip() if k == "image": out["image"] = v elif k == "resolution": try: out["resolution"] = float(v) except ValueError: pass elif k == "origin": try: nums = [float(x) for x in v.strip("[]").split(",")] out["origin"] = {"x": nums[0], "y": nums[1], "yaw": nums[2] if len(nums) > 2 else 0.0} except Exception: pass return out def _read_pgm(path: Path) -> Optional[Dict[str, Any]]: """Parse a binary PGM (P5): returns {width, height, maxval, pixels(bytes)}.""" try: data = path.read_bytes() if not data.startswith(b"P5"): return None # tokenize header (magic, width, height, maxval), skipping comments tokens: List[bytes] = [] i = 2 while len(tokens) < 3 and i < len(data): c = data[i:i + 1] if c in b" \t\r\n": i += 1 elif c == b"#": i = data.index(b"\n", i) + 1 else: j = i while j < len(data) and data[j:j + 1] not in b" \t\r\n": j += 1 tokens.append(data[i:j]) i = j w, h, maxval = int(tokens[0]), int(tokens[1]), int(tokens[2]) pixels = data[i + 1: i + 1 + w * h] if len(pixels) < w * h: return None return {"width": w, "height": h, "maxval": maxval, "pixels": pixels} except Exception: return None def _pgm_to_png_b64(pgm: Dict[str, Any]) -> str: """Grayscale 8-bit PNG from parsed PGM — pure stdlib (zlib + struct).""" import struct import zlib w, h, pixels = pgm["width"], pgm["height"], pgm["pixels"] def chunk(tag: bytes, body: bytes) -> bytes: return (struct.pack(">I", len(body)) + tag + body + struct.pack(">I", zlib.crc32(tag + body) & 0xFFFFFFFF)) ihdr = struct.pack(">IIBBBBB", w, h, 8, 0, 0, 0, 0) # 8-bit grayscale raw = b"".join(b"\x00" + pixels[y * w:(y + 1) * w] for y in range(h)) png = (b"\x89PNG\r\n\x1a\n" + chunk(b"IHDR", ihdr) + chunk(b"IDAT", zlib.compress(raw, 6)) + chunk(b"IEND", b"")) return base64.b64encode(png).decode("ascii") def _discover_slam_sets(cfg: Config) -> List[MapArtifact]: """Find slam_toolbox / map_server map sets: .yaml + .pgm (+ optional .posegraph/.data) under the maps roots and maps_slam/.""" roots = [cfg.maps_dir, cfg.maps_dir / cfg.robot, cfg.maps_dir / "maps_slam"] roots += list(cfg.extra_map_dirs) # Nav2/Pudu map dirs mounted via EXTRA_MAP_DIRS seen: set = set() out: List[MapArtifact] = [] for root in roots: if not root.exists(): continue for y in sorted(root.glob("*.yaml")): meta = _parse_map_yaml(y) img = meta.get("image", "") pgm = (y.parent / img) if img else y.with_suffix(".pgm") if not pgm.exists(): pgm = y.with_suffix(".pgm") if not pgm.exists(): continue # yaml without a raster — not a map set rp = str(y.resolve()) if rp in seen: continue seen.add(rp) files: Dict[str, Path] = {"yaml": y, "pgm": pgm} for ext in ("posegraph", "data"): p = y.with_suffix("." + ext) if p.exists(): files[ext] = p size = sum(p.stat().st_size for p in files.values()) mtime = max(int(p.stat().st_mtime) for p in files.values()) out.append(MapArtifact(path=y, name=y.name, stem=y.stem, size=size, mtime=mtime, fmt="slam_toolbox", files=files)) # Pudu/Nav2 converter emits a plain map + a keepout-BAKED twin (obstacles baked # in for Foxy, which has no KeepoutFilter). The baked one is the deploy map — drop # the redundant plain twin so the fleet server gets one canonical map, not two. baked = {m.stem[: -len("_keepout_baked")] for m in out if m.stem.endswith("_keepout_baked")} out = [m for m in out if m.stem not in baked] return out def _map_key(stem: str) -> str: stem = Path(stem).name if stem.endswith(".db"): stem = stem[:-3] return "".join(c for c in stem if c.isalnum() or c in "_-.") def _read_json(path: Path, default: Any) -> Any: try: return json.loads(path.read_text() or "") except Exception: return default def _yaw_from_pose(pose: Dict[str, Any]) -> float: if "qw" in pose or "qz" in pose: qx = float(pose.get("qx", 0.0)); qy = float(pose.get("qy", 0.0)) qz = float(pose.get("qz", 0.0)); qw = float(pose.get("qw", 1.0)) return math.atan2(2.0 * (qw * qz + qx * qy), 1.0 - 2.0 * (qy * qy + qz * qz)) return float(pose.get("yaw", 0.0)) def _places_files_for(cfg: Config, stem: str) -> List[Path]: out: List[Path] = [] key = _map_key(stem) if cfg.web_data_dir: out.append(cfg.web_data_dir / cfg.robot / "places" / f"{key}.json") if cfg.legacy_places: out.append(cfg.legacy_places) return out def load_points(cfg: Config, stem: str) -> List[Dict[str, Any]]: for pf in _places_files_for(cfg, stem): data = _read_json(pf, None) if pf.exists() else None if isinstance(data, dict) and data: pts: List[Dict[str, Any]] = [] for name, pose in data.items(): if not isinstance(pose, dict): continue try: pts.append({ "name": name, "type": str(pose.get("type", "waypoint")), "x": float(pose["x"]), "y": float(pose["y"]), "yaw": round(_yaw_from_pose(pose), 4), }) except (KeyError, TypeError, ValueError): continue return pts return [] def discover_maps(cfg: Config) -> List[MapArtifact]: roots = [cfg.maps_dir / cfg.robot, cfg.maps_dir] meta: Dict[str, Any] = {} meta_file = cfg.maps_dir / cfg.robot / "maps_meta.json" if meta_file.exists(): meta = _read_json(meta_file, {}) or {} seen: set = set() out: List[MapArtifact] = [] for root in roots: if not root.exists(): continue for p in sorted(root.glob("*.db")): rp = str(p.resolve()) if rp in seen: continue seen.add(rp) st = p.stat() out.append(MapArtifact( path=p, name=p.name, stem=p.stem, size=st.st_size, mtime=int(st.st_mtime), description=(meta.get(p.name) or {}).get("description", ""), )) # slam_toolbox map sets (office.yaml + office.pgm …) live alongside out.extend(_discover_slam_sets(cfg)) out.sort(key=lambda m: m.mtime, reverse=True) return out def _active_map_name(cfg: Config) -> Optional[str]: if not cfg.web_nav3_url: return None try: r = requests.get(cfg.web_nav3_url + "/api/status", headers={"X-Robot-Name": cfg.robot}, timeout=min(cfg.http_timeout, 5)) r.raise_for_status() am = (r.json() or {}).get("active_map") return _map_key(am) if am else None except requests.RequestException: return None def select_maps(cfg: Config, maps: List[MapArtifact]) -> List[MapArtifact]: if not maps: return [] if cfg.map_select == "newest": return maps[:1] if cfg.map_select == "active": active = _active_map_name(cfg) if active: picked = [m for m in maps if _map_key(m.stem) == active] if picked: return picked return maps[:1] return maps # "all" def _state_file(cfg: Config) -> Path: return cfg.state_dir / "uploaded.json" def load_state(cfg: Config) -> Dict[str, str]: return _read_json(_state_file(cfg), {}) if _state_file(cfg).exists() else {} def save_state(cfg: Config, state: Dict[str, str]) -> None: try: cfg.state_dir.mkdir(parents=True, exist_ok=True) _state_file(cfg).write_text(json.dumps(state, indent=2)) except Exception as e: log.warning("could not persist map state: %s", e) def build_meta(cfg: Config, m: MapArtifact) -> Dict[str, Any]: return { "sn": cfg.sn, "name": m.stem, "file": m.name, "format": m.fmt, "size_bytes": m.size, "sha256": m.sha256, "mtime": m.mtime, "description": m.description, "points": m.points, } def _upload_slam_map(cfg: Config, m: MapArtifact, session: requests.Session) -> bool: """slam_toolbox map → the spec's image JSON: PNG (from the pgm) + resolution + origin + width/height + points. This is what the dashboard renders.""" url = cfg.map_url() ymeta = _parse_map_yaml(m.files["yaml"]) pgm = _read_pgm(m.files["pgm"]) if pgm is None: log.error("map %s: cannot parse %s (not binary P5?)", m.stem, m.files["pgm"].name) return False body = build_meta(cfg, m) body.update({ "resolution": ymeta.get("resolution"), "origin": ymeta.get("origin"), "width": pgm["width"], "height": pgm["height"], "image_base64": _pgm_to_png_b64(pgm), }) try: resp = session.post(url, json=body, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) except requests.RequestException as e: log.error("map upload %s FAILED (transport): %s", m.stem, e) return False if not resp.ok: log.error("map upload %s FAILED: HTTP %s %s", m.stem, resp.status_code, resp.text[:300]) return False log.info("map uploaded: %s (slam_toolbox %dx%d @ %sm, %d points) -> HTTP %s", m.stem, pgm["width"], pgm["height"], ymeta.get("resolution"), len(m.points), resp.status_code) return True def upload_map(cfg: Config, m: MapArtifact, session: requests.Session) -> bool: if m.fmt == "slam_toolbox": return _upload_slam_map(cfg, m, session) url = cfg.map_url() meta = build_meta(cfg, m) try: if cfg.map_upload_mode == "base64json": body = dict(meta) body["db_base64"] = base64.b64encode(m.path.read_bytes()).decode("ascii") resp = session.post(url, json=body, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) else: # multipart (default) with m.path.open("rb") as fh: files = {"db": (m.name, fh, "application/octet-stream")} data = {"meta": json.dumps(meta)} resp = session.post(url, files=files, data=data, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) except requests.RequestException as e: log.error("map upload %s FAILED (transport): %s", m.name, e) return False if not resp.ok: log.error("map upload %s FAILED: HTTP %s %s", m.name, resp.status_code, resp.text[:300]) return False log.info("map uploaded: %s (%.2f MB, %d points) -> HTTP %s", m.name, m.size / 1024 / 1024, len(m.points), resp.status_code) return True # Shared map status — SHOWN in every telemetry post ("map" field). _MAP_STATUS_LOCK = threading.Lock() _MAP_STATUS: Dict[str, Any] = { "uploaded": False, "state": "pending", "maps_found": 0, "last_map": None, "error": None, "checked_ts": None, } def _set_map_status(**kw: Any) -> None: with _MAP_STATUS_LOCK: _MAP_STATUS.update(kw) _MAP_STATUS["checked_ts"] = int(time.time()) def get_map_status() -> Dict[str, Any]: with _MAP_STATUS_LOCK: return dict(_MAP_STATUS) def map_sync_once(cfg: Config, session: requests.Session, force: bool = False, dry_run: bool = False) -> int: """One map pass: scan the Sanad dashboard maps and upload anything new. Always updates the shared map status (visible in telemetry).""" try: maps = select_maps(cfg, discover_maps(cfg)) except Exception as e: _set_map_status(state="failed", uploaded=False, error=f"map scan failed: {e}") return 0 if not maps: _set_map_status(state="no_map", uploaded=False, maps_found=0, last_map=None, error=f"no saved map found in Sanad dashboard " f"(maps_dir={cfg.maps_dir}, robot={cfg.robot})") return 0 state = load_state(cfg) uploaded = failed = unstable = too_large = current = 0 last_err: Optional[str] = None now = time.time() for m in maps: prev = state.get(str(m.path.resolve())) if not force and prev == m.fingerprint(): current += 1 continue # already uploaded this exact content — one-time rule # stability guard: a db modified in the last 120 s is still being # written (active mapping) — wait until it settles before uploading if not force and (now - m.mtime) < 120: log.info("map %s still changing (mapping in progress) — waiting to settle", m.name) unstable += 1 continue # server rejects bodies over ~8 MB (client_max_body_size) — don't burn # bandwidth on uploads that will 413. Raster (slam_toolbox) maps are tiny. if m.fmt != "slam_toolbox" and (m.size / 1048576) > cfg.map_max_upload_mb: log.warning("map %s is %.0f MB — exceeds server upload cap (~%.0f MB), skipping " "(export a raster map or raise the server limit)", m.name, m.size / 1048576, cfg.map_max_upload_mb) too_large += 1 continue m.sha256 = (_sha256_set(list(m.files.values())) if m.fmt == "slam_toolbox" else _sha256(m.path)) m.points = load_points(cfg, m.stem) if dry_run: log.info("[dry-run] would upload map %s (%.2f MB, %d points)", m.name, m.size / 1024 / 1024, len(m.points)) continue if upload_map(cfg, m, session): state[str(m.path.resolve())] = m.fingerprint() save_state(cfg, state) uploaded += 1 else: failed += 1 last_err = f"upload failed for {m.name} (see agent log)" if failed: _set_map_status(state="failed", uploaded=False, maps_found=len(maps), last_map=maps[0].stem, error=last_err) elif uploaded or current: # at least one map is on the server (just now or previously); note skips note = None if too_large: note = f"{too_large} map(s) skipped: exceed server upload cap (~{cfg.map_max_upload_mb:.0f} MB)" elif unstable: note = "newer map still being written (mapping in progress)" _set_map_status(state="uploaded", uploaded=True, maps_found=len(maps), last_map=maps[0].stem, error=note) elif too_large: _set_map_status(state="failed", uploaded=False, maps_found=len(maps), last_map=maps[0].stem, error=f"map exceeds server upload cap (~{cfg.map_max_upload_mb:.0f} MB) — " "export a raster map or raise the server limit") elif unstable: # newest content is still being written (active mapping) — be honest _set_map_status(state="pending", uploaded=False, maps_found=len(maps), last_map=maps[0].stem, error="map still being written (mapping in progress) — " "will upload when it settles") else: _set_map_status(state="uploaded", uploaded=True, maps_found=len(maps), last_map=maps[0].stem, error=None) return uploaded def map_loop(cfg: Config, session: requests.Session) -> None: while True: try: map_sync_once(cfg, session) except Exception as e: log.exception("map pass failed: %s", e) time.sleep(cfg.map_poll_interval) # --------------------------------------------------------------------------- # # logs + alerts (spec: POST /{sn}/logs periodically, POST /{sn}/alert on events) # --------------------------------------------------------------------------- # class _RingLogHandler(logging.Handler): """Buffers the agent's own log lines so they can be shipped to the server.""" def __init__(self, maxlen: int = 400): super().__init__(level=logging.INFO) from collections import deque self._buf: Any = deque(maxlen=maxlen) self._blk = threading.Lock() def emit(self, record: logging.LogRecord) -> None: try: with self._blk: self._buf.append(self.format(record)) except Exception: pass def drain(self) -> List[str]: with self._blk: lines = list(self._buf) self._buf.clear() return lines def requeue(self, lines: List[str]) -> None: """Put unshipped lines back (front of the ring) so they retry next cycle instead of being lost — bounded by maxlen, oldest evicted first.""" with self._blk: self._buf.extendleft(reversed(lines)) _LOG_RING = _RingLogHandler() # shipped-status shown in every telemetry post ("logs" / "alerts" fields) _LOGS_STAT: Dict[str, Any] = {"last_sent": None, "lines_sent": 0, "ok": None} _ALERTS_STAT: Dict[str, Any] = {"sent": 0, "last": None, "last_time": None, "ok": None} # start times ("started_at" = this run, "last_start" = previous run) _STARTED: Dict[str, Any] = {"now": None, "prev": None, "mono": time.monotonic()} def _init_start_times(cfg: Config) -> None: """Record this agent start; remember the previous one (persisted in STATE_DIR).""" f = cfg.state_dir / "agent_state.json" prev = (_read_json(f, {}) or {}).get("started_at") now_s = _now_str() try: cfg.state_dir.mkdir(parents=True, exist_ok=True) f.write_text(json.dumps({"started_at": now_s})) except Exception as e: log.debug("could not persist start time: %s", e) _STARTED.update(now=now_s, prev=prev, mono=time.monotonic()) class ProjectLogTail: """Tails the robot's main PROJECT logs (e.g. the sanadr1 / sanad-p4 app) and feeds them into the shipped log lines, labeled "[-logs] …". Sources, in priority order: PROJECT_LOG_PATH explicit log file (or dir -> newest *.log) via /host PROJECT_LOG_CONTAINER a docker container name; "auto" (default) scans the host's docker metadata (/host/var/lib/docker) for a RUNNING Sanad project (sanadr1, sanad-p4, sanad*) Reads the container's json-log through the read-only /:/host mount — no docker socket needed, read-only, cannot disturb the project.""" KNOWN = ("sanadr1", "sanad-p4", "sanadv3", "sanad") def __init__(self, cfg: Config): self.label: Optional[str] = None self._cur: Optional[Path] = None self._pos = 0 self._backfill = max(0, cfg.project_log_backfill) self._exclude = None if cfg.project_log_exclude: try: import re self._exclude = re.compile(cfg.project_log_exclude) except Exception: self._exclude = None self.active = False try: self._resolve(cfg) except Exception as e: log.debug("project-log resolve failed: %s", e) if self.active: log.info("project logs: sharing '%s' (%s)", self.label, self._cur) else: log.info("project logs: none found (project_logs=null)") def _resolve(self, cfg: Config) -> None: # explicit file/dir if cfg.project_log_path: p = Path(cfg.project_log_path) if p.is_dir(): logs = sorted(p.glob("*.log"), key=lambda f: f.stat().st_mtime, reverse=True) p = logs[0] if logs else None if p and p.exists(): self._start(p, cfg.project_log_label or f"{p.stem}-logs") return # docker container json-log via /host base = Path("/host/var/lib/docker/containers") if not base.exists(): return want = cfg.project_log_container candidates: List[Any] = [] for cf in base.glob("*/config.v2.json"): try: d = json.loads(cf.read_text()) except Exception: continue name = (d.get("Name") or "").lstrip("/") running = bool((d.get("State") or {}).get("Running")) lp = d.get("LogPath") or "" if not name or not lp: continue if want != "auto": if name == want: candidates.append((0, name, lp, running)) elif running and "sanad" in name.lower() and not name.startswith("sanad-api"): # rank known Sanad projects first rank = self.KNOWN.index(name) if name in self.KNOWN else len(self.KNOWN) candidates.append((rank, name, lp, running)) if not candidates: return candidates.sort(key=lambda c: c[0]) _, name, lp, _ = candidates[0] p = Path("/host" + lp) if not lp.startswith("/host") else Path(lp) if p.exists(): self._start(p, cfg.project_log_label or f"{name}-logs") def _start(self, p: Path, label: str) -> None: self._cur = p self.label = label size = p.stat().st_size self._pos = size # new lines ship from here on self._pending: List[str] = [] # backfill: the last N RELEVANT lines (post-filter) from the last ~1 MB, # so the noise (access-log spam) doesn't eat the history window if self._backfill and size: try: take = min(size, 8 * 1024 * 1024) with p.open("rb") as f: f.seek(size - take) tail = f.read(take) raw = tail.decode("utf-8", "replace").splitlines() if take < size and raw: raw = raw[1:] # first line is a partial record — drop it for ln in raw: ln = ln.strip() if ln.startswith("{"): try: ln = (json.loads(ln).get("log") or "").rstrip() except Exception: pass if not ln: continue if self._exclude is not None and self._exclude.search(ln): continue self._pending.append(f"[{label}] {ln}") self._pending = self._pending[-self._backfill:] except Exception: self._pending = [] self.active = True def poll(self) -> List[str]: """New lines since last poll (docker json-log unwrapped), labeled. The startup backfill (last N relevant lines) is returned on first call.""" if not self.active or self._cur is None: return [] pend, self._pending = self._pending, [] out: List[str] = [] try: st = self._cur.stat() if st.st_size < self._pos: # log rotated self._pos = 0 if st.st_size > self._pos: with self._cur.open("rb") as f: f.seek(self._pos) chunk = f.read(min(st.st_size - self._pos, 256 * 1024)) self._pos = f.tell() for ln in chunk.decode("utf-8", "replace").splitlines(): ln = ln.strip() if ln.startswith("{"): try: ln = (json.loads(ln).get("log") or "").rstrip() except Exception: pass if not ln: continue if self._exclude is not None and self._exclude.search(ln): continue # access-log noise — not in the project's log panel out.append(f"[{self.label}] {ln}") except Exception as e: log.debug("project-log poll failed: %s", e) return pend + out[-100:] # backfill first, then ≤100 new/cycle _PROJECT_TAIL: Optional[ProjectLogTail] = None def ship_logs(cfg: Config, session: requests.Session) -> None: """POST buffered agent log lines (+ project logs) to /{sn}/logs. Best-effort: failures are logged at DEBUG only (below the ring's level -> no feedback loop).""" lines = _LOG_RING.drain() if _PROJECT_TAIL is not None and _PROJECT_TAIL.active: lines.extend(_PROJECT_TAIL.poll()) if not lines: return try: r = session.post(cfg.logs_url(), json={"sn": cfg.sn, "name": cfg.name, "lines": lines, "ts": int(time.time())}, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) if r.ok: _LOGS_STAT.update(last_sent=_now_str(), ok=True) _LOGS_STAT["lines_sent"] += len(lines) log.debug("logs shipped: %d lines -> HTTP %s", len(lines), r.status_code) else: _LOGS_STAT["ok"] = False _LOG_RING.requeue(lines) # retry next cycle (server keeps 500ing) log.debug("logs ship failed: HTTP %s (%d lines requeued)", r.status_code, len(lines)) except requests.RequestException as e: _LOGS_STAT["ok"] = False _LOG_RING.requeue(lines) log.debug("logs ship failed (%d lines requeued): %s", len(lines), e) def logs_loop(cfg: Config, session: requests.Session) -> None: while True: time.sleep(cfg.logs_interval) try: ship_logs(cfg, session) except Exception: pass # --------------------------------------------------------------------------- # # remote dashboard — discover the Sanad web UI and register its URL for the # fleet dashboard to embed ("full dashboard through the fleet API"). No changes # to the Sanad app: we only probe its port and POST the URL to /{sn}/remote. # --------------------------------------------------------------------------- # _REMOTE_STAT: Dict[str, Any] = {"url": None, "port": None, "kind": None, "ok": None, "ssh": None, "ssh_ok": None} def _primary_ip() -> str: import socket try: s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) s.connect(("1.1.1.1", 80)) ip = s.getsockname()[0] s.close() return ip except Exception: return "127.0.0.1" def discover_dashboard(cfg: Config) -> Optional[Dict[str, Any]]: """Find the Sanad dashboard: probe each candidate port on localhost; the one that returns 200 with a 'Sanad'/dashboard page wins. Returns {url, port}.""" if cfg.remote_url: # explicit URL (e.g. a tunnel) overrides return {"url": cfg.remote_url, "port": None} host = cfg.remote_host or _primary_ip() # the LAN IP the fleet can reach for tok in cfg.remote_ports.split(","): tok = tok.strip() if not tok.isdigit(): continue port = int(tok) try: r = requests.get(f"http://127.0.0.1:{port}/", timeout=2) except requests.RequestException: continue if r.status_code == 200 and ("sanad" in r.text.lower() or "dashboard" in r.text.lower()): return {"url": f"http://{host}:{port}", "port": port} return None def _post_remote(cfg: Config, session: requests.Session, body: Dict[str, Any]) -> bool: try: r = session.post(cfg.remote_url_ep(), json=body, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) if not r.ok: log.warning("remote(%s) register failed: HTTP %s %s", body.get("kind"), r.status_code, r.text[:120]) return bool(r.ok) except requests.RequestException as e: log.debug("remote(%s) register failed: %s", body.get("kind"), e) return False def register_remote(cfg: Config, session: requests.Session) -> None: if not cfg.remote_enable: return host = cfg.remote_host or _primary_ip() # 1) the Sanad dashboard (kind=web) — the exact UI, for the fleet to embed d = discover_dashboard(cfg) if d: ok = _post_remote(cfg, session, { "sn": cfg.sn, "name": cfg.name, "kind": cfg.remote_kind, "url": d["url"], "label": f"{cfg.name} — Sanad Dashboard", "port": d["port"], "ts": int(time.time())}) _REMOTE_STAT.update(url=d["url"], port=d["port"], kind=cfg.remote_kind, ok=ok) if ok: log.info("remote dashboard registered: %s -> HTTP 200", d["url"]) else: _REMOTE_STAT.update(url=None, port=None, ok=None) # 2) SSH access (kind=ssh) — "ssh unitree@" if cfg.ssh_enable: cmd = f"ssh {cfg.ssh_user}@{host}" if cfg.ssh_port != 22: cmd += f" -p {cfg.ssh_port}" ok = _post_remote(cfg, session, { "sn": cfg.sn, "name": cfg.name, "kind": "ssh", "url": f"ssh://{cfg.ssh_user}@{host}:{cfg.ssh_port}", # URL-valid form "command": cmd, # "ssh unitree@" "host": host, "port": cfg.ssh_port, "user": cfg.ssh_user, "label": f"{cfg.name} — SSH", "ts": int(time.time())}) _REMOTE_STAT.update(ssh=cmd, ssh_ok=ok) if ok: log.info("remote SSH registered: %s -> HTTP 200", cmd) def remote_loop(cfg: Config, session: requests.Session) -> None: while True: try: register_remote(cfg, session) except Exception as e: log.debug("remote loop: %s", e) time.sleep(cfg.remote_interval) _ALERT_SEEN: set = set() def _post_alert(cfg: Config, session: requests.Session, text: str) -> bool: """POST a single alert string to /{sn}/alert; updates the alert status.""" try: r = session.post(cfg.alert_url(), json={"sn": cfg.sn, "name": cfg.name, "alert": text, "message": text, "ts": int(time.time())}, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) _ALERTS_STAT.update(last=text, last_time=_now_str(), ok=bool(r.ok)) if r.ok: _ALERTS_STAT["sent"] += 1 log.info("alert sent: %s -> HTTP %s", text[:120], r.status_code) return bool(r.ok) except requests.RequestException as e: _ALERTS_STAT.update(last=text, last_time=_now_str(), ok=False) log.debug("alert send failed (%s): %s", text[:60], e) return False def send_alerts(cfg: Config, session: requests.Session, faults: List[str]) -> None: """POST each NEW fault (rising edge) to /{sn}/alert. Faults are strings. Includes LOW_BATTERY (<= LOW_SOC, default 50%), MOTOR_OVERTEMP, COMMS_STALE.""" global _ALERT_SEEN current = set(faults) new = current - _ALERT_SEEN _ALERT_SEEN = current for f in sorted(new): _post_alert(cfg, session, f) class LogAlertScanner: """Scans the robot's project logs (sanadr1) for error/billing patterns and fires an alert on each NEW signature (deduped with a cooldown). Catches the Gemini 'credits depleted' billing error and any ERROR/Traceback the app logs.""" def __init__(self, cfg: Config): import re self._path: Optional[Path] = None self._pos = 0 self._cooldown = cfg.alert_log_cooldown self._seen: Dict[str, float] = {} self._patterns: List[Any] = [] self._pending: List[Any] = [] for entry in cfg.alert_log_patterns.split(";;"): if "=" in entry: code, rx = entry.split("=", 1) try: # case-SENSITIVE (uppercase log levels); use (?i) inline for text self._patterns.append((code.strip(), re.compile(rx))) except re.error as e: log.warning("bad alert pattern %s: %s", code, e) if _PROJECT_TAIL is not None and _PROJECT_TAIL.active and _PROJECT_TAIL._cur: self._path = _PROJECT_TAIL._cur try: size = self._path.stat().st_size self._pos = size # ongoing reads start at EOF # STARTUP BACKFILL: scan a large tail once (access-logs bury real # errors) so a currently-active error (e.g. Gemini billing) alerts # right away. Only the matching lines are kept — cheap. back = min(size, cfg.alert_backfill_bytes) with self._path.open("rb") as f: f.seek(size - back) raw = f.read(back).decode("utf-8", "replace").splitlines() if back < size and raw: raw = raw[1:] # first line is partial matches = [] for ln in raw: ln = ln.strip() if ln.startswith("{"): try: ln = (json.loads(ln).get("log") or "").rstrip() except Exception: pass if not ln: continue for code, rx in self._patterns: if rx.search(ln): matches.append((code, ln)) break # keep the most-recent unique signature per (code+line), cap 20 seen = set() uniq = [] for code, ln in reversed(matches): sig = code + ":" + re.sub(r"\d+", "#", ln)[:120] if sig in seen: continue seen.add(sig) uniq.append((code, ln)) self._pending = list(reversed(uniq))[-20:] except Exception: self._pos = 0 if self._path and self._patterns: log.info("log-alert scan on %s (%d patterns, %d backfilled)", self._path.name, len(self._patterns), len(self._pending)) def scan(self, cfg: Config, session: requests.Session) -> None: if not self._path or not self._patterns: return import re now0 = time.monotonic() if self._pending: # flush startup backfill first pend, self._pending = self._pending, [] for code, ln in pend: sig = code + ":" + re.sub(r"\d+", "#", ln)[:120] self._seen[sig] = now0 _post_alert(cfg, session, f"{code}: {ln[:220]}") try: st = self._path.stat() if st.st_size < self._pos: # rotated self._pos = 0 if st.st_size <= self._pos: return with self._path.open("rb") as f: f.seek(self._pos) chunk = f.read(min(st.st_size - self._pos, 512 * 1024)) self._pos = f.tell() except Exception as e: log.debug("alert scan read failed: %s", e) return now = time.monotonic() for ln in chunk.decode("utf-8", "replace").splitlines(): ln = ln.strip() if ln.startswith("{"): try: ln = (json.loads(ln).get("log") or "").rstrip() except Exception: pass if not ln: continue for code, rx in self._patterns: if rx.search(ln): sig = code + ":" + re.sub(r"\d+", "#", ln)[:120] # dedup key if now - self._seen.get(sig, -1e9) > self._cooldown: self._seen[sig] = now _post_alert(cfg, session, f"{code}: {ln[:220]}") break _LOG_ALERTS: Optional[LogAlertScanner] = None def alert_scan_loop(cfg: Config, session: requests.Session) -> None: while True: time.sleep(cfg.alert_scan_interval) try: if _LOG_ALERTS is not None: _LOG_ALERTS.scan(cfg, session) except Exception as e: log.debug("alert scan loop: %s", e) # --------------------------------------------------------------------------- # # telemetry assembly # --------------------------------------------------------------------------- # def derive_faults(cfg: Config, snap: Dict[str, Any]) -> List[str]: # STRINGS, not objects — the fleet ingest 500s on fault objects. faults: List[str] = [] bms = snap.get("bms") if bms and bms.get("soc", 100) <= cfg.low_soc: faults.append(f"LOW_BATTERY: battery {bms['soc']}% (warning)") temps = snap.get("temps") or [] if temps and max(temps) >= cfg.motor_temp_max: faults.append(f"MOTOR_OVERTEMP: motor temp {max(temps):.0f}C (warning)") if snap.get("low_age") is not None and snap["low_age"] > 3.0: faults.append(f"COMMS_STALE: no rt/lowstate for {snap['low_age']:.0f}s (critical)") return faults def derive_status(cfg: Config, snap: Dict[str, Any], fsm: Optional[int]) -> str: base = _FSM_STATUS.get(fsm) if fsm is not None else None bms = snap.get("bms") charging = bool(bms and bms.get("current_a", 0.0) > 0.05) alive = snap.get("low_age") is not None and snap["low_age"] <= 3.0 if not alive and bms is None: return "offline" if charging: return "charging" if snap.get("max_dq", 0.0) > 0.15: return "moving" return base or "idle" def build_telemetry(cfg: Config, mac: str, reader: Optional[DDSReader], pos: Optional[RosbridgePosition], sim: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: if sim is not None: snap = {"bms": {"soc": sim["battery"], "current_a": 0.5 if sim["charging"] else -0.3, "voltage_v": 47.5, "temp_c": 36, "soh": 100, "cycles": 45}, "low_age": 0.1, "temps": [sim.get("temp", 45)], "max_dq": sim.get("max_dq", 0.0), "xy": sim.get("position")} fsm = sim.get("fsm") else: snap = reader.snapshot() if reader else {"bms": None, "low_age": None, "temps": [], "max_dq": 0.0, "xy": None} fsm = reader.fsm_id() if (reader and cfg.read_fsm) else None bms = snap.get("bms") battery = bms["soc"] if bms else None charging = bool(bms and bms.get("current_a", 0.0) > 0.05) status = derive_status(cfg, snap, fsm) faults = derive_faults(cfg, snap) battery_detail = None if bms: battery_detail = {"voltage_v": bms.get("voltage_v"), "current_a": bms.get("current_a"), "temp_c": bms.get("temp_c"), "soh": bms.get("soh"), "cycles": bms.get("cycles")} temps = snap.get("temps") or [] motor_temp = ({"max": round(max(temps), 1), "avg": round(sum(temps) / len(temps), 1), "min": round(min(temps), 1)} if temps else None) position = snap.get("xy") if position is None and pos is not None: position = pos.get() return { "sn": cfg.sn, "name": cfg.name, # friendly display name (e.g. go2_77) "mac": mac, "brand": cfg.brand, "type": cfg.robot_type, # humanoid "model": cfg.model, # g1 "software": read_software(cfg), # ros/os/kernel/arch/python/agent "firmware": {**read_firmware_static(), **(snap.get("fw") or {})}, # board/l4t/kernel + live robot/bms fw versions "battery": battery, # null = couldn't read (heartbeat) "charging": charging, "battery_detail": battery_detail, "motor_temp": motor_temp, # null = not receiving "storage": read_storage(cfg), "status": status, "position": position, # null when no odom/localization source "control": read_control(cfg), # loco mode (zero_torque/damp/lock/running) + switchable set "faults": faults, "map": get_map_status(), # SHOWS whether the saved map made it to the server "logs": dict(_LOGS_STAT), # log-shipping status (last_sent, lines_sent, ok) "project_logs": (_PROJECT_TAIL.label if _PROJECT_TAIL is not None and _PROJECT_TAIL.active else None), # e.g. "sanadr1-logs"; null = no project found "remote": (dict(_REMOTE_STAT) if _REMOTE_STAT.get("url") else None), # Sanad dashboard URL registered for the fleet UI "alerts": dict(_ALERTS_STAT), # alert status (sent, last, last_time, ok) "time": _now_str(), # full date+time of this post "started_at": _STARTED["now"], # when this agent run started "last_start": _STARTED["prev"], # previous agent start (null on first ever) "uptime_s": int(time.monotonic() - _STARTED["mono"]), "ts": int(time.time()), } def post_telemetry(cfg: Config, payload: Dict[str, Any], session: requests.Session) -> bool: try: r = session.post(cfg.telemetry_url(), json=payload, headers=cfg.auth_headers(), timeout=cfg.http_timeout, verify=cfg.verify_tls) except requests.RequestException as e: log.error("telemetry POST failed (transport): %s", e) return False if not r.ok: log.error("telemetry POST failed: HTTP %s %s", r.status_code, r.text[:200]) return False mp = payload.get("map") or {} log.info("telemetry ok: battery=%s charging=%s status=%s pos=%s faults=%d map=%s -> HTTP %s", payload["battery"], payload["charging"], payload["status"], payload["position"], len(payload["faults"]), mp.get("state"), r.status_code) return True def _sim_state(i: int) -> Dict[str, Any]: charging = (i % 6) in (0, 1) battery = max(5, 90 - (i % 40)) moving = (i % 3) == 2 and not charging return {"battery": battery, "charging": charging, "temp": 45 + (i % 10), "max_dq": 0.4 if moving else 0.0, "fsm": 200 if moving else 4, "position": {"x": round(1.0 + 0.1 * i, 2), "y": round(2.0 - 0.05 * i, 2)}} def cmd_list(cfg: Config) -> None: maps = discover_maps(cfg) if not maps: print(f"(no .db maps under {cfg.maps_dir} for robot '{cfg.robot}')") return print(f"{len(maps)} map(s) under {cfg.maps_dir} (robot={cfg.robot}):") for m in maps: pts = load_points(cfg, m.stem) print(f" {m.name:<28} {m.size/1024/1024:6.2f} MB {len(pts):>3} points {m.description}") # --------------------------------------------------------------------------- # # main # --------------------------------------------------------------------------- # def main(argv: Optional[List[str]] = None) -> int: ap = argparse.ArgumentParser(description="Go2 fleet agent: telemetry + map sync") ap.add_argument("--simulate", action="store_true", help="synthetic DDS state (map scan stays real)") ap.add_argument("--once", action="store_true", help="one map pass + one telemetry post, then exit") ap.add_argument("--map-only", action="store_true", help="upload discovered maps once (no DDS, no telemetry), then exit") ap.add_argument("--dry-run", action="store_true", help="build payloads, never POST") ap.add_argument("--force", action="store_true", help="re-upload maps even if unchanged") ap.add_argument("--list", action="store_true", help="list discovered maps and exit") ap.add_argument("--interval", type=float, default=None, help="override telemetry POLL_INTERVAL") ap.add_argument("-v", "--verbose", action="store_true") args = ap.parse_args(argv) logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s") # buffer our own log lines for shipping to /{sn}/logs _LOG_RING.setFormatter(logging.Formatter("%(asctime)s %(levelname)s %(name)s: %(message)s")) logging.getLogger().addHandler(_LOG_RING) _load_dotenv() cfg = Config.from_env() if args.interval is not None: cfg.poll_interval = args.interval if args.list: cmd_list(cfg) return 0 _init_start_times(cfg) global _PROJECT_TAIL, _LOG_ALERTS _PROJECT_TAIL = ProjectLogTail(cfg) _LOG_ALERTS = LogAlertScanner(cfg) # error/billing alerts from the project logs mac = read_mac(cfg.mac_interface) log.info("sanad_api_go2 — sn=%s name=%s mac=%s server=%s iface=%s pos=%s " "map_dir=%s map_every=%.0fs%s", cfg.sn, cfg.name, mac, cfg.server_url, cfg.dds_interface, cfg.position_source, cfg.maps_dir, cfg.map_poll_interval, " [SIMULATE]" if args.simulate else "") session = requests.Session() if args.map_only: # map upload only — search all map roots (incl. EXTRA_MAP_DIRS) and ship, # without opening DDS or posting telemetry (won't disturb a live feed). map_sync_once(cfg, session, force=args.force, dry_run=args.dry_run) return 0 reader = None pos = None if not args.simulate: reader = DDSReader(cfg) if cfg.position_source == "rosbridge": pos = RosbridgePosition(cfg) time.sleep(1.0) tick = 0 def one_telemetry() -> None: nonlocal tick sim = _sim_state(tick) if args.simulate else None payload = build_telemetry(cfg, mac, reader, pos, sim=sim) if args.dry_run: log.info("[dry-run] %s", json.dumps(payload)) else: post_telemetry(cfg, payload, session) send_alerts(cfg, session, payload.get("faults") or []) tick += 1 if args.once or args.dry_run: # one map pass first so the telemetry "map" field reflects it map_sync_once(cfg, session, force=args.force, dry_run=args.dry_run) if cfg.remote_enable and not args.dry_run: register_remote(cfg, session) elif cfg.remote_enable: d = discover_dashboard(cfg) if d: _REMOTE_STAT.update(url=d["url"], port=d["port"], kind=cfg.remote_kind, ok=None) if args.once: one_telemetry() ship_logs(cfg, session) return 0 for _ in range(3): one_telemetry() time.sleep(min(cfg.poll_interval, 1.0)) return 0 # loop mode: map sync + log shipping + remote registration in bg threads threading.Thread(target=map_loop, args=(cfg, session), daemon=True).start() threading.Thread(target=logs_loop, args=(cfg, session), daemon=True).start() threading.Thread(target=alert_scan_loop, args=(cfg, session), daemon=True).start() if cfg.remote_enable: threading.Thread(target=remote_loop, args=(cfg, session), daemon=True).start() log.info("telemetry every %.1fs; map check every %.0fs; logs every %.0fs (Ctrl-C to stop)", cfg.poll_interval, cfg.map_poll_interval, cfg.logs_interval) while True: try: one_telemetry() except Exception as e: log.exception("telemetry tick failed: %s", e) try: time.sleep(cfg.poll_interval) except KeyboardInterrupt: log.info("stopped") return 0 if __name__ == "__main__": sys.exit(main())