From 2f18ea7d77bf13031e072c2b5c6834de49337f7d Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 08:25:23 +0000 Subject: [PATCH 1/3] Harden the feed download against paging, empty and pinless responses Four fixes from a review of pipeline.py, each with a regression test that failed first. download_cases pages by key, OBJECTID > , instead of by resultOffset. By offset, one case deleted mid-download shifted a live case out of the next page and stamped it vanished for that build, and a server maxRecordCount of 995 lost 5 rows a page (20 of 4,300), both under FEED_COUNT_TOLERANCE. check_download_complete refuses an empty feed while the DB holds Open cases not yet vanished. 0 of 0 passed the tolerance and would have stamped all 498 open cases on the 2026-09-23 release vanished. A partial purge is still stamped, as data-quality.md says it should be. map_cases no longer KeyErrors on a feature with no geometry, which ArcGIS omits for a null shape and which crashed every build. The coordinates are REAL NOT NULL and key the geocode cache, and dropping the constraint is not an additive migration, so restore_pins gives such a case the pin the DB last stored for it; a case never pinned is set aside and printed until the feed pins it. The DB therefore never holds NULL coordinates and geocoding and the site are unchanged. backfill_reduced_pressure reads a NULL feed flag as 0. On the release DB the 7 NULL-flag rows do not match the pattern, so no published number moves. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_0168hLW2X3mkV3Jhm26LJSwQ --- notes/data-quality.md | 10 +++ src/uisce/pipeline.py | 85 +++++++++++++++++++---- tests/test_pipeline.py | 153 +++++++++++++++++++++++++++++++++-------- 3 files changed, 207 insertions(+), 41 deletions(-) diff --git a/notes/data-quality.md b/notes/data-quality.md index 206eefb..6aab8c1 100644 --- a/notes/data-quality.md +++ b/notes/data-quality.md @@ -290,6 +290,16 @@ 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 have stamped every open case vanished (498 on the 2026-09-23 release). `run` now also +refuses a download when the feed reports 0 cases or returns none while the DB holds `Open` +cases not yet vanished. 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%. + 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..455b609 100644 --- a/src/uisce/pipeline.py +++ b/src/uisce/pipeline.py @@ -35,6 +35,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 +182,42 @@ def feed_count(session): return data["count"] -def check_download_complete(features, expected): +def live_open_cases(db_path=DB_PATH): + if not db_path.exists(): + return 0 + with sqlite3.connect(f"file:{db_path}?mode=ro", uri=True) as conn: + columns = {row[1] for row in conn.execute("PRAGMA table_info(cases)")} + if not columns: + return 0 + live = " AND vanished_at IS NULL" if "vanished_at" in columns else "" + return conn.execute(f"SELECT COUNT(*) FROM cases WHERE status = 'Open'{live}").fetchone()[0] + + +def check_download_complete(features, expected, live_open=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, which would stamp every open case vanished + if live_open and (expected == 0 or not features): + raise RuntimeError( + f"the feed is empty but the DB holds {live_open} open cases; " + "refusing to stamp them all vanished" + ) 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,7 +227,7 @@ 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: @@ -220,7 +239,7 @@ def download_cases(session): if not data.get("exceededTransferLimit", False): break - offset += ARCGIS_PAGE_SIZE + last_id = features[-1]["attributes"]["OBJECTID"] time.sleep(ARCGIS_PAGE_SLEEP) print(f"Done: {len(all_features)} records") @@ -264,12 +283,16 @@ 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; restore_pins settles those + geometry = case.get("geometry") or {} + if geometry.get("x") is None or geometry.get("y") is None: + mapped_case.update(dict.fromkeys(COORD_COLUMNS)) + else: + lon, lat = transformer.transform(geometry["x"], 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) if mapped_case["county"] == "Dnegal": mapped_case["county"] = "Donegal" @@ -284,6 +307,36 @@ def map_cases(cases_to_map): return all_cases, skipped +def restore_pins(cases, db_path=DB_PATH): + """Give a case the feed served with no geometry the pin it last had. + + The coordinates are NOT NULL and key the geocode cache, so a case never + pinned cannot be stored; it is set aside, and returned by id, until the feed + pins it. Returns (cases, unplaced_ids). + """ + missing = [c["id"] for c in cases if c["full_lat"] is None] + if not missing: + return cases, [] + stored = {} + if db_path.exists(): + with sqlite3.connect(f"file:{db_path}?mode=ro", uri=True) as conn: + 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 +925,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 +994,16 @@ 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, live_open_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: + print(f"Skipped {len(unplaced)} cases the feed has never pinned: {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 d0c186a..55f8bb6 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -1,4 +1,5 @@ import json +import re import sqlite3 import pytest @@ -96,6 +97,42 @@ 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 + + +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_epoch_ms_to_iso_none_passthrough(): assert _epoch_ms_to_iso(None) is None @@ -120,59 +157,115 @@ 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): + self.ids = list(ids) + self.after_first_page = after_first_page + self.cap = cap + 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"]) + rest = [i for i in sorted(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_feed_is_refused_while_the_db_holds_live_open_cases(self): + with pytest.raises(RuntimeError, match="empty"): + pipeline.check_download_complete([], 0, live_open=498) + pipeline.check_download_complete([], 0, live_open=0) + + def test_live_open_counts_open_cases_the_feed_still_serves(self, tmp_path): + db_path = tmp_path / "u.db" + assert pipeline.live_open_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.live_open_cases(db_path) == 1 + 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: @@ -1009,3 +1102,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 From 0d8f148115255a1cb31ae0075cc3335d8e5a5f44 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 10:07:49 +0000 Subject: [PATCH 2/3] Harden the pinless, empty and out-of-order paths found in review From a code review of this PR: - An empty point written as "NaN" passed the None check, became a nan coordinate and crashed the build on a NOT NULL column. Coordinates that are not finite numbers now count as missing. - restore_pins assumed a cases table whenever the DB file existed; geocode_all can create the file first. Both readers go through one helper that answers "no table yet". - The empty-feed guard counted open rows, but the vanish stamp touches every row not yet vanished, closed ones included (4,306 on the release, not 498). It now counts those, and its dead `not features` branch and misleading message are gone. - Key paging trusted the server's order. A page that is not ascending now fails the build, and a fake feed stored out of order makes dropping orderByFields fail a test. - A case the feed has never pinned is surfaced as a ::warning:: on the run rather than a log line, and the decision is in data-quality.md. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_0168hLW2X3mkV3Jhm26LJSwQ --- notes/data-quality.md | 25 +++++++++---- src/uisce/pipeline.py | 79 +++++++++++++++++++++++++++++------------- tests/test_pipeline.py | 44 ++++++++++++++++++----- 3 files changed, 108 insertions(+), 40 deletions(-) diff --git a/notes/data-quality.md b/notes/data-quality.md index 6de21ee..01091eb 100644 --- a/notes/data-quality.md +++ b/notes/data-quality.md @@ -295,14 +295,27 @@ changing under the paging; a real purge like the one around 2026-04-20 would sti 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 have stamped every open case vanished (498 on the 2026-09-23 release). `run` now also -refuses a download when the feed reports 0 cases or returns none while the DB holds `Open` -cases not yet vanished. 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 +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 a download when the feed reports 0 cases +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%. +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 diff --git a/src/uisce/pipeline.py b/src/uisce/pipeline.py index 455b609..91f7bf5 100644 --- a/src/uisce/pipeline.py +++ b/src/uisce/pipeline.py @@ -1,5 +1,6 @@ import argparse import json +import math import os import re import sqlite3 @@ -182,28 +183,42 @@ def feed_count(session): return data["count"] -def live_open_cases(db_path=DB_PATH): +def _read_cases_table(db_path): + """A read-only connection to db_path, or None if it holds no cases table yet + (geocode_all can create the file before create_db has run).""" if not db_path.exists(): + return None + conn = sqlite3.connect(f"file:{db_path}?mode=ro", uri=True) + if not conn.execute("PRAGMA table_info(cases)").fetchall(): + conn.close() + return None + return conn + + +def unvanished_cases(db_path=DB_PATH): + """Rows load_cases would stamp vanished if the download held none of them.""" + conn = _read_cases_table(db_path) + if conn is None: return 0 - with sqlite3.connect(f"file:{db_path}?mode=ro", uri=True) as conn: + with conn: columns = {row[1] for row in conn.execute("PRAGMA table_info(cases)")} - if not columns: - return 0 - live = " AND vanished_at IS NULL" if "vanished_at" in columns else "" - return conn.execute(f"SELECT COUNT(*) FROM cases WHERE status = 'Open'{live}").fetchone()[0] + live = " WHERE vanished_at IS NULL" if "vanished_at" in columns else "" + count = conn.execute(f"SELECT COUNT(*) FROM cases{live}").fetchone()[0] + conn.close() + return count -def check_download_complete(features, expected, live_open=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, which would stamp every open case vanished - if live_open and (expected == 0 or not features): + # the tolerance passes 0 of 0, which would stamp every stored case vanished + if unvanished and expected == 0: raise RuntimeError( - f"the feed is empty but the DB holds {live_open} open cases; " - "refusing to stamp them all vanished" + f"the feed reports 0 cases ({len(features)} downloaded) but the DB holds " + f"{unvanished} not yet vanished; refusing to stamp them all" ) @@ -233,13 +248,17 @@ def download_cases(session): if not features: break + ids = [f["attributes"]["OBJECTID"] for f in features] + # key paging is only sound on ids the server returned in order + if ids != sorted(set(ids)) or ids[0] <= last_id: + raise RuntimeError(f"ArcGIS page after OBJECTID {last_id} is not in ascending order") all_features.extend(features) print(f"Fetched {len(all_features)}") if not data.get("exceededTransferLimit", False): break - last_id = features[-1]["attributes"]["OBJECTID"] + last_id = ids[-1] time.sleep(ARCGIS_PAGE_SLEEP) print(f"Done: {len(all_features)} records") @@ -283,12 +302,14 @@ 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"]) - # ArcGIS omits `geometry` for a null shape; restore_pins settles those + # ArcGIS omits `geometry` for a null shape, or writes an empty point as + # "NaN"; restore_pins settles both geometry = case.get("geometry") or {} - if geometry.get("x") is None or geometry.get("y") is None: + 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(geometry["x"], geometry["y"]) + lon, lat = transformer.transform(x, y) mapped_case["full_lat"] = lat mapped_case["full_lon"] = lon mapped_case["rounded_lat"] = round(lat, COORD_PRECISION) @@ -307,25 +328,31 @@ def map_cases(cases_to_map): return all_cases, skipped -def restore_pins(cases, db_path=DB_PATH): - """Give a case the feed served with no geometry the pin it last had. +def _coordinate(value): + try: + value = float(value) + except (TypeError, ValueError): + return None + return value if math.isfinite(value) else None - The coordinates are NOT NULL and key the geocode cache, so a case never - pinned cannot be stored; it is set aside, and returned by id, until the feed - pins it. Returns (cases, unplaced_ids). - """ + +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 = {} - if db_path.exists(): - with sqlite3.connect(f"file:{db_path}?mode=ro", uri=True) as conn: + conn = _read_cases_table(db_path) + if conn is not None: + with conn: 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} + conn.close() kept, unplaced = [], [] for case in cases: if case["full_lat"] is not None: @@ -994,7 +1021,7 @@ def run(skip_geocode=False): session = make_session() expected = feed_count(session) features = download_cases(session) - check_download_complete(features, expected, live_open_cases()) + 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)) @@ -1003,7 +1030,9 @@ def run(skip_geocode=False): print(f"Skipped {len(skipped)} cases with no usable data: {skipped}") mapped_cases, unplaced = restore_pins(mapped_cases) if unplaced: - print(f"Skipped {len(unplaced)} cases the feed has never pinned: {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 55f8bb6..eabc100 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -105,6 +105,12 @@ def test_a_feature_with_no_geometry_maps_with_no_coordinates(self): 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): @@ -133,6 +139,13 @@ def test_no_db_sets_aside_every_unpinned_case(self, tmp_path): 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" + sqlite3.connect(db_path).execute("CREATE TABLE geocode_cache (x)").connection.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 @@ -173,10 +186,11 @@ class LiveFeed: mutates it mid-download; `cap` is a server maxRecordCount below the page size requested.""" - def __init__(self, ids, after_first_page=None, cap=None): + 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): @@ -184,7 +198,9 @@ def get(self, url, params=None, timeout=None): 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"]) - rest = [i for i in sorted(self.ids) if i > after] + # 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) @@ -239,14 +255,15 @@ def test_a_short_download_is_refused_before_it_touches_the_db(self): with pytest.raises(RuntimeError, match="truncated"): pipeline.check_download_complete([{}] * 989, 1000) - def test_an_empty_feed_is_refused_while_the_db_holds_live_open_cases(self): - with pytest.raises(RuntimeError, match="empty"): - pipeline.check_download_complete([], 0, live_open=498) - pipeline.check_download_complete([], 0, live_open=0) + def test_a_feed_reporting_nothing_is_refused_while_the_db_holds_cases(self): + with pytest.raises(RuntimeError, match="reports 0 cases"): + pipeline.check_download_complete([], 0, unvanished=498) + pipeline.check_download_complete([], 0, unvanished=0) - def test_live_open_counts_open_cases_the_feed_still_serves(self, tmp_path): + 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.live_open_cases(db_path) == 0 + 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"), @@ -254,7 +271,16 @@ def test_live_open_counts_open_cases_the_feed_still_serves(self, tmp_path): pipeline.load_cases(conn, [case_record(id=2, status="Closed"), case_record(id=3, status="Open")]) conn.commit() - assert pipeline.live_open_cases(db_path) == 1 + 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_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({"count": 12097}) From aec63577fd81249a2bd65a0c84417fc29437bbc3 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 11:34:39 +0000 Subject: [PATCH 3/3] Second-pass fixes to the ingest guards From a second code review of the review-fix commit: - Read-only connections are closed on every path (a context manager over the cases table replaces `with conn:`, which never closes), the column probe runs once, and the DB URI is built with as_uri(). - The empty-feed refusal keys on an empty download, which is what the vanish stamp acts on; a 0 count with a full download now builds. - The page check is one strictly-increasing test from the last id, so a server ignoring `OBJECTID > n` is refused rather than looped over, and a feature with no OBJECTID fails with the page named. - Tests for both refusal directions and the ignored filter; the test's own sqlite connection is closed. CLAUDE.md's vanished_at row is amended and gains a row for features with no pin. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_0168hLW2X3mkV3Jhm26LJSwQ --- CLAUDE.md | 3 ++- notes/data-quality.md | 2 +- src/uisce/pipeline.py | 58 ++++++++++++++++++++++-------------------- tests/test_pipeline.py | 32 ++++++++++++++++++++--- 4 files changed, 62 insertions(+), 33 deletions(-) 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 01091eb..179ef2c 100644 --- a/notes/data-quality.md +++ b/notes/data-quality.md @@ -296,7 +296,7 @@ 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 a download when the feed reports 0 cases +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 diff --git a/src/uisce/pipeline.py b/src/uisce/pipeline.py index 91f7bf5..5494aa5 100644 --- a/src/uisce/pipeline.py +++ b/src/uisce/pipeline.py @@ -6,8 +6,10 @@ 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 @@ -183,29 +185,28 @@ def feed_count(session): return data["count"] -def _read_cases_table(db_path): - """A read-only connection to db_path, or None if it holds no cases table yet - (geocode_all can create the file before create_db has run).""" +@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(): - return None - conn = sqlite3.connect(f"file:{db_path}?mode=ro", uri=True) - if not conn.execute("PRAGMA table_info(cases)").fetchall(): + 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() - return None - return conn def unvanished_cases(db_path=DB_PATH): """Rows load_cases would stamp vanished if the download held none of them.""" - conn = _read_cases_table(db_path) - if conn is None: - return 0 - with conn: - columns = {row[1] for row in conn.execute("PRAGMA table_info(cases)")} + 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 "" - count = conn.execute(f"SELECT COUNT(*) FROM cases{live}").fetchone()[0] - conn.close() - return count + return conn.execute(f"SELECT COUNT(*) FROM cases{live}").fetchone()[0] def check_download_complete(features, expected, unvanished=0): @@ -214,11 +215,11 @@ def check_download_complete(features, expected, unvanished=0): f"downloaded {len(features)} cases but the feed reports {expected}; " "refusing to build from a truncated download" ) - # the tolerance passes 0 of 0, which would stamp every stored case vanished - if unvanished and expected == 0: + # the tolerance passes 0 of 0, and an empty download stamps every stored case vanished + if unvanished and not features: raise RuntimeError( - f"the feed reports 0 cases ({len(features)} downloaded) but the DB holds " - f"{unvanished} not yet vanished; refusing to stamp them all" + f"the download is empty (the feed reports {expected}) but the DB holds " + f"{unvanished} cases not yet vanished; refusing to stamp them all" ) @@ -248,10 +249,13 @@ def download_cases(session): if not features: break - ids = [f["attributes"]["OBJECTID"] for f in features] - # key paging is only sound on ids the server returned in order - if ids != sorted(set(ids)) or ids[0] <= last_id: - raise RuntimeError(f"ArcGIS page after OBJECTID {last_id} is not in ascending order") + 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)}") @@ -343,16 +347,14 @@ def restore_pins(cases, db_path=DB_PATH): if not missing: return cases, [] stored = {} - conn = _read_cases_table(db_path) - if conn is not None: - with conn: + 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} - conn.close() kept, unplaced = [], [] for case in cases: if case["full_lat"] is not None: diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index b3435d6..58b029f 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -1,6 +1,7 @@ import json import re import sqlite3 +from contextlib import closing import pytest import requests @@ -142,7 +143,9 @@ def test_no_db_sets_aside_every_unpinned_case(self, tmp_path): 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" - sqlite3.connect(db_path).execute("CREATE TABLE geocode_cache (x)").connection.commit() + 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 @@ -255,11 +258,18 @@ def test_a_short_download_is_refused_before_it_touches_the_db(self): with pytest.raises(RuntimeError, match="truncated"): pipeline.check_download_complete([{}] * 989, 1000) - def test_a_feed_reporting_nothing_is_refused_while_the_db_holds_cases(self): - with pytest.raises(RuntimeError, match="reports 0 cases"): + 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" @@ -277,6 +287,22 @@ 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"):