1314 lines
51 KiB
Python
1314 lines
51 KiB
Python
#!/usr/bin/env python3
|
|
"""sanad_api_go2 — Go2 fleet TELEMETRY agent.
|
|
|
|
Pushes the Unitree **Go2**'s live status to the YS Lootah fleet server:
|
|
|
|
POST {SERVER_URL}/api/v1/fleet/ingest/telemetry (Bearer device token)
|
|
body: { "sn", "mac", "battery", "charging", "status", "position":{x,y}, "faults":[] }
|
|
|
|
Same shape/cadence as the R1 agent. The DIFFERENCE is the DDS family:
|
|
|
|
* Go2 uses **unitree_go** (not unitree_hg).
|
|
* Go2 battery lives INSIDE LowState_ as a nested ``bms_state`` (soc/current) —
|
|
there is NO separate rt/lf/bmsstate topic like the G1/R1. So this agent reads
|
|
battery straight from rt/lowstate.bms_state.
|
|
* position (optional) from SportModeState (rt/lf/sportmodestate) or rosbridge /odom.
|
|
|
|
⚠️ UNVERIFIED ON HARDWARE: written from the unitree_go SDK layout but not yet run
|
|
on a real Go2. Confirm the bms current sign (charging) and sportmodestate fields
|
|
on the robot. Degrades to heartbeats if unitree_sdk2py is unavailable; --simulate
|
|
tests the upload path without a robot. SAFETY: never commands motion (read-only).
|
|
|
|
CONFIG — environment (see .env.example)
|
|
---------------------------------------
|
|
SERVER_URL, DEVICE_TOKEN fleet base URL + bearer token (required)
|
|
SN fleet id default go2_0000
|
|
DDS_INTERFACE robot network iface for DDS default eth0
|
|
DDS_DOMAIN DDS domain id default 0
|
|
MAC_INTERFACE iface whose MAC to report default = DDS_INTERFACE
|
|
GO2_POSITION_SOURCE none | sportmode | rosbridge default none
|
|
ROSBRIDGE_URL ws://127.0.0.1:9090 (position) default ws://127.0.0.1:9090
|
|
LOW_SOC / MOTOR_TEMP_MAX fault thresholds default 15 / 85
|
|
POLL_INTERVAL seconds between posts default 2
|
|
VERIFY_TLS / HTTP_TIMEOUT TLS verify (1) / timeout (10)
|
|
|
|
CLI: --simulate | --once | --dry-run | --interval N | -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")
|
|
|
|
|
|
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")
|
|
|
|
|
|
|
|
@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
|
|
dds_interface: str
|
|
dds_domain: int
|
|
mac_interface: str
|
|
position_source: str
|
|
rosbridge_url: str
|
|
low_soc: int
|
|
motor_temp_max: float
|
|
poll_interval: float
|
|
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
|
|
# 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
|
|
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),
|
|
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", "15")),
|
|
motor_temp_max=float(_env("MOTOR_TEMP_MAX", "85")),
|
|
poll_interval=float(_env("POLL_INTERVAL", "2")),
|
|
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")),
|
|
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")),
|
|
verify_tls=_env_bool("VERIFY_TLS", True),
|
|
http_timeout=float(_env("HTTP_TIMEOUT", "10")),
|
|
)
|
|
|
|
def telemetry_url(self) -> str:
|
|
return self.server_url + self.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 auth_headers(self) -> Dict[str, str]:
|
|
return {"Authorization": f"Bearer {self.device_token}"}
|
|
|
|
|
|
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))
|
|
|
|
|
|
_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
|
|
|
|
|
|
class DDSReader:
|
|
"""Subscribes rt/lowstate (unitree_go LowState_) and (optionally) sportmodestate.
|
|
Battery comes from the nested LowState_.bms_state. Passive reads only."""
|
|
|
|
def __init__(self, cfg: Config):
|
|
self.cfg = cfg
|
|
self._lock = threading.Lock()
|
|
self._bms: Optional[Dict[str, Any]] = None
|
|
self._low_ts = 0.0
|
|
self._sport_ts = 0.0
|
|
self._temps: List[float] = []
|
|
self._max_dq = 0.0
|
|
self._xy: Optional[Dict[str, float]] = 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 == "sportmode":
|
|
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_sport, 10)
|
|
self.ok = True
|
|
log.info("DDS up: domain=%d iface=%s (rt/lowstate; battery from bms_state)",
|
|
self.cfg.dds_domain, self.cfg.dds_interface)
|
|
except Exception as e:
|
|
log.warning("DDS init failed (%s) — heartbeat mode", e)
|
|
|
|
def _on_low(self, msg) -> None:
|
|
try:
|
|
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
|
|
# Go2: pack voltage is LowState_.power_v (V); pack temp from the
|
|
# BMS NTC sensors (bq_ntc / mcu_ntc, °C). All defensive getattrs.
|
|
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_vals = []
|
|
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_vals.extend(v for v in vals if -40 <= v <= 150)
|
|
if ntc_vals:
|
|
temp_c = max(ntc_vals)
|
|
except Exception:
|
|
temp_c = None
|
|
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),
|
|
}
|
|
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_sport(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)}
|
|
self._sport_ts = time.monotonic()
|
|
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,
|
|
}
|
|
|
|
|
|
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) — 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.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
|
|
size = p.stat().st_size
|
|
self._pos = size # default: only NEW lines ship
|
|
# backfill: start N lines before EOF so recent history ships at startup
|
|
if self._backfill and size:
|
|
try:
|
|
take = min(size, 512 * 1024)
|
|
with p.open("rb") as f:
|
|
f.seek(size - take)
|
|
tail = f.read(take)
|
|
parts = tail.splitlines(keepends=True)[-self._backfill:]
|
|
self._pos = size - sum(len(x) for x in parts)
|
|
except Exception:
|
|
self._pos = size
|
|
self.label = label
|
|
self.active = True
|
|
|
|
def poll(self) -> List[str]:
|
|
"""New lines since last poll (docker json-log unwrapped), labeled."""
|
|
if not self.active or self._cur is None:
|
|
return []
|
|
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 ln:
|
|
out.append(f"[{self.label}] {ln}")
|
|
except Exception as e:
|
|
log.debug("project-log poll failed: %s", e)
|
|
return out[-100:] # cap per 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
|
|
|
|
|
|
_ALERT_SEEN: set = set()
|
|
|
|
|
|
def send_alerts(cfg: Config, session: requests.Session, faults: List[str]) -> None:
|
|
"""POST each NEW fault (rising edge) to /{sn}/alert. Faults are strings."""
|
|
global _ALERT_SEEN
|
|
current = set(faults)
|
|
new = current - _ALERT_SEEN
|
|
_ALERT_SEEN = current
|
|
for f in sorted(new):
|
|
try:
|
|
r = session.post(cfg.alert_url(),
|
|
json={"sn": cfg.sn, "name": cfg.name,
|
|
"alert": f, "message": f, "ts": int(time.time())},
|
|
headers=cfg.auth_headers(),
|
|
timeout=cfg.http_timeout, verify=cfg.verify_tls)
|
|
_ALERTS_STAT.update(last=f, last_time=_now_str(), ok=bool(r.ok))
|
|
if r.ok:
|
|
_ALERTS_STAT["sent"] += 1
|
|
log.info("alert sent: %s -> HTTP %s", f, r.status_code)
|
|
except requests.RequestException as e:
|
|
_ALERTS_STAT.update(last=f, last_time=_now_str(), ok=False)
|
|
log.debug("alert send failed (%s): %s", f, e)
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# telemetry assembly
|
|
# --------------------------------------------------------------------------- #
|
|
def derive_faults(cfg: Config, snap: Dict[str, Any]) -> List[Dict[str, Any]]:
|
|
# 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]) -> str:
|
|
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 "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": 28.6, "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")}
|
|
else:
|
|
snap = reader.snapshot() if reader else {"bms": None, "low_age": None, "temps": [], "max_dq": 0.0, "xy": 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)
|
|
faults = derive_faults(cfg, snap)
|
|
|
|
# Battery detail (voltage / current / pack temp / health / cycles).
|
|
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")}
|
|
|
|
# Motor temperature stats; null = temps not receiving.
|
|
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_XX)
|
|
"mac": mac,
|
|
"brand": cfg.brand,
|
|
"type": cfg.robot_type, # humanoid | dog
|
|
"model": cfg.model, # r1 | g1 | go2
|
|
"battery": battery, "charging": charging,
|
|
"battery_detail": battery_detail,
|
|
"motor_temp": motor_temp, # null = not receiving
|
|
"storage": read_storage(cfg),
|
|
"status": status,
|
|
"position": position, "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
|
|
"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
|
|
log.info("telemetry ok: battery=%s charging=%s status=%s pos=%s faults=%d -> HTTP %s",
|
|
payload["battery"], payload["charging"], payload["status"],
|
|
payload["position"], len(payload["faults"]), 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,
|
|
"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 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}")
|
|
|
|
|
|
def main(argv: Optional[List[str]] = None) -> int:
|
|
ap = argparse.ArgumentParser(description="Go2 fleet telemetry agent")
|
|
ap.add_argument("--simulate", action="store_true")
|
|
ap.add_argument("--once", action="store_true")
|
|
ap.add_argument("--dry-run", action="store_true")
|
|
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)
|
|
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")
|
|
_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
|
|
_PROJECT_TAIL = ProjectLogTail(cfg)
|
|
|
|
mac = read_mac(cfg.mac_interface)
|
|
log.info("sanad_api_go2 telemetry — sn=%s mac=%s server=%s iface=%s domain=%d%s",
|
|
cfg.sn, mac, cfg.server_url, cfg.dds_interface, cfg.dds_domain,
|
|
" [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() -> 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:
|
|
map_sync_once(cfg, session, force=args.force, dry_run=args.dry_run)
|
|
if args.once:
|
|
one(); ship_logs(cfg, session); return 0
|
|
if args.dry_run:
|
|
for _ in range(3):
|
|
one(); time.sleep(min(cfg.poll_interval, 1.0))
|
|
return 0
|
|
|
|
threading.Thread(target=map_loop, args=(cfg, session), daemon=True).start()
|
|
threading.Thread(target=logs_loop, args=(cfg, session), daemon=True).start()
|
|
log.info("loop every %.1fs; map every %.0fs; logs every %.0fs (Ctrl-C to stop)",
|
|
cfg.poll_interval, cfg.map_poll_interval, cfg.logs_interval)
|
|
while True:
|
|
try:
|
|
one()
|
|
except Exception as e:
|
|
log.exception("tick failed: %s", e)
|
|
try:
|
|
time.sleep(cfg.poll_interval)
|
|
except KeyboardInterrupt:
|
|
log.info("stopped"); return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|