agi_fleet/agent/sanad_api_x2.py
2026-08-04 15:14:59 +04:00

2200 lines
90 KiB
Python

#!/usr/bin/env python3
"""sanad_api_x2 — AGIBOT X2 fleet agent: TELEMETRY + MAP sync in ONE service.
The single X2 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
robot's saved nav maps (pgm+yaml sets and RTAB-Map .db) 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 (AGIBOT X2 — pluggable, configured in .env)
--------------------------------------------------------
backend : X2_SOURCE = auto | http | ros2 | aimrt | none
battery / charging : X2_FIELD_SOC/VOLTAGE/CURRENT/TEMP/SOH/CYCLES
faults / motor temp: X2_FIELD_TEMPS (+ staleness of the last successful read)
position {x,y} : X2_FIELD_X / X2_FIELD_Y
status : derived (charging/moving/idle/offline)
storage : host disk via the read-only /:/host mount (+ data dir size)
maps : MAPS_DIR/**.pgm+yaml and *.db, places under
DATA_DIR/<robot>/places/<map>.json — unchanged from g1
Nothing above is hard-coded: run tools/probe_x2.sh on the robot to discover the
real topics/fields, then set them in .env. No code change, no image rebuild.
Read-only toward the robot: never commands motion. Degrades to heartbeats when
no state source is reachable; --simulate fakes only the state side (the map scan
stays real).
CONFIG — environment (see .env.example). Key vars:
SERVER_URL, DEVICE_TOKEN, SN (required at install), ROBOT_NAME,
X2_SOURCE + X2_TOPIC_*/X2_FIELD_* (state source), POLL_INTERVAL (2),
MAC_INTERFACE (which NIC's MAC to report), 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 platform
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_x2")
# --------------------------------------------------------------------------- #
# 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 / state source
net_interface: str
domain_id: int
mac_interface: str
read_fsm: bool
position_source: str
rosbridge_url: str
low_soc: int
motor_temp_max: float
alert_log_patterns: str
alert_log_cooldown: float
alert_scan_interval: float
alert_backfill_bytes: int
poll_interval: float
telemetry_endpoint: str
# map sync
robot: str
maps_dir: Path
extra_map_dirs: List[Path] # extra roots to scan for pgm+yaml sets (e.g. Nav2/Pudu maps)
web_data_dir: Optional[Path]
legacy_places: Optional[Path]
web_nav3_url: str
map_select: str
map_upload_mode: str
map_endpoint_tmpl: str
map_poll_interval: float
map_max_upload_mb: float
state_dir: Path
# logs + alerts
alert_endpoint: str
logs_endpoint: str
logs_interval: float
# remote dashboard (register the Sanad UI URL for the fleet to embed)
remote_enable: bool
remote_endpoint: str
remote_kind: str
remote_host: str
remote_ports: str
remote_url: str
remote_interval: float
ssh_enable: bool
ssh_user: str
ssh_port: int
control_url: str
control_enable: bool
# project logs (e.g. the robot's Sanad app) shipped alongside agent logs
project_log_container: str
project_log_path: str
project_log_label: str
project_log_backfill: int
project_log_exclude: str
ros_distro: str
# transport
verify_tls: bool
http_timeout: float
@classmethod
def from_env(cls) -> "Config":
server = _env("SERVER_URL").rstrip("/")
token = _env("DEVICE_TOKEN")
missing = [n for n, v in (("SERVER_URL", server), ("DEVICE_TOKEN", token)) if not v]
if missing:
raise SystemExit(f"[config] missing required env: {', '.join(missing)}")
iface = _env("X2_INTERFACE", "eth0")
data_dir = _env("DATA_DIR")
legacy = _env("LEGACY_PLACES")
return cls(
server_url=server,
device_token=token,
sn=_env("SN", "x2_0000"),
name=_env("ROBOT_NAME", "") or _env("SN", "x2_0000"),
brand=_env("ROBOT_BRAND", "agibot"),
robot_type=_env("ROBOT_TYPE", "humanoid"),
model=_env("ROBOT_MODEL", "x2"),
storage_path=_env("STORAGE_PATH", ""),
data_path=_env("STORAGE_DATA_PATH", ""),
net_interface=iface,
domain_id=int(_env("ROS_DOMAIN_ID", "0")),
mac_interface=_env("MAC_INTERFACE", iface),
read_fsm=_env_bool("X2_READ_FSM", False),
position_source=_env("X2_POSITION_SOURCE", "none").lower(),
rosbridge_url=_env("ROSBRIDGE_URL", "ws://127.0.0.1:9090"),
low_soc=int(_env("LOW_SOC", "50")),
motor_temp_max=float(_env("MOTOR_TEMP_MAX", "85")),
# log-driven alerts: "CODE=regex" entries separated by ";;" (regex may
# contain '|'). Scanned against the robot's project logs (sanadr1).
# NOTE: matching is CASE-SENSITIVE (log levels are uppercase); use an
# inline (?i) prefix for case-insensitive text (Gemini messages).
alert_log_patterns=_env("ALERT_LOG_PATTERNS",
"GEMINI_BILLING=(?i)prepayment credits.{0,40}deplet|Please go to AI Studio|Failed to connect to Gemini"
";;ROBOT_ERROR=\\bERROR\\b|\\bCRITICAL\\b|^Traceback"),
alert_log_cooldown=float(_env("ALERT_LOG_COOLDOWN", "300")), # per-signature re-alert gap
alert_scan_interval=float(_env("ALERT_SCAN_INTERVAL", "10")),
alert_backfill_bytes=int(_env("ALERT_BACKFILL_BYTES", str(8 * 1024 * 1024))),
poll_interval=float(_env("POLL_INTERVAL", "2")),
telemetry_endpoint=_env("TELEMETRY_ENDPOINT", "/api/v1/fleet/ingest/telemetry"),
robot=_env("ROBOT", "sanad"),
maps_dir=Path(_env("MAPS_DIR", "/data/maps")),
# colon-separated extra roots (mounted Nav2/Pudu map dirs). Any *.yaml+*.pgm
# set found here is rendered to PNG and uploaded like a slam_toolbox map.
extra_map_dirs=[Path(p) for p in _env("EXTRA_MAP_DIRS", "").split(":") if p.strip()],
web_data_dir=Path(data_dir) if data_dir else None,
legacy_places=Path(legacy) if legacy else None,
web_nav3_url=_env("WEB_NAV3_URL", "").rstrip("/"),
map_select=_env("MAP_SELECT", "all").lower(),
map_upload_mode=_env("MAP_UPLOAD_MODE", "multipart").lower(),
map_endpoint_tmpl=_env("MAP_ENDPOINT", "/api/v1/fleet/ingest/{sn}/map"),
map_poll_interval=float(_env("MAP_POLL_INTERVAL", "30")),
map_max_upload_mb=float(_env("MAP_MAX_UPLOAD_MB", "7")),
state_dir=Path(_env("STATE_DIR", "/data/state")),
alert_endpoint=_env("ALERT_ENDPOINT", "/api/v1/fleet/ingest/{sn}/alert"),
logs_endpoint=_env("LOGS_ENDPOINT", "/api/v1/fleet/ingest/{sn}/logs"),
logs_interval=float(_env("LOGS_INTERVAL", "60")),
remote_enable=_env_bool("REMOTE_ENABLE", True),
remote_endpoint=_env("REMOTE_ENDPOINT", "/api/v1/fleet/ingest/{sn}/remote"),
remote_kind=_env("REMOTE_KIND", "web"),
remote_host=_env("REMOTE_HOST", ""),
remote_ports=_env("REMOTE_PORTS", "8001,8014,8011,8012,8013,8000,8080"),
remote_url=_env("REMOTE_URL", ""), # explicit URL (e.g. a public tunnel) wins
remote_interval=float(_env("REMOTE_INTERVAL", "60")),
ssh_enable=_env_bool("SSH_REGISTER", True),
ssh_user=_env("SSH_USER", ""),
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"
# platform.uname(), not os.uname(): identical .release/.machine on the robot,
# but os.uname() is Unix-only and this must also import on a Windows
# workstation (dry-run / --simulate during development).
u = platform.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. The state source adds live 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"] = platform.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
# --------------------------------------------------------------------------- #
# state source (telemetry side — degrades to heartbeats if unreachable)
# --------------------------------------------------------------------------- #
def _dig(obj: Any, path: str) -> Any:
"""Walk a dotted path over dicts OR message objects.
"pose.pose.position.x" nested attribute / key
"bms.cell_temp[0]" list index
"joints.leg[*].velocity" WILDCARD — collect that field from EVERY
element, returning a flat list
The wildcard matters because real robot payloads publish joints as an array
of objects, not an array of numbers: without it, a velocity/temperature
mapping resolves to a list of dicts and silently yields nothing.
Several wildcards compose ("a[*].b[*].c" flattens both levels).
Returns None for any missing link, so a wrong mapping degrades that one
field to null — never an exception."""
if not path:
return None
cur = obj
for part in path.split("."):
if cur is None:
return None
idx = None
star = False
if part.endswith("]") and "[" in part:
part, _, raw = part[:-1].partition("[")
if raw == "*":
star = True
else:
try:
idx = int(raw)
except ValueError:
return None
def _get(o: Any, name: str) -> Any:
if not name:
return o
return o.get(name) if isinstance(o, dict) else getattr(o, name, None)
if isinstance(cur, list) and part:
# a previous [*] already fanned out — keep mapping across the list
cur = [_get(o, part) for o in cur]
else:
cur = _get(cur, part)
if star:
if cur is None:
return None
try:
cur = list(cur)
except TypeError:
return None
elif idx is not None:
try:
cur = cur[idx]
except Exception:
return None
if isinstance(cur, list) and cur and isinstance(cur[0], list):
cur = [v for sub in cur for v in (sub if isinstance(sub, list) else [sub])]
return cur
def _import_msg(spec: str):
""""sensor_msgs/msg/BatteryState" -> the message class."""
import importlib
parts = [p for p in spec.replace(".", "/").split("/") if p]
return getattr(importlib.import_module(".".join(parts[:-1])), parts[-1])
def _qos(depth: int = 10):
"""Subscription QoS for reading robot telemetry.
Defaults to BEST_EFFORT, which is the only setting that works against BOTH
kinds of publisher: a RELIABLE subscriber gets NOTHING from a BEST_EFFORT
publisher (incompatible), while a BEST_EFFORT subscriber happily reads from
either. The X2 publishes /aima/mc/leg_odometry as BEST_EFFORT, so the rclpy
default (RELIABLE) silently delivers zero messages — the subscription is
created, no error is raised, and the field just stays null forever.
Override with X2_ROS_QOS=reliable if a topic ever requires it."""
from rclpy.qos import HistoryPolicy, QoSProfile, ReliabilityPolicy
want = _env("X2_ROS_QOS", "best_effort").lower()
rel = (ReliabilityPolicy.RELIABLE if want in ("reliable", "rel")
else ReliabilityPolicy.BEST_EFFORT)
return QoSProfile(reliability=rel, history=HistoryPolicy.KEEP_LAST, depth=depth)
def _seq(v: Any) -> List[Any]:
if v is None:
return []
if isinstance(v, (str, bytes)):
return []
return list(v) if hasattr(v, "__iter__") else [v]
class AgiBotSource:
"""AGIBOT X2 state source — the ONE robot-specific class in this agent.
The X2 exposes no fixed topic contract, so NOTHING here is hard-coded:
choose a backend with X2_SOURCE and point it at the real interface with
env vars.
X2_SOURCE=auto (default) http if X2_STATE_URL is set, else ros2 if rclpy
imports, else aimrt if aimrt_py imports, else none
X2_SOURCE=http poll X2_STATE_URL (JSON); map fields with X2_FIELD_*
X2_SOURCE=ros2 subscribe X2_TOPIC_* (types from X2_TYPE_*, defaulting to
the standard sensor_msgs / nav_msgs ones)
X2_SOURCE=aimrt AimRT channels — see _start_aimrt() below
X2_SOURCE=none never read; always heartbeat
Whatever the backend, it fills one six-key snapshot —
{bms, state_age, temps, max_vel, xy, fw} — which is the entire seam between
this robot and the shared telemetry / fault / status / alert pipeline.
Every read is wrapped: a wrong mapping yields null fields and heartbeat mode,
never a crashed loop."""
# Defaults match the STANDARD ROS 2 messages. They are almost certainly
# right for a ros2 backend and almost certainly wrong for a vendor http
# payload — which is exactly why they are env-overridable.
_DEFAULTS = {
"soc": "percentage", # sensor_msgs/BatteryState
"voltage": "voltage",
"current": "current",
"temp": "temperature",
"soh": "",
"cycles": "",
"temps": "", # standard JointState has NO temperatures
"vel": "velocity", # sensor_msgs/JointState
"x": "pose.pose.position.x", # nav_msgs/Odometry
"y": "pose.pose.position.y",
"fsm": "",
"fw": "",
}
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._state_ts = 0.0
self._temps: List[float] = []
self._max_vel = 0.0
self._xy: Optional[Dict[str, float]] = None
self._fw: Dict[str, Any] = {}
self._fsm: Optional[int] = None
self._loco = None
self._stop = False
self.backend = "none"
self.ok = False
# field map: X2_FIELD_SOC, X2_FIELD_VOLTAGE, X2_FIELD_X, … (empty value
# = that field is unavailable on this robot and reports null)
self._map = {k: _env("X2_FIELD_" + k.upper(), d) for k, d in self._DEFAULTS.items()}
self._soc_scale = _env("X2_SOC_SCALE", "auto").lower()
def _f(name: str, default: str) -> float:
try:
return float(_env(name, default) or default)
except ValueError:
log.warning("bad %s — using %s", name, default)
return float(default)
# +1 = positive current means CHARGING (the ROS BatteryState
# convention). Set -1 if the X2 reports the opposite sign.
self._cur_sign = _f("X2_CURRENT_SIGN", "1")
# Unit scaling: the telemetry schema is VOLTS and AMPS. BMS firmware
# commonly publishes mV/mA instead — set 0.001 for those.
self._v_scale = _f("X2_VOLTAGE_SCALE", "1")
self._i_scale = _f("X2_CURRENT_SCALE", "1")
self._start()
# ---------------- backend selection ----------------
def _start(self) -> None:
want = (_env("X2_SOURCE", "auto").lower() or "auto")
order = ["http", "ros2", "aimrt"] if want == "auto" else [want]
for b in order:
try:
if b == "none":
break
if b == "http" and self._start_http():
self.backend = "http"
break
if b == "ros2" and self._start_ros2():
self.backend = "ros2"
break
if b == "aimrt" and self._start_aimrt():
self.backend = "aimrt"
break
except Exception as e:
log.warning("X2 source %r failed to start (%s)", b, e)
if self.backend == "none":
log.warning("no X2 state source active (X2_SOURCE=%s) — telemetry runs in "
"heartbeat mode. Set X2_STATE_URL (http) or X2_TOPIC_* (ros2); "
"run tools/probe_x2.sh on the robot to find the real names.", want)
else:
self.ok = True
log.info("X2 source up: backend=%s", self.backend)
def _start_http(self) -> bool:
"""Poll a vendor state endpoint returning one JSON object."""
url = _env("X2_STATE_URL")
if not url:
return False
self._http_url = url
try:
self._http_period = max(0.2, float(_env("X2_HTTP_INTERVAL", "1") or "1"))
except ValueError:
self._http_period = 1.0
threading.Thread(target=self._http_loop, daemon=True).start()
log.info("X2 http source: %s every %.1fs", url, self._http_period)
return True
def _http_loop(self) -> None:
while not self._stop:
try:
r = requests.get(self._http_url, timeout=4)
if r.ok:
self._ingest_all(r.json())
else:
log.debug("X2 http state: HTTP %s", r.status_code)
except Exception as e:
log.debug("X2 http poll failed: %s", e)
time.sleep(self._http_period)
def _start_ros2(self) -> bool:
"""Subscribe the configured ROS 2 topics. Types are resolved by name, so
a custom vendor message works as long as its package is on PYTHONPATH."""
try:
import rclpy
from rclpy.node import Node
except Exception as e:
log.debug("rclpy unavailable (%s)", e)
return False
wanted = [
(_env("X2_TOPIC_BATTERY", "/battery_state"),
_env("X2_TYPE_BATTERY", "sensor_msgs/msg/BatteryState"), self._ingest_battery),
(_env("X2_TOPIC_JOINTS", "/joint_states"),
_env("X2_TYPE_JOINTS", "sensor_msgs/msg/JointState"), self._ingest_joints),
(_env("X2_TOPIC_ODOM", "/odom"),
_env("X2_TYPE_ODOM", "nav_msgs/msg/Odometry"), self._ingest_odom),
]
if not any(t for t, _, _ in wanted):
return False
rclpy.init(args=None)
node = Node("sanad_api_x2")
n = 0
for topic, spec, cb in wanted:
if not topic or not spec:
continue
try:
cls = _import_msg(spec)
except Exception as e:
log.warning("X2 ros2: cannot import %s for %s (%s) — skipped", spec, topic, e)
continue
node.create_subscription(cls, topic, cb, _qos())
log.info("X2 ros2: subscribed %s (%s)", topic, spec)
n += 1
if not n:
try:
rclpy.shutdown()
except Exception:
pass
return False
self._node = node
threading.Thread(target=lambda: rclpy.spin(node), daemon=True).start()
return True
def _start_aimrt(self) -> bool:
"""AgiBot's own runtime (the X1 stack is built on it).
NOT implemented natively, on purpose: the X2's AimRT channel names and
message definitions are not public, and guessing them would produce an
agent that imports cleanly, builds, posts perfect-looking heartbeats
with battery:null forever, and only reveals the mistake on the robot.
The supported path today is AimRT's ROS 2 plugin: enable it on the robot
and run this agent with X2_SOURCE=ros2 against the bridged topics. If you
get the native channel definitions from AgiBot, implement them here — the
only contract to satisfy is calling self._ingest_battery / _ingest_joints
/ _ingest_odom with the incoming messages."""
try:
import aimrt_py # noqa: F401
except Exception as e:
log.debug("aimrt_py unavailable (%s)", e)
return False
log.warning("aimrt_py is installed, but no native X2 channel binding is built "
"(channel/message definitions are not public). Enable AimRT's ROS 2 "
"plugin and set X2_SOURCE=ros2 with X2_TOPIC_* instead.")
return False
# ---------------- field extraction ----------------
def _num(self, obj: Any, key: str) -> Optional[float]:
v = _dig(obj, self._map.get(key, ""))
try:
return float(v) if v is not None else None
except (TypeError, ValueError):
return None
def _ingest_battery(self, msg: Any) -> None:
try:
soc = self._num(msg, "soc")
if soc is None:
return # no SOC = no battery record; downstream expects an int
if self._soc_scale == "fraction" or (self._soc_scale == "auto" and 0.0 < soc <= 1.0):
soc *= 100.0 # ROS BatteryState.percentage is 0..1; vendors often use 0..100
cur = self._num(msg, "current")
volt = self._num(msg, "voltage")
temp = self._num(msg, "temp")
soh = self._num(msg, "soh")
cyc = self._num(msg, "cycles")
rec = {
"soc": max(0, min(100, int(round(soc)))),
"current_a": round(cur * self._i_scale * self._cur_sign, 2) if cur is not None else 0.0,
"voltage_v": round(volt * self._v_scale, 1) if volt is not None else None,
"temp_c": int(round(temp)) if temp is not None and -40 <= temp <= 150 else None,
"soh": int(round(soh)) if soh is not None else 0,
"cycles": int(round(cyc)) if cyc is not None else 0,
}
with self._lock:
self._bms = rec
self._bms_ts = self._state_ts = time.monotonic()
except Exception:
pass
def _ingest_joints(self, msg: Any) -> None:
try:
temps: List[float] = []
for x in _seq(_dig(msg, self._map.get("temps", ""))):
try:
f = float(x)
except (TypeError, ValueError):
continue
if 0 < f <= 200: # 0 = slot not reporting, same rule as g1
temps.append(f)
max_vel = 0.0
for x in _seq(_dig(msg, self._map.get("vel", ""))):
try:
max_vel = max(max_vel, abs(float(x)))
except (TypeError, ValueError):
pass
with self._lock:
self._temps = temps
self._max_vel = max_vel
self._state_ts = time.monotonic()
except Exception:
pass
def _ingest_odom(self, msg: Any) -> None:
try:
x, y = self._num(msg, "x"), self._num(msg, "y")
if x is None or y is None:
return
with self._lock:
self._xy = {"x": round(x, 3), "y": round(y, 3)}
self._state_ts = time.monotonic()
except Exception:
pass
def _ingest_all(self, d: Any) -> None:
"""One combined payload (http backend) feeds all three extractors."""
self._ingest_battery(d)
self._ingest_joints(d)
self._ingest_odom(d)
try:
fsm = _dig(d, self._map.get("fsm", ""))
fw = _dig(d, self._map.get("fw", ""))
with self._lock:
if fsm is not None:
try:
self._fsm = int(fsm)
except (TypeError, ValueError):
pass
if isinstance(fw, dict):
self._fw.update({str(k): str(v) for k, v in fw.items()})
except Exception:
pass
# ---------------- the snapshot contract ----------------
def snapshot(self) -> Dict[str, Any]:
with self._lock:
now = time.monotonic()
return {
"bms": dict(self._bms) if self._bms else None,
"state_age": (now - self._state_ts) if self._state_ts else None,
"temps": list(self._temps),
"max_vel": self._max_vel,
"xy": dict(self._xy) if self._xy else None,
"fw": dict(self._fw),
}
def fsm_id(self) -> Optional[int]:
"""Only the http backend can carry one (X2_FIELD_FSM); there is no
known X2 loco RPC. Read-only either way — never commands motion."""
with self._lock:
return self._fsm
# AGIBOT X2 publishes no documented loco FSM id scheme. control.mode comes
# from the read-only CONTROL_STATUS_URL when the robot exposes one; fill
# these in once the real ids are known (empty = mode reported as unknown).
_FSM_STATUS: Dict[int, str] = {}
# Control-panel mode labels + the switchable set (fsm_id -> friendly mode).
_CONTROL_MODES: Dict[int, str] = {}
_CONTROL_SWITCHABLE: List[str] = []
_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 robot control, 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
class Ros2Position:
"""Position from a ROS 2 odometry topic, INDEPENDENT of X2_SOURCE.
The X2's battery comes cleanly from the Control Dash over http, but that
payload carries no odometry — while ROS publishes it on
/aima/mc/leg_odometry (nav_msgs/Odometry). Those are two different
transports, so position is its own source rather than part of the state
backend: X2_SOURCE=http can run with X2_POSITION_SOURCE=ros2 at the
same time.
Enable with X2_POSITION_SOURCE=ros2 (topic: X2_TOPIC_ODOM, type:
X2_TYPE_ODOM, fields: X2_FIELD_X / X2_FIELD_Y). Requires rclpy — the
systemd unit must source the ROS overlay first. If rclpy is missing or the
topic never publishes, position simply stays null; nothing else is affected."""
def __init__(self, cfg: Config):
self.cfg = cfg
self._xy: Optional[Dict[str, float]] = None
self._lock = threading.Lock()
self.ok = False
# read here, NOT as class attributes: the class body runs at import,
# which is before main() calls _load_dotenv().
self._px = _env("X2_FIELD_X", "") or "pose.pose.position.x"
self._py = _env("X2_FIELD_Y", "") or "pose.pose.position.y"
topic = _env("X2_TOPIC_ODOM", "/odom")
spec = _env("X2_TYPE_ODOM", "nav_msgs/msg/Odometry")
if not topic or not spec:
log.warning("X2_POSITION_SOURCE=ros2 but X2_TOPIC_ODOM/X2_TYPE_ODOM is empty")
return
try:
import rclpy
from rclpy.node import Node
except Exception as e:
log.warning("rclpy unavailable (%s) — position stays null. The unit must "
"source the ROS overlay before starting the agent.", e)
return
try:
cls = _import_msg(spec)
except Exception as e:
log.warning("cannot import %s (%s) — position stays null", spec, e)
return
try:
# AgiBotSource may already have started rclpy when X2_SOURCE=ros2.
if not rclpy.ok():
rclpy.init(args=None)
node = Node("sanad_api_x2_pos")
node.create_subscription(cls, topic, self._on_odom, _qos())
self._node = node
threading.Thread(target=lambda: rclpy.spin(node), daemon=True).start()
self.ok = True
log.info("X2 position: ros2 %s (%s)", topic, spec)
except Exception as e:
log.warning("ros2 position init failed (%s) — position stays null", e)
def _on_odom(self, msg: Any) -> None:
try:
x = _dig(msg, self._px)
y = _dig(msg, self._py)
if x is None or y is None:
return
with self._lock:
self._xy = {"x": round(float(x), 3), "y": round(float(y), 3)}
except Exception:
pass
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"]
roots += list(cfg.extra_map_dirs) # Nav2/Pudu map dirs mounted via EXTRA_MAP_DIRS
seen: set = set()
out: List[MapArtifact] = []
for root in roots:
if not root.exists():
continue
for y in sorted(root.glob("*.yaml")):
meta = _parse_map_yaml(y)
img = meta.get("image", "")
pgm = (y.parent / img) if img else y.with_suffix(".pgm")
if not pgm.exists():
pgm = y.with_suffix(".pgm")
if not pgm.exists():
continue # yaml without a raster — not a map set
rp = str(y.resolve())
if rp in seen:
continue
seen.add(rp)
files: Dict[str, Path] = {"yaml": y, "pgm": pgm}
for ext in ("posegraph", "data"):
p = y.with_suffix("." + ext)
if p.exists():
files[ext] = p
size = sum(p.stat().st_size for p in files.values())
mtime = max(int(p.stat().st_mtime) for p in files.values())
out.append(MapArtifact(path=y, name=y.name, stem=y.stem,
size=size, mtime=mtime,
fmt="slam_toolbox", files=files))
# Pudu/Nav2 converter emits a plain map + a keepout-BAKED twin (obstacles baked
# in for Foxy, which has no KeepoutFilter). The baked one is the deploy map — drop
# the redundant plain twin so the fleet server gets one canonical map, not two.
baked = {m.stem[: -len("_keepout_baked")] for m in out if m.stem.endswith("_keepout_baked")}
out = [m for m in out if m.stem not in baked]
return out
def _map_key(stem: str) -> str:
stem = Path(stem).name
if stem.endswith(".db"):
stem = stem[:-3]
return "".join(c for c in stem if c.isalnum() or c in "_-.")
def _read_json(path: Path, default: Any) -> Any:
try:
return json.loads(path.read_text() or "")
except Exception:
return default
def _yaw_from_pose(pose: Dict[str, Any]) -> float:
if "qw" in pose or "qz" in pose:
qx = float(pose.get("qx", 0.0)); qy = float(pose.get("qy", 0.0))
qz = float(pose.get("qz", 0.0)); qw = float(pose.get("qw", 1.0))
return math.atan2(2.0 * (qw * qz + qx * qy),
1.0 - 2.0 * (qy * qy + qz * qz))
return float(pose.get("yaw", 0.0))
def _places_files_for(cfg: Config, stem: str) -> List[Path]:
out: List[Path] = []
key = _map_key(stem)
if cfg.web_data_dir:
out.append(cfg.web_data_dir / cfg.robot / "places" / f"{key}.json")
if cfg.legacy_places:
out.append(cfg.legacy_places)
return out
def load_points(cfg: Config, stem: str) -> List[Dict[str, Any]]:
for pf in _places_files_for(cfg, stem):
data = _read_json(pf, None) if pf.exists() else None
if isinstance(data, dict) and data:
pts: List[Dict[str, Any]] = []
for name, pose in data.items():
if not isinstance(pose, dict):
continue
try:
pts.append({
"name": name,
"type": str(pose.get("type", "waypoint")),
"x": float(pose["x"]),
"y": float(pose["y"]),
"yaw": round(_yaw_from_pose(pose), 4),
})
except (KeyError, TypeError, ValueError):
continue
return pts
return []
def discover_maps(cfg: Config) -> List[MapArtifact]:
roots = [cfg.maps_dir / cfg.robot, cfg.maps_dir]
meta: Dict[str, Any] = {}
meta_file = cfg.maps_dir / cfg.robot / "maps_meta.json"
if meta_file.exists():
meta = _read_json(meta_file, {}) or {}
seen: set = set()
out: List[MapArtifact] = []
for root in roots:
if not root.exists():
continue
for p in sorted(root.glob("*.db")):
rp = str(p.resolve())
if rp in seen:
continue
seen.add(rp)
st = p.stat()
out.append(MapArtifact(
path=p, name=p.name, stem=p.stem,
size=st.st_size, mtime=int(st.st_mtime),
description=(meta.get(p.name) or {}).get("description", ""),
))
# slam_toolbox map sets (office.yaml + office.pgm …) live alongside
out.extend(_discover_slam_sets(cfg))
out.sort(key=lambda m: m.mtime, reverse=True)
return out
def _active_map_name(cfg: Config) -> Optional[str]:
if not cfg.web_nav3_url:
return None
try:
r = requests.get(cfg.web_nav3_url + "/api/status",
headers={"X-Robot-Name": cfg.robot},
timeout=min(cfg.http_timeout, 5))
r.raise_for_status()
am = (r.json() or {}).get("active_map")
return _map_key(am) if am else None
except requests.RequestException:
return None
def select_maps(cfg: Config, maps: List[MapArtifact]) -> List[MapArtifact]:
if not maps:
return []
if cfg.map_select == "newest":
return maps[:1]
if cfg.map_select == "active":
active = _active_map_name(cfg)
if active:
picked = [m for m in maps if _map_key(m.stem) == active]
if picked:
return picked
return maps[:1]
return maps # "all"
def _state_file(cfg: Config) -> Path:
return cfg.state_dir / "uploaded.json"
def load_state(cfg: Config) -> Dict[str, str]:
return _read_json(_state_file(cfg), {}) if _state_file(cfg).exists() else {}
def save_state(cfg: Config, state: Dict[str, str]) -> None:
try:
cfg.state_dir.mkdir(parents=True, exist_ok=True)
_state_file(cfg).write_text(json.dumps(state, indent=2))
except Exception as e:
log.warning("could not persist map state: %s", e)
def build_meta(cfg: Config, m: MapArtifact) -> Dict[str, Any]:
return {
"sn": cfg.sn,
"name": m.stem,
"file": m.name,
"format": m.fmt,
"size_bytes": m.size,
"sha256": m.sha256,
"mtime": m.mtime,
"description": m.description,
"points": m.points,
}
def _upload_slam_map(cfg: Config, m: MapArtifact, session: requests.Session) -> bool:
"""slam_toolbox map → the spec's image JSON: PNG (from the pgm) + resolution
+ origin + width/height + points. This is what the dashboard renders."""
url = cfg.map_url()
ymeta = _parse_map_yaml(m.files["yaml"])
pgm = _read_pgm(m.files["pgm"])
if pgm is None:
log.error("map %s: cannot parse %s (not binary P5?)", m.stem, m.files["pgm"].name)
return False
body = build_meta(cfg, m)
body.update({
"resolution": ymeta.get("resolution"),
"origin": ymeta.get("origin"),
"width": pgm["width"],
"height": pgm["height"],
"image_base64": _pgm_to_png_b64(pgm),
})
try:
resp = session.post(url, json=body, headers=cfg.auth_headers(),
timeout=cfg.http_timeout, verify=cfg.verify_tls)
except requests.RequestException as e:
log.error("map upload %s FAILED (transport): %s", m.stem, e)
return False
if not resp.ok:
log.error("map upload %s FAILED: HTTP %s %s", m.stem, resp.status_code, resp.text[:300])
return False
log.info("map uploaded: %s (slam_toolbox %dx%d @ %sm, %d points) -> HTTP %s",
m.stem, pgm["width"], pgm["height"], ymeta.get("resolution"),
len(m.points), resp.status_code)
return True
def upload_map(cfg: Config, m: MapArtifact, session: requests.Session) -> bool:
if m.fmt == "slam_toolbox":
return _upload_slam_map(cfg, m, session)
url = cfg.map_url()
meta = build_meta(cfg, m)
try:
if cfg.map_upload_mode == "base64json":
body = dict(meta)
body["db_base64"] = base64.b64encode(m.path.read_bytes()).decode("ascii")
resp = session.post(url, json=body, headers=cfg.auth_headers(),
timeout=cfg.http_timeout, verify=cfg.verify_tls)
else: # multipart (default)
with m.path.open("rb") as fh:
files = {"db": (m.name, fh, "application/octet-stream")}
data = {"meta": json.dumps(meta)}
resp = session.post(url, files=files, data=data,
headers=cfg.auth_headers(),
timeout=cfg.http_timeout, verify=cfg.verify_tls)
except requests.RequestException as e:
log.error("map upload %s FAILED (transport): %s", m.name, e)
return False
if not resp.ok:
log.error("map upload %s FAILED: HTTP %s %s", m.name, resp.status_code, resp.text[:300])
return False
log.info("map uploaded: %s (%.2f MB, %d points) -> HTTP %s",
m.name, m.size / 1024 / 1024, len(m.points), resp.status_code)
return True
# Shared map status — SHOWN in every telemetry post ("map" field).
_MAP_STATUS_LOCK = threading.Lock()
_MAP_STATUS: Dict[str, Any] = {
"uploaded": False, "state": "pending", "maps_found": 0,
"last_map": None, "error": None, "checked_ts": None,
}
def _set_map_status(**kw: Any) -> None:
with _MAP_STATUS_LOCK:
_MAP_STATUS.update(kw)
_MAP_STATUS["checked_ts"] = int(time.time())
def get_map_status() -> Dict[str, Any]:
with _MAP_STATUS_LOCK:
return dict(_MAP_STATUS)
def map_sync_once(cfg: Config, session: requests.Session,
force: bool = False, dry_run: bool = False) -> int:
"""One map pass: scan the Sanad dashboard maps and upload anything new.
Always updates the shared map status (visible in telemetry)."""
try:
maps = select_maps(cfg, discover_maps(cfg))
except Exception as e:
_set_map_status(state="failed", uploaded=False, error=f"map scan failed: {e}")
return 0
if not maps:
_set_map_status(state="no_map", uploaded=False, maps_found=0, last_map=None,
error=f"no saved map found in Sanad dashboard "
f"(maps_dir={cfg.maps_dir}, robot={cfg.robot})")
return 0
state = load_state(cfg)
uploaded = failed = unstable = too_large = current = 0
last_err: Optional[str] = None
now = time.time()
for m in maps:
prev = state.get(str(m.path.resolve()))
if not force and prev == m.fingerprint():
current += 1
continue # already uploaded this exact content — one-time rule
# stability guard: a db modified in the last 120 s is still being
# written (active mapping) — wait until it settles before uploading
if not force and (now - m.mtime) < 120:
log.info("map %s still changing (mapping in progress) — waiting to settle", m.name)
unstable += 1
continue
# server rejects bodies over ~8 MB (client_max_body_size) — don't burn
# bandwidth on uploads that will 413. Raster (slam_toolbox) maps are tiny.
if m.fmt != "slam_toolbox" and (m.size / 1048576) > cfg.map_max_upload_mb:
log.warning("map %s is %.0f MB — exceeds server upload cap (~%.0f MB), skipping "
"(export a raster map or raise the server limit)",
m.name, m.size / 1048576, cfg.map_max_upload_mb)
too_large += 1
continue
m.sha256 = (_sha256_set(list(m.files.values()))
if m.fmt == "slam_toolbox" else _sha256(m.path))
m.points = load_points(cfg, m.stem)
if dry_run:
log.info("[dry-run] would upload map %s (%.2f MB, %d points)",
m.name, m.size / 1024 / 1024, len(m.points))
continue
if upload_map(cfg, m, session):
state[str(m.path.resolve())] = m.fingerprint()
save_state(cfg, state)
uploaded += 1
else:
failed += 1
last_err = f"upload failed for {m.name} (see agent log)"
if failed:
_set_map_status(state="failed", uploaded=False, maps_found=len(maps),
last_map=maps[0].stem, error=last_err)
elif uploaded or current:
# at least one map is on the server (just now or previously); note skips
note = None
if too_large:
note = f"{too_large} map(s) skipped: exceed server upload cap (~{cfg.map_max_upload_mb:.0f} MB)"
elif unstable:
note = "newer map still being written (mapping in progress)"
_set_map_status(state="uploaded", uploaded=True, maps_found=len(maps),
last_map=maps[0].stem, error=note)
elif too_large:
_set_map_status(state="failed", uploaded=False, maps_found=len(maps),
last_map=maps[0].stem,
error=f"map exceeds server upload cap (~{cfg.map_max_upload_mb:.0f} MB) — "
"export a raster map or raise the server limit")
elif unstable:
# newest content is still being written (active mapping) — be honest
_set_map_status(state="pending", uploaded=False, maps_found=len(maps),
last_map=maps[0].stem,
error="map still being written (mapping in progress) — "
"will upload when it settles")
else:
_set_map_status(state="uploaded", uploaded=True, maps_found=len(maps),
last_map=maps[0].stem, error=None)
return uploaded
def map_loop(cfg: Config, session: requests.Session) -> None:
while True:
try:
map_sync_once(cfg, session)
except Exception as e:
log.exception("map pass failed: %s", e)
time.sleep(cfg.map_poll_interval)
# --------------------------------------------------------------------------- #
# logs + alerts (spec: POST /{sn}/logs periodically, POST /{sn}/alert on events)
# --------------------------------------------------------------------------- #
class _RingLogHandler(logging.Handler):
"""Buffers the agent's own log lines so they can be shipped to the server."""
def __init__(self, maxlen: int = 400):
super().__init__(level=logging.INFO)
from collections import deque
self._buf: Any = deque(maxlen=maxlen)
self._blk = threading.Lock()
def emit(self, record: logging.LogRecord) -> None:
try:
with self._blk:
self._buf.append(self.format(record))
except Exception:
pass
def drain(self) -> List[str]:
with self._blk:
lines = list(self._buf)
self._buf.clear()
return lines
def requeue(self, lines: List[str]) -> None:
"""Put unshipped lines back (front of the ring) so they retry next cycle
instead of being lost — bounded by maxlen, oldest evicted first."""
with self._blk:
self._buf.extendleft(reversed(lines))
_LOG_RING = _RingLogHandler()
# shipped-status shown in every telemetry post ("logs" / "alerts" fields)
_LOGS_STAT: Dict[str, Any] = {"last_sent": None, "lines_sent": 0, "ok": None}
_ALERTS_STAT: Dict[str, Any] = {"sent": 0, "last": None, "last_time": None, "ok": None}
# start times ("started_at" = this run, "last_start" = previous run)
_STARTED: Dict[str, Any] = {"now": None, "prev": None, "mono": time.monotonic()}
def _init_start_times(cfg: Config) -> None:
"""Record this agent start; remember the previous one (persisted in STATE_DIR)."""
f = cfg.state_dir / "agent_state.json"
prev = (_read_json(f, {}) or {}).get("started_at")
now_s = _now_str()
try:
cfg.state_dir.mkdir(parents=True, exist_ok=True)
f.write_text(json.dumps({"started_at": now_s}))
except Exception as e:
log.debug("could not persist start time: %s", e)
_STARTED.update(now=now_s, prev=prev, mono=time.monotonic())
class ProjectLogTail:
"""Tails the robot's main PROJECT logs (e.g. the sanadr1 / sanad-p4 app)
and feeds them into the shipped log lines, labeled "[<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 <user>@<ip>". Skipped when the login user
# isn't known: registering a guessed username sends the fleet UI a command
# that silently fails for whoever tries it.
if cfg.ssh_enable and cfg.ssh_user:
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 <user>@<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: Dict[str, float] = {} # fault CODE -> monotonic ts of last alert
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 _fault_code(f: str) -> str:
"""Dedup key for a fault string: the CODE before the first ':'."""
return (f.split(":", 1)[0].strip() or f)
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.
Deduped on the fault CODE, never on the whole string: every fault embeds a
LIVE number ("battery 49%", "motor temp 87C", "no robot state for 12s") that
changes almost every tick, so string-dedup re-fired the same fault every
POLL_INTERVAL for as long as it lasted — COMMS_STALE counts up by 1 s
forever, which meant an alert POST every 2 s until the state source came back.
A code alerts on its rising edge, then at most once per ALERT_LOG_COOLDOWN
while it persists (so a long outage still re-asserts, but doesn't flood).
Clearing the fault drops the code, so the next occurrence alerts again."""
now = time.monotonic()
current: Dict[str, str] = {}
for f in faults:
current.setdefault(_fault_code(f), f)
for code in [c for c in _ALERT_SEEN if c not in current]:
del _ALERT_SEEN[code] # cleared — next occurrence is a rising edge again
for code, text in sorted(current.items()):
last = _ALERT_SEEN.get(code)
if last is None or (now - last) > cfg.alert_log_cooldown:
_ALERT_SEEN[code] = now
_post_alert(cfg, session, text)
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("state_age") is not None and snap["state_age"] > 3.0:
faults.append(f"COMMS_STALE: no robot state for {snap['state_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("state_age") is not None and snap["state_age"] <= 3.0
if not alive and bms is None:
return "offline"
if charging:
return "charging"
if snap.get("max_vel", 0.0) > 0.15:
return "moving"
return base or "idle"
def build_telemetry(cfg: Config, mac: str, reader: Optional[AgiBotSource],
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},
"state_age": 0.1, "temps": [sim.get("temp", 45)], "max_vel": sim.get("max_vel", 0.0),
"xy": sim.get("position")}
fsm = sim.get("fsm")
else:
snap = reader.snapshot() if reader else {"bms": None, "state_age": None, "temps": [], "max_vel": 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. x2_10)
"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_vel": 0.4 if moving else 0.0, "fsm": None,
"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="AGIBOT X2 fleet agent: telemetry + map sync")
ap.add_argument("--simulate", action="store_true", help="synthetic robot state (map scan stays real)")
ap.add_argument("--once", action="store_true", help="one map pass + one telemetry post, then exit")
ap.add_argument("--map-only", action="store_true", help="upload discovered maps once (no state source, no telemetry), then exit")
ap.add_argument("--dry-run", action="store_true", help="build payloads, never POST")
ap.add_argument("--force", action="store_true", help="re-upload maps even if unchanged")
ap.add_argument("--list", action="store_true", help="list discovered maps and exit")
ap.add_argument("--interval", type=float, default=None, help="override telemetry POLL_INTERVAL")
ap.add_argument("-v", "--verbose", action="store_true")
args = ap.parse_args(argv)
logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO,
format="%(asctime)s %(levelname)s %(name)s: %(message)s")
# buffer our own log lines for shipping to /{sn}/logs
_LOG_RING.setFormatter(logging.Formatter("%(asctime)s %(levelname)s %(name)s: %(message)s"))
logging.getLogger().addHandler(_LOG_RING)
_load_dotenv()
cfg = Config.from_env()
if args.interval is not None:
cfg.poll_interval = args.interval
if args.list:
cmd_list(cfg)
return 0
_init_start_times(cfg)
global _PROJECT_TAIL, _LOG_ALERTS
_PROJECT_TAIL = ProjectLogTail(cfg)
_LOG_ALERTS = LogAlertScanner(cfg) # error/billing alerts from the project logs
mac = read_mac(cfg.mac_interface)
log.info("sanad_api_x2 — 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.net_interface,
cfg.position_source, cfg.maps_dir, cfg.map_poll_interval,
" [SIMULATE]" if args.simulate else "")
session = requests.Session()
if args.map_only:
# map upload only — search all map roots (incl. EXTRA_MAP_DIRS) and ship,
# without opening a state source or posting telemetry (won't disturb a live feed).
map_sync_once(cfg, session, force=args.force, dry_run=args.dry_run)
return 0
reader = None
pos = None
if not args.simulate:
reader = AgiBotSource(cfg)
if cfg.position_source == "rosbridge":
pos = RosbridgePosition(cfg)
elif cfg.position_source == "ros2":
# independent of X2_SOURCE: the Control Dash has no odometry, ROS does
pos = Ros2Position(cfg)
time.sleep(1.0)
tick = 0
def one_telemetry() -> None:
nonlocal tick
sim = _sim_state(tick) if args.simulate else None
payload = build_telemetry(cfg, mac, reader, pos, sim=sim)
if args.dry_run:
log.info("[dry-run] %s", json.dumps(payload))
else:
post_telemetry(cfg, payload, session)
send_alerts(cfg, session, payload.get("faults") or [])
tick += 1
if args.once or args.dry_run:
# one map pass first so the telemetry "map" field reflects it
map_sync_once(cfg, session, force=args.force, dry_run=args.dry_run)
if cfg.remote_enable and not args.dry_run:
register_remote(cfg, session)
elif cfg.remote_enable:
d = discover_dashboard(cfg)
if d:
_REMOTE_STAT.update(url=d["url"], port=d["port"], kind=cfg.remote_kind, ok=None)
if args.once:
one_telemetry()
ship_logs(cfg, session)
return 0
for _ in range(3):
one_telemetry()
time.sleep(min(cfg.poll_interval, 1.0))
return 0
# loop mode: map sync + log shipping + remote registration in bg threads
threading.Thread(target=map_loop, args=(cfg, session), daemon=True).start()
threading.Thread(target=logs_loop, args=(cfg, session), daemon=True).start()
threading.Thread(target=alert_scan_loop, args=(cfg, session), daemon=True).start()
if cfg.remote_enable:
threading.Thread(target=remote_loop, args=(cfg, session), daemon=True).start()
log.info("telemetry every %.1fs; map check every %.0fs; logs every %.0fs (Ctrl-C to stop)",
cfg.poll_interval, cfg.map_poll_interval, cfg.logs_interval)
while True:
try:
one_telemetry()
except Exception as e:
log.exception("telemetry tick failed: %s", e)
try:
time.sleep(cfg.poll_interval)
except KeyboardInterrupt:
log.info("stopped")
return 0
if __name__ == "__main__":
sys.exit(main())