fleet/PIPELINE.md

19 KiB
Raw Blame History

Sanad Fleet — Data Pipeline

How robot state becomes a record on the YS Lootah fleet dashboard, end to end. This is the reference for what flows, in what shape, when, and how failures are handled. For install/ops, see README.md.


1. Big picture

Sanad Fleet — architecture and data pipeline

        ROBOT (edge)                             │        YS LOOTAH (cloud)
                                                 │
  ┌───────────────┐   read     ┌──────────────┐  │   HTTPS POST     ┌──────────────┐
  │ robot sources │ ─────────▶ │ sanad_api_*  │ ─┼───────────────▶ │ fleet server │
  │ DDS / files   │            │ (Docker)     │  │  Bearer token    │ (ingest API) │
  └───────────────┘            └──────────────┘  │                  └──────┬───────┘
        ▲                            │           │                         │
        │ passive, read-only         │ systemd   │                         ▼
        │ (never commands motion)    │ user svc  │                  ┌──────────────┐
                                     ▼           │                  │  dashboard   │
                              stays "online"     │                  │  storage     │
                              via heartbeat      │                  │  alerts      │
                                                 │                  └──────────────┘
  • Direction: outbound only. The robot opens no inbound ports; every call is an HTTPS POST to the fleet server.
  • Transport: HTTP/1.1 + JSON (telemetry) or multipart/form-data (map).
  • Auth: Authorization: Bearer <device_token> on every request.
  • Identity: each robot is keyed by sn (fleet id, e.g. r1_82); the hardware mac rides along in telemetry.

2. API contract

The full spec (from sanad-tasks-en.html) defines these ingest endpoints. This repo currently implements the bold ones; the rest are documented for later.

endpoint method agent status
/api/v1/fleet/ingest/telemetry POST G1, R1, Go2 ~2 s; 27 fields incl. software/firmware/control + all status mirrors
/api/v1/fleet/ingest/{sn}/map POST G1, R1, Go2 once per content; rtabmap .db (≤~7 MB) + slam_toolbox/Nav2/Pudu pgm+yaml→PNG
/api/v1/fleet/ingest/{sn}/alert POST G1, R1, Go2 each NEW fault, rising edge, string body
/api/v1/fleet/ingest/{sn}/logs POST G1, R1, Go2 agent + [sanadr1-logs] project lines every LOGS_INTERVAL (60 s)
/api/v1/fleet/ingest/{sn}/remote POST G1, R1, Go2 registers web (Sanad dashboard URL) + ssh (ssh unitree@<ip>) every 60 s
/api/v1/fleet/ingest/{sn}/commands GET channel works; no command executor built (motion = deliberately out)
/api/v1/fleet/ingest/commands/{id}/ack POST

Auth header (all): Authorization: Bearer <device_token>.


3. Telemetry pipeline (all agents)

Goal: a fresh status record every ~2 s; never go dark.

  DDS topics                    agent (2 s loop)                         server
  ──────────                    ────────────────                         ──────
  rt/lowstate  ─┐   callbacks   ┌─ snapshot() ─┐   build_telemetry()   ┌ POST
  (LowState_)   ├─────────────▶ │ battery      ├──────────────────────▶│ /ingest/
  rt/lf/bmsstate│  (background  │ temps/dq     │  derive status+faults │ telemetry
  (BmsState_)  ─┘   threads)    │ liveness/age │  read mac             └ (JSON)
                                └──────────────┘
                                       │ no data?
                                       └────────▶ heartbeat (battery:null, offline)

3.1 Sources per robot

field G1 / R1 (unitree_hg) Go2 (unitree_go)
battery (0100) rt/lf/bmsstateBmsState_.soc rt/lowstateLowState_.bms_state.soc
charging BmsState_.current > +0.05 A bms_state.current > +0.05 A
faults[] rt/lowstate motor temps + staleness same
status derived (see 3.3); optional loco FSM GET 7001 derived
position G1: rt/lf/odommodestate · R1: optional rosbridge /odom optional rosbridge /odom or rt/lf/sportmodestate
control Sanad GET /api/controller/statusread-only loco mode (G1 200/4/2/702, R1 0/1/4/811) same path, Go2 labels
mac NIC (/sys/class/net/<iface>/address) NIC

DDS is initialized once (ChannelFactoryInitialize(domain, interface)); each topic has a subscriber whose callback updates a locked in-memory snapshot. The loop reads the snapshot — it never blocks on the network.

3.2 Payload

POST /api/v1/fleet/ingest/telemetry
Authorization: Bearer <device_token>
Content-Type: application/json

{ // ── identity ──
  "sn": "E39N4000Q6D7E70F",              // the robot's REAL Unitree serial (server key)
  "name": "r1_82",                        // friendly display name
  "mac": "4c:bb:47:51:25:9a",
  "brand": "unitree", "type": "humanoid", "model": "r1",   // type: humanoid|dog

  // ── software / firmware ──
  "software": { "ros": "foxy", "os": "Ubuntu 20.04.5 LTS", "os_version": "20.04",
                "kernel": "5.10.104-tegra", "arch": "aarch64",
                "python": "3.10.20", "agent": "sanad_api_r1 2026.07.13" },
  "firmware": { "board": "NVIDIA Orin NX Developer Kit", "l4t": "R35.3.1",
                "kernel": "5.10.104-tegra", "robot": "0.0", "bms": "0.44" },

  // ── power / health ──
  "battery": 62, "charging": false,
  "battery_detail": { "voltage_v": 34.5, "current_a": -3.08, "temp_c": 42,
                      "soh": 99, "cycles": 10 },
  "motor_temp": { "max": 49.0, "avg": 37.8, "min": 32.0 },   // null = not receiving
  "storage": { "total_gb": 98.2, "free_gb": 58.5, "used_percent": 36.2,
               "data_kb": 770.0 },        // data_kb only when STORAGE_DATA_PATH set

  // ── state ──
  "status": "idle",
  "position": { "x": 12.4, "y": 3.1 },   // or null when no localization source
  "control": { "fsm_id": 0, "mode": "zero_torque",          // zero_torque|damp|lock|running
               "armed": false, "walk_ready": false, "teleop_active": false,
               "switchable_modes": ["zero_torque","damp","lock","running"],
               "remote_switch_enabled": false },            // read-only; no remote motion
  "faults": [],                           // strings, e.g. "LOW_BATTERY: battery 12% (warning)"

  // ── sub-system status (each stream's health, mirrored here) ──
  "map":    { "uploaded": true, "state": "uploaded", "maps_found": 2,
              "last_map": "rtabmap", "error": null },
  "logs":   { "last_sent": "…", "lines_sent": 105, "ok": true },
  "project_logs": "sanadr1-logs",         // the Sanad app whose logs are tailed (null=none)
  "remote": { "url": "http://10.255.254.82:8001", "port": 8001, "kind": "web",
              "ok": true, "ssh": "ssh unitree@10.255.254.82", "ssh_ok": true },
  "alerts": { "sent": 0, "last": null, "last_time": null, "ok": null },

  // ── time ──
  "time": "2026-07-13 17:00:22+04:00",    // full local datetime (TZ_OFFSET_HOURS, default +4)
  "started_at": "2026-07-13 17:00:21+04:00",
  "last_start": "2026-07-13 16:57:43+04:00",   // previous agent start (persisted)
  "uptime_s": 625,
  "ts": 1731000000 }

Sub-objects map / logs / alerts / remote / control are status mirrors — the real work happens on their own endpoints (below); telemetry just always shows whether each is healthy so the fleet UI never has to poll them separately.

3.3 Status derivation

if no rt/lowstate for >3 s AND never saw battery   → "offline"
elif charging (bms current > +0.05 A)              → "charging"
elif max |joint velocity| > 0.15 rad/s             → "moving"
elif R1_READ_FSM and FSM id known                  → FSM label (811 ready / 4 standing / 1 damping / 0 zero_torque)
else                                               → "idle"

3.4 Faults

Faults ride inside telemetry faults[] as strings (the ingest 500s on fault objects) — and each new fault also fires the /alert endpoint (rising edge):

code trigger severity
LOW_BATTERY soc ≤ LOW_SOC (default 50) warning
MOTOR_OVERTEMP any motor temp ≥ MOTOR_TEMP_MAX (default 85 °C) warning
COMMS_STALE no rt/lowstate for > 3 s critical
GEMINI_BILLING / ROBOT_ERROR project-log scan (see §5) critical

3.5 Heartbeat & failure handling

  • No DDS / no state → still POST with battery:null, status:"offline" so the dashboard shows the robot as reachable (the spec's "send a heartbeat" rule).
  • unitree_sdk2py missing → the agent logs a warning and runs in heartbeat mode (no crash).
  • POST fails (transport or non-2xx) → logged, loop continues; next tick retries in POLL_INTERVAL seconds. No back-pressure, no queue (latest state wins).
  • Every tick is wrapped — one bad read can't kill the loop.

4. Map pipeline (all agents)

Goal: keep the server's copy of every saved nav map current, sending each map once per content change.

  files on disk                    agent (30 s scan)                        server
  ─────────────                    ─────────────────                        ──────
  MAPS_DIR/<robot>/*.db        ─┐  discover ┌ fingerprint fmt:size:mtime ┐ changed? ┌ POST
  MAPS_DIR/**.yaml + .pgm       ├─────────▶ │ .db      → sha256+multipart├──yes────▶│ /ingest/
  EXTRA_MAP_DIRS  (Nav2/Pudu)  ─┘           │ pgm+yaml → PNG image JSON  │          │ {sn}/map
                                            └────────────────────────────┘          └
                                                     │ unchanged
                                                     └────────▶ skip (state/uploaded.json)

4.1 Formats & discovery

Two map formats are understood:

  • RTAB-Map .db (VSLAM) — sent as-is (multipart), with maps_meta.json description and places → points[]. Skipped when larger than MAP_MAX_UPLOAD_MB (7 — the server 413s at ~8 MB) with a map.error note. A 120 s stability guard skips a .db still being written (active mapping → map.state: "pending").
  • slam_toolbox / Nav2 / Pudu sets<stem>.yaml + <stem>.pgm (+ optional .posegraph/.data) → uploaded as the spec's image JSON (pure-stdlib PGM→PNG + resolution + origin + size). Small; always uploads.

Search roots (each existing dir is scanned): MAPS_DIR, MAPS_DIR/<robot>, MAPS_DIR/maps_slam, plus every dir in EXTRA_MAP_DIRS (colon-separated — the installer detects the robot's *_nav2_docker/maps dir and mounts it at /data/nav2_maps).

Pudu twin dedup: the Pudu converter emits a plain map and a keepout-baked twin; when both <stem>.yaml and <stem>_keepout_baked.yaml exist, only the baked one uploads — the server gets one canonical map.

Places → points: web/data/<robot>/places/<map>.json ({ "<name>": {x, y, z, qx, qy, qz, qw} }) → {name, type:"waypoint", x, y, yaw} (yaw from the quaternion).

4.2 Change detection

  1. Cheap pre-check: fmt : size : mtime vs STATE_DIR/uploaded.json.
  2. If different, compute sha256 (of the .db, or of the rendered image JSON).
  3. Upload; on success, record the new fingerprint. A restart re-reads state → no redundant re-upload.

MAP_SELECT chooses scope: all (default), active (via web_nav3 GET /api/status), or newest.

4.3 Wire formats

Image JSON (slam_toolbox / Nav2 / Pudu sets — the spec's map body):

POST /api/v1/fleet/ingest/E39N4000Q6D7E70F/map
{ "sn": "E39N4000Q6D7E70F", "name": "map_keepout_baked", "format": "image",
  "resolution": 0.05, "origin": [-5.85, -11.8, 0.0],
  "width": 178, "height": 304,
  "image_base64": "<PNG bytes, base64>",
  "points": [ {"name":"dock","type":"waypoint","x":1.2,"y":3.4,"yaw":0.0} ] }

Multipart (RTAB-Map .db, default MAP_UPLOAD_MODE=multipart) — file part db + form field meta:

POST /api/v1/fleet/ingest/g1_7892/map
Authorization: Bearer <device_token>
Content-Type: multipart/form-data

  db   = <floor-1.db bytes>        (application/octet-stream)
  meta = { "sn":"g1_7892", "name":"floor-1", "file":"floor-1.db",
           "format":"rtabmap_db", "size_bytes":6994944, "sha256":"…",
           "mtime":1731000000, "description":"ground floor",
           "points":[ {"name":"dock","type":"waypoint","x":1.2,"y":3.4,"yaw":0.0} ] }

base64json (set MAP_UPLOAD_MODE=base64json) — same fields as JSON with the .db as db_base64.

Server note: raster (pgm+yaml) maps use the spec's image body and display on the dashboard. The raw .db path additionally requires the server to accept a format:"rtabmap_db" artifact — and is size-capped (~8 MB), so a big VSLAM .db is skipped; export a raster with tests/map_export_once.py instead.

4.4 One-shot upload (--map-only)

--map-only runs one discovery + upload pass with no DDS and no telemetry — push a new map without disturbing a live feed:

docker run --rm --network host --env-file .env \
    -e EXTRA_MAP_DIRS=/data/nav2_maps -v <host-maps>:/data/nav2_maps:ro \
    sanad-api-<type>:latest --map-only --force

5. Logs, alerts & remote registration

  sanad app json-log ──tail──▶ ProjectLogTail ──┐
  agent's own lines  ──ring───────────────────── ├─▶ every 60 s ─▶ POST /{sn}/logs
                                                 ┘    (failed ships requeue, last ~400 lines)

  project log ──scan 10 s──▶ LogAlertScanner ──new signature──▶ POST /{sn}/alert
  faults[]    ──rising edge──────────────────────────────────▶ POST /{sn}/alert

  localhost ports 8001,8014,8011-8013,8000,8080 ──probe──▶ dashboard URL
      └─▶ every 60 s ─▶ POST /{sn}/remote  (kind "web" + kind "ssh")
  • Logs — the agent tails the robot's Sanad container json-log via the read-only /host mount (no docker socket): auto-discovers a running sanad* container, unwraps the json-log, labels lines [sanadr1-logs], filters uvicorn access-log noise, backfills the last PROJECT_LOG_BACKFILL (100) relevant lines on start.
  • AlertsALERT_LOG_PATTERNS (CODE=regex split by ;;, case-sensitive; (?i) inline where needed) scan the project log every 10 s; each new signature (code + digit-stripped line) alerts once per ALERT_LOG_COOLDOWN (300 s). On startup the last ALERT_BACKFILL_BYTES (8 MiB) are scanned so an already-active error still alerts. Defaults catch Gemini billing and any ERROR/CRITICAL/Traceback.
  • Remote — probes REMOTE_PORTS on localhost for the Sanad dashboard page and registers {kind:"web", url:"http://<lan-ip>:<port>"} plus {kind:"ssh", command:"ssh unitree@<ip>"} — no change to the Sanad app. ⚠ LAN-only URL unless REMOTE_URL pins a public tunnel.

6. Timing

stream cadence trigger
telemetry every POLL_INTERVAL (2 s) timer
map scan every MAP_POLL_INTERVAL (30 s) uploads only on content change
alert scan every ALERT_SCAN_INTERVAL (10 s) new log signature / rising-edge fault
log ship every LOGS_INTERVAL (60 s) ring buffer + project tail
remote register every REMOTE_INTERVAL (60 s) re-asserts the dashboard/ssh entry
DDS reads continuous (subscriber callbacks) firmware publish rate
heartbeat same as telemetry cadence when state is unreadable

7. Where it runs

Deploy — workstation to robot over SSH

robot host
└─ systemd --user
   └─ sanad-api-<type>.service   (Restart=always, enabled at boot via linger)
      └─ docker start -a sanad-api-<type>
         └─ container (--network host)
            └─ python -u sanad_api_<type>.py   ← the loop above
  • --network host: DDS multicast visibility (R1/Go2) and localhost access to web_nav3:8765 (G1); also the robot's real NIC MAC.
  • Single owner: the container is docker created without a docker restart policy; systemd owns start/stop/restart.

8. End-to-end sequences

Telemetry tick

loop        DDS(sub)      agent            fleet server
 │  (2 s)      │            │                   │
 │────────────▶│ snapshot   │                   │
 │             │───────────▶│ build payload     │
 │             │            │──── POST JSON ────▶│  (Bearer)
 │             │            │◀──── 200 OK ───────│
 │             │            │ log "telemetry ok" │

Map change

scan(30s)    disk          agent            fleet server
 │────────────▶│ discover   │                   │
 │             │───────────▶│ fingerprint       │
 │             │            │ changed → sha256   │
 │             │            │ load points        │
 │             │            │── POST multipart ─▶│  (db + meta)
 │             │            │◀──── 200 OK ───────│
 │             │            │ save state         │

9. Verified behavior (real hardware)

Confirmed against the live fleet during bring-up:

  • R1 is LIVE on production (https://eco.yslootahrobotics.com, sn E39N4000Q6D7E70F): telemetry every 2 s (HTTP 200), real battery draining live, control.mode tracks the actual loco FSM (lock/zero_torque observed).
  • Maps: R1's VSLAM raster export uploaded and displayed on the dashboard; the converted Pudu office map (map_keepout_baked, 178×304 @ 0.05 m) uploaded via --map-only → HTTP 200. Oversized .db correctly skipped (413 cap).
  • Alerts verified live: LOW_BATTERY fired at 50/49/48 %…, GEMINI_BILLING and ROBOT_ERROR fired from the real sanadr1 log — no false positives after the case-sensitivity fix.
  • Remote: R1 registered http://10.255.254.82:8001 (web) + ssh unitree@… → HTTP 200.
  • Logs: shipped labeled [sanadr1-logs] lines; requeue held them through the server-side /logs 500 until the handler was fixed (then stored:105).
  • G1 uploaded a map (multipart, format:rtabmap_db, points with correct yaw); G1 map + R1 telemetry landed in the same window — multi-robot works concurrently. ⚠ G1 currently offline (battery) + deploy hold.
  • Go2: pipeline code + image verified (builds, unitree_go imports, simulate payloads correct); not yet run on a real Go2.