#!/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 import re 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("PRAGMA journal_mode=WAL") conn.execute("PRAGMA wal_autocheckpoint=500") 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 ) """) 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, probe_drops INTEGER, deauth_drops INTEGER, assoc_drops INTEGER ) """) conn.execute(""" CREATE TABLE IF NOT EXISTS deauth_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, received_at TEXT NOT NULL, node_id TEXT, node_ts INTEGER, subtype TEXT, src TEXT, dst TEXT, bssid TEXT, reason INTEGER, rssi INTEGER ) """) conn.execute(""" CREATE TABLE IF NOT EXISTS assoc_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, received_at TEXT NOT NULL, node_id TEXT, node_ts INTEGER, subtype TEXT, src TEXT, bssid TEXT, ssid TEXT, rssi INTEGER ) """) conn.execute(""" CREATE TABLE IF NOT EXISTS ble_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, received_at TEXT NOT NULL, node_id TEXT, node_ts INTEGER, mac TEXT, addr_type TEXT, name TEXT, rssi INTEGER, mfr_id INTEGER, mfr_data TEXT ) """) # Safe migration: add ble_drops to heartbeat table if not present yet try: conn.execute("ALTER TABLE heartbeat_events ADD COLUMN ble_drops INTEGER") conn.commit() except sqlite3.OperationalError: pass # Indexes for common query patterns conn.executescript(""" CREATE INDEX IF NOT EXISTS idx_beacon_received_at ON beacon_events (received_at); CREATE INDEX IF NOT EXISTS idx_beacon_node_id ON beacon_events (node_id); CREATE INDEX IF NOT EXISTS idx_beacon_bssid ON beacon_events (bssid); CREATE INDEX IF NOT EXISTS idx_probe_received_at ON probe_events (received_at); CREATE INDEX IF NOT EXISTS idx_probe_node_id ON probe_events (node_id); CREATE INDEX IF NOT EXISTS idx_probe_src_mac ON probe_events (src_mac); CREATE INDEX IF NOT EXISTS idx_heartbeat_node_id ON heartbeat_events (node_id); CREATE INDEX IF NOT EXISTS idx_deauth_received_at ON deauth_events (received_at); CREATE INDEX IF NOT EXISTS idx_deauth_node_id ON deauth_events (node_id); CREATE INDEX IF NOT EXISTS idx_deauth_bssid ON deauth_events (bssid); CREATE INDEX IF NOT EXISTS idx_deauth_src ON deauth_events (src); CREATE INDEX IF NOT EXISTS idx_deauth_dst ON deauth_events (dst); CREATE INDEX IF NOT EXISTS idx_assoc_received_at ON assoc_events (received_at); CREATE INDEX IF NOT EXISTS idx_assoc_node_id ON assoc_events (node_id); CREATE INDEX IF NOT EXISTS idx_assoc_src ON assoc_events (src); CREATE INDEX IF NOT EXISTS idx_assoc_bssid ON assoc_events (bssid); CREATE INDEX IF NOT EXISTS idx_ble_received_at ON ble_events (received_at); CREATE INDEX IF NOT EXISTS idx_ble_node_id ON ble_events (node_id); CREATE INDEX IF NOT EXISTS idx_ble_mac ON ble_events (mac); CREATE INDEX IF NOT EXISTS idx_ble_mfr_id ON ble_events (mfr_id); """) 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) 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"), )) 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, probe_drops, deauth_drops, assoc_drops, ble_drops) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( received_at, ev.get("node_id"), ev.get("ts"), ev.get("uptime_ms"), ev.get("free_heap"), ev.get("wifi_rssi"), ev.get("probe_drops"), ev.get("deauth_drops"), ev.get("assoc_drops"), ev.get("ble_drops"), )) 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"), )) # ─── Validation ────────────────────────────────────────────────────────────── _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 store_deauth(conn: sqlite3.Connection, ev: dict, received_at: str): conn.execute(""" INSERT INTO deauth_events (received_at, node_id, node_ts, subtype, src, dst, bssid, reason, rssi) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( received_at, ev.get("node_id"), ev.get("ts"), ev.get("type"), ev.get("src"), ev.get("dst"), ev.get("bssid"), ev.get("reason"), ev.get("rssi"), )) def validate_deauth(ev: dict) -> str | None: for field in ("node_id", "src", "dst", "bssid", "reason", "rssi"): if ev.get(field) is None: return f"missing field '{field}'" if not _check_mac(ev["src"]): return f"bad src format: {ev['src']}" if not _check_mac(ev["dst"]): return f"bad dst format: {ev['dst']}" if not _check_mac(ev["bssid"]): return f"bad bssid format: {ev['bssid']}" if not isinstance(ev["reason"], int) or not (0 <= ev["reason"] <= 65535): return f"invalid reason code: {ev['reason']}" if not _check_rssi(ev["rssi"]): return f"rssi out of range: {ev['rssi']}" return None def store_assoc(conn: sqlite3.Connection, ev: dict, received_at: str): conn.execute(""" INSERT INTO assoc_events (received_at, node_id, node_ts, subtype, src, bssid, ssid, rssi) VALUES (?, ?, ?, ?, ?, ?, ?, ?) """, ( received_at, ev.get("node_id"), ev.get("ts"), ev.get("type"), ev.get("src"), ev.get("bssid"), ev.get("ssid", ""), ev.get("rssi"), )) def validate_assoc(ev: dict) -> str | None: for field in ("node_id", "src", "bssid", "rssi"): if ev.get(field) is None: return f"missing field '{field}'" if not _check_mac(ev["src"]): return f"bad src format: {ev['src']}" if not _check_mac(ev["bssid"]): return f"bad bssid format: {ev['bssid']}" if not _check_rssi(ev["rssi"]): return f"rssi out of range: {ev['rssi']}" return None def store_ble(conn: sqlite3.Connection, ev: dict, received_at: str): mfr_id = ev.get("mfr_id") conn.execute(""" INSERT INTO ble_events (received_at, node_id, node_ts, mac, addr_type, name, rssi, mfr_id, mfr_data) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( received_at, ev.get("node_id"), ev.get("ts"), ev.get("mac"), "public" if ev.get("addr_type") == 0 else "random", ev.get("name", ""), ev.get("rssi"), mfr_id if mfr_id is not None and mfr_id >= 0 else None, ev.get("mfr_data", ""), )) def validate_ble(ev: dict) -> str | None: for field in ("node_id", "mac", "rssi"): if ev.get(field) is None: return f"missing field '{field}'" if not _check_mac(ev["mac"]): return f"bad mac format: {ev['mac']}" if not _check_rssi(ev["rssi"]): return f"rssi out of range: {ev['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 ───────────────────────────────────────────────────────── MAGENTA = "\033[95m" 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 == "ble_adv": err = validate_ble(ev) if err: print(f"[DROP] {addr[0]} ble_adv — {err}") return store_ble(conn, ev, received_at) node = ev.get("node_id", "?") mac = ev.get("mac", "?") name = ev.get("name") or "" rssi = ev.get("rssi", 0) mfr_id = ev.get("mfr_id", -1) at = "pub" if ev.get("addr_type") == 0 else "rnd" mfr_str = f"mfr=0x{mfr_id:04X}" if mfr_id is not None and mfr_id >= 0 else "mfr=—" print(f"{MAGENTA}[{received_at}] {node} BLE {mac} ({at}) \"{name}\" {rssi:>4}dBm {mfr_str}{RESET}") elif pkt_type in ("assoc", "reassoc"): err = validate_assoc(ev) if err: print(f"[DROP] {addr[0]} {pkt_type} — {err}") return store_assoc(conn, ev, received_at) node = ev.get("node_id", "?") src = ev.get("src", "?") bssid = ev.get("bssid", "?") ssid = ev.get("ssid") or "" rssi = ev.get("rssi", 0) print(f"\033[94m[{received_at}] {node} {pkt_type.upper():<8} {src} → {bssid} \"{ssid}\" {rssi:>4}dBm{RESET}") elif pkt_type in ("deauth", "disassoc"): err = validate_deauth(ev) if err: print(f"[DROP] {addr[0]} {pkt_type} — {err}") return store_deauth(conn, ev, received_at) node = ev.get("node_id", "?") src = ev.get("src", "?") dst = ev.get("dst", "?") bssid = ev.get("bssid", "?") reason = ev.get("reason", 0) rssi = ev.get("rssi", 0) print(f"\033[91m[{received_at}] {node} {pkt_type.upper():<8} {src} → {dst} bssid={bssid} reason={reason} {rssi:>4}dBm{RESET}") elif 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 "" 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 "" 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}") conn.commit() # ─── 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()