diff --git a/work/meshai/adapter_config/defaults.py b/work/meshai/adapter_config/defaults.py index b1dec70..f91c07a 100644 --- a/work/meshai/adapter_config/defaults.py +++ b/work/meshai/adapter_config/defaults.py @@ -85,7 +85,8 @@ REGISTRY: dict[tuple[str, str], dict[str, Any]] = { }, # ================================================================= - # USGS_QUAKE -- 6 settings (regional geography + 3 mag floors + PAGER set) + # USGS_QUAKE -- 7 settings (regional geography + 3 mag floors + PAGER set + # + freshness gate) # ================================================================= ("usgs_quake", "regional_centroid"): { "default": [44.36, -114.61], # quake_handler.py:36-37 (Idaho centroid) @@ -117,6 +118,11 @@ REGISTRY: dict[tuple[str, str], dict[str, Any]] = { "type": "float", "description": "Magnitude floor for the visual escalation emoji.", }, + ("usgs_quake", "freshness_seconds"): { + "default": 3600, + "type": "int", + "description": "Staleness gate for earthquake_event broadcasts (age = now - quake origin time), read by the dispatcher instead of the generic per-toggle freshness. The upstream feed is a rolling PAST-DAY feed, so quakes are routinely detected 10min-20h after origin; 1 hour keeps quakes broadcast-worthy (still situationally relevant) while excluding the day-feed's stale backlog. 0 = disabled (not recommended -- would readmit day-old quakes).", + }, # ================================================================= # SWPC -- 3 settings (three storm-tier broadcast floors) diff --git a/work/meshai/env/usgs_quake.py b/work/meshai/env/usgs_quake.py index 943e1af..264841d 100644 --- a/work/meshai/env/usgs_quake.py +++ b/work/meshai/env/usgs_quake.py @@ -246,7 +246,14 @@ class USGSQuakeAdapter: expires=evt.get("expires"), lat=lat, lon=lon, - region=evt.get("region"), + # Deliberately NOT passing region=evt.get("region") here: the + # raw dict's region is a fixed config default (self._region, + # "magic_valley"), which doesn't match any region_routes cell + # key. Event.region is documented as "set by region tagger" + # (notifications/events.py) -- leaving it unset lets + # CoverageFilter's geometry-based tagging (event.lat/lon vs + # the named coverage areas) populate the real region name + # from the quake's actual coordinates instead. group_key=event_id, inhibit_keys=[event_id], data=canonical_data, diff --git a/work/meshai/notifications/pipeline/dispatcher.py b/work/meshai/notifications/pipeline/dispatcher.py index aafc5f6..3034af2 100644 --- a/work/meshai/notifications/pipeline/dispatcher.py +++ b/work/meshai/notifications/pipeline/dispatcher.py @@ -359,6 +359,18 @@ class Dispatcher: # v0.6-3b: fire toggle uses wfigs adapter_config freshness (0 = disabled) if fam == "fire": freshness_s = int(adapter_config.wfigs.freshness_seconds) + elif event.category == "earthquake_event": + # The upstream USGS feed is a rolling PAST-DAY feed + # (2.5_day.geojson), so a quake's age-at-ingest routinely exceeds + # the generic 600s toggle window even for a genuinely first-seen + # event -- every native quake was being silently dropped here + # (measured age-at-ingest 548s-72169s across 12 sampled quakes, + # confirmed against quake_events: 12/12 first-sighted rows with + # last_broadcast_at=NULL). Scoped to the earthquake_event + # CATEGORY, not the "seismic" family/toggle, because that family + # also carries stream_flood_warning/stream_high_water (hydro), + # which must keep using the generic per-toggle freshness_seconds. + freshness_s = int(adapter_config.usgs_quake.freshness_seconds) else: freshness_s = int(getattr(tog, "freshness_seconds", 600) or 600) if event.timestamp and freshness_s > 0: diff --git a/work/tests/test_adapter_config_api.py b/work/tests/test_adapter_config_api.py index 8e66188..19bbe4e 100644 --- a/work/tests/test_adapter_config_api.py +++ b/work/tests/test_adapter_config_api.py @@ -86,6 +86,7 @@ def test_per_adapter_list(client): "regional_centroid", "regional_radius_mi", "broadcast_pager_alerts", "global_mag_floor", "regional_mag_floor", "escalate_mag_floor", + "freshness_seconds", } diff --git a/work/tests/test_adapter_usgs_quake.py b/work/tests/test_adapter_usgs_quake.py index 7074983..d3ee539 100644 --- a/work/tests/test_adapter_usgs_quake.py +++ b/work/tests/test_adapter_usgs_quake.py @@ -200,12 +200,32 @@ def test_populates_core_fields(adapter): assert event.source == "usgs_quake" assert event.lat == 42.61 assert event.lon == -114.48 - assert event.region == "magic_valley" + # to_event() deliberately does NOT propagate the raw dict's "region" + # (a fixed config default that never matches region_routes cell keys) + # onto the Event -- Event.region is left unset so CoverageFilter's + # geometry-based tagger (lat/lon vs named coverage areas) can set the + # real region name downstream. See test_region_left_unset_for_tagger + # and tests/test_coverage_area.py for the tagging behavior itself. + assert event.region is None assert event.expires == evt["expires"] assert event.timestamp == evt["quake_time"] assert event.id +def test_region_left_unset_for_tagger(adapter): + """Regression: to_event() must never propagate the adapter's fixed + config region (e.g. "magic_valley") onto the Event. That value never + matches a region_routes cell key ("East Idaho"/"SC Idaho"/"SW Idaho"), + so pre-setting it would silently break region-routed delivery. Leaving + event.region/regions unset lets CoverageFilter's geometry tagger (which + only fires when `not event.regions`) stamp the real area name from the + quake's actual lat/lon.""" + evt = make_quake_event(lat=44.2, lon=-114.0, region="magic_valley") + event = adapter.to_event(evt) + assert event.region is None + assert event.regions == [] + + # ============================================================ # DEFENSIVE TESTS # ============================================================ diff --git a/work/tests/test_quake_freshness_and_region_tagging.py b/work/tests/test_quake_freshness_and_region_tagging.py new file mode 100644 index 0000000..ab0166f --- /dev/null +++ b/work/tests/test_quake_freshness_and_region_tagging.py @@ -0,0 +1,212 @@ +"""Regression tests for the usgs_quake dispatcher-freshness + region-tagging fix. + +Problem (see meshai/notifications/pipeline/dispatcher.py:359-374 and +meshai/env/usgs_quake.py): the usgs_quake adapter polls a rolling PAST-DAY +USGS feed (2.5_day.geojson), so a genuinely first-seen quake is routinely +already 10min-20h old by the time it's detected. The generic per-toggle +freshness_seconds (600s) silently dropped essentially every native quake +before broadcast -- confirmed against production: quake_events had 12 +first-sighting rows (decide() ran, magnitude/region gate passed, DB row +inserted) with last_broadcast_at=NULL on every single one (commit() never +reached because the staleness filter dropped the event downstream first). + +Fix 1 (dispatcher.py): earthquake_event gets its own adapter_config-backed +freshness override (adapter_config.usgs_quake.freshness_seconds, default +3600s), mirroring the existing wfigs/"fire" override -- but scoped to the +CATEGORY, not the "seismic" family/toggle, because stream_flood_warning / +stream_high_water (hydro) also live under toggle="seismic" and must keep +using the generic per-toggle freshness unchanged. + +Fix 2 (env/usgs_quake.py): to_event() no longer passes the adapter's fixed +config region ("magic_valley", which never matches a region_routes cell +key) onto the Event. Event.region is left unset so CoverageFilter's +geometry-based tagger (lat/lon vs named coverage areas) can stamp the real +region name. +""" + +from __future__ import annotations + +import asyncio +import time +from unittest.mock import MagicMock + +import pytest + +from meshai.config import Config +from meshai.coverage_area import MonitoringArea +from meshai.env.usgs_quake import USGSQuakeAdapter +from meshai.notifications.events import make_event +from meshai.notifications.pipeline.coverage_filter import CoverageFilter +from meshai.notifications.pipeline.dispatcher import Dispatcher + + +# -------------------------------------------------------------------------- +# Shared dispatcher-test helpers (mirror tests/test_v052_dispatcher.py) +# -------------------------------------------------------------------------- + +class RecChannel: + def __init__(self, rec): + self.rec = rec + + async def deliver(self, payload, rule): + self.rec.append({"name": rule.name, "message": payload.message}) + return True + + +def _make_dispatcher(cfg): + rec: list = [] + d = Dispatcher(cfg, lambda rule, conn: RecChannel(rec), connector=None) + return d, rec + + +def _dispatch_one(cfg, event): + d, rec = _make_dispatcher(cfg) + asyncio.run(d.dispatch(event)) + return d, rec + + +def _quake_cfg(): + """Config with the seismic toggle enabled, cold-start grace disabled.""" + cfg = Config() + cfg.notifications.rules = [] + cfg.notifications.cold_start_grace_seconds = 0 + t = cfg.notifications.toggles["seismic"] + t.enabled = True + t.min_severity = "routine" + t.regions = [] + t.severity_channels = { + "routine": ["mesh_broadcast"], + "priority": ["mesh_broadcast"], + "immediate": ["mesh_broadcast"], + } + t.cooldown_seconds = 0 + # Generic per-toggle freshness -- deliberately tight (600s, the old + # default) so these tests prove the earthquake_event category is NOT + # using this value anymore. + t.freshness_seconds = 600 + return cfg + + +def _quake_event(age_seconds: float, event_id="us6000abcd", lat=42.6, lon=-114.5): + return make_event( + source="usgs_quake", + category="earthquake_event", + severity="routine", + title=f"M3.0 -- {event_id}", + timestamp=time.time() - age_seconds, + lat=lat, + lon=lon, + group_key=event_id, + inhibit_keys=[event_id], + ) + + +# ============================================================ +# (a) fresh quake passes under the new per-adapter override +# ============================================================ + +def test_fresh_quake_passes_staleness_gate_with_override(): + """A 30-minute-old quake (1800s) exceeds the generic 600s toggle window + but must pass under the adapter_config.usgs_quake.freshness_seconds + override (default 3600s).""" + cfg = _quake_cfg() + event = _quake_event(age_seconds=1800) + d, rec = _dispatch_one(cfg, event) + assert len(rec) == 1, "30-min-old quake must broadcast under the 3600s override" + assert d.dispatch_stats()["stale_dropped"] == 0 + + +# ============================================================ +# (b) old quake (12h) is still dropped +# ============================================================ + +def test_old_quake_still_dropped_by_override(): + """A 12-hour-old quake (43200s) must still be dropped -- the override + widens the window, it does not disable the staleness check.""" + cfg = _quake_cfg() + event = _quake_event(age_seconds=12 * 3600) + d, rec = _dispatch_one(cfg, event) + assert rec == [], "12h-old quake must still be dropped as stale" + assert d.dispatch_stats()["stale_dropped"] == 1 + + +def test_hydro_seismic_sibling_unaffected_by_quake_override(): + """Guardrail: stream_flood_warning shares toggle='seismic' with + earthquake_event but must keep using the GENERIC per-toggle freshness + (600s here), not the quake-specific 3600s override. A 1800s-old hydro + event must still be dropped.""" + cfg = _quake_cfg() + event = make_event( + source="usgs", category="stream_flood_warning", severity="priority", + title="Snake River nr Twin Falls 12.8 ft", + timestamp=time.time() - 1800, + ) + d, rec = _dispatch_one(cfg, event) + assert rec == [], "hydro sibling must NOT inherit the quake freshness override" + assert d.dispatch_stats()["stale_dropped"] == 1 + + +# ============================================================ +# (c) region-tagging from real lat/lon +# ============================================================ + +SW_IDAHO = MonitoringArea(name="SW Idaho", west=-117.993408, south=42.05, + east=-115.389404, north=44.331707) +SC_IDAHO = MonitoringArea(name="SC Idaho", west=-115.389404, south=42.05, + east=-112.8, north=44.331707) +EAST_IDAHO = MonitoringArea(name="East Idaho", west=-112.8, south=41.9, + east=-110.9, north=45.331707) +PROD_AREAS = [SW_IDAHO, SC_IDAHO, EAST_IDAHO] + + +def _quake_adapter(): + cfg = MagicMock() + cfg.feed_url = "https://example.test/feed.geojson" + cfg.min_magnitude = 2.5 + cfg.bbox = [-115.5, 42.0, -110.0, 45.2] + cfg.region = "magic_valley" + cfg.tick_seconds = 300 + return USGSQuakeAdapter(cfg) + + +def test_quake_event_region_tags_to_matching_coverage_area(): + """A quake with real Idaho lat/lon, run through the actual adapter's + to_event() and then the real CoverageFilter (named exactly like prod's + SW/SC/East Idaho areas), must end up tagged with a region name that + matches a region_routes cell key -- NOT the adapter's stale + "magic_valley" default.""" + adapter = _quake_adapter() + raw = { + "source": "usgs_quake", + "event_id": "us7000example", + "event_type": "Earthquake", + "severity": "routine", + "headline": "M2.7 -- 10 km N of Twin Falls, ID", + "magnitude": 2.7, + "place": "10 km N of Twin Falls, ID", + "depth_km": 8.0, + "sig": 80, + "url": "https://earthquake.usgs.gov/x", + "region": "magic_valley", + "lat": 42.61, # Twin Falls -- falls inside SC Idaho + "lon": -114.48, + "quake_time": time.time(), + "expires": time.time() + 86400, + "fetched_at": time.time(), + } + event = adapter.to_event(raw) + assert event is not None + assert event.region is None # adapter no longer pre-stamps a region + + received: list = [] + flt = CoverageFilter(next_handler=received.append, areas=PROD_AREAS, enabled=True) + flt.handle(event) + + assert len(received) == 1, "in-region quake must pass the coverage gate" + tagged = received[0] + assert tagged.region == "SC Idaho", ( + f"expected geometry tagging to set 'SC Idaho', got {tagged.region!r} " + "(pre-fix this would have stayed 'magic_valley' and never matched a " + "region_routes cell key)" + ) + assert tagged.regions == ["SC Idaho"]