2200 lines
90 KiB
Python
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())
|