#!/usr/bin/env python3 """Busfleet Hub — fleet registration + dashboard server. Listens on :8080. Routes: GET / — fleet dashboard (HTML) GET /health — health check GET /api/register/ — 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 """ import json import os import socket import subprocess import sys from http.server import HTTPServer, BaseHTTPRequestHandler REGISTER_SH = "/opt/busfleet-hub/register.sh" 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): try: s = socket.create_connection(("127.0.0.1", port), timeout=2) s.close() return True except Exception: 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.""" def do_GET(self): # Fleet dashboard (homepage) if self.path in ("/", "/index.html"): self.serve_file(DASHBOARD, "text/html; charset=utf-8") return # Health check if self.path == "/health": self.send_json(200, '{"status":"ok","service":"busfleet-hub"}') return # Device registration if self.path.startswith("/api/register/"): device_id = self.path.split("/api/register/", 1)[1].strip("/") if not device_id: self.send_json_error(400, "Missing device ID") return self.call_register_script(device_id) return # 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): # Key authorization if self.path.startswith("/api/authorize-key"): content_length = int(self.headers.get("Content-Length", 0)) body = self.rfile.read(content_length).decode() if content_length else "{}" try: data = json.loads(body) except json.JSONDecodeError: self.send_json_error(400, "Invalid JSON body") return device_id = data.get("device_id", "").strip() pubkey = data.get("pubkey", "").strip() if not device_id or not pubkey: self.send_json_error(400, "Missing device_id or pubkey") return if not pubkey.startswith(("ssh-", "ecdsa-", "sk-")) or len(pubkey) < 80: self.send_json_error(400, "Invalid SSH public key format") return try: result = subprocess.run( [REGISTER_SH, "--authorize-key", device_id, pubkey], capture_output=True, text=True, timeout=15, ) status = 200 if result.returncode == 0 else 500 self.send_json(status, result.stdout.strip()) except subprocess.TimeoutExpired: self.send_json_error(504, "Key authorization timed out") except Exception as exc: self.send_json_error(500, str(exc)) return self.send_json_error(404, "Not found") # ── Helpers ──────────────────────────────────────────────────── def serve_file(self, path, content_type): """Serve a static file.""" if not os.path.isfile(path): self.send_json_error(404, "File not found") return with open(path, "rb") as f: data = f.read() self.send_response(200) self.send_header("Content-Type", content_type) self.send_header("Content-Length", str(len(data))) self.end_headers() self.wfile.write(data) def call_register_script(self, device_id): """Run register.sh and return its output.""" try: result = subprocess.run( [REGISTER_SH, device_id], capture_output=True, text=True, timeout=15, ) status = 200 if result.returncode == 0 else 500 self.send_json(status, result.stdout.strip()) except subprocess.TimeoutExpired: self.send_json_error(504, "Registration timed out") except Exception as exc: self.send_json_error(500, str(exc)) def serve_fleet_status(self): """Build and return fleet status from port-registry.json + live tunnel checks.""" devices = [] try: with open(REGISTRY) as f: registry = json.load(f) except Exception: registry = {} for dev_id, entry in registry.items(): port = entry.get("tunnel_port", 0) tunnel_up = _port_is_open(port) if port else False 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, }) # Sort: tunnel_up first, then by device_id devices.sort(key=lambda d: (not d["tunnel_up"], d["device_id"])) self.send_json(200, json.dumps({ "hub": HUB_IP, "device_count": len(devices), "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) self.send_header("Content-Type", "application/json") self.send_header("Access-Control-Allow-Origin", "*") self.end_headers() self.wfile.write(body.encode() if isinstance(body, str) else body) def send_json_error(self, status, message): """Send a JSON error response.""" self.send_json(status, json.dumps({"error": message})) def log_message(self, fmt, *args): 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) print(f"Busfleet Hub listening on :{port}", flush=True) print(f" Dashboard: http://{HUB_IP}:{port}/", flush=True) try: server.serve_forever() except KeyboardInterrupt: print("\nShutting down.", flush=True) server.shutdown() if __name__ == "__main__": main()