diff --git a/CLAUDE.md b/CLAUDE.md index 1b145ad..4975b57 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -47,7 +47,8 @@ with the evidence that closed them. | Duration outliers are categorical, not statistical. The 14-day cap is a backstop, not the outlier strategy. | data-quality.md — "Duration outliers are categorical" | | "We are investigating" reference pairing works but rescues almost nothing — not worth building. | data-quality.md — "'We are investigating' notices" (corrected 2026-07-20) | | `closed_at` is a floor: short-lived cases are never observed open. Twice-daily builds are the settled cadence. | data-quality.md — "`closed_at` is a floor" (re-measured 2026-07-31) | -| A case the feed drops while `Open` is stamped `vanished_at` (schema v4) and is closed with no signal on the site, never `closed_at`. The stamp is safe only behind the feed-count guard (`FEED_COUNT_TOLERANCE`), which refuses a short download before anything touches the DB. | data-quality.md - "Cases that vanish from the feed" (2026-09-05) | +| A case the feed drops while `Open` is stamped `vanished_at` (schema v4) and is closed with no signal on the site, never `closed_at`; the stamp touches closed rows too. It is safe only behind the feed-count guard (`FEED_COUNT_TOLERANCE`), which refuses a short download before anything touches the DB, and behind the empty-download refusal. Paging is by `OBJECTID`, refused unless each page is strictly ascending. | data-quality.md - "Cases that vanish from the feed" (2026-09-05, amended 2026-09-24) | +| A feature with no pin (no `geometry`, or `"NaN"`) keeps the pin the DB last stored, or is set aside with a `::warning::` until the feed pins it. Nullable coordinates were rejected: not an additive migration. | data-quality.md - "A feature with no pin" (2026-09-24) | | **A case is open only while nothing its own text has ended.** `is_open(row, now)` reads `status`, `vanished_at` and a passed *observed* end, decided once in `resolve_case` and carried on `Case.is_open` for every surface that says open. The close date follows the same reading: the notice's own completion, else `closed_at` (`closed_on`, 2026-09-24). The feed closes a case a median 72h after the notice reports completion; 216 of 562 `Open` cases were past one, 0 of 7,667 completions were ever followed up. Scheduled ends do not close a case for display. | statuspage-methodology.md - "The notice's own completion closes it" (2026-09-05) | | gemma-4-12b-qat over qwen3.5-9b for end-time extraction; prompt version is at v3. | model-and-runtime-benchmarks.md, end-time-eval.md | | Geography is CSO Census settlements, not the feed's `location` string (3,866 distinct values, fragments badly, carries no population). | statuspage-methodology.md — "The county drill-down" (2026-07-25) | diff --git a/notes/data-quality.md b/notes/data-quality.md index fb92bad..179ef2c 100644 --- a/notes/data-quality.md +++ b/notes/data-quality.md @@ -294,6 +294,29 @@ and refuses a download more than 1% short (`FEED_COUNT_TOLERANCE`). The 1% is fo changing under the paging; a real purge like the one around 2026-04-20 would still be stamped, which is right: that is what happened. +*2026-09-24:* the tolerance passes an empty feed, 0 downloaded of 0 reported, and that build +would stamp vanished every row not yet vanished, closed ones included (4,306 on the 2026-09-23 +release, 498 of them open). `run` now also refuses an empty download +while the DB holds rows not yet vanished; the guard counts what the stamp touches, not the +open ones only. No real purge has emptied the feed; the 2026-08-10 one left 3,044 cases, so a +partial purge is still stamped as before. The same review moved `download_cases` from +`resultOffset` to `OBJECTID > ` paging: by offset, one case deleted during the +download pushed a live case out of the next page and stamped it vanished, and a server +`maxRecordCount` below the 1,000 asked for dropped the difference at every page, both inside +the 1%. Key paging is only sound on pages returned in `OBJECTID` order, so a page that is not +ascending fails the build rather than skip or repeat rows. + +### A feature with no pin (2026-09-24) + +ArcGIS omits `geometry` for a null shape, or writes an empty point as `"NaN"`, and one such +feature crashed every build: the coordinates are `NOT NULL` in `cases` and key the geocode +cache. Making them nullable was rejected, because that is not an additive migration. +`restore_pins` instead gives a case the pin the DB last stored for it, and sets aside a case +the DB has never seen pinned, with a `::warning::` on the Actions run naming its id: such a +case is missing from the site, health notices included, until the feed pins it. It is not in +the DB, so it cannot be stamped vanished. None of the 13,588 cases on the 2026-09-23 release +lacked a pin. + The first v4 build will stamp all 9,053 (verified on a copy of the release: the stamp is idempotent across builds and clears when a case returns), and `create_db` prints the count stamped on every build from now on, so the next purge is in the build log the day it happens. diff --git a/src/uisce/pipeline.py b/src/uisce/pipeline.py index a1ca5af..5494aa5 100644 --- a/src/uisce/pipeline.py +++ b/src/uisce/pipeline.py @@ -1,12 +1,15 @@ import argparse import json +import math import os import re import sqlite3 import time from collections import Counter +from contextlib import contextmanager from dataclasses import dataclass from datetime import datetime, timezone +from itertools import pairwise from pathlib import Path import requests @@ -35,6 +38,7 @@ LOCATIONIQ_REVERSE_URL = "https://us1.locationiq.com/v1/reverse" LOCATIONIQ_GEOCODE_SLEEP = 1 COORD_PRECISION = 4 # ~10 meter +COORD_COLUMNS = ("full_lat", "full_lon", "rounded_lat", "rounded_lon") USABLE_CASE_THRESHOLD_FIELDS = ["TITLE", "DESCRIPTION"] @@ -181,24 +185,55 @@ def feed_count(session): return data["count"] -def check_download_complete(features, expected): +@contextmanager +def _cases_table(db_path): + """(read-only connection, cases columns), or (None, set()) while db_path holds + no cases table (geocode_all can create the file before create_db has run).""" + if not db_path.exists(): + yield None, set() + return + conn = sqlite3.connect(f"{Path(db_path).resolve().as_uri()}?mode=ro", uri=True) + try: + columns = {row[1] for row in conn.execute("PRAGMA table_info(cases)")} + yield (conn if columns else None), columns + finally: + conn.close() + + +def unvanished_cases(db_path=DB_PATH): + """Rows load_cases would stamp vanished if the download held none of them.""" + with _cases_table(db_path) as (conn, columns): + if conn is None: + return 0 + live = " WHERE vanished_at IS NULL" if "vanished_at" in columns else "" + return conn.execute(f"SELECT COUNT(*) FROM cases{live}").fetchone()[0] + + +def check_download_complete(features, expected, unvanished=0): if len(features) < expected * (1 - FEED_COUNT_TOLERANCE): raise RuntimeError( f"downloaded {len(features)} cases but the feed reports {expected}; " "refusing to build from a truncated download" ) + # the tolerance passes 0 of 0, and an empty download stamps every stored case vanished + if unvanished and not features: + raise RuntimeError( + f"the download is empty (the feed reports {expected}) but the DB holds " + f"{unvanished} cases not yet vanished; refusing to stamp them all" + ) def download_cases(session): all_features = [] - offset = 0 + # By key, not offset: a row deleted mid-download shifts every later offset, + # and a server maxRecordCount below the page size drops rows at each page. + last_id = -1 while True: params = { - "where": "1=1", + "where": f"OBJECTID > {last_id}", "outFields": "*", "orderByFields": "OBJECTID", - "resultOffset": offset, "resultRecordCount": ARCGIS_PAGE_SIZE, "f": "json", } @@ -208,19 +243,26 @@ def download_cases(session): data = resp.json() if "error" in data: - raise RuntimeError(f"ArcGIS error at offset {offset}: {data['error']}") + raise RuntimeError(f"ArcGIS error after OBJECTID {last_id}: {data['error']}") features = data.get("features", []) if not features: break + ids = [(f.get("attributes") or {}).get("OBJECTID") for f in features] + # key paging is only sound on ids strictly above the last page's, in order + if None in ids or not all(a < b for a, b in pairwise([last_id, *ids])): + raise RuntimeError( + f"ArcGIS page after OBJECTID {last_id} is not strictly ascending: " + f"{ids[:5]}{'...' if len(ids) > 5 else ''}" + ) all_features.extend(features) print(f"Fetched {len(all_features)}") if not data.get("exceededTransferLimit", False): break - offset += ARCGIS_PAGE_SIZE + last_id = ids[-1] time.sleep(ARCGIS_PAGE_SLEEP) print(f"Done: {len(all_features)} records") @@ -264,12 +306,18 @@ def map_cases(cases_to_map): mapped_case["start_date"] = _epoch_ms_to_iso(mapped_case["start_date"]) mapped_case["end_date"] = _epoch_ms_to_iso(mapped_case["end_date"]) - lon, lat = transformer.transform(case["geometry"]["x"], case["geometry"]["y"]) - mapped_case["full_lat"] = lat - mapped_case["full_lon"] = lon - - mapped_case["rounded_lat"] = round(lat, COORD_PRECISION) - mapped_case["rounded_lon"] = round(lon, COORD_PRECISION) + # ArcGIS omits `geometry` for a null shape, or writes an empty point as + # "NaN"; restore_pins settles both + geometry = case.get("geometry") or {} + x, y = _coordinate(geometry.get("x")), _coordinate(geometry.get("y")) + if x is None or y is None: + mapped_case.update(dict.fromkeys(COORD_COLUMNS)) + else: + lon, lat = transformer.transform(x, y) + mapped_case["full_lat"] = lat + mapped_case["full_lon"] = lon + mapped_case["rounded_lat"] = round(lat, COORD_PRECISION) + mapped_case["rounded_lon"] = round(lon, COORD_PRECISION) if mapped_case["county"] == "Dnegal": mapped_case["county"] = "Donegal" @@ -284,6 +332,40 @@ def map_cases(cases_to_map): return all_cases, skipped +def _coordinate(value): + try: + value = float(value) + except (TypeError, ValueError): + return None + return value if math.isfinite(value) else None + + +def restore_pins(cases, db_path=DB_PATH): + """(cases, unplaced ids): a pinless case gets its stored pin or is set aside. + See notes/data-quality.md, "A feature with no pin" (2026-09-24).""" + missing = [c["id"] for c in cases if c["full_lat"] is None] + if not missing: + return cases, [] + stored = {} + with _cases_table(db_path) as (conn, _): + if conn is not None: + marks = ", ".join("?" * len(missing)) + rows = conn.execute( + f"SELECT id, {', '.join(COORD_COLUMNS)} FROM cases WHERE id IN ({marks})", + missing, + ) + stored = {row[0]: dict(zip(COORD_COLUMNS, row[1:])) for row in rows} + kept, unplaced = [], [] + for case in cases: + if case["full_lat"] is not None: + kept.append(case) + elif case["id"] in stored: + kept.append(case | stored[case["id"]]) + else: + unplaced.append(case["id"]) + return kept, unplaced + + def _epoch_ms_to_iso(ms): if ms is None: return None @@ -872,7 +954,8 @@ def backfill_reduced_pressure(conn): every other backfill. """ rows = conn.execute( - "SELECT id, description FROM cases WHERE description IS NOT NULL AND NOT reduced_pressure" + "SELECT id, description FROM cases " + "WHERE description IS NOT NULL AND COALESCE(reduced_pressure, 0) = 0" ).fetchall() updates = [ (case_id,) for case_id, description in rows if _PRESSURE_ONLY.search(description) @@ -940,13 +1023,18 @@ def run(skip_geocode=False): session = make_session() expected = feed_count(session) features = download_cases(session) - check_download_complete(features, expected) + check_download_complete(features, expected, unvanished_cases()) CASES_RAW_PATH.parent.mkdir(parents=True, exist_ok=True) CASES_RAW_PATH.write_text(json.dumps(features, indent=2)) mapped_cases, skipped = map_cases(read_arcgis_cases()) if skipped: print(f"Skipped {len(skipped)} cases with no usable data: {skipped}") + mapped_cases, unplaced = restore_pins(mapped_cases) + if unplaced: + # surfaced on the Actions run: these are live notices missing from the site + print(f"::warning::{len(unplaced)} cases the feed has never pinned are left out " + f"until it does: {unplaced}") CASES_MAPPED_PATH.parent.mkdir(parents=True, exist_ok=True) CASES_MAPPED_PATH.write_text(json.dumps(mapped_cases, indent=2)) diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index 08e5d8d..58b029f 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -1,5 +1,7 @@ import json +import re import sqlite3 +from contextlib import closing import pytest import requests @@ -96,6 +98,57 @@ def test_normalises_empty_strings_to_none(self): assert mapped[0]["work_type"] is None assert mapped[0]["status"] is None + def test_a_feature_with_no_geometry_maps_with_no_coordinates(self): + # ArcGIS omits `geometry` entirely for a null shape + feature = make_feature({"DESCRIPTION": "text"}) + del feature["geometry"] + mapped, skipped = map_cases([feature]) + assert skipped == [] + assert [mapped[0][c] for c in pipeline.COORD_COLUMNS] == [None] * 4 + + def test_an_empty_point_written_as_nan_maps_with_no_coordinates(self): + feature = make_feature({"DESCRIPTION": "text"}) + feature["geometry"] = {"x": "NaN", "y": "NaN"} + mapped, _ = map_cases([feature]) + assert [mapped[0][c] for c in pipeline.COORD_COLUMNS] == [None] * 4 + + +class TestRestorePins: + def _db(self, tmp_path): + db_path = tmp_path / "u.db" + cases = [case_record( + id=1, full_lat=53.123456, full_lon=-6.5, rounded_lat=53.1235, rounded_lon=-6.5)] + skip_geocoding(cases, db_path) + pipeline.create_db(cases, db_path) + return db_path + + def test_a_known_case_keeps_its_last_pin(self, tmp_path): + cases = [case_record(id=1, title="t")] + kept, unplaced = pipeline.restore_pins(cases, self._db(tmp_path)) + assert unplaced == [] + assert [kept[0][c] for c in pipeline.COORD_COLUMNS] == [53.123456, -6.5, 53.1235, -6.5] + + def test_a_case_never_pinned_is_set_aside(self, tmp_path): + pinned = case_record(id=3, full_lat=1.0, full_lon=2.0, rounded_lat=1.0, rounded_lon=2.0) + kept, unplaced = pipeline.restore_pins( + [case_record(id=2), pinned], self._db(tmp_path)) + assert unplaced == [2] + assert kept == [pinned] + + def test_no_db_sets_aside_every_unpinned_case(self, tmp_path): + kept, unplaced = pipeline.restore_pins([case_record(id=2)], tmp_path / "none.db") + assert (kept, unplaced) == ([], [2]) + assert not (tmp_path / "none.db").exists() + + def test_a_db_with_no_cases_table_yet_sets_them_aside(self, tmp_path): + # geocode_all creates the file before create_db has run + db_path = tmp_path / "u.db" + with closing(sqlite3.connect(db_path)) as conn: + conn.execute("CREATE TABLE geocode_cache (x)") + conn.commit() + assert pipeline.restore_pins([case_record(id=2)], db_path) == ([], [2]) + assert pipeline.unvanished_cases(db_path) == 0 + def test_epoch_ms_to_iso_none_passthrough(): assert _epoch_ms_to_iso(None) is None @@ -120,59 +173,151 @@ def json(self): class FakeSession: - """Serves canned ArcGIS pages keyed by resultOffset.""" + """Serves one canned ArcGIS response per request, in order.""" - def __init__(self, pages): - self.pages = pages - self.offsets_requested = [] + def __init__(self, *responses): + self.responses = list(responses) self.params = [] def get(self, url, params=None, timeout=None): - offset = params.get("resultOffset") - self.offsets_requested.append(offset) self.params.append(params) - return FakeResponse(self.pages[offset]) + return FakeResponse(self.responses.pop(0)) + + +class LiveFeed: + """An ArcGIS layer that honours `OBJECTID > n` paging. `after_first_page` + mutates it mid-download; `cap` is a server maxRecordCount below the page + size requested.""" + + def __init__(self, ids, after_first_page=None, cap=None, ordered=True): + self.ids = list(ids) + self.after_first_page = after_first_page + self.cap = cap + self.ordered = ordered + self.pages = 0 + + def get(self, url, params=None, timeout=None): + if params.get("returnCountOnly"): + return FakeResponse({"count": len(self.ids)}) + after = int(re.fullmatch(r"OBJECTID > (-?\d+)", params["where"]).group(1)) + n = min(params["resultRecordCount"], self.cap or params["resultRecordCount"]) + # storage order unless asked, and a server that ignores the ask + by_key = params.get("orderByFields") == "OBJECTID" and self.ordered + rest = [i for i in (sorted(self.ids) if by_key else self.ids) if i > after] + self.pages += 1 + if self.pages == 1 and self.after_first_page: + self.after_first_page(self) + return FakeResponse({ + "features": [{"attributes": {"OBJECTID": i}} for i in rest[:n]], + "exceededTransferLimit": len(rest) > n, + }) class TestDownloadCases: - def test_paginates_until_transfer_limit_clear(self, monkeypatch): + @pytest.fixture(autouse=True) + def _no_sleep(self, monkeypatch): monkeypatch.setattr(pipeline, "ARCGIS_PAGE_SLEEP", 0) - page_size = pipeline.ARCGIS_PAGE_SIZE - session = FakeSession( - { - 0: {"features": [{"id": 1}, {"id": 2}], "exceededTransferLimit": True}, - page_size: {"features": [{"id": 3}]}, - } - ) - features = download_cases(session) + def _ids(self, features): + return [f["attributes"]["OBJECTID"] for f in features] - assert session.offsets_requested == [0, page_size] - assert features == [{"id": 1}, {"id": 2}, {"id": 3}] + def test_pages_by_key_until_transfer_limit_clear(self, monkeypatch): + monkeypatch.setattr(pipeline, "ARCGIS_PAGE_SIZE", 2) + feed = LiveFeed([3, 7, 9]) + assert self._ids(download_cases(feed)) == [3, 7, 9] + assert feed.pages == 2 def test_stops_on_empty_page(self): - session = FakeSession({0: {"features": []}}) - assert download_cases(session) == [] + assert download_cases(FakeSession({"features": []})) == [] + + def _vanished_live(self, feed): + conn = make_cases_table(sqlite3.connect(":memory:")) + pipeline.load_cases(conn, [case_record(id=i, status="Open") for i in feed.ids], + now="2026-09-01T00:00:00+00:00") + features = download_cases(feed) + pipeline.check_download_complete(features, len(feed.ids)) + pipeline.load_cases(conn, [case_record(id=i, status="Open") + for i in self._ids(features)], + now="2026-09-02T00:00:00+00:00") + return sorted( + i for (i,) in conn.execute("SELECT id FROM cases WHERE vanished_at IS NOT NULL") + if i in set(feed.ids) + ) + + def test_a_deletion_mid_download_loses_no_live_case(self): + # by offset, deleting case 5 during page one shifted case 1001 out of page two + feed = LiveFeed(range(1, 4301), after_first_page=lambda f: f.ids.remove(5)) + assert self._vanished_live(feed) == [] + + def test_a_server_page_cap_below_the_page_size_loses_no_case(self): + # by offset, 5 rows a page fell between pages: 20 live cases, under the tolerance + assert self._vanished_live(LiveFeed(range(1, 4301), cap=995)) == [] def test_a_short_download_is_refused_before_it_touches_the_db(self): pipeline.check_download_complete([{}] * 990, 1000) with pytest.raises(RuntimeError, match="truncated"): pipeline.check_download_complete([{}] * 989, 1000) + def test_an_empty_download_is_refused_while_the_db_holds_cases(self): + with pytest.raises(RuntimeError, match="download is empty"): + pipeline.check_download_complete([], 0, unvanished=498) + pipeline.check_download_complete([], 0, unvanished=0) + + def test_a_full_download_builds_even_if_the_count_says_zero(self): + pipeline.check_download_complete([{}] * 1000, 0, unvanished=498) + + def test_an_empty_download_against_a_nonzero_count_is_a_truncation(self): + with pytest.raises(RuntimeError, match="truncated"): + pipeline.check_download_complete([], 1000, unvanished=498) + + def test_the_guard_counts_every_row_the_stamp_would_touch(self, tmp_path): + # closed rows are stamped vanished too, so they count + db_path = tmp_path / "u.db" + assert pipeline.unvanished_cases(db_path) == 0 + assert not db_path.exists() + conn = make_cases_table(sqlite3.connect(db_path)) + pipeline.load_cases(conn, [case_record(id=1, status="Open"), + case_record(id=2, status="Closed")]) + pipeline.load_cases(conn, [case_record(id=2, status="Closed"), + case_record(id=3, status="Open")]) + conn.commit() + assert pipeline.unvanished_cases(db_path) == 2 + + def test_the_download_asks_for_key_order(self, monkeypatch): + monkeypatch.setattr(pipeline, "ARCGIS_PAGE_SIZE", 2) + assert self._ids(download_cases(LiveFeed([9, 3, 7]))) == [3, 7, 9] + + def test_a_server_that_ignores_the_key_filter_is_refused(self, monkeypatch): + # re-serves the first page whatever `OBJECTID > n` says + monkeypatch.setattr(pipeline, "ARCGIS_PAGE_SIZE", 2) + served = [] + + def first_page_again(url, params=None, timeout=None): + served.append(params["where"]) + assert len(served) < 5, "kept paging over the same rows" + return FakeResponse({"features": [{"attributes": {"OBJECTID": i}} for i in (3, 7)], + "exceededTransferLimit": True}) + + feed = LiveFeed([3, 7, 9]) + feed.get = first_page_again + with pytest.raises(RuntimeError, match="strictly ascending"): + download_cases(feed) + + def test_a_page_out_of_order_is_refused(self, monkeypatch): + monkeypatch.setattr(pipeline, "ARCGIS_PAGE_SIZE", 2) + with pytest.raises(RuntimeError, match="ascending"): + download_cases(LiveFeed([9, 3, 7], ordered=False)) + def test_the_feed_count_is_read_from_the_count_endpoint(self): - session = FakeSession({None: {"count": 12097}}) + session = FakeSession({"count": 12097}) assert pipeline.feed_count(session) == 12097 assert session.params[0]["returnCountOnly"] == "true" def test_raises_on_arcgis_error_payload(self): # ArcGIS reports errors in a 200 response body, not an HTTP status - session = FakeSession({0: {"error": {"code": 400, "message": "bad"}}}) - try: + session = FakeSession({"error": {"code": 400, "message": "bad"}}) + with pytest.raises(RuntimeError, match="ArcGIS error"): download_cases(session) - except RuntimeError as e: - assert "ArcGIS error" in str(e) - else: - raise AssertionError("expected RuntimeError") class GeocodeSession: @@ -1015,3 +1160,9 @@ def test_the_flag_is_only_ever_set_never_cleared(self): conn.execute("UPDATE cases SET reduced_pressure = 1") assert pipeline.backfill_reduced_pressure(conn) == 0 assert conn.execute("SELECT reduced_pressure FROM cases").fetchone()[0] == 1 + + def test_a_null_feed_flag_is_backfilled_like_a_zero(self): + conn = self._conn("Works may cause low pressure to Ballyduff and surrounding areas") + conn.execute("UPDATE cases SET reduced_pressure = NULL") + assert pipeline.backfill_reduced_pressure(conn) == 1 + assert conn.execute("SELECT reduced_pressure FROM cases").fetchone()[0] == 1