fix(wzdx): persist current work zones to traffic_events via dedicated ingest (#109)

The WZDx daily summary + DM query count from traffic_events, but work
zones weren't landing there: wzdx rode the generic _delta_emit path,
which silent-seeds the seen-set and returns before the decider's INSERT
on the cold-start first poll, so the current zone set never persisted
(summary would count ~0). Add a dedicated _ingest_wzdx (mirroring the
fires ingest) that UPSERTs every current coalesced zone into
traffic_events each poll (persist-only, last_broadcast_at=NULL, no emit,
no broadcast) and reconciles zones that drop out of the feed (never wipes
on an empty/failed fetch). Per-event work-zone broadcast stays suppressed
(the decider's work_zone gate is untouched). Retargets 4 tests in
test_store_received_delta.py that used a fake 'wzdx' source to exercise
the generic gate onto a neutral routing name.

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:
malice 2026-07-09 20:08:13 -06:00 committed by GitHub
commit 77e057ae86
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 463 additions and 6 deletions

View file

@ -362,6 +362,19 @@ class EnvironmentalStore:
# 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 == "wzdx":
# Native WZDx work zones DELIBERATELY bypass the generic
# `_delta_emit` path: that path silent-seeds every zone on the cold-
# start first poll and `return`s BEFORE the incident decider's
# `INSERT INTO traffic_events` ever runs, so the current active
# work-zone set (the ~127 coalesced zones) would never persist —
# only later-newly-appearing zones would trickle in, leaving the
# daily summary + DM query counting ~0. Instead `_ingest_wzdx`
# UPSERTS the CURRENT coalesced set into traffic_events every poll
# (persist-only, NEVER emit/decide/broadcast — mirrors the fires /
# generic persist-every-item pattern), and reconciles removals, so
# the table always equals the current active set the summary reads.
self._ingest_wzdx(adapter)
elif name == "avalanche":
# Avalanche: re-emit on danger_level rise (Update:) not just new
# events. The rise is a legitimate CONTENT change, so it passes
@ -551,6 +564,127 @@ class EnvironmentalStore:
if events:
self._fires_seeded = True
def _ingest_wzdx(self, adapter) -> None:
"""Native WZDx work-zone ingest — persist-only current-state mirror.
The daily summary (``wzdx_summary.fire_once``) and the DM detail
(``env_reporter.build_work_zones_detail``) both read the CURRENT active
work-zone set straight from ``traffic_events``:
SELECT ... FROM traffic_events
WHERE source='wzdx' AND (end_at IS NULL OR end_at >= now)
Left on the generic ``_delta_emit`` path, wzdx never populates that
table on a cold start: ``_delta_emit`` silent-seeds every zone into the
received-delta seen-set and ``return``s on the first poll, BEFORE the
incident decider's ``INSERT INTO traffic_events`` runs — so the current
coalesced set (the ~127 zones) is never persisted; only zones that
FIRST appear on a LATER poll would trickle in. Result: the summary/DM
count ~0 instead of the real active total.
This dedicated ingest fixes the gap the way ``_ingest_fires`` /
``_persist_generic`` do UPSERT EVERY current item on EVERY poll
but scoped to persistence only:
* UPSERT each CURRENT coalesced work zone into ``traffic_events``
(INSERT on a new ``external_id`` with ``first_seen_at`` = now and
``last_broadcast_at`` = NULL; UPDATE ``last_seen_at`` + the current
fields, preserving ``first_seen_at``), using the SAME column
set/types the incident decider writes so the summary/DM queries
work unchanged. ``last_broadcast_at`` is left NULL forever (never
armed) so nothing here ever looks "already broadcast".
* RECONCILE removals: after upserting this poll's set, DELETE any
``source='wzdx'`` rows whose ``external_id`` is NOT in the current
set, so the table == the CURRENT active set (stale zones would
otherwise inflate the count). Reconciliation runs ONLY when the
poll actually returned 1 zone an empty/failed fetch never wipes
the existing rows.
Persist ONLY: it does NOT emit to the pipeline, run the decider, or
broadcast. Per-event work-zone broadcast stays suppressed exactly as
before (the incident decider's ``work_zone`` ``broadcast=False`` gate
remains as defense-in-depth for any wzdx event that reaches it another
way). ``self._events`` is not touched here the summary/DM read the
durable table, not the in-memory cache.
"""
try:
from meshai.persistence import get_db
conn = get_db()
except Exception as e:
logger.warning("wzdx ingest skipped (DB unavailable): %s", e)
return
now = int(time.time())
events = adapter.get_events()
current_ext_ids: list = []
for evt in events:
try:
external_id = evt.get("external_id")
if not external_id:
continue # no stable identity → cannot dedup/persist
n = evt.get("normalized") or {}
road = n.get("road")
direction = n.get("direction")
mile_start = n.get("mile_start")
mile_end = n.get("mile_end")
sub_type = n.get("sub_type")
impact = n.get("impact")
lat = evt.get("lat")
lon = evt.get("lon")
start_at = evt.get("start_at")
end_at = evt.get("end_at")
current_ext_ids.append(external_id)
# UPSERT the current coalesced zone. Composite PK
# (source, external_id): INSERT on first sight (first_seen_at =
# now, last_broadcast_at = NULL → never armed), else UPDATE the
# current fields + last_seen_at while PRESERVING first_seen_at.
# end_at is refreshed from the feed on every poll so the
# summary's not-expired filter (end_at IS NULL OR end_at>=now)
# stays correct. Same columns/types the incident decider uses.
conn.execute(
"INSERT INTO traffic_events("
"source, external_id, road, direction, "
"mile_start, mile_end, county, state, lat, lon, "
"sub_type, impact, start_at, end_at, "
"first_seen_at, last_seen_at, last_broadcast_at"
") VALUES ('wzdx',?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) "
"ON CONFLICT(source, external_id) DO UPDATE SET "
"road=excluded.road, direction=excluded.direction, "
"mile_start=excluded.mile_start, mile_end=excluded.mile_end, "
"lat=excluded.lat, lon=excluded.lon, "
"sub_type=excluded.sub_type, impact=excluded.impact, "
"start_at=excluded.start_at, end_at=excluded.end_at, "
"last_seen_at=excluded.last_seen_at",
(
external_id, road, direction,
mile_start, mile_end, None, None, lat, lon,
sub_type, impact, start_at, end_at,
now, now, None,
),
)
except Exception:
logger.exception(
"wzdx ingest failed for %s", evt.get("external_id", "?"))
# RECONCILE removals — only when this poll actually returned a set.
# A zone that dropped out of the feed is deleted so the table equals
# the CURRENT active set; an empty/failed fetch (no current ext_ids)
# is left untouched so stale-but-only-momentarily-missing rows and a
# transient outage never wipe the summary's data.
if current_ext_ids:
try:
placeholders = ",".join("?" for _ in current_ext_ids)
conn.execute(
"DELETE FROM traffic_events WHERE source='wzdx' "
f"AND external_id NOT IN ({placeholders})",
current_ext_ids,
)
except Exception:
logger.exception("wzdx ingest reconcile-delete failed")
def _ingest_generic(self, adapter) -> None:
"""Native generic-source ingest — COLD-START-SILENT per source.

View file

@ -245,6 +245,25 @@ def _insert_quake(event_ids: list[str]) -> None:
)
# NOTE ON ADAPTER NAME vs EVENT SOURCE
# ------------------------------------
# These durable-pre-seed / received-delta-gate tests exercise the store's
# GENERIC emit path (store._delta_emit + _seed_from_persistent), which keys the
# seen-set and the "seeded" marker on the event's SOURCE ("wzdx"/"511"/
# "usgs_quake"), NOT on the adapter's registration name. The name only selects
# the _ingest routing branch. Native ``wzdx`` now has a DEDICATED persist-only
# ingest (store._ingest_wzdx) that upserts traffic_events and NEVER emits, so a
# ``_FakeWZDx`` injected under the literal name "wzdx" would take that path and
# broadcast nothing. To keep testing the GENERIC gate with an external_id-keyed
# source, these tests inject the ``source='wzdx'`` fake under the neutral
# routing name ``_GENERIC_NAME`` (the else-branch → _delta_emit), while every
# assertion still references the durable SOURCE "wzdx". Real wzdx persistence is
# covered separately in test_store_wzdx_persist.py. The name is arbitrary as
# long as it is NOT one of the _ingest special branches (swpc/ducting/nifc/
# generic_http/wzdx/avalanche); anything else routes to the generic else path.
_GENERIC_NAME = "wzdx_gate"
def _build_store(adapter_name: str, adapter):
"""Construct a store (runs the durable pre-seed against the current DB),
then inject a fake adapter. Insert persistent rows BEFORE calling this."""
@ -264,7 +283,7 @@ def test_persistent_preseed_known_suppressed_new_emitted():
_insert_traffic(known)
adapter = _FakeWZDx()
store, captured = _build_store("wzdx", adapter)
store, captured = _build_store(_GENERIC_NAME, adapter)
# wzdx was pre-seeded from the durable table AND marked seeded (baseline).
assert "wzdx" in store._seeded
@ -285,7 +304,7 @@ def test_persistent_preseed_cross_tick_staging_no_leak():
# seen in-process. Without the durable record it would leak.
_insert_traffic(["A", "B"]) # both already RECEIVED (durable)
adapter = _FakeWZDx()
store, captured = _build_store("wzdx", adapter)
store, captured = _build_store(_GENERIC_NAME, adapter)
adapter.set_batch(["A"])
store.refresh() # tick 1: only A present
@ -321,7 +340,7 @@ def test_incremental_empty_first_tick_then_only_new_broadcasts():
# C — genuinely never received — broadcasts.
_insert_traffic(["A", "B"])
adapter = _FakeWZDx()
store, captured = _build_store("wzdx", adapter)
store, captured = _build_store(_GENERIC_NAME, adapter)
adapter.set_batch([]) # empty first tick
store.refresh()
@ -339,14 +358,14 @@ def test_restart_against_same_persistent_db_never_rebroadcasts():
_insert_traffic(backlog)
a1 = _FakeWZDx()
store1, cap1 = _build_store("wzdx", a1)
store1, cap1 = _build_store(_GENERIC_NAME, a1)
a1.set_batch(backlog)
store1.refresh()
assert cap1 == [], "process 1: durable backlog is silent"
# RESTART: brand-new store, same persistent DB → pre-seed reloads.
a2 = _FakeWZDx()
store2, cap2 = _build_store("wzdx", a2)
store2, cap2 = _build_store(_GENERIC_NAME, a2)
a2.set_batch(backlog)
store2.refresh()
assert cap2 == [], "restart must NEVER re-broadcast the durable backlog"
@ -379,7 +398,7 @@ def test_no_durable_rows_falls_back_to_silent_first_poll():
# the first real batch is silently seeded (never leaked) — the fresh-start
# safety property.
adapter = _FakeWZDx()
store, captured = _build_store("wzdx", adapter)
store, captured = _build_store(_GENERIC_NAME, adapter)
assert "wzdx" not in store._seeded, "0 durable rows → not pre-marked seeded"
adapter.set_batch(["A", "B"])

View file

@ -0,0 +1,304 @@
"""Persistence tests for the dedicated native WZDx ingest path.
The daily summary (wzdx_summary.fire_once) and DM detail
(env_reporter.build_work_zones_detail) both read the CURRENT active work-zone
set straight from traffic_events:
SELECT ... FROM traffic_events
WHERE source='wzdx' AND (end_at IS NULL OR end_at >= now)
Before the fix, native wzdx rode the generic ``_delta_emit`` path, which
silent-seeds every zone on the cold-start first poll and ``return``s BEFORE the
incident decider's ``INSERT INTO traffic_events`` runs -- so the current
coalesced set never persisted and the summary counted ~0. ``store._ingest_wzdx``
now UPSERTS the current coalesced set into traffic_events every poll
(persist-only, never emit/decide/broadcast) and reconciles removals, so the
table always equals the current active set the summary reads.
These tests drive the REAL EnvironmentalStore + EventBus with a fake WZDx
adapter whose per-poll coalesced set we control, then assert directly against
traffic_events AND against the bus (nothing must ever be dispatched).
"""
from __future__ import annotations
from meshai.env.store import EnvironmentalStore
from meshai.config import EnvironmentalConfig
from meshai.notifications.pipeline.bus import EventBus
from meshai.notifications.events import make_event
from meshai.persistence import get_db
class _FakeWZDx:
"""Native-WZDx stand-in whose coalesced current set the test controls.
Its stored-event dicts carry the SAME shape env/wzdx.py emits: a coalesced
``external_id`` (``wzdx_<road>|<lat>|<lon>|<sub_type>``), ``lat``/``lon``,
``start_at``/``end_at`` epochs, and a ``normalized`` dict (the
_parse_wzdx_federal output: road/direction/mile_start/mile_end/sub_type/
impact). ``get_events()`` returns exactly the current coalesced set --
which is what ``_ingest_wzdx`` upserts + reconciles against.
"""
def __init__(self):
self._batch: list[dict] = []
self.emitted = False # set True if to_event() is ever called
def set_zones(self, zones: list[dict]) -> None:
"""zones: list of {ext, road, lat, lon, sub_type, impact, end_at,
start_at?, direction?} -> stored-event dicts like env/wzdx builds."""
batch = []
for z in zones:
ext = z["ext"]
n = {
"road": z.get("road"),
"direction": z.get("direction"),
"mile_start": z.get("mile_start"),
"mile_end": z.get("mile_end"),
"sub_type": z.get("sub_type"),
"impact": z.get("impact"),
}
batch.append({
"source": "wzdx",
"event_id": ext,
"external_id": ext,
"lat": z.get("lat"),
"lon": z.get("lon"),
"start_at": z.get("start_at"),
"end_at": z.get("end_at"),
"severity": "priority" if z.get("impact") == "full_closure" else "routine",
"normalized": n,
"fetched_at": 0,
})
self._batch = batch
def set_raw(self, batch: list[dict]) -> None:
"""Set the raw stored-event batch verbatim (edge cases)."""
self._batch = batch
def tick(self) -> bool:
# Real wzdx.tick() returns True even on a 0-event (registry-only) tick.
return True
def get_events(self) -> list:
return list(self._batch)
def to_event(self, raw_evt: dict):
# If _ingest_wzdx ever emitted, it would call this. It must NOT.
self.emitted = True
return make_event(source="wzdx", category="work_zone",
severity="routine", title=raw_evt.get("external_id"),
summary=raw_evt.get("external_id"),
group_key=raw_evt.get("external_id"))
def _build_store(adapter):
"""Construct a store (runs the durable pre-seed against the fresh migrated
DB the conftest points MESHAI_DB_PATH at), then inject the fake adapter.
Captures everything that reaches the EventBus so we can assert silence."""
bus = EventBus()
captured: list = []
bus.subscribe(lambda e: captured.append(e))
store = EnvironmentalStore(EnvironmentalConfig(), event_bus=bus)
store._adapters["wzdx"] = adapter
return store, captured
def _wzdx_rows():
conn = get_db()
return conn.execute(
"SELECT external_id, source, road, sub_type, impact, lat, lon, "
"end_at, first_seen_at, last_seen_at, last_broadcast_at "
"FROM traffic_events WHERE source='wzdx' "
"ORDER BY external_id"
).fetchall()
def _summary_visible_count(now: int):
"""Count exactly what wzdx_summary.fire_once / build_work_zones_detail see:
source='wzdx' AND not-expired."""
conn = get_db()
return conn.execute(
"SELECT COUNT(*) FROM traffic_events "
"WHERE source='wzdx' AND (end_at IS NULL OR end_at >= ?)",
(now,),
).fetchone()[0]
ZONES3 = [
{"ext": "wzdx_US-20|43.600|-116.200|construction work",
"road": "US-20", "lat": 43.6, "lon": -116.2,
"sub_type": "construction work", "impact": "partial",
"end_at": None},
{"ext": "wzdx_I-84|43.500|-116.400|construction work",
"road": "I-84", "lat": 43.5, "lon": -116.4,
"sub_type": "construction work", "impact": "full_closure",
"end_at": 9_000_000_000},
{"ext": "wzdx_ID-55|44.000|-116.000|maintenance",
"road": "ID-55", "lat": 44.0, "lon": -116.0,
"sub_type": "maintenance", "impact": "partial",
"end_at": 9_000_000_000},
]
def test_first_poll_persists_current_set_and_broadcasts_nothing():
# THE GAP FIX: the cold-start first poll must PERSIST every current
# coalesced zone into traffic_events (source='wzdx') and emit NOTHING.
adapter = _FakeWZDx()
store, captured = _build_store(adapter)
adapter.set_zones(ZONES3)
store.refresh() # first (cold-start) poll
rows = _wzdx_rows()
exts = {r["external_id"] for r in rows}
assert exts == {z["ext"] for z in ZONES3}, (
"first poll must persist ALL current coalesced zones (the gap fix)")
# No broadcast / no emit occurred.
assert captured == [], "persist-only: NOTHING may reach the pipeline bus"
assert adapter.emitted is False, "to_event()/emit must never be called"
# Every row is source='wzdx' with last_broadcast_at NULL (never armed).
for r in rows:
assert r["source"] == "wzdx"
assert r["last_broadcast_at"] is None
# The summary/DM query now sees the real count, not ~0.
assert _summary_visible_count(now=0) == 3
def test_columns_match_summary_and_dm_queries():
# The persisted columns the summary (road/lat/lon/sub_type/impact/end_at)
# and DM (road/direction/sub_type/impact/lat/lon/end_at) read must be
# populated from the coalesced zone, not left NULL.
adapter = _FakeWZDx()
store, _ = _build_store(adapter)
adapter.set_zones([ZONES3[1]]) # the full_closure I-84 zone
store.refresh()
r = _wzdx_rows()[0]
assert r["road"] == "I-84"
assert r["sub_type"] == "construction work"
assert r["impact"] == "full_closure"
assert abs(r["lat"] - 43.5) < 1e-9
assert abs(r["lon"] - (-116.4)) < 1e-9
assert r["end_at"] == 9_000_000_000
def test_subsequent_poll_reconciles_removed_zone():
# A zone that DROPS OUT of the feed on a later poll must be REMOVED so the
# table equals the CURRENT active set (else stale rows inflate the count).
adapter = _FakeWZDx()
store, captured = _build_store(adapter)
adapter.set_zones(ZONES3)
store.refresh()
assert len(_wzdx_rows()) == 3
# Next poll: US-20 dropped out; I-84 + ID-55 remain.
adapter.set_zones([ZONES3[1], ZONES3[2]])
store.refresh()
exts = {r["external_id"] for r in _wzdx_rows()}
assert exts == {ZONES3[1]["ext"], ZONES3[2]["ext"]}, (
"the dropped zone must be reconciled out of traffic_events")
assert captured == [], "reconcile is persist-only; still no broadcast"
def test_empty_or_failed_fetch_does_not_wipe_existing_rows():
# An empty/failed fetch (get_events() == []) must NOT delete existing rows
# -- a transient upstream outage must never zero out the summary's data.
adapter = _FakeWZDx()
store, _ = _build_store(adapter)
adapter.set_zones(ZONES3)
store.refresh()
assert len(_wzdx_rows()) == 3
adapter.set_raw([]) # empty/failed poll
store.refresh()
assert len(_wzdx_rows()) == 3, (
"an empty fetch must NEVER wipe the existing active set")
def test_upsert_preserves_first_seen_at_and_refreshes_end_at():
# A zone seen across polls keeps its first_seen_at but refreshes
# last_seen_at + end_at (so the summary's not-expired filter tracks the
# feed's latest end date).
adapter = _FakeWZDx()
store, _ = _build_store(adapter)
z = dict(ZONES3[0]); z["end_at"] = 1000
adapter.set_zones([z])
store.refresh()
r1 = _wzdx_rows()[0]
first_seen = r1["first_seen_at"]
assert r1["end_at"] == 1000
# Same zone reappears with a LATER end_at.
z2 = dict(ZONES3[0]); z2["end_at"] = 5000
adapter.set_zones([z2])
store.refresh()
r2 = _wzdx_rows()[0]
assert r2["first_seen_at"] == first_seen, "first_seen_at must be preserved"
assert r2["end_at"] == 5000, "end_at must refresh from the feed"
def test_expiry_end_at_preserved_for_not_expired_filter():
# end_at must be persisted verbatim so the summary's not-expired filter
# (end_at IS NULL OR end_at >= now) correctly includes/excludes zones.
now = 1_000_000
adapter = _FakeWZDx()
store, _ = _build_store(adapter)
zones = [
{"ext": "wzdx_open|1.0|1.0|x", "road": "OPEN", "lat": 1.0, "lon": 1.0,
"sub_type": "x", "impact": "partial", "end_at": None}, # open-ended
{"ext": "wzdx_future|2.0|2.0|x", "road": "FUT", "lat": 2.0, "lon": 2.0,
"sub_type": "x", "impact": "partial", "end_at": now + 10_000}, # not expired
{"ext": "wzdx_past|3.0|3.0|x", "road": "PAST", "lat": 3.0, "lon": 3.0,
"sub_type": "x", "impact": "partial", "end_at": now - 10_000}, # expired
]
adapter.set_zones(zones)
store.refresh()
# All 3 persisted (ingest does not itself drop expired rows) ...
assert len(_wzdx_rows()) == 3
# ... but the summary's not-expired filter counts only the open + future.
assert _summary_visible_count(now=now) == 2
def test_id_less_zone_is_skipped_not_fatal():
# A stored event with no external_id cannot be keyed; it is skipped (never
# persisted, never raises), and does not participate in reconciliation.
adapter = _FakeWZDx()
store, _ = _build_store(adapter)
good = ZONES3[0]
adapter.set_zones([good])
store.refresh()
assert len(_wzdx_rows()) == 1
# Poll with the good zone plus an id-less junk event.
adapter.set_zones([good])
junk = {"source": "wzdx", "event_id": None, "external_id": None,
"lat": 5.0, "lon": 5.0, "normalized": {}, "fetched_at": 0}
adapter._batch.append(junk)
store.refresh()
rows = _wzdx_rows()
assert {r["external_id"] for r in rows} == {good["ext"]}, (
"id-less zone skipped; good zone still persisted, reconcile intact")
def test_bulk_current_set_persists_all_like_the_real_127():
# Scale check mirroring the real ~127-zone active set: every coalesced
# zone in a large current set persists and is visible to the summary.
adapter = _FakeWZDx()
store, captured = _build_store(adapter)
zones = [
{"ext": f"wzdx_R{i}|{40.0 + i * 0.001:.3f}|-116.000|maintenance",
"road": f"R{i}", "lat": 40.0 + i * 0.001, "lon": -116.0,
"sub_type": "maintenance", "impact": "partial", "end_at": None}
for i in range(127)
]
adapter.set_zones(zones)
store.refresh()
assert len(_wzdx_rows()) == 127
assert _summary_visible_count(now=0) == 127, (
"the full active set is counted, not ~0")
assert captured == [], "bulk persist still broadcasts nothing"