mirror of
https://github.com/zvx-echo6/meshai.git
synced 2026-08-26 17:31:34 +00:00
feat(persistence): make satpass/avalanche/ducting LLM-queryable (#74)
Close the LLM data gaps: add build_satpass_detail (satpass_events was written but had no reader), and give avalanche + ducting durable tables (v24/v25) with native writers + env_reporter readers so the mesh LLM can answer avalanche, satellite-pass, and RF-propagation questions. Persistence-only; no broadcast/ gating changes. Co-authored-by: Matt Johnson <mj@k7zvx.com> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
d479ca537a
commit
3fb4e6e65c
9 changed files with 570 additions and 8 deletions
61
work/meshai/env/avalanche.py
vendored
61
work/meshai/env/avalanche.py
vendored
|
|
@ -191,11 +191,72 @@ class AvalancheAdapter:
|
||||||
|
|
||||||
self._is_loaded = True
|
self._is_loaded = True
|
||||||
|
|
||||||
|
# Persistence add (LLM-queryable): mirror the WFIGS fire path and write
|
||||||
|
# the current advisories into the durable avalanche_events table so the
|
||||||
|
# mesh LLM (env_reporter.build_avalanche_detail) can answer avalanche
|
||||||
|
# questions and state survives a restart. Persistence-only -- does NOT
|
||||||
|
# touch broadcast/gating (that stays in _delta_emit / to_event).
|
||||||
|
self._persist_events()
|
||||||
|
|
||||||
if changed:
|
if changed:
|
||||||
logger.info(f"Avalanche advisories updated: {len(new_events)} active zones")
|
logger.info(f"Avalanche advisories updated: {len(new_events)} active zones")
|
||||||
|
|
||||||
return changed
|
return changed
|
||||||
|
|
||||||
|
def _persist_events(self) -> None:
|
||||||
|
"""UPSERT the current in-memory advisories into avalanche_events.
|
||||||
|
|
||||||
|
Best-effort: a DB error is logged and swallowed so persistence never
|
||||||
|
breaks polling. Keyed by the stable adapter event_id
|
||||||
|
("avy_{center}_{zone}") -- an existing zone is updated in place.
|
||||||
|
"""
|
||||||
|
if not self._events:
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
from meshai.persistence import get_db
|
||||||
|
conn = get_db()
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("avalanche persist skipped (DB unavailable): %s", e)
|
||||||
|
return
|
||||||
|
|
||||||
|
now = int(time.time())
|
||||||
|
for evt in self._events:
|
||||||
|
try:
|
||||||
|
event_id = evt.get("event_id")
|
||||||
|
if not event_id:
|
||||||
|
continue
|
||||||
|
expires = evt.get("expires")
|
||||||
|
expires_at = int(expires) if expires is not None else None
|
||||||
|
exists = conn.execute(
|
||||||
|
"SELECT 1 FROM avalanche_events WHERE event_id=?",
|
||||||
|
(event_id,),
|
||||||
|
).fetchone()
|
||||||
|
if exists is None:
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO avalanche_events(event_id, center_id, "
|
||||||
|
"zone_name, danger_level, danger_name, travel_advice, "
|
||||||
|
"lat, lon, expires_at, first_seen_at, last_event_at) "
|
||||||
|
"VALUES (?,?,?,?,?,?,?,?,?,?,?)",
|
||||||
|
(event_id, evt.get("center_id"), evt.get("zone_name"),
|
||||||
|
evt.get("danger_level"), evt.get("danger_name"),
|
||||||
|
evt.get("travel_advice"), evt.get("lat"),
|
||||||
|
evt.get("lon"), expires_at, now, now),
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE avalanche_events SET center_id=?, zone_name=?, "
|
||||||
|
"danger_level=?, danger_name=?, travel_advice=?, "
|
||||||
|
"lat=?, lon=?, expires_at=?, last_event_at=? "
|
||||||
|
"WHERE event_id=?",
|
||||||
|
(evt.get("center_id"), evt.get("zone_name"),
|
||||||
|
evt.get("danger_level"), evt.get("danger_name"),
|
||||||
|
evt.get("travel_advice"), evt.get("lat"),
|
||||||
|
evt.get("lon"), expires_at, now, event_id),
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.exception(
|
||||||
|
"avalanche persist failed for %s", evt.get("event_id", "?"))
|
||||||
|
|
||||||
def _compute_centroid(self, geom) -> tuple:
|
def _compute_centroid(self, geom) -> tuple:
|
||||||
"""Compute centroid from GeoJSON geometry."""
|
"""Compute centroid from GeoJSON geometry."""
|
||||||
if not geom:
|
if not geom:
|
||||||
|
|
|
||||||
42
work/meshai/env/ducting.py
vendored
42
work/meshai/env/ducting.py
vendored
|
|
@ -265,6 +265,48 @@ class DuctingAdapter:
|
||||||
|
|
||||||
self._update_events()
|
self._update_events()
|
||||||
|
|
||||||
|
# Persistence add (LLM-queryable): after the tier is committed, write
|
||||||
|
# the single current assessment into the durable ducting_events table
|
||||||
|
# so the mesh LLM (env_reporter.build_ducting_detail) can answer
|
||||||
|
# RF-propagation questions and the assessment survives a restart.
|
||||||
|
# Persistence-only -- does NOT touch broadcast/gating.
|
||||||
|
self._persist_status()
|
||||||
|
|
||||||
|
def _persist_status(self) -> None:
|
||||||
|
"""UPSERT the current single-point assessment into ducting_events.
|
||||||
|
|
||||||
|
Ducting is a single-point periodic assessment, so we keep ONE durable
|
||||||
|
current row per location (id = "ducting_{lat}_{lon}"), refreshed each
|
||||||
|
poll. Best-effort: a DB error is logged and swallowed so persistence
|
||||||
|
never breaks polling.
|
||||||
|
"""
|
||||||
|
if not self._status:
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
from meshai.persistence import get_db
|
||||||
|
conn = get_db()
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("ducting persist skipped (DB unavailable): %s", e)
|
||||||
|
return
|
||||||
|
|
||||||
|
try:
|
||||||
|
loc = f"{round(self._lat, 2)}_{round(self._lon, 2)}"
|
||||||
|
row_id = f"ducting_{loc}"
|
||||||
|
assessed_at = int(self._status.get("fetched_at") or time.time())
|
||||||
|
conn.execute(
|
||||||
|
"INSERT OR REPLACE INTO ducting_events(id, lat, lon, "
|
||||||
|
"condition, tier, min_gradient, duct_base_m, duct_thickness_m, "
|
||||||
|
"assessment, assessed_at) VALUES (?,?,?,?,?,?,?,?,?,?)",
|
||||||
|
(row_id, self._lat, self._lon,
|
||||||
|
self._status.get("condition"), self._status.get("tier"),
|
||||||
|
self._status.get("min_gradient"),
|
||||||
|
self._status.get("duct_base_m"),
|
||||||
|
self._status.get("duct_thickness_m"),
|
||||||
|
self._status.get("assessment"), assessed_at),
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.exception("ducting persist failed")
|
||||||
|
|
||||||
# Tier model (enhancement strength, ascending):
|
# Tier model (enhancement strength, ascending):
|
||||||
# normal < super_refraction < duct < surface_duct
|
# normal < super_refraction < duct < surface_duct
|
||||||
_TIER_ORDER = ["normal", "super_refraction", "duct", "surface_duct"]
|
_TIER_ORDER = ["normal", "super_refraction", "duct", "surface_duct"]
|
||||||
|
|
|
||||||
|
|
@ -384,6 +384,106 @@ class EnvReporter:
|
||||||
|
|
||||||
return ("\n".join(lines) if lines else "")[:_block_cap()]
|
return ("\n".join(lines) if lines else "")[:_block_cap()]
|
||||||
|
|
||||||
|
def build_satpass_detail(self, *, hours: int = 24,
|
||||||
|
limit: int = 8,
|
||||||
|
now: Optional[int] = None) -> str:
|
||||||
|
"""Upcoming satellite passes from satpass_events (AOS in the near
|
||||||
|
future). Data is written by the native satpass path; this is the
|
||||||
|
missing LLM reader so the bot can answer satellite-pass questions.
|
||||||
|
"""
|
||||||
|
if not self._adapter_included("satpass"):
|
||||||
|
return ""
|
||||||
|
now = now if now is not None else int(time.time())
|
||||||
|
try: conn = self._conn_factory()
|
||||||
|
except Exception: return ""
|
||||||
|
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT sat_name, norad_id, observer, max_elevation, aos_at, los_at "
|
||||||
|
"FROM satpass_events WHERE aos_at IS NOT NULL "
|
||||||
|
"AND aos_at >= ? AND aos_at <= ? "
|
||||||
|
"ORDER BY aos_at ASC LIMIT ?",
|
||||||
|
(now, now + hours * 3600, limit),
|
||||||
|
).fetchall()
|
||||||
|
if not rows:
|
||||||
|
return ""
|
||||||
|
lines = [f"UPCOMING SATELLITE PASSES (next {hours}h):"]
|
||||||
|
for r in rows:
|
||||||
|
name = r["sat_name"] or (
|
||||||
|
f"NORAD {r['norad_id']}" if r["norad_id"] is not None else "sat?")
|
||||||
|
aos = _fmt_epoch(r["aos_at"]) if r["aos_at"] else "?"
|
||||||
|
el = (f"max el {int(r['max_elevation'])}°"
|
||||||
|
if r["max_elevation"] is not None else "el ?")
|
||||||
|
obs = f" @ {r['observer']}" if r["observer"] else ""
|
||||||
|
lines.append(f" - {name}: AOS {aos}, {el}{obs}")
|
||||||
|
return "\n".join(lines)[:_block_cap()]
|
||||||
|
|
||||||
|
def build_avalanche_detail(self, *, limit: int = 10,
|
||||||
|
now: Optional[int] = None) -> str:
|
||||||
|
"""Current (non-expired) avalanche advisories from avalanche_events.
|
||||||
|
NAADS danger scale 1-5. Off-season the table is empty -> empty block.
|
||||||
|
"""
|
||||||
|
if not self._adapter_included("avalanche"):
|
||||||
|
return ""
|
||||||
|
now = now if now is not None else int(time.time())
|
||||||
|
try: conn = self._conn_factory()
|
||||||
|
except Exception: return ""
|
||||||
|
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT zone_name, center_id, danger_level, danger_name, "
|
||||||
|
"travel_advice, expires_at FROM avalanche_events "
|
||||||
|
"WHERE expires_at IS NULL OR expires_at >= ? "
|
||||||
|
"ORDER BY danger_level DESC, zone_name ASC LIMIT ?",
|
||||||
|
(now, limit),
|
||||||
|
).fetchall()
|
||||||
|
if not rows:
|
||||||
|
return ""
|
||||||
|
lines = ["AVALANCHE ADVISORIES (avalanche.org, NAADS 1-5):"]
|
||||||
|
for r in rows:
|
||||||
|
zone = r["zone_name"] or "(zone?)"
|
||||||
|
dl = r["danger_level"]
|
||||||
|
level = r["danger_name"] or "?"
|
||||||
|
dl_str = f"{level} ({dl})" if dl is not None else level
|
||||||
|
center = f" [{r['center_id']}]" if r["center_id"] else ""
|
||||||
|
advice = (r["travel_advice"] or "")[:90]
|
||||||
|
advice_str = f" -- {advice}" if advice else ""
|
||||||
|
lines.append(f" - {zone}{center}: {dl_str}{advice_str}")
|
||||||
|
return "\n".join(lines)[:_block_cap()]
|
||||||
|
|
||||||
|
def build_ducting_detail(self, *, now: Optional[int] = None) -> str:
|
||||||
|
"""Latest tropospheric ducting / RF-propagation assessment from
|
||||||
|
ducting_events (single current row per location, most recent first).
|
||||||
|
Lets the LLM answer band-opening / extended-UHF-range questions.
|
||||||
|
"""
|
||||||
|
if not self._adapter_included("ducting"):
|
||||||
|
return ""
|
||||||
|
try: conn = self._conn_factory()
|
||||||
|
except Exception: return ""
|
||||||
|
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT lat, lon, condition, tier, min_gradient, duct_base_m, "
|
||||||
|
"duct_thickness_m, assessment, assessed_at FROM ducting_events "
|
||||||
|
"ORDER BY assessed_at DESC LIMIT 3",
|
||||||
|
).fetchall()
|
||||||
|
if not rows:
|
||||||
|
return ""
|
||||||
|
lines = ["RF PROPAGATION (tropospheric ducting assessment):"]
|
||||||
|
for r in rows:
|
||||||
|
tier = r["tier"] or r["condition"] or "?"
|
||||||
|
assessment = r["assessment"] or ""
|
||||||
|
mg = (f", min M-gradient {r['min_gradient']}/km"
|
||||||
|
if r["min_gradient"] is not None else "")
|
||||||
|
base = r["duct_base_m"]
|
||||||
|
thick = r["duct_thickness_m"]
|
||||||
|
duct = ""
|
||||||
|
if base is not None and thick:
|
||||||
|
duct = f", duct base {int(base)} m ~{int(thick)} m thick"
|
||||||
|
loc = (f"{r['lat']:.2f},{r['lon']:.2f}"
|
||||||
|
if r["lat"] is not None else "loc?")
|
||||||
|
when = _fmt_epoch(r["assessed_at"])
|
||||||
|
assess_str = f": {assessment}" if assessment else ""
|
||||||
|
lines.append(f" - {loc} [{tier}]{assess_str}{mg}{duct} (as of {when})")
|
||||||
|
return "\n".join(lines)[:_block_cap()]
|
||||||
|
|
||||||
def build_drop_audit(self, *, hours: int = 1) -> str:
|
def build_drop_audit(self, *, hours: int = 1) -> str:
|
||||||
"""Why-was-X-dropped: event_log handled=0 grouped by source+reason
|
"""Why-was-X-dropped: event_log handled=0 grouped by source+reason
|
||||||
+ the dispatcher_state cumulative counters."""
|
+ the dispatcher_state cumulative counters."""
|
||||||
|
|
@ -434,6 +534,9 @@ class EnvReporter:
|
||||||
self.build_traffic_detail(now=now),
|
self.build_traffic_detail(now=now),
|
||||||
self.build_gauges_detail(now=now),
|
self.build_gauges_detail(now=now),
|
||||||
self.build_swpc_detail(now=now),
|
self.build_swpc_detail(now=now),
|
||||||
|
self.build_satpass_detail(now=now),
|
||||||
|
self.build_avalanche_detail(now=now),
|
||||||
|
self.build_ducting_detail(now=now),
|
||||||
]
|
]
|
||||||
return "\n\n".join(p for p in parts if p)
|
return "\n\n".join(p for p in parts if p)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -30,7 +30,7 @@ logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
DEFAULT_DB_PATH = "/data/meshai.sqlite"
|
DEFAULT_DB_PATH = "/data/meshai.sqlite"
|
||||||
MESHAI_DB_PATH_ENV = "MESHAI_DB_PATH"
|
MESHAI_DB_PATH_ENV = "MESHAI_DB_PATH"
|
||||||
SCHEMA_VERSION = 23
|
SCHEMA_VERSION = 25
|
||||||
SCHEMA_META_TABLE = "schema_meta"
|
SCHEMA_META_TABLE = "schema_meta"
|
||||||
MIGRATIONS_DIR = Path(__file__).parent / "migrations"
|
MIGRATIONS_DIR = Path(__file__).parent / "migrations"
|
||||||
|
|
||||||
|
|
|
||||||
27
work/meshai/persistence/migrations/v24.sql
Normal file
27
work/meshai/persistence/migrations/v24.sql
Normal file
|
|
@ -0,0 +1,27 @@
|
||||||
|
-- v24 avalanche durable hazard table (LLM-queryable).
|
||||||
|
--
|
||||||
|
-- The native env/avalanche.py adapter previously kept advisories in memory
|
||||||
|
-- only (self._events), so the mesh LLM (env_reporter) was blind to avalanche
|
||||||
|
-- danger and state was lost on restart. This table gives each active zone a
|
||||||
|
-- durable row keyed by the adapter's stable event_id ("avy_{center}_{zone}").
|
||||||
|
--
|
||||||
|
-- Persistence-only: does NOT affect broadcast/gating (that stays driven by
|
||||||
|
-- the pipeline + adapter_config). NAADS danger scale 1-5. IF NOT EXISTS so a
|
||||||
|
-- fresh install and a re-run are both no-ops.
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS avalanche_events (
|
||||||
|
event_id TEXT PRIMARY KEY, -- avy_{center_id}_{zone}
|
||||||
|
center_id TEXT,
|
||||||
|
zone_name TEXT,
|
||||||
|
danger_level INTEGER, -- NAADS 1-5 (-1/0 = no rating)
|
||||||
|
danger_name TEXT,
|
||||||
|
travel_advice TEXT,
|
||||||
|
lat REAL,
|
||||||
|
lon REAL,
|
||||||
|
expires_at INTEGER, -- epoch seconds; end-of-day validity
|
||||||
|
first_seen_at INTEGER NOT NULL,
|
||||||
|
last_event_at INTEGER
|
||||||
|
);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_avalanche_expires ON avalanche_events(expires_at);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_avalanche_center ON avalanche_events(center_id);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_avalanche_danger ON avalanche_events(danger_level);
|
||||||
27
work/meshai/persistence/migrations/v25.sql
Normal file
27
work/meshai/persistence/migrations/v25.sql
Normal file
|
|
@ -0,0 +1,27 @@
|
||||||
|
-- v25 tropospheric ducting durable assessment table (LLM-queryable).
|
||||||
|
--
|
||||||
|
-- The native env/ducting.py adapter previously kept its RF-propagation
|
||||||
|
-- assessment in memory only (self._status / self._events), so the mesh LLM
|
||||||
|
-- (env_reporter) could not answer ducting / band-opening questions and the
|
||||||
|
-- assessment was lost on restart. Ducting is a single-point periodic
|
||||||
|
-- assessment, so we keep one durable current row per assessment location
|
||||||
|
-- (id = "ducting_{lat}_{lon}") that the writer UPSERTs each poll.
|
||||||
|
--
|
||||||
|
-- Persistence-only: does NOT affect broadcast/gating (tier hysteresis +
|
||||||
|
-- Event emission stay in the adapter/pipeline). IF NOT EXISTS so a fresh
|
||||||
|
-- install and a re-run are both no-ops.
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS ducting_events (
|
||||||
|
id TEXT PRIMARY KEY, -- ducting_{round(lat,2)}_{round(lon,2)}
|
||||||
|
lat REAL,
|
||||||
|
lon REAL,
|
||||||
|
condition TEXT, -- normal | super_refraction | surface_duct | elevated_duct
|
||||||
|
tier TEXT, -- normal | super_refraction | duct | surface_duct
|
||||||
|
min_gradient REAL, -- min modified-refractivity gradient (M-units/km)
|
||||||
|
duct_base_m REAL,
|
||||||
|
duct_thickness_m REAL,
|
||||||
|
assessment TEXT, -- human-readable summary line
|
||||||
|
assessed_at INTEGER NOT NULL -- epoch seconds of the assessment
|
||||||
|
);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_ducting_assessed ON ducting_events(assessed_at);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_ducting_tier ON ducting_events(tier);
|
||||||
301
work/tests/test_env_reporter_llm_gaps.py
Normal file
301
work/tests/test_env_reporter_llm_gaps.py
Normal file
|
|
@ -0,0 +1,301 @@
|
||||||
|
"""LLM-persistence-gap tests: satpass reader, avalanche + ducting
|
||||||
|
durable tables (writer -> table -> env_reporter reader).
|
||||||
|
|
||||||
|
Uses the autouse conftest fixture which points MESHAI_DB_PATH at a fresh
|
||||||
|
tmp file and runs init_db (so all migrations, now through v25, apply).
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import time
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from meshai.env.avalanche import AvalancheAdapter
|
||||||
|
from meshai.env.ducting import DuctingAdapter
|
||||||
|
from meshai.notifications.env_reporter import EnvReporter
|
||||||
|
from meshai.persistence import get_db
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def reporter():
|
||||||
|
return EnvReporter()
|
||||||
|
|
||||||
|
|
||||||
|
# ============================================================================
|
||||||
|
# migrations: the two new tables exist on a fresh boot
|
||||||
|
# ============================================================================
|
||||||
|
|
||||||
|
|
||||||
|
def test_new_tables_exist_after_init():
|
||||||
|
conn = get_db()
|
||||||
|
names = {r["name"] for r in conn.execute(
|
||||||
|
"SELECT name FROM sqlite_master WHERE type='table'").fetchall()}
|
||||||
|
assert "avalanche_events" in names
|
||||||
|
assert "ducting_events" in names
|
||||||
|
|
||||||
|
|
||||||
|
# ============================================================================
|
||||||
|
# 1. satpass reader (data already persists via the native satpass path)
|
||||||
|
# ============================================================================
|
||||||
|
|
||||||
|
|
||||||
|
def _seed_satpass(conn, *, event_id, sat_name, aos_at, norad_id=25544,
|
||||||
|
observer="Boise", max_elevation=45.0, los_at=None):
|
||||||
|
now = int(time.time())
|
||||||
|
conn.execute(
|
||||||
|
"INSERT OR REPLACE INTO satpass_events(event_id, norad_id, sat_name, "
|
||||||
|
"observer, max_elevation, aos_at, los_at, payload_json, first_seen_at) "
|
||||||
|
"VALUES (?,?,?,?,?,?,?,?,?)",
|
||||||
|
(event_id, norad_id, sat_name, observer, max_elevation, aos_at,
|
||||||
|
los_at or (aos_at + 600), "{}", now),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_satpass_detail_empty_when_no_passes(reporter):
|
||||||
|
assert reporter.build_satpass_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_satpass_detail_renders_upcoming(reporter):
|
||||||
|
conn = get_db()
|
||||||
|
now = int(time.time())
|
||||||
|
_seed_satpass(conn, event_id="P1", sat_name="ISS (ZARYA)",
|
||||||
|
aos_at=now + 1800, max_elevation=52.0, observer="Boise")
|
||||||
|
text = reporter.build_satpass_detail()
|
||||||
|
assert "UPCOMING SATELLITE PASSES" in text
|
||||||
|
assert "ISS (ZARYA)" in text
|
||||||
|
assert "max el 52" in text
|
||||||
|
assert "Boise" in text
|
||||||
|
|
||||||
|
|
||||||
|
def test_satpass_detail_ignores_past_passes(reporter):
|
||||||
|
conn = get_db()
|
||||||
|
now = int(time.time())
|
||||||
|
# A pass whose AOS already happened must not appear.
|
||||||
|
_seed_satpass(conn, event_id="P_old", sat_name="NOAA 19",
|
||||||
|
aos_at=now - 3600)
|
||||||
|
assert reporter.build_satpass_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_satpass_detail_meta_off(reporter):
|
||||||
|
conn = get_db()
|
||||||
|
now = int(time.time())
|
||||||
|
_seed_satpass(conn, event_id="P1", sat_name="ISS", aos_at=now + 600)
|
||||||
|
conn.execute(
|
||||||
|
"INSERT OR REPLACE INTO adapter_meta(adapter, include_in_llm_context, "
|
||||||
|
"updated_at) VALUES ('satpass', 0, ?)", (time.time(),))
|
||||||
|
assert reporter.build_satpass_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
# ============================================================================
|
||||||
|
# 2. avalanche: adapter writer -> avalanche_events -> reader
|
||||||
|
# ============================================================================
|
||||||
|
|
||||||
|
|
||||||
|
def _avy_config():
|
||||||
|
cfg = MagicMock()
|
||||||
|
cfg.center_ids = ["SNFAC"]
|
||||||
|
cfg.tick_seconds = 1800
|
||||||
|
cfg.season_months = [12, 1, 2, 3, 4]
|
||||||
|
return cfg
|
||||||
|
|
||||||
|
|
||||||
|
def test_avalanche_writer_persists_row_then_reader_reads_it(reporter):
|
||||||
|
adapter = AvalancheAdapter(_avy_config())
|
||||||
|
now = time.time()
|
||||||
|
# Synthetic assessment mirroring _fetch()'s stored-event dict shape.
|
||||||
|
adapter._events = [{
|
||||||
|
"source": "avalanche",
|
||||||
|
"event_id": "avy_SNFAC_banner_summit",
|
||||||
|
"center_id": "SNFAC",
|
||||||
|
"zone_name": "Banner Summit",
|
||||||
|
"danger_level": 4,
|
||||||
|
"danger_name": "High",
|
||||||
|
"travel_advice": "Very dangerous avalanche conditions.",
|
||||||
|
"lat": 44.3,
|
||||||
|
"lon": -115.2,
|
||||||
|
"expires": now + 6 * 3600,
|
||||||
|
"fetched_at": now,
|
||||||
|
}]
|
||||||
|
adapter._persist_events()
|
||||||
|
|
||||||
|
conn = get_db()
|
||||||
|
row = conn.execute(
|
||||||
|
"SELECT * FROM avalanche_events WHERE event_id=?",
|
||||||
|
("avy_SNFAC_banner_summit",)).fetchone()
|
||||||
|
assert row is not None
|
||||||
|
assert row["danger_level"] == 4
|
||||||
|
assert row["zone_name"] == "Banner Summit"
|
||||||
|
|
||||||
|
text = reporter.build_avalanche_detail()
|
||||||
|
assert "AVALANCHE ADVISORIES" in text
|
||||||
|
assert "Banner Summit" in text
|
||||||
|
assert "High (4)" in text
|
||||||
|
|
||||||
|
|
||||||
|
def test_avalanche_writer_upserts_in_place(reporter):
|
||||||
|
adapter = AvalancheAdapter(_avy_config())
|
||||||
|
now = time.time()
|
||||||
|
base = {
|
||||||
|
"source": "avalanche",
|
||||||
|
"event_id": "avy_SNFAC_z1",
|
||||||
|
"center_id": "SNFAC",
|
||||||
|
"zone_name": "Zone One",
|
||||||
|
"danger_level": 2,
|
||||||
|
"danger_name": "Moderate",
|
||||||
|
"travel_advice": "",
|
||||||
|
"lat": 44.0,
|
||||||
|
"lon": -115.0,
|
||||||
|
"expires": now + 6 * 3600,
|
||||||
|
"fetched_at": now,
|
||||||
|
}
|
||||||
|
adapter._events = [dict(base)]
|
||||||
|
adapter._persist_events()
|
||||||
|
# danger rises -> re-persist should update the same row, not duplicate.
|
||||||
|
adapter._events = [dict(base, danger_level=4, danger_name="High")]
|
||||||
|
adapter._persist_events()
|
||||||
|
|
||||||
|
conn = get_db()
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT danger_level FROM avalanche_events WHERE event_id=?",
|
||||||
|
("avy_SNFAC_z1",)).fetchall()
|
||||||
|
assert len(rows) == 1
|
||||||
|
assert rows[0]["danger_level"] == 4
|
||||||
|
|
||||||
|
|
||||||
|
def test_avalanche_detail_excludes_expired(reporter):
|
||||||
|
adapter = AvalancheAdapter(_avy_config())
|
||||||
|
now = time.time()
|
||||||
|
adapter._events = [{
|
||||||
|
"source": "avalanche",
|
||||||
|
"event_id": "avy_SNFAC_old",
|
||||||
|
"center_id": "SNFAC",
|
||||||
|
"zone_name": "Stale Zone",
|
||||||
|
"danger_level": 3,
|
||||||
|
"danger_name": "Considerable",
|
||||||
|
"travel_advice": "",
|
||||||
|
"lat": 44.0,
|
||||||
|
"lon": -115.0,
|
||||||
|
"expires": now - 3600, # already expired
|
||||||
|
"fetched_at": now,
|
||||||
|
}]
|
||||||
|
adapter._persist_events()
|
||||||
|
assert reporter.build_avalanche_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_avalanche_detail_empty_when_no_rows(reporter):
|
||||||
|
assert reporter.build_avalanche_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
# ============================================================================
|
||||||
|
# 3. ducting: adapter writer -> ducting_events -> reader
|
||||||
|
# ============================================================================
|
||||||
|
|
||||||
|
|
||||||
|
def _ducting_config():
|
||||||
|
cfg = MagicMock()
|
||||||
|
cfg.latitude = 43.6
|
||||||
|
cfg.longitude = -116.2
|
||||||
|
cfg.tick_seconds = 10800
|
||||||
|
return cfg
|
||||||
|
|
||||||
|
|
||||||
|
def test_ducting_writer_persists_row_then_reader_reads_it(reporter):
|
||||||
|
adapter = DuctingAdapter(_ducting_config())
|
||||||
|
adapter._status = {
|
||||||
|
"condition": "surface_duct",
|
||||||
|
"tier": "surface_duct",
|
||||||
|
"min_gradient": -120.0,
|
||||||
|
"duct_base_m": 110,
|
||||||
|
"duct_thickness_m": 650,
|
||||||
|
"assessment": "Ducting -- extended UHF range likely",
|
||||||
|
"fetched_at": time.time(),
|
||||||
|
}
|
||||||
|
adapter._persist_status()
|
||||||
|
|
||||||
|
conn = get_db()
|
||||||
|
row = conn.execute(
|
||||||
|
"SELECT * FROM ducting_events WHERE id=?",
|
||||||
|
("ducting_43.6_-116.2",)).fetchone()
|
||||||
|
assert row is not None
|
||||||
|
assert row["tier"] == "surface_duct"
|
||||||
|
assert row["min_gradient"] == -120.0
|
||||||
|
|
||||||
|
text = reporter.build_ducting_detail()
|
||||||
|
assert "RF PROPAGATION" in text
|
||||||
|
assert "surface_duct" in text
|
||||||
|
assert "extended UHF range" in text
|
||||||
|
|
||||||
|
|
||||||
|
def test_ducting_writer_keeps_single_current_row(reporter):
|
||||||
|
adapter = DuctingAdapter(_ducting_config())
|
||||||
|
adapter._status = {
|
||||||
|
"condition": "normal", "tier": "normal", "min_gradient": 118.0,
|
||||||
|
"duct_base_m": None, "duct_thickness_m": None,
|
||||||
|
"assessment": "Normal propagation", "fetched_at": time.time(),
|
||||||
|
}
|
||||||
|
adapter._persist_status()
|
||||||
|
adapter._status = dict(adapter._status, condition="super_refraction",
|
||||||
|
tier="super_refraction", min_gradient=40.0,
|
||||||
|
assessment="Enhanced range possible",
|
||||||
|
fetched_at=time.time() + 1)
|
||||||
|
adapter._persist_status()
|
||||||
|
|
||||||
|
conn = get_db()
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT tier FROM ducting_events WHERE id=?",
|
||||||
|
("ducting_43.6_-116.2",)).fetchall()
|
||||||
|
assert len(rows) == 1
|
||||||
|
assert rows[0]["tier"] == "super_refraction"
|
||||||
|
|
||||||
|
|
||||||
|
def test_ducting_detail_empty_when_no_rows(reporter):
|
||||||
|
assert reporter.build_ducting_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_ducting_detail_meta_off(reporter):
|
||||||
|
adapter = DuctingAdapter(_ducting_config())
|
||||||
|
adapter._status = {
|
||||||
|
"condition": "surface_duct", "tier": "surface_duct",
|
||||||
|
"min_gradient": -120.0, "duct_base_m": 110, "duct_thickness_m": 650,
|
||||||
|
"assessment": "Ducting", "fetched_at": time.time(),
|
||||||
|
}
|
||||||
|
adapter._persist_status()
|
||||||
|
conn = get_db()
|
||||||
|
conn.execute(
|
||||||
|
"INSERT OR REPLACE INTO adapter_meta(adapter, include_in_llm_context, "
|
||||||
|
"updated_at) VALUES ('ducting', 0, ?)", (time.time(),))
|
||||||
|
assert reporter.build_ducting_detail() == ""
|
||||||
|
|
||||||
|
|
||||||
|
# ============================================================================
|
||||||
|
# build_all wires the three new blocks
|
||||||
|
# ============================================================================
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_all_includes_new_blocks(reporter):
|
||||||
|
conn = get_db()
|
||||||
|
now = int(time.time())
|
||||||
|
_seed_satpass(conn, event_id="P1", sat_name="ISS", aos_at=now + 1200)
|
||||||
|
|
||||||
|
avy = AvalancheAdapter(_avy_config())
|
||||||
|
avy._events = [{
|
||||||
|
"source": "avalanche", "event_id": "avy_SNFAC_z",
|
||||||
|
"center_id": "SNFAC", "zone_name": "Banner", "danger_level": 4,
|
||||||
|
"danger_name": "High", "travel_advice": "", "lat": 44.3, "lon": -115.2,
|
||||||
|
"expires": time.time() + 6 * 3600, "fetched_at": time.time(),
|
||||||
|
}]
|
||||||
|
avy._persist_events()
|
||||||
|
|
||||||
|
duct = DuctingAdapter(_ducting_config())
|
||||||
|
duct._status = {
|
||||||
|
"condition": "surface_duct", "tier": "surface_duct",
|
||||||
|
"min_gradient": -120.0, "duct_base_m": 110, "duct_thickness_m": 650,
|
||||||
|
"assessment": "Ducting", "fetched_at": time.time(),
|
||||||
|
}
|
||||||
|
duct._persist_status()
|
||||||
|
|
||||||
|
text = reporter.build_all()
|
||||||
|
assert "UPCOMING SATELLITE PASSES" in text
|
||||||
|
assert "AVALANCHE ADVISORIES" in text
|
||||||
|
assert "RF PROPAGATION" in text
|
||||||
|
|
@ -14,8 +14,8 @@ from meshai.persistence.observer_locations import (
|
||||||
|
|
||||||
# -- schema / migration -------------------------------------------------------
|
# -- schema / migration -------------------------------------------------------
|
||||||
|
|
||||||
def test_schema_version_is_23():
|
def test_schema_version_is_25():
|
||||||
assert SCHEMA_VERSION == 23
|
assert SCHEMA_VERSION == 25
|
||||||
|
|
||||||
|
|
||||||
def test_observer_locations_table_exists():
|
def test_observer_locations_table_exists():
|
||||||
|
|
@ -26,11 +26,11 @@ def test_observer_locations_table_exists():
|
||||||
assert "observer_locations" in tables
|
assert "observer_locations" in tables
|
||||||
|
|
||||||
|
|
||||||
def test_schema_meta_at_23():
|
def test_schema_meta_at_current():
|
||||||
conn = get_db()
|
conn = get_db()
|
||||||
row = conn.execute(
|
row = conn.execute(
|
||||||
"SELECT value FROM schema_meta WHERE key='version'").fetchone()
|
"SELECT value FROM schema_meta WHERE key='version'").fetchone()
|
||||||
assert int(row["value"]) == 23
|
assert int(row["value"]) == 25
|
||||||
|
|
||||||
|
|
||||||
# -- accessors ----------------------------------------------------------------
|
# -- accessors ----------------------------------------------------------------
|
||||||
|
|
|
||||||
|
|
@ -99,8 +99,9 @@ def _ingest_envelope(norad_id=25544, observer="Boise", max_el=72.5,
|
||||||
# ── schema / migration ───────────────────────────────────────────────
|
# ── schema / migration ───────────────────────────────────────────────
|
||||||
|
|
||||||
def test_schema_version_is_current():
|
def test_schema_version_is_current():
|
||||||
# Bumped to 23 by the native-satpass observer_locations migration (v23).
|
# Bumped to 25 by the LLM-persistence migrations (v24 avalanche_events,
|
||||||
assert SCHEMA_VERSION == 23
|
# v25 ducting_events); was 23 at the native-satpass observer_locations (v23).
|
||||||
|
assert SCHEMA_VERSION == 25
|
||||||
|
|
||||||
|
|
||||||
def test_v22_migration_applies_and_adds_due_at_column(tmp_path, monkeypatch):
|
def test_v22_migration_applies_and_adds_due_at_column(tmp_path, monkeypatch):
|
||||||
|
|
@ -113,7 +114,7 @@ def test_v22_migration_applies_and_adds_due_at_column(tmp_path, monkeypatch):
|
||||||
close_thread_connection()
|
close_thread_connection()
|
||||||
conn = init_db()
|
conn = init_db()
|
||||||
row = conn.execute("SELECT value FROM schema_meta WHERE key='version'").fetchone()
|
row = conn.execute("SELECT value FROM schema_meta WHERE key='version'").fetchone()
|
||||||
assert int(row["value"]) == 23
|
assert int(row["value"]) == 25
|
||||||
cols = {r["name"] for r in conn.execute("PRAGMA table_info(satpass_pending)")}
|
cols = {r["name"] for r in conn.execute("PRAGMA table_info(satpass_pending)")}
|
||||||
assert "due_at" in cols
|
assert "due_at" in cols
|
||||||
close_thread_connection()
|
close_thread_connection()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue