#!/usr/bin/env python3
"""Aethelos Privacy Guard — local engine.

Zero extra packages. Binds 127.0.0.1 only. Reads this device's sockets.
Optional tcpdump (Termux / Linux) for live headers. Optional PostgreSQL
via AETHELOS_DATABASE_URL; otherwise a JSON file beside this script.

The command board sends Approve / Kick here. Firewall drops run on THIS
host only (ss -K, iptables/nft/ufw, netsh). No remote payload, no Tor,
no third-party disk.
"""
from __future__ import annotations

import argparse
import ipaddress
import json
import os
import re
import shutil
import socket
import subprocess
import sys
import threading
import uuid
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from typing import Any
from urllib.parse import unquote, urlparse

HOST = "127.0.0.1"
PORT = int(os.environ.get("AETHELOS_GUARD_PORT", "8787"))
STORE_PATH = Path(
    os.environ.get("AETHELOS_GUARD_STORE", Path(__file__).resolve().with_name("guard-store.json"))
)
ALLOW_ORIGIN = os.environ.get("AETHELOS_ORIGIN", "*")
DEVICE_ID = os.environ.get("AETHELOS_DEVICE_ID") or socket.gethostname() or "local"
SCHEMA_PATH = Path(__file__).resolve().with_name("guard-schema.sql")
LOCK = threading.RLock()

SCHEMA_SQL = r"""
CREATE TABLE IF NOT EXISTS devices (
  device_id TEXT PRIMARY KEY, name TEXT NOT NULL DEFAULT 'local',
  type TEXT NOT NULL DEFAULT 'hub', ip TEXT NOT NULL DEFAULT '',
  online BOOLEAN NOT NULL DEFAULT TRUE, last_seen TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE TABLE IF NOT EXISTS connections (
  id TEXT PRIMARY KEY, device_id TEXT NOT NULL DEFAULT 'local',
  remote_ip TEXT NOT NULL, remote_port INTEGER NOT NULL DEFAULT 0,
  local_port INTEGER NOT NULL DEFAULT 0, app TEXT NOT NULL DEFAULT 'unknown',
  proto TEXT NOT NULL DEFAULT 'tcp', status TEXT NOT NULL DEFAULT 'pending',
  first_seen TIMESTAMPTZ NOT NULL DEFAULT NOW(), last_seen TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS connections_status_idx ON connections (status);
CREATE TABLE IF NOT EXISTS alerts (
  id TEXT PRIMARY KEY, severity TEXT NOT NULL DEFAULT 'INFO',
  message TEXT NOT NULL, device_id TEXT NOT NULL DEFAULT '',
  conn_id TEXT NOT NULL DEFAULT '', timestamp TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE TABLE IF NOT EXISTS firewall_rules (
  id TEXT PRIMARY KEY, remote_ip TEXT NOT NULL, remote_port INTEGER NOT NULL DEFAULT 0,
  action TEXT NOT NULL DEFAULT 'drop', permanent BOOLEAN NOT NULL DEFAULT TRUE,
  reason TEXT NOT NULL DEFAULT '', applied BOOLEAN NOT NULL DEFAULT FALSE,
  notes TEXT NOT NULL DEFAULT '', created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE TABLE IF NOT EXISTS packets (
  id BIGSERIAL PRIMARY KEY, seen_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
  src TEXT NOT NULL, dst TEXT NOT NULL, sport INTEGER NOT NULL DEFAULT 0,
  dport INTEGER NOT NULL DEFAULT 0, proto TEXT NOT NULL DEFAULT 'tcp',
  bytes INTEGER NOT NULL DEFAULT 0, app TEXT NOT NULL DEFAULT ''
);
"""

DUMP_LINE = re.compile(
    r"""(?:IP6?)\s+(\S+)\s+>\s+(\S+)(?::\s+(\w+))?""",
    re.IGNORECASE,
)


def _db_url() -> str:
    """Local Postgres only. Platform cloud URLs are ignored unless AETHELOS_USE_PG=1."""
    explicit = (os.environ.get("AETHELOS_DATABASE_URL") or "").strip()
    if explicit:
        return explicit
    url = (os.environ.get("DATABASE_URL") or "").strip()
    if not url:
        return ""
    if os.environ.get("AETHELOS_USE_PG", "").lower() in {"1", "true", "yes"}:
        return url
    lowered = url.lower()
    if any(h in lowered for h in ("127.0.0.1", "localhost", "/var/run/postgresql", "@postgres:")):
        return url
    return ""


def now() -> str:
    return datetime.now(timezone.utc).replace(microsecond=0).isoformat()


def new_id(prefix: str = "id") -> str:
    return f"{prefix}-{uuid.uuid4().hex[:12]}"


def valid_ip(raw: str) -> str | None:
    try:
        return str(ipaddress.ip_address((raw or "").strip().strip("[]")))
    except ValueError:
        return None


def valid_port(raw: Any) -> int:
    try:
        n = int(raw)
    except (TypeError, ValueError):
        return 0
    return n if 0 <= n <= 65535 else 0


def is_loopback(ip: str) -> bool:
    try:
        return ipaddress.ip_address(ip).is_loopback
    except ValueError:
        return ip in {"127.0.0.1", "::1", "localhost"}


def split_hostport(token: str) -> tuple[str, int]:
    token = token.strip().rstrip(":")
    token = re.sub(r":\s*(Flags|seq|ack|win|length|cksum).*$", "", token, flags=re.I)
    if token.startswith("["):
        host, _, rest = token[1:].partition("]")
        port = rest.lstrip(".").lstrip(":")
        return host, valid_port(port)
    if token.count(":") > 1:
        host, _, port_s = token.rpartition(".")
        if not host:
            host, _, port_s = token.rpartition(":")
        return host, valid_port(port_s)
    host, _, port_s = token.rpartition(".")
    if not host:
        host, _, port_s = token.rpartition(":")
    return host.strip("[]"), valid_port(port_s)


def analyze(conns: list[dict]) -> dict:
    total = len(conns) or 1
    unique = {f"{c.get('remote_ip')}:{c.get('remote_port')}" for c in conns}
    ports = [int(c.get("remote_port") or 0) for c in conns]
    newish = sum(1 for c in conns if c.get("status") == "pending")
    port_entropy = (len(set(ports)) / len(ports)) if ports else 0.0
    score = min(1.0, (newish / total) * 0.5 + port_entropy * 0.3 + min(len(unique), 20) / 40)
    anomalies: list[dict] = []
    if newish > 2:
        anomalies.append(
            {
                "type": "new_peer_drift",
                "severity": "HIGH" if score > 0.7 else "MEDIUM",
                "message": f"{newish} pending/new peers detected (Relativistic Self-Monitoring)",
            }
        )
    if port_entropy > 0.6 and len(unique) > 3:
        anomalies.append(
            {
                "type": "forensic_geometry_port_scan",
                "severity": "HIGH",
                "message": "High port diversity + multiple peers — possible reconnaissance",
            }
        )
    rec = "block" if score >= 0.75 else "monitor" if score > 0.4 else "approve"
    return {
        "anomaly_score": round(score, 3),
        "recommendation": rec,
        "peer_count": len(unique),
        "new_pending": newish,
        "port_entropy": round(port_entropy, 2),
        "anomalies": anomalies,
        "protocol_notes": "Local capture only. Approve lets a peer stay. Kick drops it on this host.",
    }


class EngineStore:
    def __init__(self) -> None:
        self.kind = "json"
        self._pg = None
        self._psycopg = None
        self.url = _db_url()
        self.data: dict[str, Any] = {
            "devices": {},
            "connections": {},
            "alerts": [],
            "blocked": [],
            "approved": [],
            "rules": [],
            "packets": [],
        }
        if self.url:
            self._open_pg()
        if self.kind == "json":
            self._load_json()
        self._touch_device(DEVICE_ID, "This device", "hub", "127.0.0.1")

    def _open_pg(self) -> None:
        for name in ("psycopg2", "psycopg"):
            try:
                mod = __import__(name)
                conn = mod.connect(self.url)
                conn.autocommit = True
                self._psycopg = mod
                self._pg = conn
                self.kind = "postgres"
                self._exec_script(SCHEMA_SQL)
                if SCHEMA_PATH.is_file():
                    self._exec_script(SCHEMA_PATH.read_text(encoding="utf-8"))
                return
            except Exception:
                self._pg = None
                continue

    def _exec_script(self, sql: str) -> None:
        if not self._pg:
            return
        cur = self._pg.cursor()
        try:
            cur.execute(sql)
        finally:
            cur.close()

    def _q(self, sql: str, params: tuple = ()) -> list[tuple]:
        if not self._pg:
            return []
        cur = self._pg.cursor()
        try:
            cur.execute(sql, params)
            if cur.description:
                return list(cur.fetchall())
            return []
        finally:
            cur.close()

    def _load_json(self) -> None:
        if STORE_PATH.is_file():
            try:
                raw = json.loads(STORE_PATH.read_text(encoding="utf-8"))
                if isinstance(raw, dict):
                    for k in self.data:
                        if k in raw:
                            self.data[k] = raw[k]
            except (OSError, json.JSONDecodeError):
                pass

    def _save_json(self) -> None:
        STORE_PATH.parent.mkdir(parents=True, exist_ok=True)
        tmp = STORE_PATH.with_suffix(".tmp")
        tmp.write_text(json.dumps(self.data, indent=2), encoding="utf-8")
        tmp.replace(STORE_PATH)

    def _touch_device(self, device_id: str, name: str, kind: str, ip: str) -> None:
        if self.kind == "postgres":
            self._q(
                """INSERT INTO devices (device_id, name, type, ip, online, last_seen)
                   VALUES (%s,%s,%s,%s, TRUE, NOW())
                   ON CONFLICT (device_id) DO UPDATE SET
                     last_seen = NOW(), online = TRUE, ip = EXCLUDED.ip""",
                (device_id, name, kind, ip),
            )
            return
        dev = self.data["devices"].get(device_id) or {
            "device_id": device_id,
            "name": name,
            "type": kind,
            "ip": ip,
            "online": True,
        }
        dev.update(
            {"name": name, "type": kind, "ip": ip or dev.get("ip", ""), "online": True, "last_seen": now()}
        )
        self.data["devices"][device_id] = dev

    def upsert_conn(self, rec: dict) -> dict:
        cid = rec["id"]
        status = rec.get("status") or "pending"
        if self.kind == "postgres":
            rows = self._q("SELECT status, first_seen FROM connections WHERE id = %s", (cid,))
            if rows:
                status = rows[0][0]
                first = rows[0][1]
            else:
                blocked = self._q(
                    "SELECT 1 FROM firewall_rules WHERE remote_ip = %s AND action = 'drop' LIMIT 1",
                    (rec["remote_ip"],),
                )
                approved = self._q(
                    "SELECT 1 FROM connections WHERE remote_ip = %s AND status = 'approved' LIMIT 1",
                    (rec["remote_ip"],),
                )
                if blocked:
                    status = "blocked"
                elif approved:
                    status = "approved"
                first = now()
            first_s = first if isinstance(first, str) else now()
            self._q(
                """INSERT INTO connections
                   (id, device_id, remote_ip, remote_port, local_port, app, proto, status, first_seen, last_seen)
                   VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,NOW())
                   ON CONFLICT (id) DO UPDATE SET last_seen = NOW(), app = EXCLUDED.app""",
                (
                    cid,
                    rec.get("device_id", DEVICE_ID),
                    rec["remote_ip"],
                    rec["remote_port"],
                    rec.get("local_port", 0),
                    rec.get("app", "unknown"),
                    rec.get("proto", "tcp"),
                    status,
                    first_s,
                ),
            )
            rec = dict(rec)
            rec["status"] = status
            return rec
        existing = self.data["connections"].get(cid)
        peer = f"{rec['remote_ip']}:{rec['remote_port']}"
        if existing:
            existing["last_seen"] = now()
            existing["app"] = rec.get("app") or existing.get("app") or "unknown"
            rec = existing
        else:
            if rec["remote_ip"] in self.data["blocked"] or peer in self.data["blocked"]:
                status = "blocked"
            elif rec["remote_ip"] in self.data["approved"] or peer in self.data["approved"]:
                status = "approved"
            rec = {
                "id": cid,
                "device_id": rec.get("device_id") or DEVICE_ID,
                "remote_ip": rec["remote_ip"],
                "remote_port": rec["remote_port"],
                "local_port": rec.get("local_port") or 0,
                "app": rec.get("app") or "unknown",
                "proto": rec.get("proto") or "tcp",
                "status": status,
                "first_seen": now(),
                "last_seen": now(),
            }
            self.data["connections"][cid] = rec
            if status == "pending":
                self._alert("HIGH", f"New connection on {rec['device_id']}: {peer} via {rec['app']}", rec)
        self._save_json()
        return rec

    def _alert(self, severity: str, message: str, conn: dict | None = None) -> dict:
        item = {
            "id": new_id("a"),
            "severity": severity,
            "message": message,
            "device_id": (conn or {}).get("device_id", ""),
            "conn_id": (conn or {}).get("id", ""),
            "timestamp": now(),
        }
        if self.kind == "postgres":
            self._q(
                """INSERT INTO alerts (id, severity, message, device_id, conn_id, timestamp)
                   VALUES (%s,%s,%s,%s,%s,NOW())""",
                (item["id"], severity, message, item["device_id"], item["conn_id"]),
            )
            self._q(
                "DELETE FROM alerts WHERE id NOT IN (SELECT id FROM alerts ORDER BY timestamp DESC LIMIT 200)"
            )
            return item
        self.data["alerts"].insert(0, item)
        self.data["alerts"] = self.data["alerts"][:200]
        return item

    def list_devices(self) -> list[dict]:
        if self.kind == "postgres":
            rows = self._q("SELECT device_id, name, type, ip, online, last_seen FROM devices")
            return [
                {
                    "device_id": r[0],
                    "name": r[1],
                    "type": r[2],
                    "ip": r[3],
                    "online": bool(r[4]),
                    "last_seen": r[5].isoformat() if hasattr(r[5], "isoformat") else str(r[5]),
                }
                for r in rows
            ]
        return list(self.data["devices"].values())

    def list_conns(self) -> list[dict]:
        if self.kind == "postgres":
            rows = self._q(
                """SELECT id, device_id, remote_ip, remote_port, local_port, app, proto, status, first_seen, last_seen
                   FROM connections ORDER BY last_seen DESC LIMIT 400"""
            )
            out = []
            for r in rows:
                out.append(
                    {
                        "id": r[0],
                        "device_id": r[1],
                        "remote_ip": r[2],
                        "remote_port": int(r[3]),
                        "local_port": int(r[4] or 0),
                        "app": r[5],
                        "proto": r[6],
                        "status": r[7],
                        "first_seen": r[8].isoformat() if hasattr(r[8], "isoformat") else str(r[8]),
                        "last_seen": r[9].isoformat() if hasattr(r[9], "isoformat") else str(r[9]),
                    }
                )
            return out
        return sorted(self.data["connections"].values(), key=lambda c: c.get("last_seen", ""), reverse=True)

    def list_alerts(self) -> list[dict]:
        if self.kind == "postgres":
            rows = self._q(
                "SELECT id, severity, message, device_id, conn_id, timestamp FROM alerts ORDER BY timestamp DESC LIMIT 80"
            )
            return [
                {
                    "id": r[0],
                    "severity": r[1],
                    "message": r[2],
                    "device_id": r[3],
                    "conn_id": r[4],
                    "timestamp": r[5].isoformat() if hasattr(r[5], "isoformat") else str(r[5]),
                }
                for r in rows
            ]
        return list(self.data["alerts"][:80])

    def list_rules(self) -> list[dict]:
        if self.kind == "postgres":
            rows = self._q(
                """SELECT id, remote_ip, remote_port, action, permanent, reason, applied, notes, created_at
                   FROM firewall_rules ORDER BY created_at DESC LIMIT 200"""
            )
            return [
                {
                    "id": r[0],
                    "remote_ip": r[1],
                    "remote_port": int(r[2]),
                    "action": r[3],
                    "permanent": bool(r[4]),
                    "reason": r[5],
                    "applied": bool(r[6]),
                    "notes": r[7],
                    "created_at": r[8].isoformat() if hasattr(r[8], "isoformat") else str(r[8]),
                }
                for r in rows
            ]
        return list(self.data["rules"][:200])

    def _find_conn(self, cid: str, ip: str, port: int) -> dict | None:
        if cid:
            if self.kind == "postgres":
                rows = self._q(
                    """SELECT id, device_id, remote_ip, remote_port, app, proto, status
                       FROM connections WHERE id = %s""",
                    (cid,),
                )
                if rows:
                    r = rows[0]
                    return {
                        "id": r[0],
                        "device_id": r[1],
                        "remote_ip": r[2],
                        "remote_port": int(r[3]),
                        "app": r[4],
                        "proto": r[5],
                        "status": r[6],
                    }
            found = self.data["connections"].get(cid)
            if found:
                return found
        if ip:
            for c in self.list_conns():
                if c["remote_ip"] == ip and (not port or c["remote_port"] == port):
                    return c
        return None

    def approve(self, cid: str, ip: str, port: int) -> dict:
        conn = self._find_conn(cid, ip, port)
        ip = (conn or {}).get("remote_ip") or ip
        port = int((conn or {}).get("remote_port") or port or 0)
        if not conn and ip:
            conn = self.upsert_conn(
                {
                    "id": cid or f"{DEVICE_ID}:{ip}:{port}",
                    "device_id": DEVICE_ID,
                    "remote_ip": ip,
                    "remote_port": port,
                    "app": "manual",
                    "proto": "tcp",
                    "status": "approved",
                }
            )
        peer = f"{ip}:{port}"
        if self.kind == "postgres":
            if conn:
                self._q("UPDATE connections SET status = 'approved' WHERE id = %s", (conn["id"],))
            self._alert("INFO", f"Approved {peer}", conn)
            conn = self._find_conn((conn or {}).get("id", cid), ip, port) or conn
        else:
            if conn:
                conn["status"] = "approved"
                self.data["connections"][conn["id"]] = conn
            if ip and ip not in self.data["approved"]:
                self.data["approved"].append(ip)
            if peer not in self.data["approved"]:
                self.data["approved"].append(peer)
            self._alert("INFO", f"Approved {peer}", conn)
            self._save_json()
        return {"ok": True, "action": "approve", "id": (conn or {}).get("id", cid), "ip": ip, "port": port}

    def kick(self, cid: str, ip: str, port: int, reason: str, permanent: bool = True) -> dict:
        conn = self._find_conn(cid, ip, port)
        ip = valid_ip((conn or {}).get("remote_ip") or ip) or ""
        port = int((conn or {}).get("remote_port") or port or 0)
        if not ip:
            return {"ok": False, "error": "remote_ip required"}
        peer = f"{ip}:{port}"
        notes = apply_drop(ip, port)
        applied = any(n.endswith(":ok") for n in notes)
        rule = {
            "id": new_id("fw"),
            "remote_ip": ip,
            "remote_port": port,
            "action": "drop",
            "permanent": permanent,
            "reason": reason or "Kick from command board",
            "applied": applied,
            "notes": "; ".join(notes),
            "created_at": now(),
        }
        if self.kind == "postgres":
            if conn:
                self._q("UPDATE connections SET status = 'blocked' WHERE id = %s", (conn["id"],))
            self._q(
                """INSERT INTO firewall_rules
                   (id, remote_ip, remote_port, action, permanent, reason, applied, notes, created_at)
                   VALUES (%s,%s,%s,'drop',%s,%s,%s,%s,NOW())""",
                (rule["id"], ip, port, permanent, rule["reason"], applied, rule["notes"]),
            )
            self._alert("HIGH", f"AETHELOS+FIREWALL block {peer}: {rule['reason']}", conn)
        else:
            if conn:
                conn["status"] = "blocked"
                self.data["connections"][conn["id"]] = conn
            if ip not in self.data["blocked"]:
                self.data["blocked"].append(ip)
            if peer not in self.data["blocked"]:
                self.data["blocked"].append(peer)
            self.data["rules"].insert(0, rule)
            self.data["rules"] = self.data["rules"][:200]
            self._alert("HIGH", f"AETHELOS+FIREWALL block {peer}: {rule['reason']}", conn)
            self._save_json()
        return {
            "ok": True,
            "action": "kick",
            "id": (conn or {}).get("id", cid),
            "ip": ip,
            "port": port,
            "applied": applied,
            "notes": notes,
            "rule": rule,
        }

    def add_packet(self, pkt: dict) -> dict | None:
        src = valid_ip(str(pkt.get("src") or ""))
        dst = valid_ip(str(pkt.get("dst") or pkt.get("remote_ip") or ""))
        if not dst or is_loopback(dst):
            return None
        dport = valid_port(pkt.get("dport") or pkt.get("remote_port") or 0)
        sport = valid_port(pkt.get("sport") or pkt.get("local_port") or 0)
        proto = str(pkt.get("proto") or "tcp")[:8]
        app = str(pkt.get("app") or "vpn")[:40]
        device_id = str(pkt.get("device_id") or DEVICE_ID)[:80]
        rec = {
            "id": f"{device_id}:{dst}:{dport}",
            "device_id": device_id,
            "remote_ip": dst,
            "remote_port": dport,
            "local_port": sport,
            "app": app,
            "proto": proto,
            "status": "pending",
        }
        if self.kind == "postgres":
            self._q(
                """INSERT INTO packets (src, dst, sport, dport, proto, bytes, app)
                   VALUES (%s,%s,%s,%s,%s,%s,%s)""",
                (src or "", dst, sport, dport, proto, int(pkt.get("bytes") or 0), app),
            )
            self._q("DELETE FROM packets WHERE id < (SELECT COALESCE(MAX(id),0) - 2000 FROM packets)")
        else:
            self.data["packets"].insert(
                0,
                {
                    "src": src or "",
                    "dst": dst,
                    "sport": sport,
                    "dport": dport,
                    "proto": proto,
                    "bytes": int(pkt.get("bytes") or 0),
                    "app": app,
                    "seen_at": now(),
                },
            )
            self.data["packets"] = self.data["packets"][:500]
        self._touch_device(device_id, device_id, "android" if app == "vpn" else "hub", src or "")
        return self.upsert_conn(rec)

    def snapshot(self) -> dict:
        conns = self.list_conns()
        return {
            "ok": True,
            "engine": "aethelos-guard",
            "store": self.kind,
            "device_id": DEVICE_ID,
            "devices": self.list_devices(),
            "conns": conns,
            "connections": conns,
            "alerts": self.list_alerts(),
            "blocked": [r["remote_ip"] for r in self.list_rules() if r.get("action") == "drop"]
            if self.kind == "postgres"
            else list(self.data["blocked"]),
            "approved": [] if self.kind == "postgres" else list(self.data["approved"]),
            "rules": self.list_rules(),
            "anomalies": analyze(conns),
        }


STORE: EngineStore | None = None


def store() -> EngineStore:
    global STORE
    if STORE is None:
        STORE = EngineStore()
    return STORE


def _run(cmd: list[str], timeout: float = 4.0) -> subprocess.CompletedProcess[str] | None:
    try:
        return subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
    except (OSError, subprocess.SubprocessError):
        return None


def snapshot_sockets() -> int:
    """ss / netstat on this host only."""
    raw = ""
    if shutil.which("ss"):
        proc = _run(["ss", "-H", "-tnp", "state", "established"])
        if proc and proc.returncode == 0:
            raw = proc.stdout
    if not raw and shutil.which("netstat"):
        proc = _run(["netstat", "-tnp"])
        if proc:
            raw = proc.stdout
    added = 0
    for line in raw.splitlines():
        parts = line.split()
        if len(parts) < 4:
            continue
        peer = ""
        local = ""
        app = "host"
        if "users:(" in line or "pid=" in line:
            m = re.search(r'users:\(\("([^"]+)"', line)
            if m:
                app = m.group(1)
        if len(parts) >= 5 and ":" in parts[4]:
            local, peer = parts[3], parts[4]
        elif len(parts) >= 4 and ":" in parts[3]:
            local, peer = parts[2], parts[3]
        if not peer or peer.startswith(("*", "0.0.0.0", "[::]")):
            continue
        host, port = split_hostport(peer)
        _lip, lport = split_hostport(local) if local else ("", 0)
        ip = valid_ip(host)
        if not ip or is_loopback(ip):
            continue
        store().upsert_conn(
            {
                "id": f"{DEVICE_ID}:{ip}:{port}",
                "device_id": DEVICE_ID,
                "remote_ip": ip,
                "remote_port": port,
                "local_port": lport,
                "app": app,
                "proto": "tcp",
                "status": "pending",
            }
        )
        added += 1
    return added


def parse_tcpdump_line(line: str) -> dict | None:
    m = DUMP_LINE.search(line)
    if not m:
        return None
    src_h, src_p = split_hostport(m.group(1))
    dst_h, dst_p = split_hostport(m.group(2))
    proto = (m.group(3) or "tcp").lower()
    dst = valid_ip(dst_h)
    if not dst or is_loopback(dst):
        return None
    src = valid_ip(src_h) or ""
    return {
        "src": src,
        "dst": dst,
        "sport": src_p,
        "dport": dst_p,
        "proto": proto if proto in {"tcp", "udp", "icmp"} else "tcp",
        "app": "tcpdump",
        "device_id": DEVICE_ID,
        "bytes": 0,
    }


def tcpdump_loop(stop: threading.Event, iface: str) -> None:
    binary = shutil.which("tcpdump")
    if not binary:
        return
    cmd = [
        binary,
        "-i",
        iface,
        "-nn",
        "-l",
        "-q",
        "-t",
        "not",
        "port",
        str(PORT),
        "and",
        "not",
        "host",
        "127.0.0.1",
    ]
    try:
        proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, text=True, bufsize=1)
    except OSError:
        return
    try:
        assert proc.stdout is not None
        for line in proc.stdout:
            if stop.is_set():
                break
            pkt = parse_tcpdump_line(line)
            if pkt:
                with LOCK:
                    store().add_packet(pkt)
    finally:
        try:
            proc.terminate()
        except OSError:
            pass


def poll_loop(stop: threading.Event, interval: float) -> None:
    while not stop.is_set():
        try:
            with LOCK:
                snapshot_sockets()
        except Exception:
            pass
        stop.wait(interval)


def apply_drop(ip: str, port: int) -> list[str]:
    """Drop traffic to ip[:port] on THIS box. Never ssh, never a remote payload."""
    notes: list[str] = []
    checked = valid_ip(ip)
    if not checked:
        return ["ip:invalid"]
    if is_loopback(checked):
        return ["ip:loopback-refused"]

    if shutil.which("ss"):
        cmd = ["ss", "-K", "dst", checked]
        if port:
            cmd += ["dport", "=", str(port)]
        proc = _run(cmd)
        notes.append("ss-K:" + ("ok" if proc and proc.returncode == 0 else "skip"))

    def ipt(bin_name: str, proto: str | None = None) -> None:
        if not shutil.which(bin_name):
            return
        spec = [bin_name, "-C", "OUTPUT", "-d", checked, "-j", "DROP"]
        add = [bin_name, "-A", "OUTPUT", "-d", checked, "-j", "DROP"]
        if port and proto:
            spec = [bin_name, "-C", "OUTPUT", "-d", checked, "-p", proto, "--dport", str(port), "-j", "DROP"]
            add = [bin_name, "-A", "OUTPUT", "-d", checked, "-p", proto, "--dport", str(port), "-j", "DROP"]
        have = _run(spec)
        if have and have.returncode == 0:
            notes.append(f"{bin_name}:exists")
            return
        added = _run(add)
        notes.append(f"{bin_name}:" + ("ok" if added and added.returncode == 0 else "skip"))

    family = ipaddress.ip_address(checked).version
    table = "ip6tables" if family == 6 else "iptables"
    ipt(table)
    if port:
        ipt(table, "tcp")

    if shutil.which("nft"):
        proc = _run(["nft", "add", "rule", "inet", "filter", "output", "ip", "daddr", checked, "drop"])
        notes.append("nft:" + ("ok" if proc and proc.returncode == 0 else "skip"))

    if shutil.which("ufw"):
        args = ["ufw", "deny", "out", "to", checked]
        if port:
            args += ["port", str(port)]
        proc = _run(args)
        notes.append("ufw:" + ("ok" if proc and proc.returncode == 0 else "skip"))

    if os.name == "nt" or shutil.which("netsh"):
        name = f"Aethelos-Block-{checked}"
        proc = _run(
            [
                "netsh",
                "advfirewall",
                "firewall",
                "add",
                "rule",
                f"name={name}",
                "dir=out",
                "action=block",
                f"remoteip={checked}",
            ]
        )
        notes.append("netsh:" + ("ok" if proc and proc.returncode == 0 else "skip"))

    if not notes:
        notes.append("recorded-only")
    return notes


class Handler(BaseHTTPRequestHandler):
    def log_message(self, fmt: str, *args: Any) -> None:
        sys.stderr.write("guard " + (fmt % args) + "\n")

    def _cors(self) -> None:
        self.send_header("Access-Control-Allow-Origin", ALLOW_ORIGIN)
        self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
        self.send_header("Access-Control-Allow-Headers", "Content-Type")
        self.send_header("Access-Control-Allow-Private-Network", "true")
        self.send_header("Cache-Control", "no-store")

    def _json(self, code: int, payload: object) -> None:
        body = json.dumps(payload, default=str).encode("utf-8")
        self.send_response(code)
        self.send_header("Content-Type", "application/json")
        self._cors()
        self.send_header("Content-Length", str(len(body)))
        self.end_headers()
        self.wfile.write(body)

    def do_OPTIONS(self) -> None:
        self.send_response(204)
        self._cors()
        self.end_headers()

    def do_GET(self) -> None:
        path = urlparse(self.path).path.rstrip("/") or "/"
        s = store()
        with LOCK:
            if path in ("/health", "/"):
                self._json(
                    200,
                    {
                        "ok": True,
                        "service": "aethelos-guard",
                        "store": s.kind,
                        "device_id": DEVICE_ID,
                    },
                )
                return
            if path in (
                "/traffic",
                "/v1/traffic",
                "/api/traffic",
                "/api/guard/traffic",
                "/api/connections",
            ):
                snapshot_sockets()
                self._json(200, s.snapshot())
                return
            if path in ("/api/devices",):
                self._json(200, {"ok": True, "devices": s.list_devices()})
                return
            if path in ("/api/alerts",):
                self._json(200, {"ok": True, "alerts": s.list_alerts()})
                return
            if path in ("/api/firewall/rules", "/api/guard/rules"):
                self._json(200, {"ok": True, "rules": s.list_rules()})
                return
            if path in ("/api/anomalies",):
                self._json(200, {"ok": True, **analyze(s.list_conns())})
                return
        self._json(404, {"ok": False, "error": "not found"})

    def do_POST(self) -> None:
        n = int(self.headers.get("Content-Length") or 0)
        raw = self.rfile.read(n) if n else b"{}"
        try:
            body = json.loads(raw.decode("utf-8") or "{}")
            if not isinstance(body, dict):
                body = {}
        except json.JSONDecodeError:
            self._json(400, {"ok": False, "error": "json"})
            return
        path = urlparse(self.path).path.rstrip("/") or "/"
        cid = unquote(str(body.get("id") or body.get("conn_id") or ""))
        ip = str(body.get("ip") or body.get("remote_ip") or "")
        port = valid_port(body.get("port") or body.get("remote_port") or 0)
        reason = str(body.get("reason") or "")

        parts = path.split("/")
        if len(parts) >= 5 and parts[1] == "api" and parts[2] == "connections":
            cid = cid or unquote(parts[3])
            tail = parts[4]
            if tail == "approve":
                path = "/approve"
            elif tail in ("kick", "deny"):
                path = "/kick"

        s = store()
        with LOCK:
            if path in ("/approve", "/v1/approve", "/api/guard/approve"):
                self._json(200, s.approve(cid, ip, port))
                return
            if path in (
                "/kick",
                "/v1/kick",
                "/api/guard/kick",
                "/api/firewall/block",
                "/api/connections/kick",
            ):
                permanent = bool(body.get("permanent", True))
                self._json(200, s.kick(cid, ip, port, reason, permanent=permanent))
                return
            if path in ("/api/packets", "/api/connections/report"):
                conns = body.get("connections")
                added = 0
                if isinstance(conns, list):
                    for item in conns:
                        if isinstance(item, dict):
                            item.setdefault("device_id", body.get("device_id") or DEVICE_ID)
                            if s.add_packet(item):
                                added += 1
                elif s.add_packet(body):
                    added = 1
                self._json(200, {"ok": True, "ingested": added})
                return
        self._json(404, {"ok": False, "error": "not found"})


def self_test() -> int:
    global STORE_PATH, STORE
    import tempfile

    os.environ.pop("AETHELOS_DATABASE_URL", None)
    os.environ["AETHELOS_USE_PG"] = "0"
    STORE_PATH = Path(tempfile.mkdtemp()) / "guard-store.json"
    STORE = None
    sample = [
        "IP 192.168.1.18.54321 > 203.0.113.5.4444: tcp",
        "IP 10.0.0.8.41222 > 1.1.1.1.443: tcp",
        "IP 127.0.0.1.80 > 127.0.0.1.443: tcp",
        "IP6 2001:db8::1.443 > 2001:db8::2.51234: tcp",
    ]
    parsed = [parse_tcpdump_line(s) for s in sample]
    assert parsed[0] and parsed[0]["dst"] == "203.0.113.5" and parsed[0]["dport"] == 4444
    assert parsed[1] and parsed[1]["dport"] == 443
    assert parsed[2] is None
    s = store()
    s.add_packet(parsed[0] or {})
    s.add_packet(parsed[1] or {})
    a = s.approve("", "1.1.1.1", 443)
    k = s.kick("", "203.0.113.5", 4444, "self-test kick")
    snap = s.snapshot()
    assert a["ok"] and k["ok"]
    assert any(c["remote_ip"] == "1.1.1.1" and c["status"] == "approved" for c in snap["conns"])
    assert any(c["remote_ip"] == "203.0.113.5" and c["status"] == "blocked" for c in snap["conns"])
    notes = apply_drop("203.0.113.5", 4444)
    assert notes
    print(json.dumps({"ok": True, "store": s.kind, "conns": len(snap["conns"]), "notes": notes, "self_test": True}))
    return 0


def main(argv: list[str] | None = None) -> int:
    parser = argparse.ArgumentParser(description="Aethelos local Guard engine")
    parser.add_argument("--tcpdump", action="store_true", help="capture headers with tcpdump on this device")
    parser.add_argument("--iface", default=os.environ.get("AETHELOS_IFACE", "any"))
    parser.add_argument("--port", type=int, default=int(os.environ.get("AETHELOS_GUARD_PORT", "8787")))
    parser.add_argument("--self-test", action="store_true")
    args = parser.parse_args(argv)
    if args.self_test:
        return self_test()

    bind_port = args.port
    stop = threading.Event()
    threading.Thread(target=poll_loop, args=(stop, 3.0), name="ss-poll", daemon=True).start()
    want_dump = args.tcpdump or os.environ.get("AETHELOS_TCPDUMP", "").lower() in {"1", "true", "yes"}
    if want_dump:
        threading.Thread(target=tcpdump_loop, args=(stop, args.iface), name="tcpdump", daemon=True).start()

    httpd = ThreadingHTTPServer((HOST, bind_port), Handler)
    print(f"Aethelos Guard engine on http://{HOST}:{bind_port}  store={store().kind}  device={DEVICE_ID}")
    print("This device only. Approve / Kick from the command board land here.")
    if want_dump:
        print("tcpdump header capture armed (this host).")
    else:
        print("Sockets via ss/netstat. Pass --tcpdump for live headers (Termux: pkg install tcpdump tsu).")
    try:
        httpd.serve_forever()
    except KeyboardInterrupt:
        stop.set()
        httpd.server_close()
        return 0
    return 0


if __name__ == "__main__":
    raise SystemExit(main())
