Files
esp32-cluster/coordinator/udp_ingest.py
T

249 lines
8.3 KiB
Python
Raw Normal View History

2026-04-03 14:38:38 +03:00
#!/usr/bin/env python3
"""
WiFi Recon Coordinator — UDP Ingestor
Receives beacon and probe events from ESP32 nodes over UDP, writes to SQLite.
This process has one job: ingest. Dashboard is separate.
"""
import socket
import json
import sqlite3
import datetime
import os
UDP_IP = "0.0.0.0"
UDP_PORT = 5005
DB_PATH = os.path.join(os.path.dirname(__file__), "events.db")
# ─── Database ────────────────────────────────────────────────────────────────
def init_db(path: str) -> sqlite3.Connection:
conn = sqlite3.connect(path, check_same_thread=False)
conn.execute("""
CREATE TABLE IF NOT EXISTS beacon_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
received_at TEXT NOT NULL,
node_id TEXT,
node_ts INTEGER,
ssid TEXT,
bssid TEXT,
rssi INTEGER,
channel INTEGER,
encryption TEXT,
importance TEXT,
confidence TEXT
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS probe_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
received_at TEXT NOT NULL,
node_id TEXT,
node_ts INTEGER,
src_mac TEXT,
ssid TEXT,
rssi INTEGER,
importance TEXT
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS heartbeat_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
received_at TEXT NOT NULL,
node_id TEXT,
node_ts INTEGER,
uptime_ms INTEGER,
free_heap INTEGER,
wifi_rssi INTEGER
)
""")
conn.commit()
return conn
def store_beacon(conn: sqlite3.Connection, ev: dict, received_at: str):
conn.execute("""
INSERT INTO beacon_events
(received_at, node_id, node_ts, ssid, bssid, rssi, channel, encryption, importance, confidence)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", (
received_at,
ev.get("node_id"),
ev.get("ts"),
ev.get("ssid"),
ev.get("bssid"),
ev.get("rssi"),
ev.get("ch"),
ev.get("enc"),
ev.get("imp"),
ev.get("conf"),
))
conn.commit()
def store_heartbeat(conn: sqlite3.Connection, ev: dict, received_at: str):
conn.execute("""
INSERT INTO heartbeat_events
(received_at, node_id, node_ts, uptime_ms, free_heap, wifi_rssi)
VALUES (?, ?, ?, ?, ?, ?)
""", (
received_at,
ev.get("node_id"),
ev.get("ts"),
ev.get("uptime_ms"),
ev.get("free_heap"),
ev.get("wifi_rssi"),
))
conn.commit()
def store_probe(conn: sqlite3.Connection, ev: dict, received_at: str):
conn.execute("""
INSERT INTO probe_events
(received_at, node_id, node_ts, src_mac, ssid, rssi, importance)
VALUES (?, ?, ?, ?, ?, ?, ?)
""", (
received_at,
ev.get("node_id"),
ev.get("ts"),
ev.get("src_mac"),
ev.get("ssid"),
ev.get("rssi"),
ev.get("imp"),
))
conn.commit()
# ─── Validation ──────────────────────────────────────────────────────────────
import re
_MAC_RE = re.compile(r'^([0-9A-Fa-f]{2}:){5}[0-9A-Fa-f]{2}$')
def _check_mac(val: str) -> bool:
return isinstance(val, str) and bool(_MAC_RE.match(val))
def _check_rssi(val) -> bool:
return isinstance(val, (int, float)) and -100 <= val <= 0
def validate_beacon(ev: dict) -> str | None:
"""Returns an error string if invalid, else None."""
for field in ("node_id", "bssid", "rssi", "ch"):
if ev.get(field) is None:
return f"missing field '{field}'"
if not _check_rssi(ev["rssi"]):
return f"rssi out of range: {ev['rssi']}"
if not isinstance(ev["ch"], int) or not (1 <= ev["ch"] <= 14):
return f"channel out of range: {ev['ch']}"
if not _check_mac(ev["bssid"]):
return f"bad bssid format: {ev['bssid']}"
return None
def validate_heartbeat(ev: dict) -> str | None:
for field in ("node_id", "uptime_ms", "free_heap", "wifi_rssi"):
if ev.get(field) is None:
return f"missing field '{field}'"
if not isinstance(ev["uptime_ms"], int) or ev["uptime_ms"] < 0:
return f"invalid uptime_ms: {ev['uptime_ms']}"
if not isinstance(ev["free_heap"], int) or ev["free_heap"] < 0:
return f"invalid free_heap: {ev['free_heap']}"
if not _check_rssi(ev["wifi_rssi"]):
return f"wifi_rssi out of range: {ev['wifi_rssi']}"
return None
def validate_probe(ev: dict) -> str | None:
"""Returns an error string if invalid, else None."""
for field in ("node_id", "src_mac", "rssi"):
if ev.get(field) is None:
return f"missing field '{field}'"
if not _check_rssi(ev["rssi"]):
return f"rssi out of range: {ev['rssi']}"
if not _check_mac(ev["src_mac"]):
return f"bad src_mac format: {ev['src_mac']}"
return None
# ─── Packet handling ─────────────────────────────────────────────────────────
IMP_COLOR = {
"high": "\033[92m",
"normal": "\033[0m",
"low": "\033[90m",
}
RESET = "\033[0m"
CYAN = "\033[96m"
YELLOW = "\033[93m"
def handle_packet(data: bytes, addr: tuple, conn: sqlite3.Connection):
try:
ev = json.loads(data.decode("utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError):
print(f"[WARN] Bad packet from {addr[0]}: {data[:80]}")
return
received_at = datetime.datetime.utcnow().isoformat(timespec="seconds")
pkt_type = ev.get("type", "beacon")
if pkt_type == "heartbeat":
err = validate_heartbeat(ev)
if err:
print(f"[DROP] {addr[0]} heartbeat — {err}")
return
store_heartbeat(conn, ev, received_at)
node = ev.get("node_id", "?")
uptime = ev.get("uptime_ms", 0) // 1000
heap = ev.get("free_heap", 0) // 1024
wrssi = ev.get("wifi_rssi", 0)
print(f"{YELLOW}[{received_at}] {node} HB ↑{uptime}s {heap}K heap AP:{wrssi}dBm{RESET}")
elif pkt_type == "probe":
err = validate_probe(ev)
if err:
print(f"[DROP] {addr[0]} probe — {err} raw={data[:120]}")
return
store_probe(conn, ev, received_at)
node = ev.get("node_id", "?")
src_mac = ev.get("src_mac", "?")
ssid = ev.get("ssid") or "<wildcard>"
rssi = ev.get("rssi", 0)
print(f"{CYAN}[{received_at}] {node} PROBE {src_mac}\"{ssid}\" {rssi:>4}dBm{RESET}")
else:
err = validate_beacon(ev)
if err:
print(f"[DROP] {addr[0]} beacon — {err} raw={data[:120]}")
return
store_beacon(conn, ev, received_at)
node = ev.get("node_id", "?")
ssid = ev.get("ssid") or "<hidden>"
bssid = ev.get("bssid", "?")
rssi = ev.get("rssi", 0)
ch = ev.get("ch", 0)
enc = ev.get("enc", "?")
imp = ev.get("imp", "normal")
color = IMP_COLOR.get(imp, "")
print(f"{color}[{received_at}] {node} {ssid:<32} {bssid} ch{ch:<3} {rssi:>4}dBm {enc}{RESET}")
# ─── Main ────────────────────────────────────────────────────────────────────
def main():
conn = init_db(DB_PATH)
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.bind((UDP_IP, UDP_PORT))
print(f"[INGESTOR] Listening on UDP :{UDP_PORT}")
print(f"[INGESTOR] Storing events to {DB_PATH}")
print(f"{'─' * 80}")
try:
while True:
data, addr = sock.recvfrom(1024)
handle_packet(data, addr, conn)
except KeyboardInterrupt:
print("\n[INGESTOR] Stopped.")
finally:
conn.close()
sock.close()
if __name__ == "__main__":
main()