feat: unified fleet telemetry — canonical format, GPS/IP geo, WAN health dashboard
- Rewrite telemetry-synology.sh: canonical GL format (modem_0001/wan keys), self-contained metrics (no aiwanbal), Eyeride GPS with IP geo fallback - Add /api/fleet-telemetry to hub: joins tunnel status with fleet Postgres - Update dashboard.html: per-device WAN health bars, GPS, signal strength - Fix hub/ingest.sh normalisation: remap old wan1/wan2 → modem_0001/wan - Bump SPK version: 0.1-0001 → 0.5.0-0001 for GL parity Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+184
-1
@@ -7,6 +7,8 @@ Listens on :8080. Routes:
|
||||
GET /api/register/<DEVICE_ID> — assign tunnel port, return JSON
|
||||
GET /api/fleet — all registered devices with
|
||||
tunnel status + access links
|
||||
GET /api/fleet-telemetry — fleet status enriched with telemetry
|
||||
from fleet Postgres (WAN health, GPS)
|
||||
POST /api/authorize-key — authorize device tunnel key
|
||||
"""
|
||||
|
||||
@@ -22,6 +24,13 @@ REGISTRY = "/opt/busfleet-hub/port-registry.json"
|
||||
DASHBOARD = "/opt/busfleet-hub/dashboard.html"
|
||||
HUB_IP = "162.243.83.36"
|
||||
|
||||
# Fleet Postgres connection (for telemetry enrichment)
|
||||
FLEET_DB_HOST = os.environ.get("FLEET_DB_HOST", "167.172.237.162")
|
||||
FLEET_DB_PORT = os.environ.get("FLEET_DB_PORT", "5432")
|
||||
FLEET_DB_NAME = os.environ.get("FLEET_DB_NAME", "fleet")
|
||||
FLEET_DB_USER = os.environ.get("FLEET_DB_USER", "fleet")
|
||||
FLEET_DB_PASS = os.environ.get("FLEET_DB_PASS", "")
|
||||
|
||||
|
||||
# ── Helper: check if a TCP port is open (tunnel active) ──────────
|
||||
def _port_is_open(port):
|
||||
@@ -33,6 +42,84 @@ def _port_is_open(port):
|
||||
return False
|
||||
|
||||
|
||||
# ── Helper: query fleet Postgres for latest device_status ─────────
|
||||
def _query_fleet_telemetry():
|
||||
"""Query fleet Postgres for the latest telemetry data for all devices.
|
||||
Returns a dict keyed by device_id, or empty dict if DB unreachable."""
|
||||
try:
|
||||
import psycopg2
|
||||
conn = psycopg2.connect(
|
||||
host=FLEET_DB_HOST,
|
||||
port=FLEET_DB_PORT,
|
||||
dbname=FLEET_DB_NAME,
|
||||
user=FLEET_DB_USER,
|
||||
password=FLEET_DB_PASS,
|
||||
connect_timeout=5,
|
||||
)
|
||||
cur = conn.cursor()
|
||||
cur.execute("""
|
||||
SELECT d.device_id, d.vendor, d.model,
|
||||
st.online, st.ts, st.version, st.uptime_s,
|
||||
st.primary_member, st.active_wan,
|
||||
st.gps_lat, st.gps_lon, st.gps_fix,
|
||||
st.wan, st.starlink, st.state
|
||||
FROM devices d
|
||||
LEFT JOIN device_status st ON st.device_id = d.device_id
|
||||
ORDER BY d.device_id
|
||||
""")
|
||||
rows = cur.fetchall()
|
||||
cur.close()
|
||||
conn.close()
|
||||
|
||||
result = {}
|
||||
for row in rows:
|
||||
device_id = row[0]
|
||||
wan_data = row[12] if row[12] else {}
|
||||
# psycopg2 returns JSONB as dict
|
||||
if isinstance(wan_data, str):
|
||||
try:
|
||||
wan_data = json.loads(wan_data)
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
wan_data = {}
|
||||
|
||||
starlink_data = row[13] if row[13] else {}
|
||||
if isinstance(starlink_data, str):
|
||||
try:
|
||||
starlink_data = json.loads(starlink_data)
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
starlink_data = {}
|
||||
|
||||
state_data = row[14] if row[14] else {}
|
||||
if isinstance(state_data, str):
|
||||
try:
|
||||
state_data = json.loads(state_data)
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
state_data = {}
|
||||
|
||||
result[device_id] = {
|
||||
"vendor": row[1],
|
||||
"model": row[2],
|
||||
"online": row[3] if row[3] is not None else False,
|
||||
"ts": str(row[4]) if row[4] else None,
|
||||
"version": row[5],
|
||||
"uptime_s": row[6],
|
||||
"primary_member": row[7],
|
||||
"active_wan": row[8],
|
||||
"gps": {
|
||||
"lat": float(row[9]) if row[9] is not None else None,
|
||||
"lon": float(row[10]) if row[10] is not None else None,
|
||||
"fix": row[11],
|
||||
},
|
||||
"wan": wan_data,
|
||||
"starlink": starlink_data,
|
||||
"state": state_data,
|
||||
}
|
||||
return result
|
||||
except Exception as e:
|
||||
print(f"[hub] fleet telemetry query failed: {e}", file=sys.stderr)
|
||||
return {}
|
||||
|
||||
|
||||
class HubHandler(BaseHTTPRequestHandler):
|
||||
"""Handle hub registration, fleet status, and dashboard requests."""
|
||||
|
||||
@@ -56,11 +143,16 @@ class HubHandler(BaseHTTPRequestHandler):
|
||||
self.call_register_script(device_id)
|
||||
return
|
||||
|
||||
# Fleet status — all registered devices
|
||||
# Fleet status — basic tunnel-only (original endpoint)
|
||||
if self.path == "/api/fleet":
|
||||
self.serve_fleet_status()
|
||||
return
|
||||
|
||||
# Fleet status — enriched with telemetry from fleet DB
|
||||
if self.path == "/api/fleet-telemetry":
|
||||
self.serve_fleet_telemetry()
|
||||
return
|
||||
|
||||
self.send_json_error(404, "Not found")
|
||||
|
||||
def do_POST(self):
|
||||
@@ -161,6 +253,73 @@ class HubHandler(BaseHTTPRequestHandler):
|
||||
"devices": devices,
|
||||
}))
|
||||
|
||||
def serve_fleet_telemetry(self):
|
||||
"""Build enriched fleet status: tunnel status + telemetry from fleet DB."""
|
||||
# Read port registry (tunnel info)
|
||||
try:
|
||||
with open(REGISTRY) as f:
|
||||
registry = json.load(f)
|
||||
except Exception:
|
||||
registry = {}
|
||||
|
||||
# Query fleet Postgres for telemetry
|
||||
fleet_data = _query_fleet_telemetry()
|
||||
|
||||
devices = []
|
||||
for dev_id, entry in registry.items():
|
||||
port = entry.get("tunnel_port", 0)
|
||||
tunnel_up = _port_is_open(port) if port else False
|
||||
|
||||
# Enrich with fleet telemetry
|
||||
telem = fleet_data.get(dev_id, {})
|
||||
|
||||
# Derive platform from vendor or device_id prefix
|
||||
vendor = telem.get("vendor", "")
|
||||
if vendor == "synology" or dev_id.startswith("x"):
|
||||
platform = "synology"
|
||||
elif vendor == "glinet" or dev_id.startswith(("B", "C")):
|
||||
platform = "openwrt"
|
||||
else:
|
||||
platform = telem.get("vendor", "unknown")
|
||||
|
||||
# Format WAN summary for dashboard
|
||||
wan_summary = _format_wan_summary(telem.get("wan", {}), telem.get("primary_member"))
|
||||
|
||||
devices.append({
|
||||
"device_id": dev_id,
|
||||
"tunnel_port": port,
|
||||
"assigned": entry.get("assigned", ""),
|
||||
"tunnel_up": tunnel_up,
|
||||
"ssh_command": f"ssh -p {port} kitadmin@{HUB_IP}",
|
||||
"dashboard_url": f"http://{HUB_IP}:{port}" if tunnel_up else None,
|
||||
# Telemetry enrichment
|
||||
"platform": platform,
|
||||
"version": telem.get("version"),
|
||||
"uptime_s": telem.get("uptime_s"),
|
||||
"online_telemetry": telem.get("online", False),
|
||||
"primary_member": telem.get("primary_member"),
|
||||
"active_wan": telem.get("active_wan"),
|
||||
"gps": telem.get("gps", {}),
|
||||
"wan_summary": wan_summary,
|
||||
"starlink": telem.get("starlink"),
|
||||
"autonomous": telem.get("state", {}).get("autonomous"),
|
||||
})
|
||||
|
||||
# Sort: tunnel_up first, then by device_id
|
||||
devices.sort(key=lambda d: (not d["tunnel_up"], d["device_id"]))
|
||||
|
||||
# Summary stats
|
||||
online = sum(1 for d in devices if d["tunnel_up"])
|
||||
telemetry_online = sum(1 for d in devices if d.get("online_telemetry"))
|
||||
|
||||
self.send_json(200, json.dumps({
|
||||
"hub": HUB_IP,
|
||||
"device_count": len(devices),
|
||||
"online": online,
|
||||
"telemetry_online": telemetry_online,
|
||||
"devices": devices,
|
||||
}))
|
||||
|
||||
def send_json(self, status, body):
|
||||
"""Send a JSON response."""
|
||||
self.send_response(status)
|
||||
@@ -177,6 +336,30 @@ class HubHandler(BaseHTTPRequestHandler):
|
||||
print(f"[hub] {args[0]}", file=sys.stderr)
|
||||
|
||||
|
||||
def _format_wan_summary(wan, primary_member):
|
||||
"""Format WAN data into a concise summary for the dashboard."""
|
||||
if not wan or not isinstance(wan, dict):
|
||||
return []
|
||||
summary = []
|
||||
for member, metrics in wan.items():
|
||||
if not isinstance(metrics, dict):
|
||||
continue
|
||||
entry = {
|
||||
"member": member,
|
||||
"active": member == primary_member,
|
||||
"score": metrics.get("score"),
|
||||
"latency_ms": metrics.get("latency_ms"),
|
||||
"loss_pct": metrics.get("loss_pct"),
|
||||
"rsrp_dbm": metrics.get("rsrp_dbm"),
|
||||
"carrier": metrics.get("carrier"),
|
||||
"technology": metrics.get("technology"),
|
||||
}
|
||||
summary.append(entry)
|
||||
# Sort: active first, then by score descending
|
||||
summary.sort(key=lambda x: (not x["active"], -(x["score"] or 0)))
|
||||
return summary
|
||||
|
||||
|
||||
def main():
|
||||
port = int(sys.argv[1]) if len(sys.argv) > 1 else 8080
|
||||
server = HTTPServer(("0.0.0.0", port), HubHandler)
|
||||
|
||||
Reference in New Issue
Block a user