engineai_fleet/agent/sanad_api_eng.py
2026-08-27 16:24:12 +04:00

2679 lines
111 KiB
Python

#!/usr/bin/env python3
"""sanad_api_eng — EngineAI PM01 fleet agent: TELEMETRY + MAP sync in ONE service.
Same concept, same payload schema and same endpoints as the AGIBOT X2 agent
(agi_fleet/agent/sanad_api_x2.py) — retargeted to the EngineAI PM01 ("T800")
and extended to post the SAME data to MORE THAN ONE fleet server.
1. TELEMETRY — every ~2 s POSTs the robot's live status to EVERY enabled
fleet server:
POST {SERVER_URL}/api/v1/fleet/ingest/telemetry
{ sn, name, mac, brand, type, model, software, firmware, battery,
charging, battery_detail, motor_temp, storage, status, position,
control, faults, map, logs, project_logs, remote, alerts, time,
started_at, last_start, uptime_s, ts }
2. MAP — a background loop (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 PER SERVER (re-upload only if its content changes):
POST {SERVER_URL}/api/v1/fleet/ingest/{sn}/map
The map result is SHOWN inside every telemetry post as the "map" field, so
the server always sees whether the map made it — and why not.
3. LOGS / ALERTS / REMOTE — /{sn}/logs, /{sn}/alert, /{sn}/remote, all
fanned out to every enabled server exactly like telemetry.
MULTI-SERVER (the one structural difference from the X2 agent)
---------------------------------------------------------------
The X2 agent posts to a single SERVER_URL. This one keeps a LIST of endpoints,
each with its OWN token, and sends the identical payload to each:
SERVER_URL / DEVICE_TOKEN (primary)
SERVER_2_URL / SERVER_2_TOKEN / SERVER_2_ENABLE
SERVER_3_URL / SERVER_3_TOKEN / SERVER_3_ENABLE
A server that is unreachable, unauthorised or slow NEVER blocks the others:
each POST is independent and its result is reported per-server in the log line
and in the telemetry status blocks. Turning a second server on is an .env edit
plus a restart — never a code change.
DATA SOURCES (EngineAI PM01 — pluggable, configured in .env)
-------------------------------------------------------------
backend : ENG_SOURCE = auto | ros2 | http | none
battery / charging : /hardware/power_info interface_protocol/msg/PowerInfo
motor temperatures : /hardware/motor_debug interface_protocol/msg/MotorDebug
locomotion (moving): /hardware/joint_state interface_protocol/msg/JointState
motion mode : /motion/motion_state interface_protocol/msg/MotionState
position {x,y} : ENG_POSITION_SOURCE (none|ros2|rosbridge|http) — the PM01
publishes NO odometry, so this is null until a nav stack
is running; see docs/PM01_INTERFACE.md
storage / OS / MAC : read directly from the host
maps : MAPS_DIR/**.pgm+yaml and *.db, places under
DATA_DIR/<robot>/places/<map>.json
Nothing above is hard-coded: run tools/probe_eng.sh on the robot to discover the
real topics/fields, then set them in .env.
Read-only toward the robot: it subscribes to four topics and reads an HTTP
status page. It NEVER publishes, and never commands motion.
CLI: --simulate | --once | --dry-run | --force (map re-upload) | --list | -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_eng")
# --------------------------------------------------------------------------- #
# env helpers
# --------------------------------------------------------------------------- #
def _load_dotenv(path: str = "") -> None:
"""Load KEY=VALUE lines from .env (next to this file unless overridden)."""
p = Path(path) if path else Path(__file__).resolve().parent / ".env"
if not p.exists():
p = Path(".env")
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 _env_float(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)
def _now_str() -> str:
"""Full local date+time with UTC offset (e.g. 2026-08-27 19:20:33+04:00).
TZ_OFFSET_HOURS (default +4, Dubai) keeps clocks honest without tzdata."""
off = _env_float("TZ_OFFSET_HOURS", "4")
tz = _dt.timezone(_dt.timedelta(hours=off))
return _dt.datetime.now(tz).isoformat(sep=" ", timespec="seconds")
# --------------------------------------------------------------------------- #
# fleet endpoints — the SAME payload goes to every enabled server
# --------------------------------------------------------------------------- #
def _host_label(url: str) -> str:
"""'https://eco-dev.yslootahrobotics.com' -> 'eco-dev' (for log lines)."""
host = url.split("://", 1)[-1].split("/", 1)[0]
return host.split(".", 1)[0] or host
@dataclass
class Endpoint:
"""One fleet server: its base URL and its OWN device token.
Tokens are per-server on purpose. eco and eco-dev maintain independent
token stores, so a token minted on one is rejected (401) by the other —
sharing a single DEVICE_TOKEN across both silently 401s forever on
whichever server did not issue it."""
name: str
url: str
token: str
enabled: bool = True
def headers(self) -> Dict[str, str]:
return {"Authorization": f"Bearer {self.token}"}
def at(self, tmpl: str, sn: str = "") -> str:
return self.url + (tmpl.format(sn=sn) if "{sn}" in tmpl else tmpl)
def _endpoints_from_env() -> List[Endpoint]:
"""SERVER_URL/DEVICE_TOKEN plus SERVER_<N>_URL/SERVER_<N>_TOKEN for N=2..5.
A slot with no URL is skipped entirely. A slot with a URL but no token is
kept but DISABLED with a loud warning — that is the "server configured,
token not issued yet" state, and it must be visible rather than silent."""
out: List[Endpoint] = []
primary_url = _env("SERVER_URL").rstrip("/")
primary_tok = _env("DEVICE_TOKEN")
if primary_url:
out.append(Endpoint(name=_env("SERVER_NAME", "") or _host_label(primary_url),
url=primary_url, token=primary_tok,
enabled=bool(primary_tok) and _env_bool("SERVER_ENABLE", True)))
for n in range(2, 6):
url = _env(f"SERVER_{n}_URL").rstrip("/")
if not url:
continue
tok = _env(f"SERVER_{n}_TOKEN")
if not tok and _env_bool(f"SERVER_{n}_SHARE_TOKEN", False):
tok = primary_tok
out.append(Endpoint(name=_env(f"SERVER_{n}_NAME", "") or _host_label(url),
url=url, token=tok,
enabled=bool(tok) and _env_bool(f"SERVER_{n}_ENABLE", True)))
return out
# --------------------------------------------------------------------------- #
# config
# --------------------------------------------------------------------------- #
@dataclass
class Config:
endpoints: List[Endpoint]
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
position_source: str
rosbridge_url: str
low_soc: int
motor_temp_max: float
mos_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]
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
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 (the robot's Sanad app container)
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":
eps = _endpoints_from_env()
if not eps:
raise SystemExit("[config] missing required env: SERVER_URL (and DEVICE_TOKEN)")
for e in eps:
if not e.token:
log.warning("server %r (%s) has NO token — DISABLED. Set its "
"SERVER_*_TOKEN in .env and restart to enable it.", e.name, e.url)
elif not e.enabled:
log.warning("server %r (%s) is configured but switched OFF "
"(SERVER_*_ENABLE=0)", e.name, e.url)
if not any(e.enabled for e in eps):
raise SystemExit("[config] no fleet server is enabled — set DEVICE_TOKEN "
"(and/or SERVER_<N>_TOKEN) in .env")
iface = _env("ENG_INTERFACE", "eth0")
data_dir = _env("DATA_DIR")
legacy = _env("LEGACY_PLACES")
return cls(
endpoints=eps,
sn=_env("SN", "pm01_0000"),
name=_env("ROBOT_NAME", "") or _env("SN", "pm01_0000"),
brand=_env("ROBOT_BRAND", "engineai"),
robot_type=_env("ROBOT_TYPE", "humanoid"),
model=_env("ROBOT_MODEL", "pm01"),
storage_path=_env("STORAGE_PATH", ""),
data_path=_env("STORAGE_DATA_PATH", ""),
net_interface=iface,
domain_id=int(_env("ROS_DOMAIN_ID", "69") or "69"),
mac_interface=_env("MAC_INTERFACE", iface),
position_source=_env("ENG_POSITION_SOURCE", "none").lower(),
rosbridge_url=_env("ROSBRIDGE_URL", "ws://127.0.0.1:9090"),
low_soc=int(_env("LOW_SOC", "50") or "50"),
motor_temp_max=_env_float("MOTOR_TEMP_MAX", "85"),
mos_temp_max=_env_float("MOS_TEMP_MAX", "100"),
# log-driven alerts: "CODE=regex" entries separated by ";;" (regex may
# contain '|'). NOTE: matching is CASE-SENSITIVE (log levels are
# uppercase); use an inline (?i) prefix for case-insensitive text.
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=_env_float("ALERT_LOG_COOLDOWN", "300"),
alert_scan_interval=_env_float("ALERT_SCAN_INTERVAL", "10"),
alert_backfill_bytes=int(_env("ALERT_BACKFILL_BYTES", str(8 * 1024 * 1024))),
poll_interval=_env_float("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")),
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=_env_float("MAP_POLL_INTERVAL", "30"),
map_max_upload_mb=_env_float("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=_env_float("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", "8014,8001,9002,9003,8000,8080"),
remote_url=_env("REMOTE_URL", ""),
remote_interval=_env_float("REMOTE_INTERVAL", "60"),
ssh_enable=_env_bool("SSH_REGISTER", True),
ssh_user=_env("SSH_USER", ""),
ssh_port=int(_env("SSH_PORT", "22") or "22"),
control_url=_env("CONTROL_STATUS_URL", ""),
control_enable=_env_bool("CONTROL_ENABLE", False),
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") or "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=_env_float("HTTP_TIMEOUT", "30"),
)
def live(self) -> List[Endpoint]:
return [e for e in self.endpoints if e.enabled]
# --------------------------------------------------------------------------- #
# multi-server POST fan-out
#
# Every ingest call goes to EVERY enabled server. One server failing (down,
# 401, slow) must never stop the others or raise: each POST is independent and
# its outcome is recorded per server name.
# --------------------------------------------------------------------------- #
def _brief(text: str, limit: int = 160) -> str:
"""One-line summary of a server error body.
Error pages are HTML: dumping them raw breaks the log into a dozen lines and
those lines are then shipped to /{sn}/logs, so a single 401 becomes a dozen
junk entries per report. Collapse to one line and strip tags."""
import re
t = re.sub(r"<[^>]+>", " ", text or "")
t = re.sub(r"\s+", " ", t).strip()
return t[:limit]
# Repeated-failure throttle. A server that is down or holding a bad token would
# otherwise emit one ERROR every POLL_INTERVAL (every 2 s) forever, burying the
# rest of the log and filling the shipped-log ring. Each (what, server, code)
# signature is reported in full at most once per ERROR_LOG_COOLDOWN; the
# suppressed occurrences are COUNTED and reported on the next line that does
# print, so nothing is hidden — only de-duplicated.
_ERR_SEEN: Dict[str, List[float]] = {} # sig -> [last_logged_ts, suppressed]
_ERR_LOCK = threading.Lock()
def _should_log_error(sig: str, cooldown: float) -> Optional[int]:
"""None = suppress. Otherwise the number of occurrences suppressed since the
last time this signature was logged (0 on a first/rising-edge report)."""
now = time.monotonic()
with _ERR_LOCK:
ent = _ERR_SEEN.get(sig)
if ent is None or (now - ent[0]) > cooldown:
suppressed = int(ent[1]) if ent else 0
_ERR_SEEN[sig] = [now, 0.0]
return suppressed
ent[1] += 1
return None
def _clear_error(sig: str) -> None:
"""Forget a signature once it succeeds, so its next failure logs immediately."""
with _ERR_LOCK:
_ERR_SEEN.pop(sig, None)
def _post_each(cfg: Config, session: requests.Session, tmpl: str, *,
json_body: Optional[Dict[str, Any]] = None,
body_bytes: Optional[bytes] = None,
filename: str = "", form: Optional[Dict[str, str]] = None,
what: str = "post", quiet: bool = False) -> Dict[str, Any]:
"""POST to every enabled endpoint. Returns {server_name: {ok, code, error}}.
body_bytes+filename selects multipart upload (map .db); json_body selects a
JSON body. The bytes are held in memory rather than a file handle because a
handle can only be streamed ONCE — with several servers the second upload
would send an empty body.
quiet=True downgrades failures to DEBUG (used by the log-shipping path,
whose own ERROR lines would otherwise feed back into the ring it ships).
Every other caller reports failures, throttled per signature.
"""
results: Dict[str, Any] = {}
cooldown = _env_float("ERROR_LOG_COOLDOWN", "60")
for e in cfg.live():
url = e.at(tmpl, cfg.sn)
try:
if body_bytes is not None:
resp = session.post(url,
files={"db": (filename, body_bytes, "application/octet-stream")},
data=form or {}, headers=e.headers(),
timeout=cfg.http_timeout, verify=cfg.verify_tls)
else:
resp = session.post(url, json=json_body, headers=e.headers(),
timeout=cfg.http_timeout, verify=cfg.verify_tls)
except requests.RequestException as ex:
results[e.name] = {"ok": False, "code": None, "error": _brief(str(ex), 200)}
sig = f"{what}|{e.name}|transport"
if quiet:
log.debug("%s -> %s FAILED (transport): %s", what, e.name, ex)
else:
n = _should_log_error(sig, cooldown)
if n is not None:
log.error("%s -> %s FAILED (transport): %s%s", what, e.name, ex,
f" [+{n} more since last report]" if n else "")
continue
ok = bool(resp.ok)
results[e.name] = {"ok": ok, "code": resp.status_code,
"error": None if ok else _brief(resp.text, 200)}
sig = f"{what}|{e.name}|{resp.status_code}"
if ok:
_clear_error(sig)
continue
if quiet:
log.debug("%s -> %s FAILED: HTTP %s", what, e.name, resp.status_code)
continue
n = _should_log_error(sig, cooldown)
if n is not None:
hint = ""
if resp.status_code in (401, 403):
hint = (f" <- token rejected by {e.name}; each fleet server issues its "
f"OWN token, so a token minted elsewhere will never work here "
f"(set the token for {e.url} in .env)")
elif resp.status_code == 413:
hint = " <- body too large for the server's upload cap"
log.error("%s -> %s FAILED: HTTP %s %s%s%s", what, e.name,
resp.status_code, _brief(resp.text), hint,
f" [+{n} more since last report]" if n else "")
return results
def _fmt_results(results: Dict[str, Any]) -> str:
"""'eco-dev=200 eco=401' — one compact token per server for the log line."""
return " ".join(f"{n}={(r.get('code') if r.get('code') is not None else 'ERR')}"
for n, r in results.items()) or "(no server enabled)"
def _any_ok(results: Dict[str, Any]) -> bool:
return any(r.get("ok") for r in results.values())
# --------------------------------------------------------------------------- #
# mac + storage + software/firmware cards
# --------------------------------------------------------------------------- #
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.08.27"
_software_cache: Optional[Dict[str, Any]] = None
def read_software(cfg: Config) -> Dict[str, Any]:
"""Robot software/OS card: ROS distro, host OS, kernel, arch, python."""
global _software_cache
if _software_cache is not None:
return _software_cache
sw: Dict[str, Any] = {}
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
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
sw["os_version"] = os_ver
# 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
sw["arch"] = u.machine
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): compute board model + Jetson L4T/BSP
release + kernel + the EngineAI product tag. The state source adds live fw."""
global _firmware_cache
if _firmware_cache is not None:
return dict(_firmware_cache)
fw: Dict[str, Any] = {}
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: "# R36 (release), REVISION: 4.3, ..." -> "R36.4.3"
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
# board serial (Jetson module) — stable hardware id, useful next to SN
for p in ("/proc/device-tree/serial-number", "/host/proc/device-tree/serial-number"):
try:
s = Path(p).read_bytes().decode().strip("\x00 \n")
if s:
fw["board_serial"] = s
break
except Exception:
continue
# EngineAI product tag, e.g. "export PRODUCT=t800"
for p in ("/app/applications/install/bringup/product.env",
"/host/app/applications/install/bringup/product.env"):
try:
for line in Path(p).read_text().splitlines():
if "PRODUCT" in line and "=" in line:
fw["product"] = line.split("=", 1)[1].strip().strip('"')
break
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 app data-dir size."""
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 ROS 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
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:
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):
""""interface_protocol/msg/PowerInfo" -> 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 subscription is created, no error is raised, and the field
just stays null forever.
Override with ENG_ROS_QOS=reliable if a topic ever requires it."""
from rclpy.qos import HistoryPolicy, QoSProfile, ReliabilityPolicy
want = _env("ENG_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 _ros_shutdown() -> None:
"""Stop rclpy before the interpreter tears down.
Without this the CycloneDDS C++ worker threads are still running when Python
finalises, and the process dies with SIGABRT ("terminate called without an
active exception") AFTER doing all its work correctly. systemd would record
that as a failed exit and it would mask a genuine crash, so the exit path is
made explicit rather than left to garbage collection."""
try:
import rclpy
if rclpy.ok():
rclpy.shutdown()
time.sleep(0.2) # let spin() unblock and the DDS threads wind down
except Exception:
pass
def _exit(code: int) -> int:
"""Flush logs, stop ROS, and leave with a clean status.
os._exit skips the C++ static destructors that abort the process even after
a successful rclpy.shutdown(); the flush above it means nothing is lost."""
_ros_shutdown()
try:
logging.shutdown()
sys.stdout.flush()
sys.stderr.flush()
except Exception:
pass
os._exit(code)
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]
def _floats(v: Any, lo: float, hi: float) -> List[float]:
"""Numeric members of v inside (lo, hi]. Out-of-range slots (0 = 'not
reporting' on this hardware) are dropped, not counted as real readings."""
out: List[float] = []
for x in _seq(v):
try:
f = float(x)
except (TypeError, ValueError):
continue
if lo < f <= hi:
out.append(f)
return out
class EngineAiSource:
"""EngineAI PM01 state source — the ONE robot-specific class in this agent.
Backends:
ENG_SOURCE=auto (default) ros2 if rclpy imports, else http if
ENG_STATE_URL is set, else none (heartbeat mode)
ENG_SOURCE=ros2 subscribe ENG_TOPIC_* (types from ENG_TYPE_*)
ENG_SOURCE=http poll ENG_STATE_URL (one JSON object); same ENG_FIELD_*
ENG_SOURCE=none never read; always heartbeat
Whatever the backend, it fills one snapshot —
{bms, state_age, temps, mos_temps, max_vel, xy, motion, transitions,
motor_faults, power_err, power_enabled, fw} — which is the entire seam
between this robot and the shared telemetry / fault / status pipeline.
Every read is wrapped: a wrong mapping yields null fields and heartbeat
mode, never a crashed loop.
THROTTLING matters here. /hardware/joint_state publishes at 500 Hz and
/hardware/motor_debug at 100 Hz. Running a Python callback on every message
would burn CPU on the robot's own compute for data we resample every 2 s, so
each stream is decimated to ENG_*_MIN_PERIOD. The agent must not tax the
machine it is monitoring."""
# Defaults match the PM01's REAL interface_protocol messages.
_DEFAULTS = {
"soc": "percentage", # PowerInfo.percentage (0..100)
"voltage": "voltage", # PowerInfo.voltage (V)
"current": "current", # PowerInfo.current (A)
"current_limit": "current_limit",
"power_err": "error_code", # PowerInfo.error_code
"power_enable": "enable", # PowerInfo.enable
"temp": "", # PowerInfo carries NO pack temp
"soh": "", # not published
"cycles": "", # not published
"temps": "motor_temperature", # MotorDebug.motor_temperature[]
"mos_temps": "mos_temperature", # MotorDebug.mos_temperature[]
"motor_err": "error_code", # MotorDebug.error_code[]
"motor_offline": "offline", # MotorDebug.offline[]
"vel": "velocity", # JointState.velocity[]
"motion": "current_motion_task", # MotionState.current_motion_task
"transitions": "available_transition_motions",
"x": "pose.pose.position.x", # nav_msgs/Odometry (if one appears)
"y": "pose.pose.position.y",
"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._mos_temps: List[float] = []
self._motor_faults: List[str] = []
self._motor_offline: List[int] = []
self._power_err: Optional[int] = None
self._power_enabled: Optional[bool] = None
self._max_vel = 0.0
self._xy: Optional[Dict[str, float]] = None
self._motion: Optional[str] = None
self._transitions: List[str] = []
self._fw: Dict[str, Any] = {}
self._stop = False
self.backend = "none"
self.ok = False
self._map = {k: _env("ENG_FIELD_" + k.upper(), d) for k, d in self._DEFAULTS.items()}
# +1 = positive current means CHARGING (the ROS BatteryState convention).
# The PM01 reports the OPPOSITE: measured on the robot, PowerInfo.current
# stays positive (~2 A) while the pack drains (27% -> 26%, 54.74 -> 54.58 V),
# so the shipped default is -1 and "charging" is not inverted.
self._cur_sign = _env_float("ENG_CURRENT_SIGN", "-1")
# The telemetry schema is VOLTS and AMPS; PowerInfo already uses both.
self._v_scale = _env_float("ENG_VOLTAGE_SCALE", "1")
self._i_scale = _env_float("ENG_CURRENT_SCALE", "1")
self._soc_scale = _env("ENG_SOC_SCALE", "percent").lower()
# decimation (see class docstring) — seconds between processed messages
self._joint_period = _env_float("ENG_JOINT_MIN_PERIOD", "0.05") # 500Hz -> 20Hz
self._motor_period = _env_float("ENG_MOTOR_MIN_PERIOD", "0.2") # 100Hz -> 5Hz
self._last_joint = 0.0
self._last_motor = 0.0
self._start()
# ---------------- backend selection ----------------
def _start(self) -> None:
want = (_env("ENG_SOURCE", "auto").lower() or "auto")
order = ["ros2", "http"] if want == "auto" else [want]
for b in order:
try:
if b == "none":
break
if b == "ros2" and self._start_ros2():
self.backend = "ros2"
break
if b == "http" and self._start_http():
self.backend = "http"
break
except Exception as e:
log.warning("EngineAI source %r failed to start (%s)", b, e)
if self.backend == "none":
log.warning("no EngineAI state source active (ENG_SOURCE=%s) — telemetry runs "
"in heartbeat mode. Check that the unit sources ros_env.sh "
"(ROS_DOMAIN_ID=%s + CYCLONEDDS_URI) and run tools/probe_eng.sh "
"on the robot to confirm the topic names.", want, self.cfg.domain_id)
else:
self.ok = True
log.info("EngineAI source up: backend=%s", self.backend)
def _start_http(self) -> bool:
url = _env("ENG_STATE_URL")
if not url:
return False
self._http_url = url
self._http_period = max(0.2, _env_float("ENG_HTTP_INTERVAL", "1"))
threading.Thread(target=self._http_loop, daemon=True).start()
log.info("EngineAI 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("EngineAI http state: HTTP %s", r.status_code)
except Exception as e:
log.debug("EngineAI 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
the vendor's interface_protocol messages work as long as the overlay is
sourced (the systemd unit sources ros_env.sh, which does exactly that)."""
try:
import rclpy
from rclpy.node import Node
except Exception as e:
log.debug("rclpy unavailable (%s)", e)
return False
# (topic, type, callback, high_rate_key) — high-rate topics are
# subscribed raw so the gate can drop messages before deserialization.
wanted = [
(_env("ENG_TOPIC_POWER", "/hardware/power_info"),
_env("ENG_TYPE_POWER", "interface_protocol/msg/PowerInfo"),
self._ingest_battery, ""),
(_env("ENG_TOPIC_MOTORS", "/hardware/motor_debug"),
_env("ENG_TYPE_MOTORS", "interface_protocol/msg/MotorDebug"),
self._ingest_motors, "motor"),
(_env("ENG_TOPIC_JOINTS", "/hardware/joint_state"),
_env("ENG_TYPE_JOINTS", "interface_protocol/msg/JointState"),
self._ingest_joints, "joint"),
(_env("ENG_TOPIC_MOTION", "/motion/motion_state"),
_env("ENG_TYPE_MOTION", "interface_protocol/msg/MotionState"),
self._ingest_motion, ""),
(_env("ENG_TOPIC_ODOM", ""),
_env("ENG_TYPE_ODOM", "nav_msgs/msg/Odometry"), self._ingest_odom, ""),
]
if not any(t for t, _, _, _ in wanted):
return False
use_raw = _env_bool("ENG_RAW_SUBSCRIBE", True)
if not rclpy.ok():
rclpy.init(args=None)
node = Node("sanad_api_eng")
n = 0
for topic, spec, cb, hot in wanted:
if not topic or not spec:
continue
try:
cls = _import_msg(spec)
except Exception as e:
log.warning("EngineAI ros2: cannot import %s for %s (%s) — skipped",
spec, topic, e)
continue
raw = bool(hot) and use_raw
if raw:
try:
node.create_subscription(cls, topic, self._raw_cb(cls, hot),
_qos(), raw=True)
except TypeError:
# rclpy build without raw= support: fall back to a normal
# subscription (correct, just more CPU). Never a reason to
# lose the topic entirely.
log.info("EngineAI ros2: raw subscriptions unsupported here — "
"normal subscription for %s", topic)
raw = False
node.create_subscription(cls, topic, cb, _qos())
else:
node.create_subscription(cls, topic, cb, _qos())
note = ""
if raw:
per = self._joint_period if hot == "joint" else self._motor_period
note = " [raw, gated to ~%.0f Hz]" % (1.0 / per if per > 0 else 0)
log.info("EngineAI ros2: subscribed %s (%s)%s", topic, spec, note)
n += 1
if not n:
return False
self._node = node
threading.Thread(target=lambda: rclpy.spin(node), daemon=True).start()
return True
# ---------------- 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 _gate(self, which: str) -> bool:
"""Rate gate for the high-frequency streams. True = process this message.
Cheap and lock-free: it is the first thing every high-rate callback does,
including the raw ones, so a dropped message costs one clock read."""
now = time.monotonic()
if which == "joint":
if now - self._last_joint < self._joint_period:
return False
self._last_joint = now
return True
if now - self._last_motor < self._motor_period:
return False
self._last_motor = now
return True
def _raw_cb(self, cls: Any, which: str):
"""Wrap a high-rate subscription so the rate gate runs BEFORE the message
is deserialized.
Measured on this robot, per subscription, as % of ONE core (12 available):
all four topics 26.5 <- as first written
without /hardware/joint_state 6.5 => joint_state alone ~20%
without joint_state+motor_debug 2.2 => motor_debug ~4.3%
=> everything else ~2.2%
So the 500 Hz stream was ~75% of the agent's cost, for one number
(max |velocity|) that telemetry samples every 2 s.
raw=True hands the callback the serialized bytes, so a message the gate
drops is never deserialized. That measured 27.2% -> 23.8%: a real win,
but smaller than expected, which says most of the cost is CycloneDDS
delivering 500 msg/s and waking the executor, NOT building the Python
object. Deserialization is only the part we can avoid without changing
what DDS does.
The remaining ~20% could be removed by creating and destroying the
subscription around each sample, but that would emit DDS endpoint
discovery traffic every 2 s onto eth1 — the network this robot's MOTION
CONTROLLER lives on (peer 192.168.0.163). Trading ~2% of an application
core for periodic discovery churn on a realtime control network is the
wrong trade, so the subscription is left stable.
If the CPU ever matters more than the "moving" status, ENG_TOPIC_JOINTS=
(empty) drops this subscription entirely and the agent costs ~6%."""
from rclpy.serialization import deserialize_message
sink = self._joints_from if which == "joint" else self._motors_from
def cb(data: Any) -> None:
if not self._gate(which):
return
try:
sink(deserialize_message(data, cls))
except Exception:
pass
return cb
def _ingest_battery(self, msg: Any) -> None:
"""PowerInfo -> the battery record the telemetry schema expects."""
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
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")
lim = self._num(msg, "current_limit")
err = _dig(msg, self._map.get("power_err", ""))
en = _dig(msg, self._map.get("power_enable", ""))
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,
"current_limit_a": round(lim, 1) if lim is not None else None,
}
with self._lock:
self._bms = rec
if err is not None:
try:
self._power_err = int(err)
except (TypeError, ValueError):
pass
if en is not None:
self._power_enabled = bool(en)
self._bms_ts = self._state_ts = time.monotonic()
except Exception:
pass
def _ingest_motors(self, msg: Any) -> None:
"""MotorDebug -> per-motor temperatures + per-motor fault/offline flags.
This is the field the X2 could never fill: that robot publishes no
per-motor temperature at all, so its motor_temp was permanently null.
The PM01 publishes 25 motor temps AND 25 MOSFET temps, so both are
real readings here — reported, not invented."""
if not self._gate("motor"):
return
self._motors_from(msg)
def _motors_from(self, msg: Any) -> None:
now = time.monotonic()
try:
temps = _floats(_dig(msg, self._map.get("temps", "")), 0, 200)
mos = _floats(_dig(msg, self._map.get("mos_temps", "")), 0, 200)
faults: List[str] = []
offline: List[int] = []
for i, code in enumerate(_seq(_dig(msg, self._map.get("motor_err", "")))):
try:
c = int(code)
except (TypeError, ValueError):
continue
if c != 0:
faults.append(f"joint{i}=0x{c:x}")
for i, off in enumerate(_seq(_dig(msg, self._map.get("motor_offline", "")))):
try:
if int(off):
offline.append(i)
except (TypeError, ValueError):
continue
with self._lock:
self._temps = temps
self._mos_temps = mos
self._motor_faults = faults
self._motor_offline = offline
self._state_ts = now
except Exception:
pass
def _ingest_joints(self, msg: Any) -> None:
"""JointState -> max |velocity| across joints, which drives 'moving'."""
if not self._gate("joint"):
return
self._joints_from(msg)
def _joints_from(self, msg: Any) -> None:
now = time.monotonic()
try:
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._max_vel = max_vel
self._state_ts = now
except Exception:
pass
def _ingest_motion(self, msg: Any) -> None:
"""MotionState -> the live motion task + the transitions it allows.
Read-only: the transitions list is REPORTED so the dashboard can show
what the robot would accept. The agent never requests one."""
try:
task = _dig(msg, self._map.get("motion", ""))
trans = _seq(_dig(msg, self._map.get("transitions", "")))
with self._lock:
if task is not None:
self._motion = str(task)
self._transitions = [str(t) for t in trans]
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 every extractor."""
self._ingest_battery(d)
self._ingest_motors(d)
self._ingest_joints(d)
self._ingest_motion(d)
self._ingest_odom(d)
try:
fw = _dig(d, self._map.get("fw", ""))
if isinstance(fw, dict):
with self._lock:
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),
"mos_temps": list(self._mos_temps),
"motor_faults": list(self._motor_faults),
"motor_offline": list(self._motor_offline),
"power_err": self._power_err,
"power_enabled": self._power_enabled,
"max_vel": self._max_vel,
"xy": dict(self._xy) if self._xy else None,
"motion": self._motion,
"transitions": list(self._transitions),
"fw": dict(self._fw),
}
# --------------------------------------------------------------------------- #
# control card — the PM01 reports its own motion mode natively
# --------------------------------------------------------------------------- #
_control_cache: Dict[str, Any] = {"ts": 0.0, "data": None}
def read_control(cfg: Config, snap: Dict[str, Any]) -> Optional[Dict[str, Any]]:
"""READ-ONLY control status (no robot control, no motion).
Unlike the X2 — where the mode had to be guessed from an undocumented FSM id
and reported "unknown" — the PM01 publishes /motion/motion_state, which
carries the mode BY NAME plus the exact set of transitions it will accept.
So control.mode and control.switchable_modes are the robot's own words.
The Sanad dashboard's /api/controller/status is merged in when reachable
(arm/teleop/velocity detail), but it is never required."""
out: Dict[str, Any] = {
"mode": snap.get("motion"), # e.g. "pd_sitdown"
"switchable_modes": snap.get("transitions") or [],
"remote_switch_enabled": bool(cfg.control_enable),
"source": "ros2:/motion/motion_state" if snap.get("motion") else None,
}
now = time.monotonic()
if _control_cache["data"] is not None and now - _control_cache["ts"] < 1.5:
d = _control_cache["data"]
else:
url = cfg.control_url
if not url:
port = _REMOTE_STAT.get("port")
url = f"http://127.0.0.1:{port}/api/controller/status" if port else ""
d = None
if url:
try:
r = requests.get(url, timeout=2)
d = r.json() if r.ok else None
except Exception:
d = None
_control_cache.update(ts=now, data=d)
if isinstance(d, dict):
out.update({
"fsm_id": d.get("fsm_id"),
"fsm_mode": d.get("fsm_mode"),
"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"),
})
if out["mode"] is None and not isinstance(d, dict):
return None
return out
class RosbridgePosition:
"""Position over a rosbridge websocket (ENG_POSITION_SOURCE=rosbridge)."""
def __init__(self, cfg: Config):
self.cfg = cfg
self._xy: Optional[Dict[str, float]] = None
self._lock = threading.Lock()
self._stop = False
self._topic = _env("ENG_TOPIC_ODOM", "") or "/odom"
self._type = _env("ENG_TYPE_ODOM", "nav_msgs/Odometry").replace("/msg/", "/")
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": self._topic,
"type": self._type, "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 ENG_SOURCE.
The PM01 as shipped publishes NO odometry topic at all (see
docs/PM01_INTERFACE.md) — the motion stack is a whole-body controller, not a
navigation stack, so there is nothing to localise against until a nav
bringup is started. Position therefore reports null rather than a fabricated
origin. This class exists so that the moment an odometry topic does appear,
it is one .env line (ENG_TOPIC_ODOM=...) to start reporting it."""
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("ENG_FIELD_X", "") or "pose.pose.position.x"
self._py = _env("ENG_FIELD_Y", "") or "pose.pose.position.y"
topic = _env("ENG_TOPIC_ODOM", "")
spec = _env("ENG_TYPE_ODOM", "nav_msgs/msg/Odometry")
if not topic or not spec:
log.warning("ENG_POSITION_SOURCE=ros2 but ENG_TOPIC_ODOM is empty — the PM01 "
"publishes no odometry; position stays null until a nav stack runs")
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:
# EngineAiSource may already have started rclpy when ENG_SOURCE=ros2.
if not rclpy.ok():
rclpy.init(args=None)
node = Node("sanad_api_eng_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("EngineAI 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
class HttpPosition:
"""Position from the Sanad dashboard's nav API (ENG_POSITION_SOURCE=http).
The PM01's Sanad app exposes /api/nav/status; when its nav bringup is
running that payload carries the robot pose. While bringup is down the
endpoint answers reachable:false and this reports null — never a stale or
invented coordinate."""
def __init__(self, cfg: Config):
self.cfg = cfg
self._xy: Optional[Dict[str, float]] = None
self._lock = threading.Lock()
self._stop = False
self._url = _env("ENG_POSITION_URL", "") or "http://127.0.0.1:8014/api/nav/status"
self._px = _env("ENG_POSITION_FIELD_X", "") or "pose.x"
self._py = _env("ENG_POSITION_FIELD_Y", "") or "pose.y"
self._period = max(0.5, _env_float("ENG_POSITION_INTERVAL", "2"))
self.ok = True
threading.Thread(target=self._run, daemon=True).start()
log.info("EngineAI position: http %s (%s / %s)", self._url, self._px, self._py)
def _run(self) -> None:
while not self._stop:
try:
r = requests.get(self._url, timeout=3)
if r.ok:
d = r.json()
x, y = _dig(d, self._px), _dig(d, self._py)
if x is not None and y is not None:
with self._lock:
self._xy = {"x": round(float(x), 3), "y": round(float(y), 3)}
else:
with self._lock:
self._xy = None
except Exception as e:
log.debug("http position poll failed: %s", e)
time.sleep(self._period)
def get(self) -> Optional[Dict[str, float]]:
with self._lock:
return dict(self._xy) if self._xy else None
# --------------------------------------------------------------------------- #
# map sync (saved maps → EVERY fleet server, uploaded ONCE per content per server)
# --------------------------------------------------------------------------- #
@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)
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()
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
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."""
roots = [cfg.maps_dir, cfg.maps_dir / cfg.robot, cfg.maps_dir / "maps_slam"]
roots += list(cfg.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))
# a keepout-BAKED twin 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")
out.append(cfg.web_data_dir / "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", ""),
))
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, Dict[str, str]]:
"""Uploaded-map state, keyed PER SERVER: {server_name: {map_path: fingerprint}}.
Per-server on purpose. With one shared key, enabling a second fleet server
later would find every map already 'uploaded' and that server would never
receive them."""
raw = _read_json(_state_file(cfg), {}) if _state_file(cfg).exists() else {}
if not isinstance(raw, dict):
return {}
# migrate the flat {path: fingerprint} layout written by the X2 agent
if raw and all(isinstance(v, str) for v in raw.values()):
primary = cfg.endpoints[0].name if cfg.endpoints else "default"
return {primary: dict(raw)}
return {k: dict(v) for k, v in raw.items() if isinstance(v, dict)}
def save_state(cfg: Config, state: Dict[str, 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_map(cfg: Config, m: MapArtifact, session: requests.Session,
targets: List[Endpoint]) -> Dict[str, Any]:
"""Upload ONE map to the given servers. Returns {server: {ok, code, error}}.
The body is built once and posted to each target, so every server receives
byte-identical map content."""
if m.fmt == "slam_toolbox":
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 {e.name: {"ok": False, "code": None, "error": "unparseable pgm"} for e in targets}
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),
})
results = _post_each(_scoped(cfg, targets), session, cfg.map_endpoint_tmpl,
json_body=body, what=f"map[{m.stem}]")
if _any_ok(results):
log.info("map uploaded: %s (slam_toolbox %dx%d @ %sm, %d points) -> %s",
m.stem, pgm["width"], pgm["height"], ymeta.get("resolution"),
len(m.points), _fmt_results(results))
return results
meta = build_meta(cfg, m)
if cfg.map_upload_mode == "base64json":
body = dict(meta)
body["db_base64"] = base64.b64encode(m.path.read_bytes()).decode("ascii")
results = _post_each(_scoped(cfg, targets), session, cfg.map_endpoint_tmpl,
json_body=body, what=f"map[{m.name}]")
else: # multipart (default) — read once, post to each server
blob = m.path.read_bytes()
results = _post_each(_scoped(cfg, targets), session, cfg.map_endpoint_tmpl,
body_bytes=blob, filename=m.name,
form={"meta": json.dumps(meta)}, what=f"map[{m.name}]")
if _any_ok(results):
log.info("map uploaded: %s (%.2f MB, %d points) -> %s",
m.name, m.size / 1024 / 1024, len(m.points), _fmt_results(results))
return results
def _scoped(cfg: Config, targets: List[Endpoint]) -> Config:
"""A shallow Config view limited to `targets` (for per-server map uploads)."""
import copy
c = copy.copy(cfg)
c.endpoints = list(targets)
return c
# 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, "servers": {},
}
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 saved maps and upload anything new, PER SERVER.
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,
servers={}, error=f"no saved map found "
f"(maps_dir={cfg.maps_dir}, robot={cfg.robot})")
return 0
state = load_state(cfg)
uploaded = failed = unstable = too_large = current = 0
per_server: Dict[str, Any] = {}
last_err: Optional[str] = None
now = time.time()
for m in maps:
key = str(m.path.resolve())
# which servers still need THIS content?
targets = [e for e in cfg.live()
if force or state.get(e.name, {}).get(key) != m.fingerprint()]
if not targets:
current += 1
continue # every server already has this exact content — one-time rule
# stability guard: a map 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 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) to %s",
m.name, m.size / 1024 / 1024, len(m.points),
",".join(e.name for e in targets))
continue
results = upload_map(cfg, m, session, targets)
per_server.update({k: v.get("code") for k, v in results.items()})
got_one = False
for e in targets:
r = results.get(e.name) or {}
if r.get("ok"):
state.setdefault(e.name, {})[key] = m.fingerprint()
got_one = True
else:
last_err = f"{m.name} -> {e.name}: {r.get('error') or 'failed'}"
if got_one:
save_state(cfg, state)
uploaded += 1
if any(not (results.get(e.name) or {}).get("ok") for e in targets):
failed += 1
if failed and not uploaded:
_set_map_status(state="failed", uploaded=False, maps_found=len(maps),
last_map=maps[0].stem, servers=per_server, error=last_err)
elif uploaded or current:
note = last_err
if not note and too_large:
note = f"{too_large} map(s) skipped: exceed server upload cap (~{cfg.map_max_upload_mb:.0f} MB)"
elif not note and 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, servers=per_server, error=note)
elif too_large:
_set_map_status(state="failed", uploaded=False, maps_found=len(maps),
last_map=maps[0].stem, servers=per_server,
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:
_set_map_status(state="pending", uploaded=False, maps_found=len(maps),
last_map=maps[0].stem, servers=per_server,
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, servers=per_server, 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 (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 servers."""
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, "servers": {}}
_ALERTS_STAT: Dict[str, Any] = {"sent": 0, "last": None, "last_time": None,
"ok": None, "servers": {}}
# 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 (the Sanad 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)
PROJECT_LOG_CONTAINER a docker container name; "auto" (default) scans the
host's docker metadata for a RUNNING Sanad project
(sanad-t8, sanadr1, sanad-p4, sanad*)
Reads the container's json-log directly from the docker data dir — no docker
socket needed, read-only, cannot disturb the project. On this robot that
path is root-owned, which is one reason the agent runs as a system service."""
KNOWN = ("sanad-t8", "sanadr1", "sanad-p4", "sanadv3", "sanad")
# /host/... first (containerised agent with the read-only / mount), then the
# native path (this agent runs natively on the PM01).
DOCKER_ROOTS = ("/host/var/lib/docker/containers", "/var/lib/docker/containers")
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
self._pending: List[str] = []
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
base = None
for root in self.DOCKER_ROOTS:
if Path(root).is_dir():
base = Path(root)
break
if base is None:
return
want = cfg.project_log_container
candidates: List[Any] = []
try:
config_files = list(base.glob("*/config.v2.json"))
except PermissionError:
log.warning("cannot read %s (permission denied) — project logs disabled. "
"Run the agent as a system service (root) to ship them.", base)
return
for cf in config_files:
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 = 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(lp)
if not p.exists() and Path("/host" + lp).exists():
p = Path("/host" + 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
# backfill: the last N RELEVANT lines (post-filter) from the tail, so
# access-log noise 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 is returned on the 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 on EVERY
server. Best-effort: failures are logged at DEBUG only (below the ring's
level -> no feedback loop). Lines are requeued unless at least one server
accepted them."""
lines = _LOG_RING.drain()
if _PROJECT_TAIL is not None and _PROJECT_TAIL.active:
lines.extend(_PROJECT_TAIL.poll())
if not lines:
return
body = {"sn": cfg.sn, "name": cfg.name, "lines": lines, "ts": int(time.time())}
results = _post_each(cfg, session, cfg.logs_endpoint, json_body=body,
what="logs", quiet=True)
_LOGS_STAT["servers"] = {n: r.get("code") for n, r in results.items()}
if _any_ok(results):
_LOGS_STAT.update(last_sent=_now_str(), ok=True)
_LOGS_STAT["lines_sent"] += len(lines)
log.debug("logs shipped: %d lines -> %s", len(lines), _fmt_results(results))
else:
_LOGS_STAT["ok"] = False
_LOG_RING.requeue(lines) # retry next cycle
log.debug("logs ship failed: %s (%d lines requeued)", _fmt_results(results), len(lines))
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. No changes to the Sanad app: we only probe its port
# and POST the URL to /{sn}/remote on every server.
# --------------------------------------------------------------------------- #
_REMOTE_STAT: Dict[str, Any] = {"url": None, "port": None, "kind": None, "ok": None,
"ssh": None, "ssh_ok": None, "servers": {}}
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 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:
results = _post_each(cfg, session, cfg.remote_endpoint, json_body={
"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())}, what="remote(web)", quiet=True)
ok = _any_ok(results)
_REMOTE_STAT.update(url=d["url"], port=d["port"], kind=cfg.remote_kind, ok=ok,
servers={n: r.get("code") for n, r in results.items()})
if ok:
log.info("remote dashboard registered: %s -> %s", d["url"], _fmt_results(results))
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}"
results = _post_each(cfg, session, cfg.remote_endpoint, json_body={
"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())},
what="remote(ssh)", quiet=True)
ok = _any_ok(results)
_REMOTE_STAT.update(ssh=cmd, ssh_ok=ok)
if ok:
log.info("remote SSH registered: %s -> %s", cmd, _fmt_results(results))
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 on every server."""
body = {"sn": cfg.sn, "name": cfg.name, "alert": text,
"message": text, "ts": int(time.time())}
results = _post_each(cfg, session, cfg.alert_endpoint, json_body=body,
what="alert", quiet=True)
ok = _any_ok(results)
_ALERTS_STAT.update(last=text, last_time=_now_str(), ok=ok,
servers={n: r.get("code") for n, r in results.items()})
if ok:
_ALERTS_STAT["sent"] += 1
log.info("alert sent: %s -> %s", text[:120], _fmt_results(results))
return ok
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.
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 would re-fire the same fault every
POLL_INTERVAL for as long as it lasted.
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 for error/billing patterns and fires an
alert on each NEW signature (deduped with a cooldown)."""
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 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]:
"""Fault STRINGS, not objects — the fleet ingest 500s on fault objects.
The PM01 exposes real hardware fault channels the X2 did not: PowerInfo
carries enable+error_code, and MotorDebug carries a per-motor error_code and
offline flag. Those become first-class faults rather than being inferred."""
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)")
mos = snap.get("mos_temps") or []
if mos and max(mos) >= cfg.mos_temp_max:
faults.append(f"MOS_OVERTEMP: driver MOSFET temp {max(mos):.0f}C (warning)")
perr = snap.get("power_err")
if perr:
faults.append(f"POWER_FAULT: power_info error_code 0x{int(perr):x} (critical)")
if snap.get("power_enabled") is False:
faults.append("POWER_DISABLED: battery output reports enable=false (critical)")
mf = snap.get("motor_faults") or []
if mf:
faults.append(f"MOTOR_FAULT: {len(mf)} motor(s) reporting an error code "
f"[{', '.join(mf[:6])}] (critical)")
off = snap.get("motor_offline") or []
if off:
faults.append(f"MOTOR_OFFLINE: {len(off)} motor(s) offline "
f"[{', '.join(str(i) for i in off[:8])}] (critical)")
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
# Motion tasks that mean "powered down / not actively controlled". Everything
# else is a live controller, and then velocity decides moving vs idle.
_PASSIVE_MOTIONS = {"passive", "idle", "none", ""}
def derive_status(cfg: Config, snap: Dict[str, Any]) -> str:
"""charging | moving | idle | offline — the same vocabulary the X2 reports."""
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) > _env_float("MOVING_VEL", "0.15"):
return "moving"
return "idle"
def build_telemetry(cfg: Config, mac: str, reader: Optional[EngineAiSource],
pos: Optional[Any],
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": 54.8, "temp_c": None, "soh": 0, "cycles": 0,
"current_limit_a": 90.0},
"state_age": 0.1, "temps": [sim.get("temp", 45)],
"mos_temps": [sim.get("temp", 45) + 2], "motor_faults": [],
"motor_offline": [], "power_err": 0, "power_enabled": True,
"max_vel": sim.get("max_vel", 0.0), "xy": sim.get("position"),
"motion": sim.get("motion", "passive"), "transitions": [], "fw": {}}
else:
snap = (reader.snapshot() if reader else
{"bms": None, "state_age": None, "temps": [], "mos_temps": [],
"motor_faults": [], "motor_offline": [], "power_err": None,
"power_enabled": None, "max_vel": 0.0, "xy": None,
"motion": None, "transitions": [], "fw": {}})
bms = snap.get("bms")
battery = bms["soc"] if bms else None
charging = bool(bms and bms.get("current_a", 0.0) > 0.05)
status = derive_status(cfg, snap)
faults = derive_faults(cfg, snap)
battery_detail = 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"),
"current_limit_a": bms.get("current_limit_a")}
temps = snap.get("temps") or []
mos = snap.get("mos_temps") or []
motor_temp = None
if temps:
motor_temp = {"max": round(max(temps), 1), "avg": round(sum(temps) / len(temps), 1),
"min": round(min(temps), 1), "count": len(temps)}
if mos:
motor_temp["mos_max"] = round(max(mos), 1)
motor_temp["mos_avg"] = round(sum(mos) / len(mos), 1)
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. pm01_150)
"mac": mac,
"brand": cfg.brand, # engineai
"type": cfg.robot_type, # humanoid
"model": cfg.model, # pm01
"software": read_software(cfg), # ros/os/kernel/arch/python/agent
"firmware": {**read_firmware_static(), **(snap.get("fw") or {})},
# board/l4t/product/board_serial + live fw
"battery": battery, # null = couldn't read (heartbeat)
"charging": charging,
"battery_detail": battery_detail,
"motor_temp": motor_temp, # REAL on the PM01 (25 motors) — null = not receiving
"storage": read_storage(cfg),
"status": status,
"position": position, # null when no odom/localization source
"control": read_control(cfg, snap), # live motion task + allowed transitions
"faults": faults,
"map": get_map_status(), # SHOWS whether the saved map made it to each 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. "sanad-t8-logs"; null = no project found
"remote": (dict(_REMOTE_STAT) if _REMOTE_STAT.get("url") else None),
"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:
"""POST the identical payload to every enabled fleet server."""
results = _post_each(cfg, session, cfg.telemetry_endpoint,
json_body=payload, what="telemetry")
mp = payload.get("map") or {}
mt = payload.get("motor_temp") or {}
log.info("telemetry ok: battery=%s charging=%s status=%s mode=%s pos=%s "
"motor_max=%s faults=%d map=%s -> %s",
payload["battery"], payload["charging"], payload["status"],
(payload.get("control") or {}).get("mode"),
payload["position"], mt.get("max"), len(payload["faults"]),
mp.get("state"), _fmt_results(results))
return _any_ok(results)
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,
"motion": "rl_basic" if moving else "pd_sitdown",
"position": {"x": round(1.0 + 0.1 * i, 2), "y": round(2.0 - 0.05 * i, 2)}}
def cmd_list(cfg: Config) -> None:
maps = discover_maps(cfg)
if not maps:
print(f"(no maps under {cfg.maps_dir} for robot '{cfg.robot}')")
return
print(f"{len(maps)} map(s) under {cfg.maps_dir} (robot={cfg.robot}):")
for m in maps:
pts = load_points(cfg, m.stem)
print(f" {m.name:<28} {m.size/1024/1024:6.2f} MB {len(pts):>3} points {m.description}")
# --------------------------------------------------------------------------- #
# main
# --------------------------------------------------------------------------- #
def main(argv: Optional[List[str]] = None) -> int:
ap = argparse.ArgumentParser(description="EngineAI PM01 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, 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)
# systemd stops the unit with SIGINT (KillSignal=SIGINT) and falls back to
# SIGTERM. Handle both so a restart tears ROS down cleanly instead of
# aborting — otherwise every `systemctl restart` is recorded as a failure.
import signal
def _on_signal(sig: int, _frame: Any) -> None:
log.info("received %s — stopping", signal.Signals(sig).name)
_exit(0)
for _s in (signal.SIGINT, signal.SIGTERM):
try:
signal.signal(_s, _on_signal)
except Exception:
pass
_load_dotenv()
cfg = Config.from_env()
if args.interval is not None:
cfg.poll_interval = args.interval
if args.list:
cmd_list(cfg)
return _exit(0)
_init_start_times(cfg)
global _PROJECT_TAIL, _LOG_ALERTS
_PROJECT_TAIL = ProjectLogTail(cfg)
_LOG_ALERTS = LogAlertScanner(cfg)
mac = read_mac(cfg.mac_interface)
log.info("sanad_api_eng — sn=%s name=%s mac=%s iface=%s pos=%s map_dir=%s "
"map_every=%.0fs%s",
cfg.sn, cfg.name, mac, cfg.mac_interface, cfg.position_source,
cfg.maps_dir, cfg.map_poll_interval,
" [SIMULATE]" if args.simulate else "")
for e in cfg.endpoints:
log.info(" fleet server %-10s %s %s", e.name, e.url,
"ENABLED" if e.enabled else "disabled (no token / SERVER_*_ENABLE=0)")
session = requests.Session()
if args.map_only:
map_sync_once(cfg, session, force=args.force, dry_run=args.dry_run)
return _exit(0)
reader = None
pos = None
if not args.simulate:
reader = EngineAiSource(cfg)
if cfg.position_source == "rosbridge":
pos = RosbridgePosition(cfg)
elif cfg.position_source == "ros2":
pos = Ros2Position(cfg)
elif cfg.position_source == "http":
pos = HttpPosition(cfg)
time.sleep(1.5) # let the first ROS messages land before the first post
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, indent=2))
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 _exit(0)
for _ in range(3):
one_telemetry()
time.sleep(min(cfg.poll_interval, 1.0))
return _exit(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 to %d server(s); map check every %.0fs; logs every %.0fs",
cfg.poll_interval, len(cfg.live()), 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 _exit(0)
if __name__ == "__main__":
sys.exit(main())