mirror of
https://github.com/zvx-echo6/meshai.git
synced 2026-08-26 17:31:34 +00:00
feat(generic): config-driven REST/GeoJSON source adapter (ported from Central) (#78)
Universal, no-code data sources: one GenericHttpAdapter polls any public REST/GeoJSON feed per config.generic_sources[] — dotted-path field mapping (items/id/lat/lon/geometry/title/fields) → coverage-gated, persisted (generic_events, v26), cold-start-silent, LLM-queryable events. Ports Central's GenericHttpAdapter to meshai native. First real use case: Idaho Power outages, configured (not hardcoded) — anyone can point it at their own utility/feed. GUI editor is a follow-up. 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
38f2f828ca
commit
30212ceb18
11 changed files with 897 additions and 7 deletions
|
|
@ -847,6 +847,13 @@ class Config:
|
||||||
notifications: NotificationsConfig = field(default_factory=NotificationsConfig)
|
notifications: NotificationsConfig = field(default_factory=NotificationsConfig)
|
||||||
danger_zones: DangerZonesConfig = field(default_factory=DangerZonesConfig)
|
danger_zones: DangerZonesConfig = field(default_factory=DangerZonesConfig)
|
||||||
coverage: Coverage = field(default_factory=Coverage)
|
coverage: Coverage = field(default_factory=Coverage)
|
||||||
|
# Config-driven REST/GeoJSON sources for the universal GenericHttpAdapter.
|
||||||
|
# Each entry is a PLAIN DICT (not a nested dataclass — the adapter parses
|
||||||
|
# and validates it) so the loader passes it through verbatim in both
|
||||||
|
# directions. Keys: name, enabled, url, items_path, id_path, lat_path,
|
||||||
|
# lon_path, geometry_path, title_path, time_path, category, poll_seconds,
|
||||||
|
# severity, field_mappings:[{source_path,dest_key}], summary_template, emoji.
|
||||||
|
generic_sources: list = field(default_factory=list)
|
||||||
|
|
||||||
_config_path: Optional[Path] = field(default=None, repr=False)
|
_config_path: Optional[Path] = field(default=None, repr=False)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -57,6 +57,8 @@ SECTION_TO_FILE: dict[str, str] = {
|
||||||
"dashboard": "dashboard.yaml",
|
"dashboard": "dashboard.yaml",
|
||||||
"danger_zones": "danger_zones.yaml",
|
"danger_zones": "danger_zones.yaml",
|
||||||
"coverage": "config.yaml",
|
"coverage": "config.yaml",
|
||||||
|
# Top-level list section (like mesh_sources): the whole file IS the list.
|
||||||
|
"generic_sources": "generic_sources.yaml",
|
||||||
}
|
}
|
||||||
|
|
||||||
# Fields that should be written to local.yaml instead of domain files
|
# Fields that should be written to local.yaml instead of domain files
|
||||||
|
|
|
||||||
428
work/meshai/env/generic_http.py
vendored
Normal file
428
work/meshai/env/generic_http.py
vendored
Normal file
|
|
@ -0,0 +1,428 @@
|
||||||
|
"""Generic, config-driven REST/GeoJSON source adapter (meshai native).
|
||||||
|
|
||||||
|
ONE class, MANY sources. Every source is a plain dict in
|
||||||
|
``config.generic_sources[]`` describing how to poll a public REST or GeoJSON
|
||||||
|
endpoint and map its fields onto meshai Events — no Python required to add a
|
||||||
|
new feed. This is the UNIVERSAL way to wire any public data feed; the first
|
||||||
|
real use case (Idaho Power outages) is CONFIGURED, not hardcoded.
|
||||||
|
|
||||||
|
Ported from Central's ``GenericHttpAdapter`` (src/central/adapters/generic_http.py):
|
||||||
|
* ``_dig(obj, path)`` dotted-path walker — ported verbatim.
|
||||||
|
* settings schema (url / items_path / id_path / lat_path / lon_path /
|
||||||
|
geometry_path / title_path / time_path / severity / category /
|
||||||
|
field_mappings) — ported concept.
|
||||||
|
* ``_item_to_event`` mapping — id (required, dedup), geometry-first then
|
||||||
|
lat/lon fallback, title + field_mappings → data.
|
||||||
|
|
||||||
|
Central-coupled bits REPLACED for native meshai:
|
||||||
|
* aiohttp → stdlib ``urllib.request`` (mirrors env/usgs_quake.py).
|
||||||
|
* pydantic settings model → plain dicts (the adapter validates them).
|
||||||
|
* NATS ``domain`` / ``subject_for`` → a meshai ``category`` string.
|
||||||
|
* Central Event / Geo → meshai ``make_event`` (lat/lon carried on the
|
||||||
|
Event so the deployed Shapely CoverageFilter gates it automatically).
|
||||||
|
"""
|
||||||
|
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
|
import re
|
||||||
|
import time
|
||||||
|
from typing import TYPE_CHECKING, Optional
|
||||||
|
from urllib.error import HTTPError, URLError
|
||||||
|
from urllib.request import Request, urlopen
|
||||||
|
|
||||||
|
from meshai.notifications.events import Event, make_event
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
pass
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# Only substitute {word} tokens in summary_template so stray braces in data
|
||||||
|
# values can never blow up the formatter.
|
||||||
|
_TEMPLATE_TOKEN = re.compile(r"\{(\w+)\}")
|
||||||
|
|
||||||
|
|
||||||
|
def _dig(obj, path):
|
||||||
|
"""Walk a nested dict/list by a dotted path string.
|
||||||
|
|
||||||
|
Each segment is tried as a dict key first; if the current node is a
|
||||||
|
list/tuple the segment is coerced to an integer index. Returns None on
|
||||||
|
any miss (missing key, out-of-range index, wrong node type).
|
||||||
|
|
||||||
|
Ported verbatim from Central's generic_http._dig. Examples::
|
||||||
|
|
||||||
|
_dig({"a": {"b": 1}}, "a.b") # 1
|
||||||
|
_dig({"a": [10, 20]}, "a.1") # 20
|
||||||
|
_dig({"a": 1}, "a.b") # None
|
||||||
|
_dig(None, "anything") # None
|
||||||
|
"""
|
||||||
|
if not path:
|
||||||
|
return None
|
||||||
|
parts = path.split(".")
|
||||||
|
cur = obj
|
||||||
|
for part in parts:
|
||||||
|
if cur is None:
|
||||||
|
return None
|
||||||
|
if isinstance(cur, dict):
|
||||||
|
cur = cur.get(part)
|
||||||
|
elif isinstance(cur, (list, tuple)):
|
||||||
|
try:
|
||||||
|
cur = cur[int(part)]
|
||||||
|
except (ValueError, IndexError):
|
||||||
|
return None
|
||||||
|
else:
|
||||||
|
return None
|
||||||
|
return cur
|
||||||
|
|
||||||
|
|
||||||
|
class GenericHttpAdapter:
|
||||||
|
"""Config-driven REST/GeoJSON source adapter — one instance, many sources.
|
||||||
|
|
||||||
|
meshai registers a SINGLE ``generic_http`` adapter; it internally polls
|
||||||
|
every enabled source in ``config.generic_sources`` on that source's own
|
||||||
|
cadence. Each source dict carries (mirrors Central's settings)::
|
||||||
|
|
||||||
|
{
|
||||||
|
name, enabled, url, items_path, id_path,
|
||||||
|
lat_path, lon_path, geometry_path, title_path, time_path,
|
||||||
|
category, poll_seconds, severity,
|
||||||
|
field_mappings: [{source_path, dest_key}, ...],
|
||||||
|
summary_template, emoji,
|
||||||
|
}
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, sources: list, coverage: dict = None):
|
||||||
|
# Keep only enabled sources that carry the minimum needed to poll +
|
||||||
|
# dedup: a url and an id_path. Everything else has sane fallbacks.
|
||||||
|
self._sources = []
|
||||||
|
for s in (sources or []):
|
||||||
|
if not isinstance(s, dict):
|
||||||
|
continue
|
||||||
|
if not s.get("enabled", True):
|
||||||
|
continue
|
||||||
|
if not s.get("url") or not s.get("id_path"):
|
||||||
|
logger.warning(
|
||||||
|
"generic_http: skipping source %r (needs url + id_path)",
|
||||||
|
s.get("name", "?"))
|
||||||
|
continue
|
||||||
|
self._sources.append(s)
|
||||||
|
|
||||||
|
# Pipeline CoverageFilter (Shapely) gates lat/lon events; the adapter
|
||||||
|
# does not hard-filter. Retained for parity / future use.
|
||||||
|
self._coverage = coverage
|
||||||
|
|
||||||
|
# Per-source last-poll epoch (keyed by source name) for cadence.
|
||||||
|
self._last_poll: dict = {}
|
||||||
|
# (source_name, event_id) -> internal event dict (current snapshot).
|
||||||
|
self._events: dict = {}
|
||||||
|
self._consecutive_errors = 0
|
||||||
|
self._last_error = None
|
||||||
|
self._is_loaded = False
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# Poll cadence
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
def tick(self) -> bool:
|
||||||
|
"""Poll each source whose cadence has elapsed. Returns True if the set
|
||||||
|
of (source, event_id) keys changed this tick.
|
||||||
|
|
||||||
|
Errors are guarded PER SOURCE so one broken feed never kills the
|
||||||
|
others — its exception is logged and the remaining sources still poll.
|
||||||
|
"""
|
||||||
|
now = time.time()
|
||||||
|
changed = False
|
||||||
|
for source in self._sources:
|
||||||
|
name = source.get("name") or source.get("url")
|
||||||
|
cadence = source.get("poll_seconds") or 300
|
||||||
|
last = self._last_poll.get(name, 0.0)
|
||||||
|
if now - last < cadence:
|
||||||
|
continue
|
||||||
|
self._last_poll[name] = now
|
||||||
|
try:
|
||||||
|
if self._poll_source(source, now):
|
||||||
|
changed = True
|
||||||
|
except Exception:
|
||||||
|
logger.exception("generic_http: source %r poll failed", name)
|
||||||
|
self._last_error = f"{name}: poll error"
|
||||||
|
self._consecutive_errors += 1
|
||||||
|
if changed:
|
||||||
|
self._is_loaded = True
|
||||||
|
return changed
|
||||||
|
|
||||||
|
def _poll_source(self, source: dict, now: float) -> bool:
|
||||||
|
"""Fetch + map one source. Returns True if its id-set changed."""
|
||||||
|
name = source.get("name") or source.get("url")
|
||||||
|
raw = self._fetch(source["url"])
|
||||||
|
if raw is None:
|
||||||
|
return False
|
||||||
|
|
||||||
|
items = _dig(raw, source.get("items_path") or "features")
|
||||||
|
if not isinstance(items, list):
|
||||||
|
logger.warning(
|
||||||
|
"generic_http: source %r items_path %r did not resolve to a "
|
||||||
|
"list (got %s)", name, source.get("items_path"),
|
||||||
|
type(items).__name__)
|
||||||
|
items = []
|
||||||
|
|
||||||
|
mapped = {}
|
||||||
|
for item in items:
|
||||||
|
evt = self._item_to_event(item, source, now)
|
||||||
|
if evt is not None:
|
||||||
|
mapped[(evt["_source"], evt["event_id"])] = evt
|
||||||
|
|
||||||
|
# Replace only THIS source's slice of the snapshot; detect change on
|
||||||
|
# the id-set for this source.
|
||||||
|
old_ids = {k for k in self._events if k[0] == name}
|
||||||
|
new_ids = set(mapped.keys())
|
||||||
|
changed = old_ids != new_ids
|
||||||
|
for k in old_ids:
|
||||||
|
self._events.pop(k, None)
|
||||||
|
self._events.update(mapped)
|
||||||
|
|
||||||
|
self._consecutive_errors = 0
|
||||||
|
self._last_error = None
|
||||||
|
if changed:
|
||||||
|
logger.info("generic_http: source %r updated (%d item(s))",
|
||||||
|
name, len(mapped))
|
||||||
|
return changed
|
||||||
|
|
||||||
|
def _fetch(self, url: str):
|
||||||
|
"""GET ``url`` and return parsed JSON, or None on any error.
|
||||||
|
|
||||||
|
stdlib urllib (mirrors env/usgs_quake.py); 30s timeout; UA header.
|
||||||
|
"""
|
||||||
|
headers = {"User-Agent": "MeshAI/1.0", "Accept": "application/json"}
|
||||||
|
try:
|
||||||
|
req = Request(url, headers=headers)
|
||||||
|
with urlopen(req, timeout=30) as resp:
|
||||||
|
return json.loads(resp.read().decode("utf-8"))
|
||||||
|
except HTTPError as e:
|
||||||
|
logger.warning("generic_http HTTP error %s for %s", e.code, url)
|
||||||
|
self._last_error = f"HTTP {e.code}"
|
||||||
|
self._consecutive_errors += 1
|
||||||
|
return None
|
||||||
|
except URLError as e:
|
||||||
|
logger.warning("generic_http connection error for %s: %s",
|
||||||
|
url, e.reason)
|
||||||
|
self._last_error = str(e.reason)
|
||||||
|
self._consecutive_errors += 1
|
||||||
|
return None
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("generic_http fetch error for %s: %s", url, e)
|
||||||
|
self._last_error = str(e)
|
||||||
|
self._consecutive_errors += 1
|
||||||
|
return None
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# Item -> internal event dict
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
def _item_to_event(self, item, source: dict, now: float = None):
|
||||||
|
"""Map one source item to an internal event dict, or None to skip.
|
||||||
|
|
||||||
|
Mirrors Central's ``_item_to_event``: id is required (dedup);
|
||||||
|
geometry-first (Point → centroid) then lat/lon fallback; title +
|
||||||
|
field_mappings written into ``data``. The resolved geometry dict (when
|
||||||
|
present) is carried in ``data['geometry']`` so the coverage gate can
|
||||||
|
use it even for non-Point shapes.
|
||||||
|
"""
|
||||||
|
now = now if now is not None else time.time()
|
||||||
|
name = source.get("name") or source.get("url")
|
||||||
|
|
||||||
|
# --- id (required, dedup) ---
|
||||||
|
raw_id = _dig(item, source["id_path"])
|
||||||
|
if raw_id is None:
|
||||||
|
return None
|
||||||
|
event_id = str(raw_id)
|
||||||
|
|
||||||
|
category = source.get("category") or "generic"
|
||||||
|
severity = source.get("severity") or "routine"
|
||||||
|
|
||||||
|
# --- geo: geometry Point centroid first, else lat/lon paths ---
|
||||||
|
lat = lon = None
|
||||||
|
geometry = None
|
||||||
|
geom_path = source.get("geometry_path")
|
||||||
|
if geom_path:
|
||||||
|
geom = _dig(item, geom_path)
|
||||||
|
if isinstance(geom, dict):
|
||||||
|
geometry = geom
|
||||||
|
if geom.get("type") == "Point":
|
||||||
|
coords = geom.get("coordinates") or []
|
||||||
|
if len(coords) >= 2:
|
||||||
|
try:
|
||||||
|
lon = float(coords[0])
|
||||||
|
lat = float(coords[1])
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
lat = lon = None
|
||||||
|
if lat is None or lon is None:
|
||||||
|
lat_path = source.get("lat_path")
|
||||||
|
lon_path = source.get("lon_path")
|
||||||
|
if lat_path and lon_path:
|
||||||
|
raw_lat = _dig(item, lat_path)
|
||||||
|
raw_lon = _dig(item, lon_path)
|
||||||
|
if raw_lat is not None and raw_lon is not None:
|
||||||
|
try:
|
||||||
|
lat = float(raw_lat)
|
||||||
|
lon = float(raw_lon)
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
lat = lon = None
|
||||||
|
|
||||||
|
# --- title ---
|
||||||
|
title = None
|
||||||
|
if source.get("title_path"):
|
||||||
|
t = _dig(item, source["title_path"])
|
||||||
|
if t is not None:
|
||||||
|
title = str(t)
|
||||||
|
|
||||||
|
# --- data: title + coords + field_mappings + geometry + render hints ---
|
||||||
|
data = {}
|
||||||
|
if title is not None:
|
||||||
|
data["title"] = title
|
||||||
|
if lat is not None and lon is not None:
|
||||||
|
data["latitude"] = lat
|
||||||
|
data["longitude"] = lon
|
||||||
|
for fm in (source.get("field_mappings") or []):
|
||||||
|
if not isinstance(fm, dict):
|
||||||
|
continue
|
||||||
|
dest = fm.get("dest_key")
|
||||||
|
src_path = fm.get("source_path")
|
||||||
|
if not dest or not src_path:
|
||||||
|
continue
|
||||||
|
data[dest] = _dig(item, src_path)
|
||||||
|
if geometry is not None:
|
||||||
|
data["geometry"] = geometry
|
||||||
|
data["_source"] = name
|
||||||
|
data["_summary_template"] = source.get("summary_template")
|
||||||
|
data["_emoji"] = source.get("emoji")
|
||||||
|
|
||||||
|
# --- optional event time ---
|
||||||
|
ts = now
|
||||||
|
if source.get("time_path"):
|
||||||
|
raw_time = _dig(item, source["time_path"])
|
||||||
|
if raw_time is not None:
|
||||||
|
try:
|
||||||
|
# epoch ms (int) or ISO-8601 string both accepted.
|
||||||
|
if isinstance(raw_time, (int, float)):
|
||||||
|
ts = float(raw_time)
|
||||||
|
if ts > 1e12: # milliseconds
|
||||||
|
ts /= 1000.0
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
ts = now
|
||||||
|
|
||||||
|
return {
|
||||||
|
# store keys on evt["source"] + evt["event_id"]; namespace the
|
||||||
|
# source per configured name so each feed dedups independently.
|
||||||
|
"source": f"generic:{name}",
|
||||||
|
"_source": name, # bare configured name (table + cold-start)
|
||||||
|
"event_id": event_id,
|
||||||
|
"category": category,
|
||||||
|
"severity": severity,
|
||||||
|
"title": title,
|
||||||
|
"latitude": lat,
|
||||||
|
"longitude": lon,
|
||||||
|
"data": data,
|
||||||
|
"fetched_at": now,
|
||||||
|
"quake_time": ts, # generic event time
|
||||||
|
}
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# Internal event dict -> pipeline Event
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
def to_event(self, evt: dict) -> Optional["Event"]:
|
||||||
|
"""Build a meshai Event from an internal event dict.
|
||||||
|
|
||||||
|
lat/lon (and ``data['geometry']`` when resolved) ride on the Event so
|
||||||
|
the deployed Shapely CoverageFilter gates it with no extra wiring.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
event_id = evt.get("event_id")
|
||||||
|
if not event_id:
|
||||||
|
return None
|
||||||
|
summary = self._render(evt)
|
||||||
|
return make_event(
|
||||||
|
source="generic:" + (evt.get("_source") or "?"),
|
||||||
|
category=evt.get("category") or "generic",
|
||||||
|
severity=evt.get("severity") or "routine",
|
||||||
|
title=evt.get("title") or summary,
|
||||||
|
summary=summary,
|
||||||
|
timestamp=evt.get("fetched_at"),
|
||||||
|
expires=evt.get("expires"),
|
||||||
|
lat=evt.get("latitude"),
|
||||||
|
lon=evt.get("longitude"),
|
||||||
|
group_key=event_id,
|
||||||
|
inhibit_keys=[event_id],
|
||||||
|
data=evt.get("data") or {},
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.exception("generic_http to_event failed for %s",
|
||||||
|
evt.get("event_id"))
|
||||||
|
return None
|
||||||
|
|
||||||
|
def _render(self, evt: dict) -> str:
|
||||||
|
"""Render the mesh wire string for one event.
|
||||||
|
|
||||||
|
With ``summary_template`` set: substitute ``{key}`` tokens from the
|
||||||
|
event's data dict (title included). Missing keys render BLANK — the
|
||||||
|
formatter never crashes on an absent field.
|
||||||
|
|
||||||
|
Without a template: ``"{emoji} {title}"`` plus a compact join of the
|
||||||
|
mapped fields (``dest_key: value``). Kept short for the mesh budget.
|
||||||
|
"""
|
||||||
|
data = evt.get("data") or {}
|
||||||
|
template = data.get("_summary_template")
|
||||||
|
emoji = data.get("_emoji") or ""
|
||||||
|
title = evt.get("title") or data.get("title") or ""
|
||||||
|
|
||||||
|
if template:
|
||||||
|
mapping = {k: v for k, v in data.items() if not k.startswith("_")}
|
||||||
|
if title and "title" not in mapping:
|
||||||
|
mapping["title"] = title
|
||||||
|
|
||||||
|
def _sub(m):
|
||||||
|
v = mapping.get(m.group(1))
|
||||||
|
return "" if v is None else str(v)
|
||||||
|
|
||||||
|
return _TEMPLATE_TOKEN.sub(_sub, template).strip()
|
||||||
|
|
||||||
|
# Default render: emoji + title + compact mapped fields.
|
||||||
|
parts = []
|
||||||
|
head = f"{emoji} {title}".strip()
|
||||||
|
if head:
|
||||||
|
parts.append(head)
|
||||||
|
extras = []
|
||||||
|
for k, v in data.items():
|
||||||
|
if k.startswith("_") or k in ("title", "latitude", "longitude",
|
||||||
|
"geometry"):
|
||||||
|
continue
|
||||||
|
if v is None:
|
||||||
|
continue
|
||||||
|
extras.append(f"{k}: {v}")
|
||||||
|
if extras:
|
||||||
|
parts.append(", ".join(extras))
|
||||||
|
return " — ".join(parts) if parts else (evt.get("category") or "event")
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# State / status
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
def get_active(self) -> list:
|
||||||
|
"""Current mapped events across all sources (store ingests these)."""
|
||||||
|
return list(self._events.values())
|
||||||
|
|
||||||
|
# Alias so the store's generic get_active(source=...) style also works.
|
||||||
|
def get_events(self) -> list:
|
||||||
|
return self.get_active()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def health_status(self) -> dict:
|
||||||
|
return {
|
||||||
|
"source": "generic_http",
|
||||||
|
"is_loaded": self._is_loaded,
|
||||||
|
"last_error": str(self._last_error) if self._last_error else None,
|
||||||
|
"consecutive_errors": self._consecutive_errors,
|
||||||
|
"event_count": len(self._events),
|
||||||
|
"source_count": len(self._sources),
|
||||||
|
"last_fetch": max(self._last_poll.values()) if self._last_poll else 0.0,
|
||||||
|
}
|
||||||
112
work/meshai/env/store.py
vendored
112
work/meshai/env/store.py
vendored
|
|
@ -40,7 +40,11 @@ class EnvironmentalStore:
|
||||||
event_bus: Optional["EventBus"] = None,
|
event_bus: Optional["EventBus"] = None,
|
||||||
coverage_bbox: list = None,
|
coverage_bbox: list = None,
|
||||||
coverage_excluded: list = None,
|
coverage_excluded: list = None,
|
||||||
|
generic_sources: list = None,
|
||||||
):
|
):
|
||||||
|
# Config-driven REST/GeoJSON sources (top-level config.generic_sources)
|
||||||
|
# for the universal GenericHttpAdapter. Plain list of dicts.
|
||||||
|
self._generic_sources = generic_sources or []
|
||||||
self._adapters = {} # name -> adapter instance
|
self._adapters = {} # name -> adapter instance
|
||||||
self._failed_adapters = {} # name -> last_error string
|
self._failed_adapters = {} # name -> last_error string
|
||||||
self._events = {} # (source, event_id) -> event dict
|
self._events = {} # (source, event_id) -> event dict
|
||||||
|
|
@ -90,6 +94,13 @@ class EnvironmentalStore:
|
||||||
# LATER poll is a genuine ignition and broadcasts "New".
|
# LATER poll is a genuine ignition and broadcasts "New".
|
||||||
self._fires_seeded: bool = False
|
self._fires_seeded: bool = False
|
||||||
|
|
||||||
|
# Generic-source cold-start silent-seed gate (PER source name), mirror
|
||||||
|
# of _fires_seeded. The FIRST non-empty poll for a given generic source
|
||||||
|
# seeds every current item (records it seen + persists it) and emits
|
||||||
|
# NOTHING — that batch is pre-existing backlog. Items that first appear
|
||||||
|
# on a LATER poll broadcast. Keyed by the bare configured source name.
|
||||||
|
self._generic_seeded: set[str] = set()
|
||||||
|
|
||||||
# Create adapter instances with error isolation
|
# Create adapter instances with error isolation
|
||||||
self._register_adapter("nws", config.nws, ".nws", "NWSAlertsAdapter",
|
self._register_adapter("nws", config.nws, ".nws", "NWSAlertsAdapter",
|
||||||
lambda cfg: (cfg, self._coverage_for("nws")))
|
lambda cfg: (cfg, self._coverage_for("nws")))
|
||||||
|
|
@ -133,6 +144,22 @@ class EnvironmentalStore:
|
||||||
logger.warning("Failed to initialize firms adapter: %s", err_msg)
|
logger.warning("Failed to initialize firms adapter: %s", err_msg)
|
||||||
self._failed_adapters["firms"] = err_msg
|
self._failed_adapters["firms"] = err_msg
|
||||||
|
|
||||||
|
# Universal config-driven REST/GeoJSON sources. ONE adapter instance
|
||||||
|
# handles every source in config.generic_sources; construct it only
|
||||||
|
# when at least one source is configured. Guarded like firms so a bad
|
||||||
|
# import/config never takes the whole store down.
|
||||||
|
if self._generic_sources:
|
||||||
|
try:
|
||||||
|
from .generic_http import GenericHttpAdapter
|
||||||
|
self._generic = GenericHttpAdapter(
|
||||||
|
self._generic_sources,
|
||||||
|
self._coverage_for("generic_http"))
|
||||||
|
self._adapters["generic_http"] = self._generic
|
||||||
|
except Exception as e:
|
||||||
|
err_msg = f"{type(e).__name__}: {e}"
|
||||||
|
logger.warning("Failed to initialize generic_http adapter: %s", err_msg)
|
||||||
|
self._failed_adapters["generic_http"] = err_msg
|
||||||
|
|
||||||
_central = [n for n in ("nws", "swpc", "ducting", "fires", "avalanche", "usgs", "usgs_quake", "traffic", "roads511", "wzdx", "firms", "satpass")
|
_central = [n for n in ("nws", "swpc", "ducting", "fires", "avalanche", "usgs", "usgs_quake", "traffic", "roads511", "wzdx", "firms", "satpass")
|
||||||
if getattr(getattr(config, n, None), "feed_source", "native") == "central"]
|
if getattr(getattr(config, n, None), "feed_source", "native") == "central"]
|
||||||
if _central:
|
if _central:
|
||||||
|
|
@ -298,6 +325,10 @@ class EnvironmentalStore:
|
||||||
key = (evt["source"], evt["event_id"])
|
key = (evt["source"], evt["event_id"])
|
||||||
self._events[key] = evt
|
self._events[key] = evt
|
||||||
self._ingest_fires(adapter)
|
self._ingest_fires(adapter)
|
||||||
|
elif name == "generic_http":
|
||||||
|
# Universal config-driven sources: custom ingest with COLD-START-
|
||||||
|
# SILENT per source name + persist-every-item (see _ingest_generic).
|
||||||
|
self._ingest_generic(adapter)
|
||||||
elif name == "avalanche":
|
elif name == "avalanche":
|
||||||
# Avalanche: re-emit on danger_level rise (Update:) not just new
|
# Avalanche: re-emit on danger_level rise (Update:) not just new
|
||||||
# events. The rise is a legitimate CONTENT change, so it passes
|
# events. The rise is a legitimate CONTENT change, so it passes
|
||||||
|
|
@ -446,6 +477,87 @@ class EnvironmentalStore:
|
||||||
if events:
|
if events:
|
||||||
self._fires_seeded = True
|
self._fires_seeded = True
|
||||||
|
|
||||||
|
def _ingest_generic(self, adapter) -> None:
|
||||||
|
"""Native generic-source ingest — COLD-START-SILENT per source.
|
||||||
|
|
||||||
|
Mirrors the ``_fires_seeded`` cold-start pattern, but PER configured
|
||||||
|
source name:
|
||||||
|
|
||||||
|
1. The FIRST non-empty poll for a given source seeds every current
|
||||||
|
item into the received-delta seen-set and PERSISTS it, but emits
|
||||||
|
NOTHING — that batch is pre-existing backlog.
|
||||||
|
2. On later polls, only items whose key is not already seen (they
|
||||||
|
newly appeared upstream = "just received") broadcast via
|
||||||
|
``_emit_event`` (which runs the coverage/decider path).
|
||||||
|
3. EVERY item — seeded or new — is upserted into ``generic_events``
|
||||||
|
so the LLM (env_reporter.build_generic_detail) sees it immediately.
|
||||||
|
Dedup is by (source, event_id).
|
||||||
|
|
||||||
|
The received-delta seen-set is namespaced by the event's ``source``
|
||||||
|
(``generic:<name>``) so distinct sources never cross-contaminate.
|
||||||
|
"""
|
||||||
|
conn = None
|
||||||
|
try:
|
||||||
|
from meshai.persistence import get_db
|
||||||
|
conn = get_db()
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("generic ingest skipped persistence (DB unavailable): %s", e)
|
||||||
|
|
||||||
|
# Group this poll's events by bare source name to apply the per-source
|
||||||
|
# cold-start gate consistently across all of that source's items.
|
||||||
|
by_source: dict = {}
|
||||||
|
for evt in adapter.get_active():
|
||||||
|
by_source.setdefault(evt.get("_source") or "?", []).append(evt)
|
||||||
|
|
||||||
|
for src_name, items in by_source.items():
|
||||||
|
cold_start = src_name not in self._generic_seeded
|
||||||
|
for evt in items:
|
||||||
|
key = (evt["source"], evt["event_id"])
|
||||||
|
self._events[key] = evt # always track current state
|
||||||
|
if conn is not None:
|
||||||
|
self._persist_generic(conn, evt)
|
||||||
|
|
||||||
|
seen = self._seen.setdefault(evt["source"], set())
|
||||||
|
seen_key = self._seen_key(evt)
|
||||||
|
if cold_start:
|
||||||
|
seen.add(seen_key) # silent seed — backlog, no emit
|
||||||
|
continue
|
||||||
|
if seen_key in seen:
|
||||||
|
continue # already received (prior poll)
|
||||||
|
seen.add(seen_key)
|
||||||
|
if self._event_bus is not None and hasattr(adapter, "to_event"):
|
||||||
|
self._emit_event(adapter, evt)
|
||||||
|
# First non-empty poll for this source complete: later polls emit.
|
||||||
|
if items:
|
||||||
|
self._generic_seeded.add(src_name)
|
||||||
|
|
||||||
|
def _persist_generic(self, conn, evt: dict) -> None:
|
||||||
|
"""Upsert one generic event into ``generic_events`` (source keyed by
|
||||||
|
the bare configured name). Mirrors how _ingest_fires writes the fires
|
||||||
|
table. Never fatal — a persistence error logs and continues."""
|
||||||
|
now = int(time.time())
|
||||||
|
try:
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO generic_events(source, event_id, category, title, "
|
||||||
|
"lat, lon, severity, data_json, first_seen, last_seen) "
|
||||||
|
"VALUES (?,?,?,?,?,?,?,?,?,?) "
|
||||||
|
"ON CONFLICT(source, event_id) DO UPDATE SET "
|
||||||
|
"category=excluded.category, title=excluded.title, "
|
||||||
|
"lat=excluded.lat, lon=excluded.lon, "
|
||||||
|
"severity=excluded.severity, data_json=excluded.data_json, "
|
||||||
|
"last_seen=excluded.last_seen",
|
||||||
|
(
|
||||||
|
evt.get("_source"), evt["event_id"], evt.get("category"),
|
||||||
|
evt.get("title"), evt.get("latitude"), evt.get("longitude"),
|
||||||
|
evt.get("severity"),
|
||||||
|
json.dumps(evt.get("data", {}), default=str),
|
||||||
|
now, now,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.exception("generic persist failed for %s",
|
||||||
|
evt.get("event_id", "?"))
|
||||||
|
|
||||||
def _seed_from_persistent(self) -> None:
|
def _seed_from_persistent(self) -> None:
|
||||||
"""Pre-seed ``self._seen`` from the durable hazard tables at startup.
|
"""Pre-seed ``self._seen`` from the durable hazard tables at startup.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -655,6 +655,7 @@ class MeshAI:
|
||||||
config=env_cfg, region_anchors=region_anchors,
|
config=env_cfg, region_anchors=region_anchors,
|
||||||
coverage_bbox=coverage_bbox, event_bus=self.event_bus,
|
coverage_bbox=coverage_bbox, event_bus=self.event_bus,
|
||||||
coverage_excluded=cov.excluded_adapters,
|
coverage_excluded=cov.excluded_adapters,
|
||||||
|
generic_sources=self.config.generic_sources,
|
||||||
)
|
)
|
||||||
logger.info(f"Environmental feeds enabled ({len(self.env_store._adapters)} adapters)")
|
logger.info(f"Environmental feeds enabled ({len(self.env_store._adapters)} adapters)")
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -14,6 +14,7 @@ Read-only library; no writes to any table.
|
||||||
"""
|
"""
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
import logging
|
import logging
|
||||||
import sqlite3
|
import sqlite3
|
||||||
import time
|
import time
|
||||||
|
|
@ -275,6 +276,64 @@ class EnvReporter:
|
||||||
lines.append(f" - {mag} {place}, {depth} depth, {when}{ts}")
|
lines.append(f" - {mag} {place}, {depth} depth, {when}{ts}")
|
||||||
return "\n".join(lines)[:_block_cap()]
|
return "\n".join(lines)[:_block_cap()]
|
||||||
|
|
||||||
|
def build_generic_detail(self, *, hours: int = 24,
|
||||||
|
limit: int = 20,
|
||||||
|
now: Optional[int] = None) -> str:
|
||||||
|
"""LLM block for config-driven generic sources (env/generic_http.py).
|
||||||
|
|
||||||
|
Groups active generic_events by source, showing recent items (title +
|
||||||
|
a couple of key mapped data fields) so a mesh user can ask the LLM
|
||||||
|
about ANY configured feed. Honors the ``generic_http`` adapter's
|
||||||
|
include_in_llm_context gate.
|
||||||
|
"""
|
||||||
|
if not self._adapter_included("generic_http"):
|
||||||
|
return ""
|
||||||
|
now = now if now is not None else int(time.time())
|
||||||
|
try: conn = self._conn_factory()
|
||||||
|
except Exception: return ""
|
||||||
|
|
||||||
|
try:
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT source, event_id, category, title, lat, lon, severity, "
|
||||||
|
"data_json, last_seen FROM generic_events "
|
||||||
|
"WHERE last_seen >= ? ORDER BY source, last_seen DESC",
|
||||||
|
(now - hours * 3600,),
|
||||||
|
).fetchall()
|
||||||
|
except Exception:
|
||||||
|
return ""
|
||||||
|
if not rows:
|
||||||
|
return ""
|
||||||
|
|
||||||
|
# Bucket by source, capping items per source so no single feed hogs
|
||||||
|
# the block budget.
|
||||||
|
by_source: dict = {}
|
||||||
|
for r in rows:
|
||||||
|
by_source.setdefault(r["source"], []).append(r)
|
||||||
|
|
||||||
|
lines = [f"CONFIGURED DATA SOURCES (last {hours}h):"]
|
||||||
|
for src, items in by_source.items():
|
||||||
|
cat = items[0]["category"] or "generic"
|
||||||
|
lines.append(f"{src} ({cat}): {len(items)} item(s)")
|
||||||
|
for r in items[:limit]:
|
||||||
|
title = r["title"] or r["event_id"]
|
||||||
|
extras = []
|
||||||
|
try:
|
||||||
|
data = json.loads(r["data_json"]) if r["data_json"] else {}
|
||||||
|
except Exception:
|
||||||
|
data = {}
|
||||||
|
for k, v in data.items():
|
||||||
|
if k.startswith("_") or k in ("title", "latitude",
|
||||||
|
"longitude", "geometry"):
|
||||||
|
continue
|
||||||
|
if v is None:
|
||||||
|
continue
|
||||||
|
extras.append(f"{k}: {v}")
|
||||||
|
if len(extras) >= 3:
|
||||||
|
break
|
||||||
|
tail = (" — " + ", ".join(extras)) if extras else ""
|
||||||
|
lines.append(f" - {title}{tail}")
|
||||||
|
return "\n".join(lines)[:_block_cap()]
|
||||||
|
|
||||||
def build_traffic_detail(self, *, state: Optional[str] = "ID",
|
def build_traffic_detail(self, *, state: Optional[str] = "ID",
|
||||||
hours: int = 2,
|
hours: int = 2,
|
||||||
limit: int = 10,
|
limit: int = 10,
|
||||||
|
|
@ -537,6 +596,7 @@ class EnvReporter:
|
||||||
self.build_satpass_detail(now=now),
|
self.build_satpass_detail(now=now),
|
||||||
self.build_avalanche_detail(now=now),
|
self.build_avalanche_detail(now=now),
|
||||||
self.build_ducting_detail(now=now),
|
self.build_ducting_detail(now=now),
|
||||||
|
self.build_generic_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 = 25
|
SCHEMA_VERSION = 26
|
||||||
SCHEMA_META_TABLE = "schema_meta"
|
SCHEMA_META_TABLE = "schema_meta"
|
||||||
MIGRATIONS_DIR = Path(__file__).parent / "migrations"
|
MIGRATIONS_DIR = Path(__file__).parent / "migrations"
|
||||||
|
|
||||||
|
|
|
||||||
32
work/meshai/persistence/migrations/v26.sql
Normal file
32
work/meshai/persistence/migrations/v26.sql
Normal file
|
|
@ -0,0 +1,32 @@
|
||||||
|
-- v26 generic config-driven source events (LLM-queryable durable store).
|
||||||
|
--
|
||||||
|
-- The universal GenericHttpAdapter (env/generic_http.py) polls any public
|
||||||
|
-- REST/GeoJSON feed described in config.generic_sources[] and maps its items
|
||||||
|
-- onto meshai Events. This table is that adapter's durable store so the mesh
|
||||||
|
-- LLM (env_reporter.build_generic_detail) can answer questions about any
|
||||||
|
-- configured source, and so cold-start-silent seeding survives restarts.
|
||||||
|
--
|
||||||
|
-- One row per (source, event_id): source is the bare configured source name
|
||||||
|
-- (e.g. "idaho_power"); event_id is the item's stable id (id_path). The store
|
||||||
|
-- upserts every polled item (seeded or newly broadcast) here. data_json holds
|
||||||
|
-- the full mapped data dict (title + field_mappings + geometry, etc.).
|
||||||
|
--
|
||||||
|
-- Persistence-only: does NOT affect broadcast/gating (the received-delta gate
|
||||||
|
-- + coverage/decider path stay in the store/pipeline). IF NOT EXISTS so a
|
||||||
|
-- fresh install and a re-run are both no-ops.
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS generic_events (
|
||||||
|
source TEXT NOT NULL, -- bare configured source name
|
||||||
|
event_id TEXT NOT NULL, -- stable per-item id (id_path)
|
||||||
|
category TEXT, -- configured category (e.g. power_outage)
|
||||||
|
title TEXT, -- resolved title_path value
|
||||||
|
lat REAL,
|
||||||
|
lon REAL,
|
||||||
|
severity TEXT, -- routine | priority | immediate
|
||||||
|
data_json TEXT, -- full mapped data dict as JSON
|
||||||
|
first_seen INTEGER, -- epoch seconds first ingested
|
||||||
|
last_seen INTEGER, -- epoch seconds most recently seen
|
||||||
|
PRIMARY KEY (source, event_id)
|
||||||
|
);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_generic_last_seen ON generic_events(last_seen);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_generic_source ON generic_events(source);
|
||||||
248
work/tests/test_generic_http.py
Normal file
248
work/tests/test_generic_http.py
Normal file
|
|
@ -0,0 +1,248 @@
|
||||||
|
"""Tests for the universal config-driven GenericHttpAdapter (env/generic_http.py).
|
||||||
|
|
||||||
|
Ported-behavior coverage:
|
||||||
|
* _dig dotted-path walker (nested dict, list index, miss cases)
|
||||||
|
* _item_to_event field mapping (Idaho Power outage item) + _render template
|
||||||
|
* summary_template with a missing key does not crash
|
||||||
|
* cold-start-silent: first poll seeds + persists but broadcasts nothing
|
||||||
|
* geometry-path Point -> centroid extraction
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
|
||||||
|
from meshai.env.generic_http import GenericHttpAdapter, _dig
|
||||||
|
from meshai.env.store import EnvironmentalStore
|
||||||
|
from meshai.config import EnvironmentalConfig
|
||||||
|
from meshai.notifications.pipeline.bus import EventBus
|
||||||
|
from meshai.persistence import get_db
|
||||||
|
|
||||||
|
|
||||||
|
# --- Idaho Power source config (CONFIGURED, not hardcoded) -------------------
|
||||||
|
|
||||||
|
IDAHO_POWER_SOURCE = {
|
||||||
|
"name": "idaho_power",
|
||||||
|
"enabled": True,
|
||||||
|
"url": "https://apiedge.idahopower.com/api/Outage/GetCurrentOutageInformation",
|
||||||
|
"items_path": "object.outages",
|
||||||
|
"id_path": "omsOutageId",
|
||||||
|
"lat_path": "latitude",
|
||||||
|
"lon_path": "longitude",
|
||||||
|
"title_path": "probableCause",
|
||||||
|
"category": "power_outage",
|
||||||
|
"poll_seconds": 300,
|
||||||
|
"severity": "routine",
|
||||||
|
"field_mappings": [
|
||||||
|
{"source_path": "omsCustomerCount", "dest_key": "customers"},
|
||||||
|
{"source_path": "omsEstimatedTimeToRestorationMessage", "dest_key": "eta"},
|
||||||
|
{"source_path": "omsStatusDescription", "dest_key": "status"},
|
||||||
|
],
|
||||||
|
"summary_template": "⚡ Power out — {customers} affected, ETA {eta} ({status})",
|
||||||
|
"emoji": "⚡",
|
||||||
|
}
|
||||||
|
|
||||||
|
IDAHO_POWER_ITEM = {
|
||||||
|
"omsOutageId": "123",
|
||||||
|
"latitude": 43.6,
|
||||||
|
"longitude": -116.2,
|
||||||
|
"omsCustomerCount": 250,
|
||||||
|
"probableCause": "Equipment",
|
||||||
|
"omsEstimatedTimeToRestorationMessage": "10:30 PM",
|
||||||
|
"omsStatusDescription": "Crew assigned",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
# ===========================================================================
|
||||||
|
# _dig
|
||||||
|
# ===========================================================================
|
||||||
|
|
||||||
|
def test_dig_nested_dict():
|
||||||
|
assert _dig({"a": {"b": 1}}, "a.b") == 1
|
||||||
|
assert _dig({"object": {"outages": [1, 2]}}, "object.outages") == [1, 2]
|
||||||
|
|
||||||
|
|
||||||
|
def test_dig_list_index():
|
||||||
|
assert _dig({"a": [10, 20]}, "a.1") == 20
|
||||||
|
assert _dig({"geometry": {"coordinates": [-116.2, 43.6]}},
|
||||||
|
"geometry.coordinates.0") == -116.2
|
||||||
|
|
||||||
|
|
||||||
|
def test_dig_miss_cases():
|
||||||
|
assert _dig({"a": 1}, "a.b") is None # scalar then key
|
||||||
|
assert _dig(None, "anything") is None # None root
|
||||||
|
assert _dig({"a": [1]}, "a.5") is None # out-of-range index
|
||||||
|
assert _dig({"a": {}}, "a.missing") is None # missing key
|
||||||
|
assert _dig({"a": [1]}, "a.x") is None # non-int index on list
|
||||||
|
|
||||||
|
|
||||||
|
# ===========================================================================
|
||||||
|
# _item_to_event + _render
|
||||||
|
# ===========================================================================
|
||||||
|
|
||||||
|
def test_item_to_event_idaho_power_mapping():
|
||||||
|
adapter = GenericHttpAdapter([IDAHO_POWER_SOURCE])
|
||||||
|
evt = adapter._item_to_event(IDAHO_POWER_ITEM, IDAHO_POWER_SOURCE, now=1000.0)
|
||||||
|
|
||||||
|
assert evt is not None
|
||||||
|
assert evt["event_id"] == "123"
|
||||||
|
assert evt["category"] == "power_outage"
|
||||||
|
assert evt["latitude"] == 43.6
|
||||||
|
assert evt["longitude"] == -116.2
|
||||||
|
assert evt["title"] == "Equipment"
|
||||||
|
assert evt["data"]["customers"] == 250
|
||||||
|
assert evt["data"]["eta"] == "10:30 PM"
|
||||||
|
assert evt["data"]["status"] == "Crew assigned"
|
||||||
|
assert evt["data"]["latitude"] == 43.6
|
||||||
|
assert evt["data"]["longitude"] == -116.2
|
||||||
|
# source namespaced per configured name for independent dedup
|
||||||
|
assert evt["source"] == "generic:idaho_power"
|
||||||
|
assert evt["_source"] == "idaho_power"
|
||||||
|
|
||||||
|
|
||||||
|
def test_render_summary_template():
|
||||||
|
adapter = GenericHttpAdapter([IDAHO_POWER_SOURCE])
|
||||||
|
evt = adapter._item_to_event(IDAHO_POWER_ITEM, IDAHO_POWER_SOURCE)
|
||||||
|
assert adapter._render(evt) == (
|
||||||
|
"⚡ Power out — 250 affected, ETA 10:30 PM (Crew assigned)"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_render_template_missing_key_does_not_crash():
|
||||||
|
source = dict(IDAHO_POWER_SOURCE)
|
||||||
|
# Reference a key that no field mapping produces -> must render blank.
|
||||||
|
source["summary_template"] = "⚡ {customers} out, cause {nonexistent_key}!"
|
||||||
|
adapter = GenericHttpAdapter([source])
|
||||||
|
evt = adapter._item_to_event(IDAHO_POWER_ITEM, source)
|
||||||
|
rendered = adapter._render(evt) # must not raise
|
||||||
|
assert "250" in rendered
|
||||||
|
assert "{nonexistent_key}" not in rendered # token substituted (blank)
|
||||||
|
assert rendered == "⚡ 250 out, cause !"
|
||||||
|
|
||||||
|
|
||||||
|
def test_render_default_no_template():
|
||||||
|
source = dict(IDAHO_POWER_SOURCE)
|
||||||
|
source.pop("summary_template")
|
||||||
|
adapter = GenericHttpAdapter([source])
|
||||||
|
evt = adapter._item_to_event(IDAHO_POWER_ITEM, source)
|
||||||
|
rendered = adapter._render(evt)
|
||||||
|
assert "Equipment" in rendered # title
|
||||||
|
assert "customers: 250" in rendered # mapped field
|
||||||
|
|
||||||
|
|
||||||
|
def test_item_missing_id_is_skipped():
|
||||||
|
adapter = GenericHttpAdapter([IDAHO_POWER_SOURCE])
|
||||||
|
evt = adapter._item_to_event({"latitude": 43.6, "longitude": -116.2},
|
||||||
|
IDAHO_POWER_SOURCE)
|
||||||
|
assert evt is None
|
||||||
|
|
||||||
|
|
||||||
|
# ===========================================================================
|
||||||
|
# geometry-path Point -> centroid
|
||||||
|
# ===========================================================================
|
||||||
|
|
||||||
|
def test_geometry_point_centroid():
|
||||||
|
source = {
|
||||||
|
"name": "geo_feed",
|
||||||
|
"enabled": True,
|
||||||
|
"url": "https://example.com/feed.geojson",
|
||||||
|
"items_path": "features",
|
||||||
|
"id_path": "id",
|
||||||
|
"geometry_path": "geometry",
|
||||||
|
"title_path": "properties.name",
|
||||||
|
"category": "geo",
|
||||||
|
}
|
||||||
|
item = {
|
||||||
|
"id": "g1",
|
||||||
|
"geometry": {"type": "Point", "coordinates": [-116.2, 43.6]},
|
||||||
|
"properties": {"name": "Test Point"},
|
||||||
|
}
|
||||||
|
adapter = GenericHttpAdapter([source])
|
||||||
|
evt = adapter._item_to_event(item, source)
|
||||||
|
assert evt["latitude"] == 43.6
|
||||||
|
assert evt["longitude"] == -116.2
|
||||||
|
# the resolved geometry dict rides on data so the coverage gate can use it
|
||||||
|
assert evt["data"]["geometry"]["type"] == "Point"
|
||||||
|
|
||||||
|
|
||||||
|
def test_to_event_carries_latlon_for_coverage_gate():
|
||||||
|
adapter = GenericHttpAdapter([IDAHO_POWER_SOURCE])
|
||||||
|
evt = adapter._item_to_event(IDAHO_POWER_ITEM, IDAHO_POWER_SOURCE)
|
||||||
|
event = adapter.to_event(evt)
|
||||||
|
assert event is not None
|
||||||
|
assert event.lat == 43.6
|
||||||
|
assert event.lon == -116.2
|
||||||
|
assert event.category == "power_outage"
|
||||||
|
assert event.summary == "⚡ Power out — 250 affected, ETA 10:30 PM (Crew assigned)"
|
||||||
|
|
||||||
|
|
||||||
|
# ===========================================================================
|
||||||
|
# cold-start-silent (full store path) + persistence
|
||||||
|
# ===========================================================================
|
||||||
|
|
||||||
|
def _payload(items):
|
||||||
|
return {"object": {"outages": items, "totalCustomersAffected": sum(
|
||||||
|
i.get("omsCustomerCount", 0) for i in items)}}
|
||||||
|
|
||||||
|
|
||||||
|
def _make_store_with_generic():
|
||||||
|
bus = EventBus()
|
||||||
|
captured = []
|
||||||
|
bus.subscribe(lambda e: captured.append(e))
|
||||||
|
store = EnvironmentalStore(
|
||||||
|
EnvironmentalConfig(), event_bus=bus,
|
||||||
|
generic_sources=[dict(IDAHO_POWER_SOURCE)],
|
||||||
|
)
|
||||||
|
adapter = store._adapters["generic_http"]
|
||||||
|
return store, adapter, captured
|
||||||
|
|
||||||
|
|
||||||
|
def test_cold_start_silent_first_poll_seeds_persists_no_emit():
|
||||||
|
store, adapter, captured = _make_store_with_generic()
|
||||||
|
# Stub the network fetch with one active outage.
|
||||||
|
adapter._fetch = lambda url: _payload([IDAHO_POWER_ITEM])
|
||||||
|
|
||||||
|
store.refresh() # poll 1 == pre-existing backlog
|
||||||
|
|
||||||
|
# Nothing broadcast on the cold-start poll...
|
||||||
|
assert captured == [], "first poll must broadcast NOTHING (cold-start seed)"
|
||||||
|
# ...but the item IS persisted so the LLM sees it immediately.
|
||||||
|
row = get_db().execute(
|
||||||
|
"SELECT source, event_id, category, title, lat, lon FROM generic_events "
|
||||||
|
"WHERE source=? AND event_id=?", ("idaho_power", "123")).fetchone()
|
||||||
|
assert row is not None
|
||||||
|
assert row["category"] == "power_outage"
|
||||||
|
assert row["title"] == "Equipment"
|
||||||
|
assert abs(row["lat"] - 43.6) < 1e-6
|
||||||
|
|
||||||
|
|
||||||
|
def test_later_poll_broadcasts_newly_received_item():
|
||||||
|
store, adapter, captured = _make_store_with_generic()
|
||||||
|
adapter._fetch = lambda url: _payload([IDAHO_POWER_ITEM])
|
||||||
|
store.refresh() # poll 1 — seed silently
|
||||||
|
assert captured == []
|
||||||
|
|
||||||
|
# A genuinely NEW outage appears on a later poll -> it must broadcast.
|
||||||
|
new_item = dict(IDAHO_POWER_ITEM, omsOutageId="456", omsCustomerCount=99)
|
||||||
|
adapter._fetch = lambda url: _payload([IDAHO_POWER_ITEM, new_item])
|
||||||
|
adapter._last_poll.clear() # force cadence to elapse
|
||||||
|
store.refresh() # poll 2
|
||||||
|
|
||||||
|
assert len(captured) == 1, "only the newly-received outage broadcasts"
|
||||||
|
assert captured[0].category == "power_outage"
|
||||||
|
# both items persisted (dedup by source+event_id)
|
||||||
|
n = get_db().execute(
|
||||||
|
"SELECT COUNT(*) c FROM generic_events WHERE source=?",
|
||||||
|
("idaho_power",)).fetchone()["c"]
|
||||||
|
assert n == 2
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_generic_detail_reader():
|
||||||
|
from meshai.notifications.env_reporter import EnvReporter
|
||||||
|
store, adapter, captured = _make_store_with_generic()
|
||||||
|
adapter._fetch = lambda url: _payload([IDAHO_POWER_ITEM])
|
||||||
|
store.refresh()
|
||||||
|
|
||||||
|
text = EnvReporter().build_generic_detail()
|
||||||
|
assert "idaho_power" in text
|
||||||
|
assert "power_outage" in text
|
||||||
|
assert "Equipment" in text
|
||||||
|
|
@ -15,7 +15,7 @@ from meshai.persistence.observer_locations import (
|
||||||
# -- schema / migration -------------------------------------------------------
|
# -- schema / migration -------------------------------------------------------
|
||||||
|
|
||||||
def test_schema_version_is_25():
|
def test_schema_version_is_25():
|
||||||
assert SCHEMA_VERSION == 25
|
assert SCHEMA_VERSION == 26
|
||||||
|
|
||||||
|
|
||||||
def test_observer_locations_table_exists():
|
def test_observer_locations_table_exists():
|
||||||
|
|
@ -30,7 +30,7 @@ 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"]) == 25
|
assert int(row["value"]) == 26
|
||||||
|
|
||||||
|
|
||||||
# -- accessors ----------------------------------------------------------------
|
# -- accessors ----------------------------------------------------------------
|
||||||
|
|
|
||||||
|
|
@ -99,9 +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 25 by the LLM-persistence migrations (v24 avalanche_events,
|
# Bumped to 26 by the generic-source migration (v26 generic_events); v24
|
||||||
# v25 ducting_events); was 23 at the native-satpass observer_locations (v23).
|
# avalanche_events, v25 ducting_events; was 23 at native-satpass (v23).
|
||||||
assert SCHEMA_VERSION == 25
|
assert SCHEMA_VERSION == 26
|
||||||
|
|
||||||
|
|
||||||
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):
|
||||||
|
|
@ -114,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"]) == 25
|
assert int(row["value"]) == 26
|
||||||
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