feat(satpass): add satellite pass alert family (tier 1)

- Add satpass_handler.py with observer/elevation/NORAD filters
- Subscribe to central.sat.pass.> in consumer.py
- Add satpass toggle family to config.py
- Add sat_pass category mapping in categories.py
- Add satpass config section in adapter_config/defaults.py
- Add 9 unit tests for handler logic

Wire format: 3-line LoRa-tight (sat name, AOS/LOS times, observer)
Dedup: hourly bucket per sat/observer
Severity: 60deg+=immediate, 45deg+=priority, else routine
This commit is contained in:
Matt Johnson (via Claude) 2026-06-11 23:23:24 +00:00
commit 782899d4b1
6 changed files with 512 additions and 1 deletions

View file

@ -591,6 +591,26 @@ REGISTRY: dict[tuple[str, str], dict[str, Any]] = {
"description": "Minimum danger level to broadcast (3=Considerable, 4=High, 5=Extreme).",
},
# =================================================================
# SATPASS -- Satellite pass broadcasts
# =================================================================
("satpass", "observers"): {
"default": [],
"type": "list",
"description": "Observer location names to include (empty = all).",
},
("satpass", "min_elevation"): {
"default": 30,
"type": "int",
"description": "Minimum max elevation (degrees) to broadcast a pass.",
},
("satpass", "norad_ids"): {
"default": [],
"type": "list",
"description": "NORAD catalog IDs to include (empty = all).",
},
# =================================================================
# DASHBOARD -- UI-only settings persisted for the operator
# =================================================================
@ -734,6 +754,11 @@ ADAPTER_META: dict[str, dict[str, Any]] = {
"include_in_llm_context": False,
"description": "Operator UI preferences persisted to adapter_config (region selectors, display options).",
},
"satpass": {
"display_name": "Satellite passes",
"include_in_llm_context": True,
"description": "Regional satellite pass broadcasts (ISS, amateur radio sats).",
},
}

View file

@ -50,6 +50,7 @@ _SUBJECTS_BARE: dict[str, list[str]] = {
"traffic": ["central.traffic.>"],
"roads511": ["central.traffic.>"], # shared with traffic; sub-adapter routing
"avalanche": ["central.avy.advisory.>"],
"satpass": ["central.sat.pass.>"],
}
# Backwards-compat: keep ADAPTER_SUBJECTS importable for legacy readers/tests.
@ -207,6 +208,7 @@ CENTRAL_ADAPTER_TO_SOURCE: dict[str, str] = {
"itd_511": "roads511",
"avalanche_org": "avalanche",
"firms": "firms",
"sat_passes": "satpass",
}
# Central hierarchical category prefix -> meshai flat category.
@ -236,6 +238,7 @@ _CATEGORY_MAP: list[tuple[str, str]] = [
("incident", "road_incident"),
("closure", "road_closure"),
("traffic.", "traffic_congestion"),
("sat.", "sat_pass"),
]
@ -535,6 +538,9 @@ class CentralConsumer:
# commit #5 (env_reporter). Closes the v0.5.13
# silent-drop on central.fire.hotspot.> (audit doc
# finding #2).
elif inner.get("adapter") == "sat_passes":
from meshai.central.satpass_handler import handle_satpass
synthesized = handle_satpass(envelope, subject, data=data) or None
elif inner.get("adapter") == "firms":
from meshai.central.firms_handler import handle_firms
synthesized = handle_firms(envelope, subject, data=data) or None

View file

@ -0,0 +1,308 @@
"""v0.7 Satellite pass handler.
Broadcast regional satellite passes from Central's CENTRAL_SAT stream.
Filter criteria:
(a) Pass must be for an observer in adapter_config.satpass.observers
(empty list = all observers)
(b) Max elevation must meet adapter_config.satpass.min_elevation (default 30)
(c) Optional NORAD ID filter via adapter_config.satpass.norad_ids
Dedup bucketing: canonical event_id = {norad_id}:{observer}:{aos_bucket}
where aos_bucket = floor(aos_epoch / 3600) -- one broadcast per satellite
per observer per hour window.
Severity mapping:
4 = immediate (>= 60 deg max elevation)
3 = priority (>= 45 deg max elevation)
<= 2 = routine
Wire format (multi-line, LoRa-tight):
Line 1: satellite emoji {sat_name} Pass -- {max_el} deg max
Line 2: AOS {aos_time} . LOS {los_time}
Line 3: {observer} . {direction}
"""
from __future__ import annotations
import json
import logging
import time
from datetime import datetime, timezone
from typing import Any, Optional
from meshai.adapter_config import adapter_config
from meshai.persistence import get_db
logger = logging.getLogger(__name__)
def _now() -> int:
return int(time.time())
def _coerce_float(v) -> Optional[float]:
if v is None:
return None
if isinstance(v, (int, float)):
return float(v)
try:
return float(v)
except (TypeError, ValueError):
return None
def _coerce_int(v) -> Optional[int]:
if v is None:
return None
if isinstance(v, int):
return v
try:
return int(v)
except (TypeError, ValueError):
return None
def _parse_iso_epoch(s) -> Optional[int]:
"""Parse ISO-8601 timestamp to epoch seconds."""
if not s or not isinstance(s, str):
return None
try:
dt = datetime.fromisoformat(s.replace("Z", "+00:00"))
return int(dt.timestamp())
except Exception:
return None
def _format_time(epoch: Optional[int]) -> str:
"""Format epoch to HH:MM local time string."""
if epoch is None:
return "?"
try:
dt = datetime.fromtimestamp(epoch)
return dt.strftime("%H:%M")
except Exception:
return "?"
def _map_severity(max_el: float) -> str:
"""Map max elevation to severity word."""
if max_el >= 60:
return "immediate"
if max_el >= 45:
return "priority"
return "routine"
def _canonical_id(norad_id: int, observer: str, aos_epoch: int) -> str:
"""Generate dedup-bucketed canonical event ID.
Bucket = floor(aos_epoch / 3600) -- one broadcast per sat/observer/hour.
"""
bucket = aos_epoch // 3600
return f"{norad_id}:{observer}:{bucket}"
def handle_satpass(envelope: dict, subject: str,
data: Optional[dict] = None,
now: Optional[int] = None) -> Optional[str]:
"""Process a satellite pass event from Central.
Returns wire message string if pass should be broadcast, None otherwise.
"""
if not isinstance(envelope, dict):
return None
inner = envelope.get("data") or {}
adapter = inner.get("adapter") or ""
# Only handle sat_passes adapter
if adapter != "sat_passes":
return None
d = inner.get("data") or {}
now = now if now is not None else _now()
# Extract pass data
norad_id = _coerce_int(d.get("norad_id") or d.get("satid"))
sat_name = d.get("sat_name") or d.get("name") or f"SAT-{norad_id}"
observer = d.get("observer") or d.get("location") or "unknown"
max_el = _coerce_float(d.get("max_elevation") or d.get("maxEl"))
aos_iso = d.get("aos") or d.get("rise_time")
los_iso = d.get("los") or d.get("set_time")
direction = d.get("direction") or d.get("pass_type") or ""
if norad_id is None or max_el is None:
logger.debug("satpass_handler: missing norad_id or max_elevation")
return None
aos_epoch = _parse_iso_epoch(aos_iso)
los_epoch = _parse_iso_epoch(los_iso)
if aos_epoch is None:
logger.debug("satpass_handler: could not parse aos time")
return None
# Config filters
cfg = adapter_config.satpass
# Observer filter (empty = all)
observers = getattr(cfg, "observers", []) or []
if observers and observer not in observers:
logger.debug("satpass_handler: observer %r not in configured list", observer)
return None
# NORAD ID filter (empty = all)
norad_ids = getattr(cfg, "norad_ids", []) or []
if norad_ids and norad_id not in norad_ids:
logger.debug("satpass_handler: norad_id %d not in configured list", norad_id)
return None
# Elevation floor
min_el = float(getattr(cfg, "min_elevation", 30))
if max_el < min_el:
logger.debug("satpass_handler: max_el %.1f below floor %.1f", max_el, min_el)
return None
# Generate canonical dedup ID
event_id = _canonical_id(norad_id, observer, aos_epoch)
severity_word = _map_severity(max_el)
category_raw = inner.get("category") or "sat.pass"
try:
conn = get_db()
except Exception:
logger.exception("satpass_handler: persistence unavailable")
return None
# Check for existing broadcast
row = conn.execute(
"SELECT last_broadcast_at FROM satpass_events WHERE event_id=?",
(event_id,)).fetchone()
payload_json = None
try:
payload_json = json.dumps(d, default=str)[:8000]
except Exception:
pass
# Log the event
log_id = _log_event_returning_id(
conn, now=now, source="satpass", category=category_raw,
severity_word=severity_word, event_id_external=event_id,
subject=subject, handled=0,
table_name="satpass_events", table_pk=event_id)
if row is None:
# First time seeing this pass bucket
_upsert_satpass(conn, event_id=event_id, norad_id=norad_id,
sat_name=sat_name, observer=observer,
max_elevation=max_el, aos_at=aos_epoch,
los_at=los_epoch, payload_json=payload_json,
first_seen_at=now, set_last_broadcast=False)
wire = _render(sat_name, max_el, aos_epoch, los_epoch, observer, direction)
_attach_commit(data, event_id=event_id, event_log_row_id=log_id)
return wire
if row["last_broadcast_at"] is None:
# Seen but not yet broadcast
wire = _render(sat_name, max_el, aos_epoch, los_epoch, observer, direction)
_attach_commit(data, event_id=event_id, event_log_row_id=log_id)
return wire
# Already broadcast this pass bucket
return None
def _render(sat_name: str, max_el: float, aos_epoch: Optional[int],
los_epoch: Optional[int], observer: str, direction: str) -> str:
"""Render wire format message."""
aos_str = _format_time(aos_epoch)
los_str = _format_time(los_epoch)
line1 = f"\U0001F6F0\uFE0F {sat_name} Pass \u2014 {int(max_el)}\u00B0 max"
line2 = f"AOS {aos_str} \u00B7 LOS {los_str}"
line3 = f"{observer}"
if direction:
line3 += f" \u00B7 {direction}"
return "\n".join([line1, line2, line3])
def _upsert_satpass(conn, *, event_id, norad_id, sat_name, observer,
max_elevation, aos_at, los_at, payload_json,
first_seen_at, set_last_broadcast=False,
broadcast_at=None) -> None:
"""Insert or update satpass_events row."""
existing = conn.execute(
"SELECT 1 FROM satpass_events WHERE event_id=?", (event_id,)).fetchone()
if existing is None:
conn.execute(
"INSERT INTO satpass_events(event_id, norad_id, sat_name, observer, "
"max_elevation, aos_at, los_at, payload_json, first_seen_at, "
"last_broadcast_at) VALUES (?,?,?,?,?,?,?,?,?,?)",
(event_id, norad_id, sat_name, observer, max_elevation, aos_at,
los_at, payload_json, first_seen_at,
broadcast_at if set_last_broadcast else None))
else:
conn.execute(
"UPDATE satpass_events SET sat_name=?, max_elevation=?, "
"payload_json=? WHERE event_id=?",
(sat_name, max_elevation, payload_json, event_id))
def _attach_commit(data: Optional[dict], *, event_id: str,
event_log_row_id: Optional[int]) -> None:
"""Attach post-broadcast commit callback."""
if not isinstance(data, dict):
return
def _on_commit(committed_at: float) -> None:
try:
conn = get_db()
except Exception:
logger.exception("satpass commit: persistence unavailable")
return
conn.execute(
"UPDATE satpass_events SET last_broadcast_at=?, "
"first_broadcast_at=COALESCE(first_broadcast_at, ?) WHERE event_id=?",
(int(committed_at), int(committed_at), event_id))
if event_log_row_id is not None:
conn.execute("UPDATE event_log SET handled=1 WHERE id=?",
(int(event_log_row_id),))
data["_on_broadcast_committed"] = _on_commit
data["_broadcast_audit"] = {"table": "satpass_events", "pk": event_id}
def _log_event_returning_id(conn, *, now, source, category, severity_word,
event_id_external, subject, handled,
table_name, table_pk) -> int:
"""Insert event_log row and return its ID."""
cur = conn.execute(
"INSERT INTO event_log(received_at, source, category, severity_word, "
"event_id_external, nats_subject, handled, table_name, table_pk) "
"VALUES (?,?,?,?,?,?,?,?,?)",
(now, source, category, severity_word, event_id_external, subject,
int(bool(handled)), table_name, table_pk))
return int(cur.lastrowid)
# Schema for satpass_events table (run once at startup via persistence)
SCHEMA_SATPASS_EVENTS = """
CREATE TABLE IF NOT EXISTS satpass_events (
event_id TEXT PRIMARY KEY,
norad_id INTEGER,
sat_name TEXT,
observer TEXT,
max_elevation REAL,
aos_at INTEGER,
los_at INTEGER,
payload_json TEXT,
first_seen_at INTEGER,
first_broadcast_at INTEGER,
last_broadcast_at INTEGER
);
CREATE INDEX IF NOT EXISTS idx_satpass_norad ON satpass_events(norad_id);
CREATE INDEX IF NOT EXISTS idx_satpass_observer ON satpass_events(observer);
CREATE INDEX IF NOT EXISTS idx_satpass_aos ON satpass_events(aos_at);
"""

View file

@ -583,7 +583,7 @@ class NotificationToggle:
TOGGLE_FAMILIES = [
"mesh_health", "weather", "fire", "rf_propagation",
"mesh_health", "weather", "fire", "rf_propagation", "satpass",
"roads", "avalanche", "seismic", "tracking",
]

View file

@ -51,6 +51,7 @@ _TOGGLE_PREFIX_FALLBACK = [
("solar_radiation", "rf_propagation"),
("rf_", "rf_propagation"),
("avalanche", "avalanche"),
("sat", "satpass"),
]
@ -524,6 +525,14 @@ ALERT_CATEGORIES = {
"example_message": "⛷ Avalanche Danger CONSIDERABLE: Sawtooth Zone — dangerous conditions on steep slopes.",
"toggle": "avalanche",
},
# Satellite passes
"sat_pass": {
"name": "Satellite Pass",
"description": "Notable satellite pass visible from your location (ISS, amateur sats)",
"default_severity": "routine",
"example_message": "🛰️ ISS Pass — 75° max\nAOS 19:32 · LOS 19:38\nBoise · NW→SE",
"toggle": "satpass",
},
}

View file

@ -0,0 +1,163 @@
"""v0.7 satpass_handler tests."""
import pytest
from unittest.mock import MagicMock, patch
def _envelope(norad_id=25544, sat_name="ISS", observer="Boise",
max_el=75.0, aos="2026-06-12T03:32:00Z",
los="2026-06-12T03:38:00Z", direction="NW-SE"):
"""Build a CloudEvents envelope for a satellite pass."""
return {
"specversion": "1.0",
"type": "central.sat.pass",
"source": "central",
"id": f"pass-{norad_id}-{aos}",
"data": {
"adapter": "sat_passes",
"category": "sat.pass",
"severity": 0,
"data": {
"norad_id": norad_id,
"sat_name": sat_name,
"observer": observer,
"max_elevation": max_el,
"aos": aos,
"los": los,
"direction": direction,
}
}
}
@pytest.fixture
def mock_db():
"""Mock database connection."""
conn = MagicMock()
conn.execute.return_value.fetchone.return_value = None
conn.execute.return_value.lastrowid = 1
with patch("meshai.central.satpass_handler.get_db", return_value=conn):
yield conn
@pytest.fixture
def mock_adapter_config():
"""Mock adapter_config.satpass."""
cfg = MagicMock()
cfg.observers = [] # empty = all observers
cfg.min_elevation = 30
cfg.norad_ids = [] # empty = all satellites
with patch("meshai.central.satpass_handler.adapter_config") as mock:
mock.satpass = cfg
yield cfg
class TestSatpassHandler:
"""Tests for handle_satpass function."""
def test_high_elevation_pass_broadcasts(self, mock_db, mock_adapter_config):
"""A pass with high elevation should broadcast."""
from meshai.central.satpass_handler import handle_satpass
env = _envelope(max_el=75.0)
result = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
assert result is not None
assert "ISS Pass" in result
assert "75" in result
assert "Boise" in result
def test_low_elevation_pass_filtered(self, mock_db, mock_adapter_config):
"""A pass below min_elevation should be filtered."""
from meshai.central.satpass_handler import handle_satpass
mock_adapter_config.min_elevation = 30
env = _envelope(max_el=25.0)
result = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
assert result is None
def test_observer_filter_blocks_mismatch(self, mock_db, mock_adapter_config):
"""A pass for non-configured observer should be filtered."""
from meshai.central.satpass_handler import handle_satpass
mock_adapter_config.observers = ["Magic Valley"]
env = _envelope(observer="Boise")
result = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
assert result is None
def test_observer_filter_allows_match(self, mock_db, mock_adapter_config):
"""A pass for configured observer should broadcast."""
from meshai.central.satpass_handler import handle_satpass
mock_adapter_config.observers = ["Boise", "Magic Valley"]
env = _envelope(observer="Boise", max_el=45.0)
result = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
assert result is not None
def test_norad_id_filter(self, mock_db, mock_adapter_config):
"""NORAD ID filter should block non-matching satellites."""
from meshai.central.satpass_handler import handle_satpass
mock_adapter_config.norad_ids = [25544] # ISS only
env = _envelope(norad_id=12345, max_el=60.0)
result = handle_satpass(env, "central.sat.pass.other", data={}, now=1718163120)
assert result is None
def test_dedup_blocks_second_broadcast(self, mock_db, mock_adapter_config):
"""Second pass in same hour bucket should be deduplicated."""
from meshai.central.satpass_handler import handle_satpass
# First call returns no existing broadcast
mock_db.execute.return_value.fetchone.return_value = None
env = _envelope(max_el=60.0)
result1 = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
assert result1 is not None
# Second call simulates existing broadcast
mock_db.execute.return_value.fetchone.return_value = {"last_broadcast_at": 1718163120}
result2 = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163180)
assert result2 is None
def test_wire_format(self, mock_db, mock_adapter_config):
"""Wire format should have 3 lines with correct info."""
from meshai.central.satpass_handler import handle_satpass
env = _envelope(sat_name="ISS", max_el=75, observer="Boise", direction="NW-SE")
result = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
lines = result.split("\n")
assert len(lines) == 3
assert "ISS Pass" in lines[0]
assert "75" in lines[0]
assert "AOS" in lines[1]
assert "LOS" in lines[1]
assert "Boise" in lines[2]
def test_commit_callback_attached(self, mock_db, mock_adapter_config):
"""Broadcast should attach commit callback."""
from meshai.central.satpass_handler import handle_satpass
data = {}
env = _envelope(max_el=60.0)
result = handle_satpass(env, "central.sat.pass.iss", data=data, now=1718163120)
assert result is not None
assert "_on_broadcast_committed" in data
assert "_broadcast_audit" in data
assert data["_broadcast_audit"]["table"] == "satpass_events"
def test_wrong_adapter_ignored(self, mock_db, mock_adapter_config):
"""Envelope with wrong adapter should be ignored."""
from meshai.central.satpass_handler import handle_satpass
env = _envelope()
env["data"]["adapter"] = "some_other_adapter"
result = handle_satpass(env, "central.sat.pass.iss", data={}, now=1718163120)
assert result is None