"""Prometheus exporter over the simulator's Parquet output. Periodically scans DATA_DIR with DuckDB and exposes fleet statistics as /metrics. Zero dependencies beyond duckdb: the exposition format is plain text, served with the standard library HTTP server. """ from __future__ import annotations import os import threading import time from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path import duckdb DATA_DIR = Path(os.environ.get("DATA_DIR", "./data")) PORT = int(os.environ.get("EXPORTER_PORT", "9105")) SCAN_INTERVAL_S = float(os.environ.get("SCAN_INTERVAL_S", "5")) _lock = threading.Lock() _payload = "# swarm exporter starting\n" def _q(con: duckdb.DuckDBPyConnection, sql: str) -> list[tuple]: try: return con.sql(sql).fetchall() except duckdb.Error: return [] # partitions may not exist yet while the swarm warms up def collect() -> str: con = duckdb.connect() # union_by_name: sensors have different schemas under one dataset glob glob = lambda ds: ( # noqa: E731 f"'{DATA_DIR}/dataset={ds}/**/*.parquet', hive_partitioning=true, union_by_name=true" ) lines: list[str] = [] def metric(name: str, help_text: str, mtype: str, rows: list[str]) -> None: lines.append(f"# HELP {name} {help_text}") lines.append(f"# TYPE {name} {mtype}") lines.extend(rows) metric( "swarm_rows_total", "Telemetry rows written per drone and sensor", "gauge", [f'swarm_rows_total{{drone="{d}",sensor="{s}"}} {n}' for d, s, n in _q(con, f"SELECT drone, sensor, count(*) FROM read_parquet({glob('telemetry')}) GROUP BY 1,2")], ) metric( "swarm_detections_total", "Detection events per drone and class", "gauge", [f'swarm_detections_total{{drone="{d}",cls="{c}"}} {n}' for d, c, n in _q(con, f"SELECT drone, cls, count(*) FROM read_parquet({glob('detections')}) GROUP BY 1,2")], ) metric( "swarm_state_frames_total", "Pose broadcast frames per drone and direction", "gauge", [f'swarm_state_frames_total{{drone="{d}",direction="{dr}"}} {n}' for d, dr, n in _q(con, f"SELECT drone, direction, count(*) FROM read_parquet({glob('state')}) GROUP BY 1,2")], ) metric( "swarm_battery_pct", "Latest battery level per drone", "gauge", [f'swarm_battery_pct{{drone="{d}"}} {v}' for d, v in _q(con, f""" SELECT drone, arg_max(level_pct, ts_ns) FROM read_parquet({glob('telemetry')}) WHERE sensor='battery' GROUP BY drone""")], ) metric( "swarm_rssi_dbm", "Latest RSSI per drone-peer link", "gauge", [f'swarm_rssi_dbm{{drone="{d}",peer="{p}"}} {v}' for d, p, v in _q(con, f""" SELECT drone, peer_id, arg_max(rssi_dbm, ts_ns) FROM read_parquet({glob('telemetry')}) WHERE sensor='rssi' GROUP BY drone, peer_id""")], ) metric( "swarm_peer_distance_m", "Latest inter-drone distance estimate", "gauge", [f'swarm_peer_distance_m{{drone="{d}",peer="{p}"}} {v}' for d, p, v in _q(con, f""" SELECT drone, peer_id, arg_max(distance_m, ts_ns) FROM read_parquet({glob('telemetry')}) WHERE sensor='rssi' GROUP BY drone, peer_id""")], ) files = list(DATA_DIR.rglob("*.parquet")) metric( "swarm_parquet_bytes", "Bytes on disk per dataset", "gauge", [f'swarm_parquet_bytes{{dataset="{ds}"}} {sum(f.stat().st_size for f in files if f"dataset={ds}" in str(f))}' for ds in ("telemetry", "detections", "state")], ) metric("swarm_parquet_files", "Parquet files on disk", "gauge", [f"swarm_parquet_files {len(files)}"]) con.close() return "\n".join(lines) + "\n" def scanner() -> None: global _payload while True: started = time.monotonic() try: payload = collect() except Exception as exc: # keep serving stale metrics over dying payload = f"# collect error: {exc}\n" with _lock: _payload = payload time.sleep(max(0.5, SCAN_INTERVAL_S - (time.monotonic() - started))) class Handler(BaseHTTPRequestHandler): def do_GET(self) -> None: # noqa: N802 — http.server API if self.path != "/metrics": self.send_response(404) self.end_headers() return with _lock: body = _payload.encode() self.send_response(200) self.send_header("Content-Type", "text/plain; version=0.0.4") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def log_message(self, *_args: object) -> None: pass # scrapes every few seconds; keep the log quiet if __name__ == "__main__": threading.Thread(target=scanner, daemon=True).start() print(f"swarm exporter on :{PORT}/metrics, scanning {DATA_DIR} every {SCAN_INTERVAL_S}s") ThreadingHTTPServer(("", PORT), Handler).serve_forever()