1869 lines
75 KiB
Python
1869 lines
75 KiB
Python
#!/usr/bin/env python3
|
|
"""sanad_api_g1 — G1 fleet agent: TELEMETRY + MAP sync in ONE service.
|
|
|
|
The single G1 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 G1, unitree_hg DDS)
|
|
-----------------------------------------
|
|
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/<robot>/*.db (+ maps_meta.json), places under
|
|
DATA_DIR/<robot>/places/<map>.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_g1")
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# 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
|
|
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", "g1_0000"),
|
|
name=_env("ROBOT_NAME", "") or _env("SN", "g1_0000"),
|
|
brand=_env("ROBOT_BRAND", "unitree"),
|
|
robot_type=_env("ROBOT_TYPE", "humanoid"),
|
|
model=_env("ROBOT_MODEL", "g1"),
|
|
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("G1_READ_FSM", False),
|
|
position_source=_env("G1_POSITION_SOURCE", "odom").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")),
|
|
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/lf/bmsstate + rt/lowstate (+ rt/lf/odommodestate for position).
|
|
Passive reads; the only RPC ever issued is GET_FSM_ID."""
|
|
|
|
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] = {} # live fw versions (robot ctrl + bms)
|
|
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_hg.msg.dds_ import LowState_
|
|
try:
|
|
from unitree_sdk2py.idl.unitree_hg.msg.dds_ import BmsState_
|
|
except Exception:
|
|
BmsState_ = None
|
|
SportModeState_ = None
|
|
if self.cfg.position_source == "odom":
|
|
try:
|
|
# G1 firmware publishes odom as unitree_go SportModeState_.
|
|
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 BmsState_ is not None:
|
|
self._bms_sub = ChannelSubscriber("rt/lf/bmsstate", BmsState_)
|
|
self._bms_sub.Init(self._on_bms, 10)
|
|
else:
|
|
log.warning("BmsState_ not in this unitree_sdk2py — battery will be null")
|
|
if SportModeState_ is not None:
|
|
self._odom_sub = ChannelSubscriber("rt/lf/odommodestate", SportModeState_)
|
|
self._odom_sub.Init(self._on_odom, 10)
|
|
if self.cfg.read_fsm:
|
|
self._init_loco()
|
|
self.ok = True
|
|
log.info("DDS up: domain=%d iface=%s (bmsstate + lowstate%s)",
|
|
self.cfg.dds_domain, self.cfg.dds_interface,
|
|
" + odom" if SportModeState_ is not None else "")
|
|
except Exception as e:
|
|
log.warning("DDS init failed (%s) — heartbeat mode", e)
|
|
|
|
def _init_loco(self) -> None:
|
|
try:
|
|
from unitree_sdk2py.rpc.client import Client # type: ignore
|
|
except Exception as e:
|
|
log.warning("loco RPC client unavailable (%s) — status from BMS/motion only", e)
|
|
return
|
|
try:
|
|
c = Client("loco", 0); c.Init(); c.SetTimeout(3.0)
|
|
self._loco = c
|
|
log.info("loco FSM read enabled (GET-only, no motion)")
|
|
except Exception as e:
|
|
log.warning("loco client init failed (%s) — status from BMS/motion only", e)
|
|
self._loco = None
|
|
|
|
def _on_bms(self, msg) -> None:
|
|
try:
|
|
soc = int(getattr(msg, "soc", 0) or 0)
|
|
cur_mA = int(getattr(msg, "current", 0) or 0)
|
|
# Pack voltage: prefer bmsvoltage[0] (mV); else sum of cell voltages.
|
|
volt_mv = 0
|
|
bv = getattr(msg, "bmsvoltage", None)
|
|
try:
|
|
if bv is not None and len(bv) and int(bv[0]):
|
|
volt_mv = int(bv[0])
|
|
except Exception:
|
|
volt_mv = 0
|
|
if not volt_mv:
|
|
cv = getattr(msg, "cell_vol", None)
|
|
if cv is not None:
|
|
try:
|
|
volt_mv = int(sum(int(x) for x in cv if x))
|
|
except Exception:
|
|
volt_mv = 0
|
|
# Max plausible pack temperature (int16 °C).
|
|
temp_c = None
|
|
tt = getattr(msg, "temperature", None)
|
|
if tt is not None:
|
|
try:
|
|
vals = [int(x) for x in tt if -40 <= int(x) <= 150]
|
|
if vals:
|
|
temp_c = max(vals)
|
|
except Exception:
|
|
temp_c = None
|
|
# BMS firmware version (version_high.version_low)
|
|
try:
|
|
vh, vl = getattr(msg, "version_high", None), getattr(msg, "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_mA / 1000.0, 2),
|
|
"voltage_v": round(volt_mv / 1000.0, 1) if volt_mv else None,
|
|
"temp_c": temp_c,
|
|
"soh": int(getattr(msg, "soh", 0) or 0),
|
|
"cycles": int(getattr(msg, "cycle", 0) or 0),
|
|
}
|
|
self._bms_ts = time.monotonic()
|
|
except Exception:
|
|
pass
|
|
|
|
def _on_low(self, msg) -> None:
|
|
try:
|
|
# robot controller firmware version (LowState.version array)
|
|
try:
|
|
v = getattr(msg, "version", None)
|
|
if v is not None and "robot" not in self._fw:
|
|
vals = [int(x) for x in v] if hasattr(v, "__iter__") else [int(v)]
|
|
self._fw["robot"] = ".".join(map(str, vals))
|
|
except Exception:
|
|
pass
|
|
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]:
|
|
if not self._loco:
|
|
return None
|
|
try:
|
|
code, data = self._loco._Call(7001, "{}") # GET_FSM_ID — read-only
|
|
if code == 0 and data:
|
|
return int(json.loads(data).get("data", data)) if data.strip().startswith("{") else int(data)
|
|
except Exception as e:
|
|
log.debug("fsm read failed: %s", e)
|
|
return None
|
|
|
|
|
|
# G1 FSM ids (differ from the R1's): 200 balance/walk-ready, 4 StandUp, 2 Squat, 702 Lie2Stand.
|
|
_FSM_STATUS = {200: "ready", 4: "standing", 2: "squat", 702: "lie2stand"}
|
|
# Control-panel mode labels + the switchable set (fsm_id -> friendly mode).
|
|
_CONTROL_MODES = {200: "running", 4: "lock", 2: "squat", 702: "lie2stand", 0: "zero_torque", 1: "damp"}
|
|
_CONTROL_SWITCHABLE = ["zero_torque", "damp", "lock", "running"]
|
|
|
|
|
|
_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: <stem>.yaml + <stem>.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"]
|
|
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))
|
|
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 "[<project>-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@<ip>"
|
|
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@<ip>"
|
|
"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. g1_58)
|
|
"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="G1 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("--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_g1 — 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 "")
|
|
|
|
reader = None
|
|
pos = None
|
|
if not args.simulate:
|
|
reader = DDSReader(cfg)
|
|
if cfg.position_source == "rosbridge":
|
|
pos = RosbridgePosition(cfg)
|
|
time.sleep(1.0)
|
|
|
|
session = requests.Session()
|
|
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())
|