Commit 137b3438e9
Verified · cmc
Layout: unified · split
.github/workflows/daily-pull.yml +22 −5
| @@ -5,9 +5,11 @@ name: daily pull | |||
| 5 | 5 | ||
| 6 | on: | 6 | on: |
| 7 | schedule: | 7 | schedule: |
| 8 | # 06:00 America/Chicago in summer, 05:00 in winter. Offset from the hour | 8 | # Twice a day, because the cadence sets how much the rolling feeds drop |
| 9 | # because GitHub drops on-the-hour scheduled runs under load. | 9 | # before it is captured: a five-hour gap cost 13 Sarpy records once. |
| 10 | # Offset from the hour, GitHub drops on-the-hour runs under load. | ||
| 10 | - cron: "17 11 * * *" | 11 | - cron: "17 11 * * *" |
| 12 | - cron: "17 23 * * *" | ||
| 11 | workflow_dispatch: | 13 | workflow_dispatch: |
| 12 | inputs: | 14 | inputs: |
| 13 | bootstrap: | 15 | bootstrap: |
| @@ -34,7 +36,7 @@ env: | |||
| 34 | jobs: | 36 | jobs: |
| 35 | pull: | 37 | pull: |
| 36 | runs-on: ubuntu-latest | 38 | runs-on: ubuntu-latest |
| 37 | timeout-minutes: 30 | 39 | timeout-minutes: 45 |
| 38 | 40 | ||
| 39 | steps: | 41 | steps: |
| 40 | - uses: actions/checkout@v4 | 42 | - uses: actions/checkout@v4 |
| @@ -73,12 +75,22 @@ jobs: | |||
| 73 | sqlite3 -noheader -separator ' ' "$DB" \ | 75 | sqlite3 -noheader -separator ' ' "$DB" \ |
| 74 | "SELECT source || '+amend', COUNT(*) FROM incident_amendments | 76 | "SELECT source || '+amend', COUNT(*) FROM incident_amendments |
| 75 | GROUP BY source ORDER BY source" >> before.txt 2>/dev/null || true | 77 | GROUP BY source ORDER BY source" >> before.txt 2>/dev/null || true |
| 78 | sqlite3 -noheader -separator ' ' "$DB" \ | ||
| 79 | "SELECT source || '+raw', COUNT(*) FROM raw_records | ||
| 80 | GROUP BY source ORDER BY source" >> before.txt 2>/dev/null || true | ||
| 81 | sqlite3 -noheader -separator ' ' "$DB" \ | ||
| 82 | "SELECT source || '+raw', COUNT(*) FROM raw_records | ||
| 83 | GROUP BY source ORDER BY source" >> before.txt 2>/dev/null || true | ||
| 76 | fi | 84 | fi |
| 77 | cat before.txt | 85 | cat before.txt |
| 78 | 86 | ||
| 79 | - name: Pull feeds | 87 | - name: Pull feeds |
| 80 | run: | | 88 | run: | |
| 81 | if [ "${{ inputs.full }}" = "true" ] || [ "${{ inputs.bootstrap }}" = "true" ]; then | 89 | # A 30-day window cannot see an agency amending a record filed months |
| 90 | # ago, and OPD does exactly that, so sweep the whole feed on Sundays. | ||
| 91 | if [ "${{ inputs.full }}" = "true" ] || [ "${{ inputs.bootstrap }}" = "true" ] \ | ||
| 92 | || [ "$(date -u +%u)" = "7" ]; then | ||
| 93 | echo "full sweep" | ||
| 82 | python ingest.py --full | 94 | python ingest.py --full |
| 83 | else | 95 | else |
| 84 | python ingest.py | 96 | python ingest.py |
| @@ -92,6 +104,9 @@ jobs: | |||
| 92 | sqlite3 -noheader -separator ' ' "$DB" \ | 104 | sqlite3 -noheader -separator ' ' "$DB" \ |
| 93 | "SELECT source || '+amend', COUNT(*) FROM incident_amendments | 105 | "SELECT source || '+amend', COUNT(*) FROM incident_amendments |
| 94 | GROUP BY source ORDER BY source" >> after.txt | 106 | GROUP BY source ORDER BY source" >> after.txt |
| 107 | sqlite3 -noheader -separator ' ' "$DB" \ | ||
| 108 | "SELECT source || '+raw', COUNT(*) FROM raw_records | ||
| 109 | GROUP BY source ORDER BY source" >> after.txt | ||
| 95 | cat after.txt | 110 | cat after.txt |
| 96 | test -s after.txt || { echo "::error::archive is empty"; exit 1; } | 111 | test -s after.txt || { echo "::error::archive is empty"; exit 1; } |
| 97 | # Keyed on FILENAME, not NR == FNR: before.txt is empty on a bootstrap | 112 | # Keyed on FILENAME, not NR == FNR: before.txt is empty on a bootstrap |
| @@ -141,7 +156,9 @@ jobs: | |||
| 141 | echo "incidents holds each record as first published; every later" | 156 | echo "incidents holds each record as first published; every later" |
| 142 | echo "version the feed served is a row in incident_amendments" | 157 | echo "version the feed served is a row in incident_amendments" |
| 143 | echo "($(sqlite3 -noheader "$DB" 'SELECT COUNT(*) FROM incident_amendments') so far)." | 158 | echo "($(sqlite3 -noheader "$DB" 'SELECT COUNT(*) FROM incident_amendments') so far)." |
| 144 | echo "incidents_current is the newest version of each." | 159 | echo "incidents_current is the newest version of each, and" |
| 160 | echo "raw_records keeps the feed's own JSON for every version" | ||
| 161 | echo "so a parse can be redone against what actually arrived." | ||
| 145 | echo | 162 | echo |
| 146 | echo '```' | 163 | echo '```' |
| 147 | sqlite3 -header -column "$DB" \ | 164 | sqlite3 -header -column "$DB" \ |
README.nfo +18 −1
| @@ -40,13 +40,17 @@ USE | |||
| 40 | .venv/bin/python app.py | 40 | .venv/bin/python app.py |
| 41 | 41 | ||
| 42 | ARCHIVE | 42 | ARCHIVE |
| 43 | .github/workflows/daily-pull.yml runs the pull at 11:17 utc and | 43 | .github/workflows/daily-pull.yml runs at 11:17 and 23:17 utc and |
| 44 | keeps the database as metro.db.gz on the "archive" release, so the | 44 | keeps the database as metro.db.gz on the "archive" release, so the |
| 45 | archive does not depend on any one machine. each run restores that | 45 | archive does not depend on any one machine. each run restores that |
| 46 | asset, pulls, refuses to publish if any source came back with fewer | 46 | asset, pulls, refuses to publish if any source came back with fewer |
| 47 | rows than it started with, then uploads and fails loudly if a feed | 47 | rows than it started with, then uploads and fails loudly if a feed |
| 48 | has not moved in seven days. | 48 | has not moved in seven days. |
| 49 | 49 | ||
| 50 | sundays it sweeps every feed in full instead of the last 30 days, | ||
| 51 | because a 30-day window cannot see an agency amending a record it | ||
| 52 | filed months ago, and omaha does that. | ||
| 53 | |||
| 50 | first run: trigger it manually with bootstrap enabled, which pulls | 54 | first run: trigger it manually with bootstrap enabled, which pulls |
| 51 | every feed in full and creates the release. after that the restore | 55 | every feed in full and creates the release. after that the restore |
| 52 | step is mandatory -- a bootstrap over a live archive throws away | 56 | step is mandatory -- a bootstrap over a live archive throws away |
| @@ -61,6 +65,19 @@ ARCHIVE | |||
| 61 | 65 | ||
| 62 | 0 6 * * * cd /path/to/omaha-incidents && .venv/bin/python ingest.py | 66 | 0 6 * * * cd /path/to/omaha-incidents && .venv/bin/python ingest.py |
| 63 | 67 | ||
| 68 | RAW | ||
| 69 | raw_records keeps the feed's own json for every version of every | ||
| 70 | record, keyed the same way amendments are. a parse that turns out | ||
| 71 | wrong, or a field a feed adds later, can only be applied to history | ||
| 72 | if the bytes were kept, and the rolling feeds mean there is no | ||
| 73 | second chance to fetch them. the payloads already carry fields | ||
| 74 | ingest does not map: council bluffs response times and priority, | ||
| 75 | sarpy case status. | ||
| 76 | |||
| 77 | it costs about 0.7 mb gzipped a day and roughly triples the | ||
| 78 | database: 20 mb published without it, 52 mb with. 319 records | ||
| 79 | predate it and their raw is gone; the feeds no longer serve them. | ||
| 80 | |||
| 64 | AMENDMENTS | 81 | AMENDMENTS |
| 65 | agencies edit records after publishing them: a disposition changes, | 82 | agencies edit records after publishing them: a disposition changes, |
| 66 | a case reopens, a record is withdrawn. nothing in the incidents | 83 | a case reopens, a record is withdrawn. nothing in the incidents |
ingest.py +39 −24
| @@ -194,15 +194,18 @@ def digest(values): | |||
| 194 | def upsert(conn, rows): | 194 | def upsert(conn, rows): |
| 195 | """Insert records not seen before; file a changed record as an amendment. | 195 | """Insert records not seen before; file a changed record as an amendment. |
| 196 | 196 | ||
| 197 | Nothing in incidents is ever updated. A record whose payload differs from | 197 | Takes (values, raw) pairs, where raw is the feature exactly as the feed |
| 198 | the one on file is appended to incident_amendments, so the version the | 198 | served it. Nothing in incidents is ever updated: a record whose payload |
| 199 | agency published first stays readable next to what it published later.""" | 199 | differs from the one on file is appended to incident_amendments, so the |
| 200 | version the agency published first stays readable next to what it published | ||
| 201 | later. Every version's raw payload is kept too, so a parse can be redone | ||
| 202 | against what actually arrived.""" | ||
| 200 | now = datetime.now(LOCAL).strftime("%Y-%m-%dT%H:%M:%S") | 203 | now = datetime.now(LOCAL).strftime("%Y-%m-%dT%H:%M:%S") |
| 201 | marks = ",".join("?" * (len(COLUMNS) + 1)) | 204 | marks = ",".join("?" * (len(COLUMNS) + 2)) |
| 202 | staged = [r + (digest(r[2:]),) for r in rows] | 205 | staged = [r + (digest(r[2:]), raw) for r, raw in rows] |
| 203 | 206 | ||
| 204 | conn.execute("DROP TABLE IF EXISTS temp.incoming") | 207 | conn.execute("DROP TABLE IF EXISTS temp.incoming") |
| 205 | conn.execute(f"CREATE TEMP TABLE incoming ({','.join(COLUMNS)}, digest)") | 208 | conn.execute(f"CREATE TEMP TABLE incoming ({','.join(COLUMNS)}, digest, raw)") |
| 206 | conn.executemany(f"INSERT INTO temp.incoming VALUES ({marks})", staged) | 209 | conn.executemany(f"INSERT INTO temp.incoming VALUES ({marks})", staged) |
| 207 | conn.execute("CREATE INDEX temp.incoming_key ON incoming (source, source_key)") | 210 | conn.execute("CREATE INDEX temp.incoming_key ON incoming (source, source_key)") |
| 208 | 211 | ||
| @@ -218,6 +221,14 @@ def upsert(conn, rows): | |||
| 218 | ON o.source = i.source AND o.source_key = i.source_key | 221 | ON o.source = i.source AND o.source_key = i.source_key |
| 219 | WHERE o.digest <> i.digest""", (now,)).rowcount | 222 | WHERE o.digest <> i.digest""", (now,)).rowcount |
| 220 | 223 | ||
| 224 | # OR IGNORE keyed on the version, so a run that re-serves a known record | ||
| 225 | # stores nothing and the first full run backfills whatever is still served. | ||
| 226 | conn.execute( | ||
| 227 | """INSERT OR IGNORE INTO raw_records (source, source_key, digest, | ||
| 228 | fetched_at, payload) | ||
| 229 | SELECT source, source_key, digest, ?, raw FROM temp.incoming | ||
| 230 | WHERE raw IS NOT NULL""", (now,)) | ||
| 231 | |||
| 221 | conn.execute("DROP TABLE temp.incoming") | 232 | conn.execute("DROP TABLE temp.incoming") |
| 222 | return len(rows), amended | 233 | return len(rows), amended |
| 223 | 234 | ||
| @@ -229,9 +240,10 @@ def ingest_opd(conn, since): | |||
| 229 | occurred = local_iso(a["dteMidpoint"]) | 240 | occurred = local_iso(a["dteMidpoint"]) |
| 230 | if occurred is None: | 241 | if occurred is None: |
| 231 | continue | 242 | continue |
| 232 | rows.append(("opd", str(a["PK"]), "Omaha PD", a.get("RB"), occurred, | 243 | rows.append((("opd", str(a["PK"]), "Omaha PD", a.get("RB"), occurred, |
| 233 | a.get("NIBRSCategory"), None, None, None, 0, | 244 | a.get("NIBRSCategory"), None, None, None, 0, |
| 234 | a.get("AddressBlock"), a.get("LatBlock"), a.get("LonBlock"))) | 245 | a.get("AddressBlock"), a.get("LatBlock"), a.get("LonBlock")), |
| 246 | json.dumps(f, sort_keys=True))) | ||
| 235 | return upsert(conn, rows) | 247 | return upsert(conn, rows) |
| 236 | 248 | ||
| 237 | 249 | ||
| @@ -249,11 +261,12 @@ def ingest_sarpy(conn, since): | |||
| 249 | unmapped.add(prefix) | 261 | unmapped.add(prefix) |
| 250 | agency = f"Unmapped {prefix}" | 262 | agency = f"Unmapped {prefix}" |
| 251 | g = f.get("geometry") or {} | 263 | g = f.get("geometry") or {} |
| 252 | rows.append(("sarpy", iid, agency, iid, occurred, a.get("Category"), | 264 | rows.append((("sarpy", iid, agency, iid, occurred, a.get("Category"), |
| 253 | a.get("CadTypeDesc"), a.get("CadDisposition"), | 265 | a.get("CadTypeDesc"), a.get("CadDisposition"), |
| 254 | a.get("StatuteDesc"), | 266 | a.get("StatuteDesc"), |
| 255 | int(a.get("Category") == SARPY_STOP_CATEGORY), | 267 | int(a.get("Category") == SARPY_STOP_CATEGORY), |
| 256 | a.get("BlkAddress"), g.get("y"), g.get("x"))) | 268 | a.get("BlkAddress"), g.get("y"), g.get("x")), |
| 269 | json.dumps(f, sort_keys=True))) | ||
| 257 | result = upsert(conn, rows) | 270 | result = upsert(conn, rows) |
| 258 | if unmapped: | 271 | if unmapped: |
| 259 | print(f" unmapped IncidentId prefixes: {sorted(unmapped)}") | 272 | print(f" unmapped IncidentId prefixes: {sorted(unmapped)}") |
| @@ -270,11 +283,12 @@ def ingest_cbpd(conn, since): | |||
| 270 | g = f.get("geometry") or {} | 283 | g = f.get("geometry") or {} |
| 271 | code = a.get("incident_code") | 284 | code = a.get("incident_code") |
| 272 | # Council Bluffs withholds the street address; the point is still exact. | 285 | # Council Bluffs withholds the street address; the point is still exact. |
| 273 | rows.append(("cbpd", a["cfs_number"], "Council Bluffs PD", | 286 | rows.append((("cbpd", a["cfs_number"], "Council Bluffs PD", |
| 274 | a.get("case_number") or a.get("cfs_number"), occurred, | 287 | a.get("case_number") or a.get("cfs_number"), occurred, |
| 275 | a.get("incident_category"), code, a.get("disp_code"), None, | 288 | a.get("incident_category"), code, a.get("disp_code"), None, |
| 276 | int(code == CBPD_STOP_CODE), | 289 | int(code == CBPD_STOP_CODE), |
| 277 | None, g.get("y"), g.get("x"))) | 290 | None, g.get("y"), g.get("x")), |
| 291 | json.dumps(f, sort_keys=True))) | ||
| 278 | return upsert(conn, rows) | 292 | return upsert(conn, rows) |
| 279 | 293 | ||
| 280 | 294 | ||
| @@ -296,11 +310,12 @@ def ingest_opd_csv(conn, _since): | |||
| 296 | when = datetime.strptime(f"{date} {tm}", "%m/%d/%Y %H:%M") | 310 | when = datetime.strptime(f"{date} {tm}", "%m/%d/%Y %H:%M") |
| 297 | except ValueError: | 311 | except ValueError: |
| 298 | continue | 312 | continue |
| 299 | rows.append(("opd_csv", f"{path.stem}:{i}", "Omaha PD", rb, | 313 | rows.append((("opd_csv", f"{path.stem}:{i}", "Omaha PD", rb, |
| 300 | when.strftime("%Y-%m-%dT%H:%M:%S"), None, None, None, | 314 | when.strftime("%Y-%m-%dT%H:%M:%S"), None, None, |
| 301 | desc, 0, loc, | 315 | None, desc, 0, loc, |
| 302 | float(lat) if lat else None, | 316 | float(lat) if lat else None, |
| 303 | float(lon) if lon else None)) | 317 | float(lon) if lon else None), |
| 318 | json.dumps(r))) | ||
| 304 | return upsert(conn, rows) | 319 | return upsert(conn, rows) |
| 305 | 320 | ||
| 306 | 321 | ||
schema.sql +13
| @@ -76,6 +76,19 @@ FROM ( | |||
| 76 | ) | 76 | ) |
| 77 | WHERE rn = 1; | 77 | WHERE rn = 1; |
| 78 | 78 | ||
| 79 | -- What the feed actually served, one row per version, keyed the same way | ||
| 80 | -- incident_amendments is. A parse that turns out wrong or a field a feed adds | ||
| 81 | -- later can only be applied to history if the bytes were kept, and the feeds | ||
| 82 | -- age out, so there is no second chance to fetch them. | ||
| 83 | CREATE TABLE IF NOT EXISTS raw_records ( | ||
| 84 | source TEXT NOT NULL, | ||
| 85 | source_key TEXT NOT NULL, | ||
| 86 | digest TEXT NOT NULL, -- the version of the record this payload produced | ||
| 87 | fetched_at TEXT NOT NULL, | ||
| 88 | payload TEXT NOT NULL, -- the feature object as served, JSON | ||
| 89 | PRIMARY KEY (source, source_key, digest) | ||
| 90 | ); | ||
| 91 | |||
| 79 | -- ALPR cameras from OpenStreetMap (ODbL). first_seen/last_seen track when a node | 92 | -- ALPR cameras from OpenStreetMap (ODbL). first_seen/last_seen track when a node |
| 80 | -- entered and was last present in the Overpass result, so cameras that appear or | 93 | -- entered and was last present in the Overpass result, so cameras that appear or |
| 81 | -- are removed are visible over time. | 94 | -- are removed are visible over time. |