Commit 63f759d70a
Verified · cmc
Layout: unified · split
.github/workflows/daily-pull.yml +28
| @@ -69,6 +69,10 @@ jobs: | ||
| 69 | 69 | sqlite3 -noheader -separator ' ' "$DB" \ |
| 70 | 70 | "SELECT source, COUNT(*) FROM incidents GROUP BY source ORDER BY source" \ |
| 71 | 71 | > before.txt |
| 72 | # absent until the amendment migration has run against this archive | |
| 73 | sqlite3 -noheader -separator ' ' "$DB" \ | |
| 74 | "SELECT source || '+amend', COUNT(*) FROM incident_amendments | |
| 75 | GROUP BY source ORDER BY source" >> before.txt 2>/dev/null || true | |
| 72 | 76 | fi |
| 73 | 77 | cat before.txt |
| 74 | 78 | |
| @@ -85,6 +89,9 @@ jobs: | ||
| 85 | 89 | sqlite3 -noheader -separator ' ' "$DB" \ |
| 86 | 90 | "SELECT source, COUNT(*) FROM incidents GROUP BY source ORDER BY source" \ |
| 87 | 91 | > after.txt |
| 92 | sqlite3 -noheader -separator ' ' "$DB" \ | |
| 93 | "SELECT source || '+amend', COUNT(*) FROM incident_amendments | |
| 94 | GROUP BY source ORDER BY source" >> after.txt | |
| 88 | 95 | cat after.txt |
| 89 | 96 | test -s after.txt || { echo "::error::archive is empty"; exit 1; } |
| 90 | 97 | # Keyed on FILENAME, not NR == FNR: before.txt is empty on a bootstrap |
| @@ -101,6 +108,22 @@ jobs: | ||
| 101 | 108 | exit bad |
| 102 | 109 | }' before.txt after.txt |
| 103 | 110 | |
| 111 | # An amendment identical to its original means the digest drifted | |
| 112 | # against what SQLite stores, and every run would file the same | |
| 113 | # phantom again. Cheap to check, silent and cumulative if it happens. | |
| 114 | phantom=$(sqlite3 -noheader "$DB" " | |
| 115 | SELECT COUNT(*) FROM incident_amendments a | |
| 116 | JOIN incidents o ON o.source = a.source AND o.source_key = a.source_key | |
| 117 | WHERE a.agency IS o.agency AND a.case_id IS o.case_id | |
| 118 | AND a.occurred_at IS o.occurred_at AND a.category IS o.category | |
| 119 | AND a.call_type IS o.call_type AND a.disposition IS o.disposition | |
| 120 | AND a.offense_desc IS o.offense_desc AND a.is_stop IS o.is_stop | |
| 121 | AND a.address IS o.address AND a.lat IS o.lat AND a.lon IS o.lon") | |
| 122 | if [ "$phantom" -ne 0 ]; then | |
| 123 | echo "::error::$phantom amendments are identical to their original" | |
| 124 | exit 1 | |
| 125 | fi | |
| 126 | ||
| 104 | 127 | - name: Publish archive |
| 105 | 128 | run: | |
| 106 | 129 | sqlite3 "$DB" "VACUUM;" |
| @@ -115,6 +138,11 @@ jobs: | ||
| 115 | 138 | echo |
| 116 | 139 | echo "Updated $(date -u '+%Y-%m-%d %H:%M UTC'). Schema: schema.sql." |
| 117 | 140 | echo |
| 141 | echo "incidents holds each record as first published; every later" | |
| 142 | echo "version the feed served is a row in incident_amendments" | |
| 143 | echo "($(sqlite3 -noheader "$DB" 'SELECT COUNT(*) FROM incident_amendments') so far)." | |
| 144 | echo "incidents_current is the newest version of each." | |
| 145 | echo | |
| 118 | 146 | echo '```' |
| 119 | 147 | sqlite3 -header -column "$DB" \ |
| 120 | 148 | "SELECT agency, COUNT(*) AS rows, SUM(is_stop) AS stops, |
README.nfo +20
| @@ -61,6 +61,26 @@ ARCHIVE | ||
| 61 | 61 | |
| 62 | 62 | 0 6 * * * cd /path/to/omaha-incidents && .venv/bin/python ingest.py |
| 63 | 63 | |
| 64 | AMENDMENTS | |
| 65 | agencies edit records after publishing them: a disposition changes, | |
| 66 | a case reopens, a record is withdrawn. nothing in the incidents | |
| 67 | table is ever updated, so what an agency published first stays | |
| 68 | readable. every later version the feed serves lands in | |
| 69 | incident_amendments, and incidents_current is the newest version of | |
| 70 | each record. analysis.changed_stop_outcomes() lists stops whose | |
| 71 | disposition changed after filing. | |
| 72 | ||
| 73 | a version is keyed on the hash of its payload, so a record that | |
| 74 | reverts to a payload already on file is not recorded again. this | |
| 75 | holds the set of distinct states observed, not a strict timeline. | |
| 76 | ||
| 77 | the hash has to survive a round trip through sqlite. a lon of -96 | |
| 78 | arrives from the feed as a json int and comes back out of a REAL | |
| 79 | column as -96.0, so lat, lon and is_stop are coerced before | |
| 80 | hashing. get this wrong and every affected record is filed as | |
| 81 | amended on every run, forever. the workflow fails if any amendment | |
| 82 | is byte-identical to its original. | |
| 83 | ||
| 64 | 84 | NOTES |
| 65 | 85 | all three arcgis services return utc epochs; their where-clause |
| 66 | 86 | literals do not agree (opd and council bluffs utc, sarpy central). |
analysis.py +40 −8
| @@ -1,4 +1,8 @@ | ||
| 1 | """Queries behind the dashboard: agency activity and distance to the nearest ALPR.""" | |
| 1 | """Queries behind the dashboard: agency activity and distance to the nearest ALPR. | |
| 2 | ||
| 3 | Reads incidents_current, the newest known version of each record. The originals | |
| 4 | stay in incidents and every superseded version in incident_amendments, so a | |
| 5 | disposition an agency changed after the fact is still recoverable.""" | |
| 2 | 6 | |
| 3 | 7 | import sqlite3 |
| 4 | 8 | from pathlib import Path |
| @@ -36,8 +40,8 @@ def load_incidents(conn, agencies=None, start=None, end=None, categories=None): | ||
| 36 | 40 | params.append(f"{end}T23:59:59") |
| 37 | 41 | df = pd.read_sql_query( |
| 38 | 42 | f"""SELECT source, agency, case_id, occurred_at, category, call_type, |
| 39 | disposition, offense_desc, is_stop, address, lat, lon | |
| 40 | FROM incidents WHERE {' AND '.join(where)}""", | |
| 43 | disposition, offense_desc, is_stop, address, lat, lon, amended | |
| 44 | FROM incidents_current WHERE {' AND '.join(where)}""", | |
| 41 | 45 | conn, params=params) |
| 42 | 46 | df["occurred_at"] = pd.to_datetime(df["occurred_at"]) |
| 43 | 47 | return df |
| @@ -126,21 +130,49 @@ def camera_proximity(df, cameras, bin_m=200, max_m=2000): | ||
| 126 | 130 | return g |
| 127 | 131 | |
| 128 | 132 | |
| 133 | def amendment_history(conn, source, source_key): | |
| 134 | """Every version of one record, oldest first.""" | |
| 135 | return pd.read_sql_query( | |
| 136 | """SELECT 'original' AS version, occurred_at, category, call_type, | |
| 137 | disposition, is_stop, address, first_seen AS seen_at | |
| 138 | FROM incidents WHERE source = ? AND source_key = ? | |
| 139 | UNION ALL | |
| 140 | SELECT 'amended', occurred_at, category, call_type, | |
| 141 | disposition, is_stop, address, seen_at | |
| 142 | FROM incident_amendments WHERE source = ? AND source_key = ? | |
| 143 | ORDER BY seen_at""", | |
| 144 | conn, params=[source, source_key, source, source_key]) | |
| 145 | ||
| 146 | ||
| 147 | def changed_stop_outcomes(conn): | |
| 148 | """Stops whose disposition the agency changed after first publishing it.""" | |
| 149 | return pd.read_sql_query( | |
| 150 | """SELECT o.agency, o.case_id, o.occurred_at, | |
| 151 | o.disposition AS first_published, | |
| 152 | a.disposition AS later_published, a.seen_at | |
| 153 | FROM incident_amendments a | |
| 154 | JOIN incidents o | |
| 155 | ON o.source = a.source AND o.source_key = a.source_key | |
| 156 | WHERE o.is_stop = 1 | |
| 157 | AND IFNULL(a.disposition, '') <> IFNULL(o.disposition, '') | |
| 158 | ORDER BY a.seen_at DESC""", conn) | |
| 159 | ||
| 160 | ||
| 129 | 161 | def agency_options(conn): |
| 130 | 162 | rows = conn.execute( |
| 131 | "SELECT agency, COUNT(*) FROM incidents GROUP BY agency ORDER BY 2 DESC" | |
| 132 | ).fetchall() | |
| 163 | "SELECT agency, COUNT(*) FROM incidents_current GROUP BY agency" | |
| 164 | " ORDER BY 2 DESC").fetchall() | |
| 133 | 165 | return [a for a, _ in rows] |
| 134 | 166 | |
| 135 | 167 | |
| 136 | 168 | def category_options(conn): |
| 137 | 169 | rows = conn.execute( |
| 138 | "SELECT category, COUNT(*) FROM incidents WHERE category IS NOT NULL" | |
| 139 | " GROUP BY category ORDER BY 2 DESC").fetchall() | |
| 170 | "SELECT category, COUNT(*) FROM incidents_current" | |
| 171 | " WHERE category IS NOT NULL GROUP BY category ORDER BY 2 DESC").fetchall() | |
| 140 | 172 | return [c for c, _ in rows] |
| 141 | 173 | |
| 142 | 174 | |
| 143 | 175 | def date_bounds(conn): |
| 144 | 176 | lo, hi = conn.execute( |
| 145 | "SELECT MIN(occurred_at), MAX(occurred_at) FROM incidents").fetchone() | |
| 177 | "SELECT MIN(occurred_at), MAX(occurred_at) FROM incidents_current").fetchone() | |
| 146 | 178 | return lo[:10], hi[:10] |
app.py +1
| @@ -131,6 +131,7 @@ def update_kpis(agencies, categories, start, end, _theme): | ||
| 131 | 131 | ("Vehicle stops", f"{len(stops):,}"), |
| 132 | 132 | ("Stops per day", rate), |
| 133 | 133 | ("Agencies", f"{df['agency'].nunique()}"), |
| 134 | ("Amended since filing", f"{int(df['amended'].sum()):,}"), | |
| 134 | 135 | ] |
| 135 | 136 | return [html.Div(className="kpi", children=[html.Span(v, className="kpi-value"), |
| 136 | 137 | html.Span(k, className="kpi-label")]) |
ingest.py +87 −20
| @@ -18,6 +18,7 @@ carries its own literal timezone and page size. | ||
| 18 | 18 | |
| 19 | 19 | import argparse |
| 20 | 20 | import csv |
| 21 | import hashlib | |
| 21 | 22 | import json |
| 22 | 23 | import sqlite3 |
| 23 | 24 | import ssl |
| @@ -158,20 +159,67 @@ def local_iso(epoch_ms): | ||
| 158 | 159 | return dt.astimezone(LOCAL).strftime("%Y-%m-%dT%H:%M:%S") |
| 159 | 160 | |
| 160 | 161 | |
| 162 | # The order every row tuple is built in, and the order digest() hashes. | |
| 163 | COLUMNS = ("source", "source_key", "agency", "case_id", "occurred_at", "category", | |
| 164 | "call_type", "disposition", "offense_desc", "is_stop", "address", | |
| 165 | "lat", "lon") | |
| 166 | PAYLOAD = COLUMNS[2:] # everything the feed can change; the first two are the key | |
| 167 | ||
| 168 | ||
| 169 | # Columns whose SQLite affinity rewrites what the feed sent: a JSON lon of -96 | |
| 170 | # arrives as int and comes back out of a REAL column as -96.0. Hashing the raw | |
| 171 | # value would mark such a record amended on every run forever, so coerce to the | |
| 172 | # stored representation first. | |
| 173 | REAL_FIELDS = {"lat", "lon"} | |
| 174 | INT_FIELDS = {"is_stop"} | |
| 175 | ||
| 176 | ||
| 177 | def digest(values): | |
| 178 | """Hash of a record's payload. Feeds amend records after publishing them, so | |
| 179 | this is what tells an unchanged record from a genuinely new version. Must | |
| 180 | give the same answer for a value going into the database and coming back.""" | |
| 181 | parts = [] | |
| 182 | for name, v in zip(PAYLOAD, values): | |
| 183 | if v is None: | |
| 184 | parts.append("") | |
| 185 | elif name in REAL_FIELDS: | |
| 186 | parts.append(repr(float(v))) | |
| 187 | elif name in INT_FIELDS: | |
| 188 | parts.append(str(int(v))) | |
| 189 | else: | |
| 190 | parts.append(str(v)) | |
| 191 | return hashlib.blake2b("\x1f".join(parts).encode(), digest_size=8).hexdigest() | |
| 192 | ||
| 193 | ||
| 161 | 194 | def upsert(conn, rows): |
| 162 | conn.executemany( | |
| 163 | """INSERT INTO incidents | |
| 164 | (source, source_key, agency, case_id, occurred_at, category, | |
| 165 | call_type, disposition, offense_desc, is_stop, address, lat, lon) | |
| 166 | VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?) | |
| 167 | ON CONFLICT(source, source_key) DO UPDATE SET | |
| 168 | agency=excluded.agency, case_id=excluded.case_id, | |
| 169 | occurred_at=excluded.occurred_at, category=excluded.category, | |
| 170 | call_type=excluded.call_type, disposition=excluded.disposition, | |
| 171 | offense_desc=excluded.offense_desc, is_stop=excluded.is_stop, | |
| 172 | address=excluded.address, lat=excluded.lat, lon=excluded.lon""", | |
| 173 | rows) | |
| 174 | return len(rows) | |
| 195 | """Insert records not seen before; file a changed record as an amendment. | |
| 196 | ||
| 197 | Nothing in incidents is ever updated. A record whose payload differs from | |
| 198 | the one on file is appended to incident_amendments, so the version the | |
| 199 | agency published first stays readable next to what it published later.""" | |
| 200 | now = datetime.now(LOCAL).strftime("%Y-%m-%dT%H:%M:%S") | |
| 201 | marks = ",".join("?" * (len(COLUMNS) + 1)) | |
| 202 | staged = [r + (digest(r[2:]),) for r in rows] | |
| 203 | ||
| 204 | conn.execute("DROP TABLE IF EXISTS temp.incoming") | |
| 205 | conn.execute(f"CREATE TEMP TABLE incoming ({','.join(COLUMNS)}, digest)") | |
| 206 | conn.executemany(f"INSERT INTO temp.incoming VALUES ({marks})", staged) | |
| 207 | conn.execute("CREATE INDEX temp.incoming_key ON incoming (source, source_key)") | |
| 208 | ||
| 209 | cols = ",".join(COLUMNS) | |
| 210 | conn.execute( | |
| 211 | f"""INSERT OR IGNORE INTO incidents ({cols}, digest, first_seen) | |
| 212 | SELECT {cols}, digest, ? FROM temp.incoming""", (now,)) | |
| 213 | amended = conn.execute( | |
| 214 | f"""INSERT OR IGNORE INTO incident_amendments ({cols}, digest, seen_at) | |
| 215 | SELECT i.{', i.'.join(COLUMNS)}, i.digest, ? | |
| 216 | FROM temp.incoming i | |
| 217 | JOIN incidents o | |
| 218 | ON o.source = i.source AND o.source_key = i.source_key | |
| 219 | WHERE o.digest <> i.digest""", (now,)).rowcount | |
| 220 | ||
| 221 | conn.execute("DROP TABLE temp.incoming") | |
| 222 | return len(rows), amended | |
| 175 | 223 | |
| 176 | 224 | |
| 177 | 225 | def ingest_opd(conn, since): |
| @@ -206,10 +254,10 @@ def ingest_sarpy(conn, since): | ||
| 206 | 254 | a.get("StatuteDesc"), |
| 207 | 255 | int(a.get("Category") == SARPY_STOP_CATEGORY), |
| 208 | 256 | a.get("BlkAddress"), g.get("y"), g.get("x"))) |
| 209 | n = upsert(conn, rows) | |
| 257 | result = upsert(conn, rows) | |
| 210 | 258 | if unmapped: |
| 211 | 259 | print(f" unmapped IncidentId prefixes: {sorted(unmapped)}") |
| 212 | return n | |
| 260 | return result | |
| 213 | 261 | |
| 214 | 262 | |
| 215 | 263 | def ingest_cbpd(conn, since): |
| @@ -278,20 +326,38 @@ def ingest_alpr(conn, _since): | ||
| 278 | 326 | direction=excluded.direction, last_seen=excluded.last_seen, |
| 279 | 327 | tags=excluded.tags""", |
| 280 | 328 | rows) |
| 281 | return len(rows) | |
| 329 | return len(rows), 0 | |
| 282 | 330 | |
| 283 | 331 | |
| 284 | 332 | def migrate(conn): |
| 285 | 333 | """Add columns introduced after a database was first built. Runs before |
| 286 | schema.sql so its indexes can reference the new columns.""" | |
| 334 | schema.sql so its views and indexes can reference the new columns.""" | |
| 287 | 335 | have = {r[1] for r in conn.execute("PRAGMA table_info(incidents)")} |
| 288 | if have and "is_stop" not in have: | |
| 336 | if not have: | |
| 337 | return | |
| 338 | ||
| 339 | if "is_stop" not in have: | |
| 289 | 340 | conn.execute("ALTER TABLE incidents ADD COLUMN is_stop INTEGER NOT NULL" |
| 290 | 341 | " DEFAULT 0") |
| 291 | 342 | conn.execute("UPDATE incidents SET is_stop = 1 WHERE category = ?", |
| 292 | 343 | (SARPY_STOP_CATEGORY,)) |
| 293 | 344 | conn.commit() |
| 294 | 345 | |
| 346 | if "digest" not in have: | |
| 347 | # Backfill from the stored values, using the same function ingest uses, | |
| 348 | # so the next run sees the existing rows as unchanged rather than | |
| 349 | # amending all of them. | |
| 350 | conn.execute("ALTER TABLE incidents ADD COLUMN digest TEXT") | |
| 351 | conn.execute("ALTER TABLE incidents ADD COLUMN first_seen TEXT") | |
| 352 | payload = ", ".join(PAYLOAD) | |
| 353 | rows = conn.execute( | |
| 354 | f"SELECT source, source_key, {payload} FROM incidents").fetchall() | |
| 355 | conn.executemany( | |
| 356 | "UPDATE incidents SET digest = ? WHERE source = ? AND source_key = ?", | |
| 357 | [(digest(r[2:]), r[0], r[1]) for r in rows]) | |
| 358 | conn.commit() | |
| 359 | print(f" migrated: digested {len(rows)} existing rows") | |
| 360 | ||
| 295 | 361 | |
| 296 | 362 | SOURCES = { |
| 297 | 363 | "opd": ingest_opd, |
| @@ -324,9 +390,10 @@ def main(): | ||
| 324 | 390 | for name in sources: |
| 325 | 391 | start = time.monotonic() |
| 326 | 392 | print(f" {name}: pulling...") |
| 327 | n = SOURCES[name](conn, since) | |
| 393 | seen, amended = SOURCES[name](conn, since) | |
| 328 | 394 | conn.commit() |
| 329 | print(f" {name}: {n} rows in {time.monotonic() - start:.1f}s") | |
| 395 | note = f", {amended} amended" if amended else "" | |
| 396 | print(f" {name}: {seen} rows{note} in {time.monotonic() - start:.1f}s") | |
| 330 | 397 | conn.close() |
| 331 | 398 | |
| 332 | 399 | |
schema.sql +56
| @@ -1,5 +1,8 @@ | ||
| 1 | -- Incidents as first observed. Rows here are never updated: when a feed serves a | |
| 2 | -- changed version of a record it goes to incident_amendments instead, so the | |
| 3 | -- original survives a reclassification, a reopened case or a withdrawn record. | |
| 4 | -- occurred_at is local time (America/Chicago); the services return UTC epochs | |
| 5 | -- and ingest.py converts on the way in. | |
| 1 | 6 | CREATE TABLE IF NOT EXISTS incidents ( |
| 2 | 7 | source TEXT NOT NULL, -- opd | sarpy | cbpd | opd_csv |
| 3 | 8 | source_key TEXT NOT NULL, -- PK (opd) | IncidentId (sarpy) | cfs_number (cbpd) |
| @@ -14,6 +17,8 @@ CREATE TABLE IF NOT EXISTS incidents ( | ||
| 14 | 17 | address TEXT, |
| 15 | 18 | lat REAL, |
| 16 | 19 | lon REAL, |
| 20 | digest TEXT NOT NULL, -- hash of the payload, for change detection | |
| 21 | first_seen TEXT, -- when ingest first saw it; NULL if pre-dating | |
| 17 | 22 | PRIMARY KEY (source, source_key) |
| 18 | 23 | ); |
| 19 | 24 | |
| @@ -22,6 +27,55 @@ CREATE INDEX IF NOT EXISTS incidents_agency ON incidents (agency, occurred_at) | ||
| 22 | 27 | CREATE INDEX IF NOT EXISTS incidents_category ON incidents (category); |
| 23 | 28 | CREATE INDEX IF NOT EXISTS incidents_stop ON incidents (is_stop, occurred_at); |
| 24 | 29 | |
| 30 | -- Every distinct later version of a record, one row per version. Keyed on the | |
| 31 | -- payload digest, so a version is stored once no matter how many runs serve it. | |
| 32 | -- A record that reverts to a payload already on file is therefore not recorded | |
| 33 | -- again: this holds the set of distinct states observed, not a strict timeline. | |
| 34 | CREATE TABLE IF NOT EXISTS incident_amendments ( | |
| 35 | source TEXT NOT NULL, | |
| 36 | source_key TEXT NOT NULL, | |
| 37 | agency TEXT NOT NULL, | |
| 38 | case_id TEXT, | |
| 39 | occurred_at TEXT NOT NULL, | |
| 40 | category TEXT, | |
| 41 | call_type TEXT, | |
| 42 | disposition TEXT, | |
| 43 | offense_desc TEXT, | |
| 44 | is_stop INTEGER NOT NULL DEFAULT 0, | |
| 45 | address TEXT, | |
| 46 | lat REAL, | |
| 47 | lon REAL, | |
| 48 | digest TEXT NOT NULL, | |
| 49 | seen_at TEXT NOT NULL, -- when this version was first observed | |
| 50 | PRIMARY KEY (source, source_key, digest) | |
| 51 | ); | |
| 52 | ||
| 53 | CREATE INDEX IF NOT EXISTS amendments_key ON incident_amendments (source, source_key); | |
| 54 | CREATE INDEX IF NOT EXISTS amendments_seen ON incident_amendments (seen_at); | |
| 55 | ||
| 56 | -- The newest known version of each record: its latest amendment, or the | |
| 57 | -- original where a record has never been amended. | |
| 58 | CREATE VIEW IF NOT EXISTS incidents_current AS | |
| 59 | SELECT source, source_key, agency, case_id, occurred_at, category, call_type, | |
| 60 | disposition, offense_desc, is_stop, address, lat, lon, observed_at, | |
| 61 | amended | |
| 62 | FROM ( | |
| 63 | SELECT *, ROW_NUMBER() OVER (PARTITION BY source, source_key | |
| 64 | ORDER BY amended DESC, observed_at DESC) AS rn | |
| 65 | FROM ( | |
| 66 | SELECT source, source_key, agency, case_id, occurred_at, category, | |
| 67 | call_type, disposition, offense_desc, is_stop, address, lat, lon, | |
| 68 | first_seen AS observed_at, 0 AS amended | |
| 69 | FROM incidents | |
| 70 | UNION ALL | |
| 71 | SELECT source, source_key, agency, case_id, occurred_at, category, | |
| 72 | call_type, disposition, offense_desc, is_stop, address, lat, lon, | |
| 73 | seen_at AS observed_at, 1 AS amended | |
| 74 | FROM incident_amendments | |
| 75 | ) | |
| 76 | ) | |
| 77 | WHERE rn = 1; | |
| 78 | ||
| 25 | 79 | -- ALPR cameras from OpenStreetMap (ODbL). first_seen/last_seen track when a node |
| 26 | 80 | -- entered and was last present in the Overpass result, so cameras that appear or |
| 27 | 81 | -- are removed are visible over time. |