720 lines
22 KiB
Python
720 lines
22 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Presence daemon — Multi-signal WiFi/ARP/DHCP/BLE person presence tracking.
|
|
Passive sensors with exponential decay and presence anchoring.
|
|
Zero active probing. No tool fingerprints in output.
|
|
"""
|
|
|
|
import asyncio
|
|
import http.server
|
|
import json
|
|
import logging
|
|
import math
|
|
import os
|
|
import socket
|
|
import sqlite3
|
|
import struct
|
|
import sys
|
|
import threading
|
|
import time
|
|
import urllib.parse
|
|
import urllib.request
|
|
from dataclasses import dataclass, field
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
|
|
# --- Config (env-driven) ---
|
|
MONITOR_IFACE = os.getenv("PRESENCE_MONITOR_IFACE", "wlan1")
|
|
DATA_IFACE = os.getenv("PRESENCE_DATA_IFACE", "wlan0")
|
|
DB_PATH = os.getenv("PRESENCE_DB_PATH", "/opt/sensor/presence.db")
|
|
MATRIX_WEBHOOK = os.getenv("PRESENCE_MATRIX_WEBHOOK", "")
|
|
STATUS_PORT = int(os.getenv("PRESENCE_STATUS_PORT", "9191"))
|
|
PRESENCE_THRESHOLD = 0.40
|
|
ABSENCE_THRESHOLD = 0.10
|
|
ANCHOR_MINUTES = 45
|
|
DECAY_LAMBDA = 0.08
|
|
|
|
# Signal bump constants
|
|
PROBE_BUMP = 0.60
|
|
ARP_BUMP = 0.50
|
|
DHCP_BUMP = 0.80
|
|
BLE_BUMP = 0.40
|
|
|
|
# Logging
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [%(levelname)s] %(message)s",
|
|
handlers=[
|
|
logging.FileHandler("/var/log/sensor.log", encoding="utf-8"),
|
|
logging.StreamHandler(sys.stdout)
|
|
]
|
|
)
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Top OUI prefixes (inline, no external DB)
|
|
OUI_DICT = {
|
|
"88A29E": "Apple",
|
|
"001A7D": "Apple",
|
|
"185F3F": "Apple",
|
|
"A4C3F0": "Apple",
|
|
"0899D8": "Apple",
|
|
"004096": "Apple",
|
|
"00219B": "Apple",
|
|
"0021E9": "Apple",
|
|
"005973": "Apple",
|
|
"006377": "Apple",
|
|
"0064B9": "Apple",
|
|
"0084F3": "Apple",
|
|
"00A04D": "Apple",
|
|
"00D04B": "Apple",
|
|
"28879F": "Google",
|
|
"2887BA": "Google",
|
|
"5427EB": "Google",
|
|
"542758": "Google",
|
|
"341513": "Amazon",
|
|
"0C47C2": "Amazon",
|
|
"B827EB": "Raspberry Pi",
|
|
"2C56DC": "Raspberry Pi",
|
|
"E45F01": "Raspberry Pi",
|
|
}
|
|
|
|
# Netlink constants
|
|
NETLINK_ROUTE = 0
|
|
RTMGRP_NEIGH = 0x4
|
|
RTM_NEWNEIGH = 28
|
|
RTM_DELNEIGH = 29
|
|
NUD_REACHABLE = 0x02
|
|
NUD_STALE = 0x04
|
|
NUD_DELAY = 0x08
|
|
NUD_PROBE = 0x10
|
|
NDA_DST = 1
|
|
NDA_LLADDR = 2
|
|
|
|
|
|
@dataclass
|
|
class DeviceState:
|
|
"""In-memory device state."""
|
|
mac: str
|
|
certainty: float = 0.0
|
|
last_seen: float = field(default_factory=time.time)
|
|
last_signal: str = ""
|
|
person_id: Optional[int] = None
|
|
|
|
|
|
@dataclass
|
|
class PersonState:
|
|
"""In-memory person state."""
|
|
person_id: int
|
|
name: str
|
|
status: str = "UNKNOWN"
|
|
last_change: float = field(default_factory=time.time)
|
|
|
|
|
|
# Shared state with locks
|
|
devices = {}
|
|
devices_lock = threading.Lock()
|
|
|
|
persons = {}
|
|
persons_lock = threading.Lock()
|
|
|
|
infrastructure_ips = set()
|
|
infrastructure_lock = threading.Lock()
|
|
|
|
|
|
def lookup_oui(mac: str) -> str:
|
|
"""Look up OUI vendor from MAC (first 3 octets, inline dict)."""
|
|
if not mac or mac == "00:00:00:00:00:00":
|
|
return "unknown"
|
|
oui_prefix = mac.replace(":", "").upper()[:6]
|
|
return OUI_DICT.get(oui_prefix, "unknown")
|
|
|
|
|
|
def init_db() -> None:
|
|
"""Initialize SQLite database from schema if not exists."""
|
|
schema_path = Path(__file__).parent / "presence_schema.sql"
|
|
if not schema_path.exists():
|
|
logger.error(f"Schema file not found: {schema_path}")
|
|
return
|
|
|
|
db_path = Path(DB_PATH)
|
|
db_path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
conn = sqlite3.connect(DB_PATH, timeout=10)
|
|
try:
|
|
cursor = conn.cursor()
|
|
schema = schema_path.read_text()
|
|
cursor.executescript(schema)
|
|
conn.commit()
|
|
cursor.execute("PRAGMA journal_mode=WAL")
|
|
cursor.execute("PRAGMA synchronous=NORMAL")
|
|
cursor.execute("PRAGMA foreign_keys=ON")
|
|
conn.commit()
|
|
logger.info(f"Database initialized at {DB_PATH}")
|
|
except Exception as e:
|
|
logger.error(f"Failed to init database: {e}")
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def lookup_person(mac: str) -> Optional[int]:
|
|
"""Look up person_id for a MAC from database."""
|
|
conn = sqlite3.connect(DB_PATH, timeout=10)
|
|
try:
|
|
cursor = conn.cursor()
|
|
result = cursor.execute(
|
|
"SELECT person_id FROM devices WHERE mac = ?",
|
|
(mac.lower(),)
|
|
).fetchone()
|
|
return result[0] if result else None
|
|
except Exception as e:
|
|
logger.debug(f"lookup_person error for {mac}: {e}")
|
|
return None
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def seed_persons_from_db() -> None:
|
|
"""Load all persons from database into memory."""
|
|
conn = sqlite3.connect(DB_PATH, timeout=10)
|
|
try:
|
|
conn.row_factory = sqlite3.Row
|
|
cursor = conn.cursor()
|
|
rows = cursor.execute("SELECT id, name FROM persons").fetchall()
|
|
|
|
with persons_lock:
|
|
for row in rows:
|
|
person = PersonState(
|
|
person_id=row["id"],
|
|
name=row["name"]
|
|
)
|
|
persons[row["id"]] = person
|
|
logger.info(f"Loaded {len(rows)} persons from database")
|
|
except Exception as e:
|
|
logger.error(f"Failed to seed persons: {e}")
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def log_signal(mac: str, signal_type: str, certainty: float) -> None:
|
|
"""Insert signal into database (non-blocking via queue would be ideal, but inline for simplicity)."""
|
|
conn = sqlite3.connect(DB_PATH, timeout=10)
|
|
try:
|
|
cursor = conn.cursor()
|
|
cursor.execute(
|
|
"INSERT INTO signals (mac, signal_type, certainty, ts) VALUES (?, ?, ?, ?)",
|
|
(mac.lower(), signal_type, certainty, time.time())
|
|
)
|
|
conn.commit()
|
|
except Exception as e:
|
|
logger.debug(f"Failed to log signal: {e}")
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def log_presence_change(person_id: int, old_status: str, new_status: str) -> None:
|
|
"""Insert occupancy change into database."""
|
|
conn = sqlite3.connect(DB_PATH, timeout=10)
|
|
try:
|
|
cursor = conn.cursor()
|
|
cursor.execute(
|
|
"INSERT INTO occupancy_log (person_id, old_status, new_status, ts) VALUES (?, ?, ?, ?)",
|
|
(person_id, old_status, new_status, time.time())
|
|
)
|
|
conn.commit()
|
|
except Exception as e:
|
|
logger.debug(f"Failed to log presence change: {e}")
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def send_alert(person_name: str, status: str, signal_type: str, certainty: float) -> None:
|
|
"""Send alert to Matrix webhook if configured."""
|
|
if not MATRIX_WEBHOOK:
|
|
logger.info(f"Alert (no webhook): {person_name} — {status} ({signal_type}, certainty: {certainty:.2f})")
|
|
return
|
|
|
|
try:
|
|
action = "ARRIVED" if status == "PRESENT" else "DEPARTED"
|
|
message = f"SENSOR: {person_name} — {action} ({signal_type}, certainty: {certainty:.2f})"
|
|
payload = json.dumps({"text": message})
|
|
req = urllib.request.Request(
|
|
MATRIX_WEBHOOK,
|
|
data=payload.encode("utf-8"),
|
|
headers={"Content-Type": "application/json"}
|
|
)
|
|
with urllib.request.urlopen(req, timeout=5) as response:
|
|
logger.info(f"Alert sent: {message}")
|
|
except Exception as e:
|
|
logger.error(f"Failed to send alert: {e}")
|
|
|
|
|
|
def on_signal(mac: str, signal_type: str, bump: float) -> None:
|
|
"""
|
|
Process a signal bump for a device.
|
|
Decay old certainty, apply bump, update state.
|
|
"""
|
|
mac = mac.lower()
|
|
with devices_lock:
|
|
dev = devices.get(mac)
|
|
if dev is None:
|
|
dev = DeviceState(
|
|
mac=mac,
|
|
certainty=0.0,
|
|
last_seen=time.time(),
|
|
last_signal="",
|
|
person_id=lookup_person(mac)
|
|
)
|
|
devices[mac] = dev
|
|
|
|
# Apply exponential decay since last signal
|
|
elapsed_minutes = (time.time() - dev.last_seen) / 60.0
|
|
dev.certainty = dev.certainty * math.exp(-DECAY_LAMBDA * elapsed_minutes)
|
|
|
|
# Apply bump
|
|
dev.certainty = min(1.0, dev.certainty + bump)
|
|
dev.last_seen = time.time()
|
|
dev.last_signal = signal_type
|
|
|
|
log_signal(mac, signal_type, dev.certainty)
|
|
update_person_state(dev.person_id)
|
|
|
|
|
|
def update_person_state(person_id: Optional[int]) -> None:
|
|
"""
|
|
Compute aggregate certainty for a person, update state machine.
|
|
Fire alerts on state transitions.
|
|
"""
|
|
if person_id is None:
|
|
return
|
|
|
|
# Snapshot device state under lock
|
|
with devices_lock:
|
|
# Get all devices for this person
|
|
devices_for_person = [d for d in devices.values() if d.person_id == person_id]
|
|
if not devices_for_person:
|
|
return
|
|
|
|
# Compute aggregate certainty (max + 0.08 * second_max)
|
|
certainties = sorted([d.certainty for d in devices_for_person], reverse=True)
|
|
if len(certainties) >= 2:
|
|
agg_certainty = certainties[0] + 0.08 * certainties[1]
|
|
else:
|
|
agg_certainty = certainties[0]
|
|
|
|
# Get max last_seen time
|
|
max_last_seen = max(d.last_seen for d in devices_for_person)
|
|
|
|
# Determine new status based on snapshot
|
|
current_time = time.time()
|
|
elapsed_since_signal = (current_time - max_last_seen) / 60.0
|
|
|
|
with persons_lock:
|
|
person = persons.get(person_id)
|
|
if person is None:
|
|
return
|
|
|
|
old_status = person.status
|
|
new_status = old_status
|
|
|
|
# State machine logic
|
|
if person.status == "UNKNOWN":
|
|
if agg_certainty >= PRESENCE_THRESHOLD:
|
|
new_status = "PRESENT"
|
|
elif person.status == "PRESENT":
|
|
# Check presence anchor
|
|
if elapsed_since_signal >= ANCHOR_MINUTES:
|
|
new_status = "ABSENT"
|
|
elif person.status == "ABSENT":
|
|
if agg_certainty >= PRESENCE_THRESHOLD:
|
|
new_status = "PRESENT"
|
|
|
|
# Apply state transition if changed
|
|
if new_status != old_status:
|
|
person.status = new_status
|
|
person.last_change = current_time
|
|
log_presence_change(person_id, old_status, new_status)
|
|
|
|
# Send alert outside of lock
|
|
if new_status != old_status:
|
|
if new_status == "PRESENT":
|
|
send_alert(person.name, "PRESENT", "signal", agg_certainty)
|
|
else:
|
|
send_alert(person.name, "ABSENT", "timeout", agg_certainty)
|
|
|
|
|
|
def probe_sniffer(iface: str) -> None:
|
|
"""
|
|
WiFi probe request sniffer.
|
|
Raw 802.11 frame capture on monitor mode interface.
|
|
Extracts source MAC from probe requests.
|
|
"""
|
|
try:
|
|
s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(0x0003))
|
|
s.bind((iface, 0))
|
|
logger.info(f"Probe sniffer started on {iface}")
|
|
|
|
while True:
|
|
try:
|
|
data, _ = s.recvfrom(65535)
|
|
if len(data) < 24:
|
|
continue
|
|
|
|
# 802.11 frame: FC (2) + duration (2) + dest (6) + src (6) + BSSID (6) = 22 bytes min
|
|
fc = struct.unpack("<H", data[0:2])[0]
|
|
frame_type = (fc >> 2) & 0x3
|
|
frame_subtype = (fc >> 4) & 0xf
|
|
|
|
# Type 0 = management, Subtype 4 = probe request
|
|
if frame_type == 0 and frame_subtype == 4:
|
|
if len(data) >= 16:
|
|
# Source MAC at bytes 10-16
|
|
mac_bytes = data[10:16]
|
|
mac = ':'.join(f'{b:02x}' for b in mac_bytes)
|
|
if mac != "ff:ff:ff:ff:ff:ff" and mac != "00:00:00:00:00:00":
|
|
on_signal(mac, "probe", PROBE_BUMP)
|
|
except Exception as e:
|
|
logger.debug(f"Probe sniffer error: {e}")
|
|
time.sleep(0.1)
|
|
except Exception as e:
|
|
logger.error(f"Failed to start probe sniffer: {e}")
|
|
|
|
|
|
def arp_listener() -> None:
|
|
"""
|
|
ARP listener via netlink RTM_NEWNEIGH.
|
|
Kernel pushes neighbor discovery events.
|
|
"""
|
|
try:
|
|
s = socket.socket(socket.AF_NETLINK, socket.SOCK_RAW, NETLINK_ROUTE)
|
|
s.bind((os.getpid(), RTMGRP_NEIGH))
|
|
logger.info("ARP listener started")
|
|
|
|
while True:
|
|
try:
|
|
data = s.recv(65535)
|
|
parse_netlink(data)
|
|
except Exception as e:
|
|
logger.debug(f"ARP listener error: {e}")
|
|
time.sleep(0.1)
|
|
except Exception as e:
|
|
logger.error(f"Failed to start ARP listener: {e}")
|
|
|
|
|
|
def parse_netlink(data: bytes) -> None:
|
|
"""Parse netlink messages for neighbor events."""
|
|
try:
|
|
offset = 0
|
|
while offset < len(data):
|
|
if offset + 16 > len(data):
|
|
break
|
|
|
|
nlmsg_len, nlmsg_type, _, _, _ = struct.unpack_from('=IHHII', data, offset)
|
|
if nlmsg_len < 16:
|
|
break
|
|
|
|
# Verify complete message is available before extracting payload
|
|
if offset + nlmsg_len > len(data):
|
|
return
|
|
|
|
payload = data[offset + 16:offset + nlmsg_len]
|
|
|
|
if nlmsg_type in (RTM_NEWNEIGH, RTM_DELNEIGH):
|
|
parse_ndmsg(payload, nlmsg_type)
|
|
|
|
offset += (nlmsg_len + 3) & ~3
|
|
except Exception as e:
|
|
logger.debug(f"Netlink parse error: {e}")
|
|
|
|
|
|
def parse_ndmsg(data: bytes, msg_type: int) -> None:
|
|
"""Parse ndmsg (neighbor discovery message)."""
|
|
try:
|
|
if len(data) < 12:
|
|
return
|
|
|
|
state = struct.unpack_from('=H', data, 8)[0]
|
|
|
|
mac = None
|
|
ip = None
|
|
offset = 12
|
|
|
|
while offset + 4 <= len(data):
|
|
rta_len, rta_type = struct.unpack_from('=HH', data, offset)
|
|
if rta_len < 4:
|
|
break
|
|
|
|
# Verify complete RTA attribute is available before extracting
|
|
if offset + rta_len > len(data):
|
|
return
|
|
|
|
val = data[offset + 4:offset + rta_len]
|
|
|
|
if rta_type == NDA_LLADDR and len(val) == 6:
|
|
mac = ':'.join(f'{b:02x}' for b in val)
|
|
elif rta_type == NDA_DST:
|
|
if len(val) == 4:
|
|
ip = socket.inet_ntoa(val)
|
|
|
|
offset += (rta_len + 3) & ~3
|
|
|
|
if not mac or mac in ('00:00:00:00:00:00', 'ff:ff:ff:ff:ff:ff'):
|
|
return
|
|
|
|
# Ignore infrastructure IPs
|
|
with infrastructure_lock:
|
|
if ip and ip in infrastructure_ips:
|
|
return
|
|
|
|
if msg_type == RTM_NEWNEIGH and (state & (NUD_REACHABLE | NUD_STALE | NUD_DELAY | NUD_PROBE)):
|
|
on_signal(mac, "arp", ARP_BUMP)
|
|
except Exception as e:
|
|
logger.debug(f"ndmsg parse error: {e}")
|
|
|
|
|
|
def dhcp_sniffer(iface: str) -> None:
|
|
"""
|
|
DHCP sniffer via AF_PACKET raw socket.
|
|
Captures DHCP traffic passively, extracts client MACs.
|
|
"""
|
|
try:
|
|
s = socket.socket(socket.AF_PACKET, socket.SOCK_RAW, socket.htons(0x0003))
|
|
s.bind((iface, 0))
|
|
logger.info(f"DHCP sniffer started on {iface}")
|
|
|
|
while True:
|
|
try:
|
|
data, _ = s.recvfrom(65535)
|
|
parse_dhcp(data)
|
|
except Exception as e:
|
|
logger.debug(f"DHCP sniffer error: {e}")
|
|
time.sleep(0.1)
|
|
except Exception as e:
|
|
logger.error(f"Failed to start DHCP sniffer: {e}")
|
|
|
|
|
|
def parse_dhcp(data: bytes) -> None:
|
|
"""Parse DHCP packet and extract client MAC."""
|
|
try:
|
|
if len(data) < 14:
|
|
return
|
|
|
|
# Ethernet frame
|
|
eth_type = struct.unpack('!H', data[12:14])[0]
|
|
if eth_type != 0x0800: # not IPv4
|
|
return
|
|
|
|
# IP header
|
|
if len(data) < 23:
|
|
return
|
|
proto = data[23]
|
|
if proto != 17: # not UDP
|
|
return
|
|
|
|
# UDP ports
|
|
if len(data) < 38:
|
|
return
|
|
src_port = struct.unpack('!H', data[34:36])[0]
|
|
dst_port = struct.unpack('!H', data[36:38])[0]
|
|
if not (src_port in (67, 68) and dst_port in (67, 68)):
|
|
return
|
|
|
|
# DHCP payload (starts at byte 42)
|
|
if len(data) < 42:
|
|
return
|
|
dhcp = data[42:]
|
|
if len(dhcp) < 236:
|
|
return
|
|
|
|
# MAC address (chaddr at offset 28, 6 bytes)
|
|
if len(dhcp) < 34:
|
|
return
|
|
mac_bytes = dhcp[28:34]
|
|
mac = ':'.join(f'{b:02x}' for b in mac_bytes)
|
|
if mac in ('00:00:00:00:00:00', 'ff:ff:ff:ff:ff:ff'):
|
|
return
|
|
|
|
# Parse DHCP options (start at byte 240)
|
|
if len(dhcp) < 240:
|
|
return
|
|
|
|
msg_type = None
|
|
i = 240
|
|
while i < len(dhcp):
|
|
opt = dhcp[i]
|
|
if opt == 255:
|
|
break
|
|
if opt == 0:
|
|
i += 1
|
|
continue
|
|
if i + 1 >= len(dhcp):
|
|
break
|
|
|
|
length = dhcp[i + 1]
|
|
if i + 2 + length > len(dhcp):
|
|
break
|
|
|
|
val = dhcp[i + 2:i + 2 + length]
|
|
|
|
if opt == 53 and length == 1: # DHCP message type
|
|
msg_type = val[0]
|
|
|
|
i += 2 + length
|
|
|
|
# DISCOVER or REQUEST indicates presence
|
|
if msg_type in (1, 3):
|
|
on_signal(mac, "dhcp", DHCP_BUMP)
|
|
except Exception as e:
|
|
logger.debug(f"DHCP parse error: {e}")
|
|
|
|
|
|
def ble_scanner() -> None:
|
|
"""
|
|
BLE passive scanner.
|
|
Gracefully degrades if bleak is not available.
|
|
"""
|
|
try:
|
|
from bleak import BleakScanner
|
|
except ImportError:
|
|
logger.warning("bleak not available; BLE scanning disabled")
|
|
return
|
|
|
|
async def run_ble_scan() -> None:
|
|
"""Run BLE scan in asyncio loop."""
|
|
def on_ble_detect(device, adv_data):
|
|
"""Callback for BLE device detection."""
|
|
if device and device.address:
|
|
on_signal(device.address.lower(), "ble", BLE_BUMP)
|
|
|
|
try:
|
|
async with BleakScanner(
|
|
detection_callback=on_ble_detect,
|
|
scanning_mode="passive"
|
|
):
|
|
await asyncio.sleep(float('inf'))
|
|
except Exception as e:
|
|
logger.error(f"BLE scanner error: {e}")
|
|
|
|
try:
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
loop.run_until_complete(run_ble_scan())
|
|
except Exception as e:
|
|
logger.error(f"Failed to start BLE scanner: {e}")
|
|
|
|
|
|
def decay_loop() -> None:
|
|
"""
|
|
Background thread: apply decay to all devices every 60 seconds.
|
|
Prune devices not seen in 120+ minutes.
|
|
"""
|
|
while True:
|
|
try:
|
|
time.sleep(60)
|
|
|
|
with devices_lock:
|
|
now = time.time()
|
|
to_prune = []
|
|
|
|
for mac, dev in devices.items():
|
|
elapsed_minutes = (now - dev.last_seen) / 60.0
|
|
|
|
# Apply decay
|
|
dev.certainty = dev.certainty * math.exp(-DECAY_LAMBDA * elapsed_minutes)
|
|
|
|
# Mark for pruning if not seen in 120+ minutes
|
|
if elapsed_minutes >= 120:
|
|
to_prune.append(mac)
|
|
|
|
# Prune old devices
|
|
for mac in to_prune:
|
|
del devices[mac]
|
|
logger.debug(f"Pruned device {mac}")
|
|
|
|
# Update all person states due to decay
|
|
for person_id in set(d.person_id for d in devices.values() if d.person_id):
|
|
update_person_state(person_id)
|
|
|
|
except Exception as e:
|
|
logger.error(f"Decay loop error: {e}")
|
|
|
|
|
|
def status_handler(request, client_address, server):
|
|
"""HTTP request handler for status endpoint."""
|
|
class StatusHandler(http.server.BaseHTTPRequestHandler):
|
|
def do_GET(self):
|
|
if self.path == "/":
|
|
try:
|
|
with persons_lock:
|
|
person_list = [
|
|
{
|
|
"id": p.person_id,
|
|
"name": p.name,
|
|
"status": p.status,
|
|
"last_change": p.last_change
|
|
}
|
|
for p in persons.values()
|
|
]
|
|
response = json.dumps({"persons": person_list})
|
|
self.send_response(200)
|
|
self.send_header("Content-Type", "application/json")
|
|
self.end_headers()
|
|
self.wfile.write(response.encode("utf-8"))
|
|
except Exception as e:
|
|
logger.error(f"Status handler error: {e}")
|
|
self.send_response(500)
|
|
self.end_headers()
|
|
else:
|
|
self.send_response(404)
|
|
self.end_headers()
|
|
|
|
def log_message(self, format, *args):
|
|
"""Suppress HTTP server logs."""
|
|
pass
|
|
|
|
return StatusHandler(request, client_address, server)
|
|
|
|
|
|
def status_server() -> None:
|
|
"""Start HTTP status server on 127.0.0.1:STATUS_PORT."""
|
|
try:
|
|
server = http.server.HTTPServer(
|
|
("127.0.0.1", STATUS_PORT),
|
|
status_handler
|
|
)
|
|
logger.info(f"Status server started on http://127.0.0.1:{STATUS_PORT}")
|
|
server.serve_forever()
|
|
except Exception as e:
|
|
logger.error(f"Failed to start status server: {e}")
|
|
|
|
|
|
def main() -> None:
|
|
"""Initialize and run all sensor threads."""
|
|
init_db()
|
|
seed_persons_from_db()
|
|
|
|
threads = [
|
|
threading.Thread(target=probe_sniffer, args=(MONITOR_IFACE,), daemon=True, name="probe-sniffer"),
|
|
threading.Thread(target=arp_listener, daemon=True, name="arp-listener"),
|
|
threading.Thread(target=dhcp_sniffer, args=(DATA_IFACE,), daemon=True, name="dhcp-sniffer"),
|
|
threading.Thread(target=ble_scanner, daemon=True, name="ble-scanner"),
|
|
threading.Thread(target=decay_loop, daemon=True, name="decay-loop"),
|
|
threading.Thread(target=status_server, daemon=True, name="status-server"),
|
|
]
|
|
|
|
for t in threads:
|
|
t.start()
|
|
|
|
logger.info("Sensor daemon started")
|
|
|
|
# Handle shutdown gracefully
|
|
import signal
|
|
signal.signal(signal.SIGTERM, lambda *_: sys.exit(0))
|
|
signal.signal(signal.SIGINT, lambda *_: sys.exit(0))
|
|
|
|
for t in threads:
|
|
t.join()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|