Files
esp32-cluster/coordinator/udp_ingest.py
T

395 lines
14 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
2026-04-05 11:29:59 +03:00
import re
2026-04-03 14:38:38 +03:00
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)
2026-04-05 11:29:59 +03:00
conn.execute("PRAGMA journal_mode=WAL")
2026-04-03 14:38:38 +03:00
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 (
2026-04-05 13:06:25 +03:00
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 (
2026-04-03 14:38:38 +03:00
id INTEGER PRIMARY KEY AUTOINCREMENT,
received_at TEXT NOT NULL,
node_id TEXT,
node_ts INTEGER,
2026-04-05 13:06:25 +03:00
subtype TEXT,
src TEXT,
dst TEXT,
bssid TEXT,
reason INTEGER,
rssi INTEGER
2026-04-03 14:38:38 +03:00
)
""")
conn.execute("""
2026-04-05 13:06:25 +03:00
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,
2026-04-05 13:06:25 +03:00
ssid TEXT,
rssi INTEGER
)
""")
2026-04-05 13:06:25 +03:00
# 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);
""")
2026-04-03 14:38:38 +03:00
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
2026-04-05 13:06:25 +03:00
(received_at, node_id, node_ts, uptime_ms, free_heap, wifi_rssi,
probe_drops, deauth_drops, assoc_drops)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
2026-04-03 14:38:38 +03:00
""", (
received_at,
ev.get("node_id"),
ev.get("ts"),
ev.get("uptime_ms"),
ev.get("free_heap"),
ev.get("wifi_rssi"),
2026-04-05 13:06:25 +03:00
ev.get("probe_drops"),
ev.get("deauth_drops"),
ev.get("assoc_drops"),
2026-04-03 14:38:38 +03:00
))
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 ──────────────────────────────────────────────────────────────
_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"),
))
conn.commit()
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
2026-04-05 13:06:25 +03:00
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"),
))
conn.commit()
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
2026-04-03 14:38:38 +03:00
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")
2026-04-05 13:06:25 +03:00
if 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 "<hidden>"
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":
2026-04-03 14:38:38 +03:00
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()