diff --git a/meshai/adapter_config/defaults.py b/meshai/adapter_config/defaults.py index d96346c..63d4944 100644 --- a/meshai/adapter_config/defaults.py +++ b/meshai/adapter_config/defaults.py @@ -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).", + }, } diff --git a/meshai/central/consumer.py b/meshai/central/consumer.py index 82799a3..300b359 100644 --- a/meshai/central/consumer.py +++ b/meshai/central/consumer.py @@ -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 diff --git a/meshai/central/satpass_handler.py b/meshai/central/satpass_handler.py new file mode 100644 index 0000000..14992c2 --- /dev/null +++ b/meshai/central/satpass_handler.py @@ -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); +""" diff --git a/meshai/config.py b/meshai/config.py index 7aa291d..894e60b 100644 --- a/meshai/config.py +++ b/meshai/config.py @@ -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", ] diff --git a/meshai/notifications/categories.py b/meshai/notifications/categories.py index faf5b76..ba9bfc8 100644 --- a/meshai/notifications/categories.py +++ b/meshai/notifications/categories.py @@ -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", + }, } diff --git a/tests/test_satpass_handler.py b/tests/test_satpass_handler.py new file mode 100644 index 0000000..020028d --- /dev/null +++ b/tests/test_satpass_handler.py @@ -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