Public-data ingestion
How NYC agency data (HPD, DOB, ECB, OATH, 311, DOF) gets into Complied's own database, and why nothing in the app calls the city's APIs at browse time. For engineers working on domain/src/ingestion/, sync-orchestrator, or the ingestion CLI.
Principle
Agency data is stored locally and synced on a schedule. The map, the violations explorer and project intake read Postgres; they never live-fetch Socrata. Live calls are reserved for user-initiated interactions such as geocoding an address.
Three tables
| Table | Question it answers | Notes |
|---|---|---|
buildings | Who and where is this building? | One row per BIN, citywide. Address, BIN, BBL, coordinates, year built, unit count, plus denormalized compliance counts for map pins. |
public_events | What happened? | One row per agency violation, complaint or event, keyed (agency, source_id). Dates, status_norm, order numbers, and the raw agency payload in raw_data. |
tenant_buildings | Who cares? | A tenant's portfolio: "this tenant tracks this building for this client". |
There is no separate store for tracked buildings' events. "Tracked" and "prospect" are the same rows, distinguished by a tenant_buildings join. Individual violation rows never live on buildings; only summary counts do (a number on a map pin comes from buildings, anything a project or the deadline engine needs comes from public_events).
Supporting tables:
| Table | Role |
|---|---|
unlinked_events | Quarantine for rows that could not be tied to a building. Unique on (agency, source_id); carries the attempted BIN/BBL/address and a reason. |
sync_runs | One row per feed invocation: rows fetched/upserted/quarantined/tombstoned, outcome (success, partial, failure), trigger (cron, manual). |
feed_sync_state | Per-feed position: the forward delta_watermark and the backward backfill cursor. |
buildings and public_events carry a guard (ingestion_guard schema, see 20260929090000_protect_ingestion_tables.sql) that blocks DELETE, TRUNCATE and dropping the tables or their columns unless a transaction first sets complied.allow_ingestion_table_deletes = 'on'. Ingestion only upserts.
The five layers
| Layer | What | Cadence | Runs in |
|---|---|---|---|
| 1. Buildings identity | buildings from NYC Building Footprints (BIN, outlines) joined to PLUTO on BBL (year built, unit count) | Manual | Local CLI over a direct Postgres connection |
| 2. Citywide delta | public_events for every feed, changed rows only | Nightly 08:10 UTC, one request per feed | sync-orchestrator edge function, driven by pg_cron |
| 2 backfill | Historical rows, walking backward in date windows | Manual | ingest CLI |
| 2b. Summary refresh | Compliance counts on buildings, map grid, snapshot stamp | After every nightly delta | RPCs called by the orchestrator |
| 5. Obligations | Deadlines computed from public_events and building facts | On read | @complied/domain in the browser; no ingest job |
Layer 1 needs both datasets. PLUTO is organised by tax lot and carries no BIN; one lot can hold several buildings. Footprints supplies the BIN and geometry, PLUTO supplies year built (the pre-1978 lead-paint cutoff cannot be answered without it). The whole citywide set is loaded (about 1.1M rows), which is also what lets address-only feeds (311, OATH) resolve at all. Identity loading is too large for an edge function's time and memory budget, so it runs through db-tests/scripts/run-sync-orchestrator-locally.ts identity.
Layer 2 uses Socrata's hidden :updated_at field where a publisher updates incrementally ($where=:updated_at > '<watermark>'), instead of asking each building "anything new?". Some datasets are replaced wholesale nightly (DOB complaints reports ~100% changed daily); a feed like that is pulled in full each run and is safe to tombstone against. Check a dataset's :updated_at behaviour (changed-row count against total) before relying on it for a new feed.
Layer 2b runs refresh_building_compliance_summary() to write open/overdue counts onto buildings, refresh_map_grid() to rebuild the map's aggregate grid, and mark_map_stale() to tell the map exporter the data changed (see map.md). Order matters: the grid derives from the summary columns.
Adopting a building inserts a tenant_buildings row (trackBuilding in src/data/tenantBuildings.ts). The building row already exists, so adopting adds the tracking relationship, not an identity. Events for the building already arrived through the citywide sync.
The shared runFeed engine
Every dataset has its own field names, status vocabulary and building identifier (BIN, BBL or only an address). One shared engine does the same steps for all of them; a feed supplies only how to fetch and how to normalize.
| Piece | File | Job |
|---|---|---|
runCitywideFeed | runFeedCitywide.ts | The nightly runner: fetchSince, normalize, persist, tombstone (only against a complete snapshot), log. |
runFeed | runFeed.ts | The candidate-scoped runner (fetch scoped to a list of buildings). Never auto-creates buildings. |
persistNormalizedEvents | persistEvents.ts | The shared middle: resolve, auto-create from BIN, link units, dedupe, upsert or quarantine. The nightly sync and the backfill both call it so history never disagrees with the present. |
runBackfillChunk | runBackfill.ts | One bounded date window of history. |
resolveBuilding(s) | buildingResolver.ts | BIN, then BBL, then normalized address. Never creates, never fuzzy-matches. |
computeXxxStatus | statusNormalization.ts | Collapses each agency's status vocabulary to OPEN, CLOSED or DISMISSED; the agency-specific value is kept in status_detail. |
fetchAllPages, socrataBackfill | feeds/socrata.ts | Shared pagination (50k pages, 500k safety cap surfaced as truncated) and a declarative backfill spec. |
units.ts | units.ts | Links HPD apartment labels to units, creating one only if the label passes isCreatableUnitLabel. |
Invariants the engine enforces:
- Resolution never guesses from an address. A fuzzy address match once attached violations to the wrong building. The only auto-created building comes from a real BIN; an address-only row that does not match exactly is quarantined.
- Tombstoning needs a complete snapshot. "Not seen this run" means "gone" only when the fetch was neither incremental nor truncated. A truncated fetch logs
partialand skips the sweep; a backfill never tombstones. Tombstoned rows are set toCLOSEDwithtombstone_reason = 'NOT_IN_LATEST_FETCH'and resurrected if they reappear. - Idempotency. Every write is an upsert on
(agency, source_id). Each feed'ssource_idcarries itssourceIdPrefix, so feeds sharing an agency (three HPD feeds) never collide on upsert or tombstone. - Bounded memory. Rows are persisted in batches of 2,000; the backfill bounds each window so a run never buffers a whole dataset.
Feeds
| Feed id | Dataset | Notes |
|---|---|---|
hpd-violations | HPD violations wvxf-dwi5 | Incremental by :updated_at. Carries unit labels. |
hpd-complaints | HPD complaints ygpa-z7cr | Incremental. Rows carry bin. Carries unit labels. |
hpd-litigation | HPD litigation 59kj-x8nc | Full pull each run (about 240k rows); no backfill. |
dob-violations | DOB violations 3h2n-5cm9 plus DOB complaints eabe-havv | Two datasets under one feed id. Violations are incremental; complaints are a full nightly replace at the source and are pulled in full. |
ecb-violations | ECB violations 6bgk-3dad | Full pull each run. Compact-text dates. |
oath-hearings | OATH hearings jz4z-kudi | Full pull each run. Resolves by address and borough. Agency (FDNY, DSNY, DOHMH) is decided per row; DOB-issued tickets are excluded to avoid duplicating ECB. |
nyc-311 | 311 requests erm2-nwe9 | Incremental. Resolves by address and borough. Filtered to HPD, DOB, DEP, DSNY, FDNY, DOT. Updates several times a day. |
dof-liens | DOF tax lien sale list 9rz4-mjek | A lien-sale eligibility roster: no amounts or dates, so every row is OPEN. Resolves by BBL composed from borough, block, lot. Full pull each run; no backfill. |
Field-level detail per dataset: NYC open data reference.
Sync state: two cursors per feed
feed_sync_state holds one row per feed and answers "where is this feed up to", which sync_runs cannot.
| Value | Meaning | Moves |
|---|---|---|
delta_watermark | Everything changed at or after this moment is loaded | Forward, nightly. Written only by a run whose fetch was not truncated. |
backfill_cursor | Everything dated at or after this day is loaded | Backward, one window per chunk, until it passes backfill_floor (the dataset's earliest record, probed once). |
backfill_status is pending, running, complete, or not_applicable (feeds with no backfill spec, because the nightly full pull already covers the dataset). The two cursors start together and walk apart; once the backfill completes, only the nightly delta runs, with no special case.
A backfill chunk:
- Reads the cursor and
backfill_window_days. - Probes
count(*)for the window and halves or widens it so it lands between 20,000 and 150,000 rows. Window size therefore tunes itself per feed. - Fetches, normalizes and persists that window (no tombstoning), then writes the new cursor and row counter.
The cursor advances only after the window's rows are durably written, so a crash or Ctrl-C just reruns the same window. Direction is newest first: recent, likely-open violations arrive before decade-old closed rows. Rows with no usable date (ECB and DOB use the string "0") are swept once at the end via fetchUndated.
Scheduling and where things run
| Job | Runner | Details |
|---|---|---|
| Nightly delta | pg_cron job sync-orchestrator-nightly-delta, 08:10 UTC | Posts one request per feed with { mode: "citywide", triggeredBy: "cron", feedId }, bearer CRON_SECRET read from Vault at call time. Per-feed requests keep each run inside the edge gateway's ~150 s limit. |
| Nightly delta from the ingestion box | systemd unit complied-ingest-delta on the ingestion box, 08:10 UTC, running run-delta-daily.sh | Runs ingest delta --all --since-days=45 over a direct Postgres connection, mirroring the pg_cron job, then ingest status. It does not chain a map export. |
| Manual delta | Same function, manage_tenant caller who is also a platform operator (or service role) | A citywide sync writes reference data every tenant reads, so a tenant admin alone is refused. |
| Backfill, redrive, identity | db-tests/scripts/ CLIs | Direct Postgres connection: no wall-clock or memory ceiling. |
The cron secret is a capability scoped to citywide mode with a single feedId; it is not the service-role key. Migrations: 20260903083457_nightly_delta_sync_per_feed.sql, 20260903090000_feed_sync_state.sql.
The ingestion CLI
Run from db-tests/ as npm run ingest -- <verb>. It reads TEST_DATABASE_URL, so it writes to whatever database that points at.
| Verb | Effect |
|---|---|
status | Per feed: rows loaded, percent backfilled, cursor, watermark, last run, quarantined total. Run this first, always. |
backfill [feed...] [--all] [--chunks=N] [--max-rows=N] | Runs N windows per feed and stops; interactive picker with no feed. |
delta [feed] [--since-days=45] | Forward pull; the only verb that moves delta_watermark. |
redrive [--limit=50000] | Re-normalizes quarantined rows from raw_data, re-resolves against the current buildings table, promotes matches into public_events, and deletes the promoted quarantine rows. Safe to repeat. |
backfill-units | Links unit ids on events that landed before unit auto-linking existed. |
refresh | Runs the Layer 2b rollup and map grid rebuild so the UI reflects a session's loads. |
Operational guidance is in the hpd-ingestion-backfill skill. See local development for the warning about the local database holding the real citywide data.
Quarantine
Rows reach unlinked_events with a reason: NO_BUILDING_MATCH_FOR_BIN, _BBL, _ADDRESS, NO_RESOLUTION_KEYS, or BUILDING_CREATE_FAILED. Two exits exist:
ingest redrivepromotes rows that resolve afterbuildingshas grown (the common case for 311 and OATH).resolve_unlinked_event(p_id, p_building_id)marks one rowresolvedby hand; it is restricted to platform operators.
Adding a feed
- Create
domain/src/ingestion/feeds/<name>.tsexporting aFeedDefinition: a stableid, asourceIdPrefixunique across feeds,possibleAgencies,fetchSince(sinceIso, deps)returning{ rows, truncated, incremental }, andnormalize(raw)returning aNormalizedEvent(ornullto skip a row). UsefetchAllPagesandsocrataUrlfromfeeds/socrata.ts. - For a large dataset, add
backfill: socrataBackfill({ datasetUrl, dateField, format, extraWhere }).extraWheremust repeat the feed's permanent filter so a backfill loads exactly the slice the delta keeps current. A feed whose dataset is small omitsbackfill. - Map status vocabulary through
statusNormalization.ts; keep the agency's own wording instatusDetail. - Export it from
feeds/index.tsand add it toFEED_REGISTRYinsync-orchestrator/index.tsand to the CLI's registry indb-tests/scripts/ingest.ts. - Add a migration that seeds a
feed_sync_staterow (not_applicableif it has no backfill) and adds the feed id to theunnest(array[...])list of the nightly cron job. - Test against real rows saved from the live dataset in
domain/tests/ingestion/fixtures(fetch is injected throughFeedFetchDeps.fetchJson, so tests never call Socrata), then runredrive-style checks withingest status.
No new tables are needed per agency: a feed adds rows to public_events and, where useful, summary columns on buildings.
Tests
domain/tests/ingestion/ covers the resolver, dedupe, status mapping, the feed normalizers against real-row fixtures, the Socrata pager, runFeed and runBackfillChunk against an in-memory Db. db-tests/tests/ingestion.test.ts covers the SQL side: constraints, the normalized_address trigger, quarantine RPCs, buildings identity, the compliance summary, and ops views.