Skip to main content

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​

TableQuestion it answersNotes
buildingsWho 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_eventsWhat 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_buildingsWho 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:

TableRole
unlinked_eventsQuarantine 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_runsOne row per feed invocation: rows fetched/upserted/quarantined/tombstoned, outcome (success, partial, failure), trigger (cron, manual).
feed_sync_statePer-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​

LayerWhatCadenceRuns in
1. Buildings identitybuildings from NYC Building Footprints (BIN, outlines) joined to PLUTO on BBL (year built, unit count)ManualLocal CLI over a direct Postgres connection
2. Citywide deltapublic_events for every feed, changed rows onlyNightly 08:10 UTC, one request per feedsync-orchestrator edge function, driven by pg_cron
2 backfillHistorical rows, walking backward in date windowsManualingest CLI
2b. Summary refreshCompliance counts on buildings, map grid, snapshot stampAfter every nightly deltaRPCs called by the orchestrator
5. ObligationsDeadlines computed from public_events and building factsOn 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.

PieceFileJob
runCitywideFeedrunFeedCitywide.tsThe nightly runner: fetchSince, normalize, persist, tombstone (only against a complete snapshot), log.
runFeedrunFeed.tsThe candidate-scoped runner (fetch scoped to a list of buildings). Never auto-creates buildings.
persistNormalizedEventspersistEvents.tsThe 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.
runBackfillChunkrunBackfill.tsOne bounded date window of history.
resolveBuilding(s)buildingResolver.tsBIN, then BBL, then normalized address. Never creates, never fuzzy-matches.
computeXxxStatusstatusNormalization.tsCollapses each agency's status vocabulary to OPEN, CLOSED or DISMISSED; the agency-specific value is kept in status_detail.
fetchAllPages, socrataBackfillfeeds/socrata.tsShared pagination (50k pages, 500k safety cap surfaced as truncated) and a declarative backfill spec.
units.tsunits.tsLinks 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 partial and skips the sweep; a backfill never tombstones. Tombstoned rows are set to CLOSED with tombstone_reason = 'NOT_IN_LATEST_FETCH' and resurrected if they reappear.
  • Idempotency. Every write is an upsert on (agency, source_id). Each feed's source_id carries its sourceIdPrefix, 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 idDatasetNotes
hpd-violationsHPD violations wvxf-dwi5Incremental by :updated_at. Carries unit labels.
hpd-complaintsHPD complaints ygpa-z7crIncremental. Rows carry bin. Carries unit labels.
hpd-litigationHPD litigation 59kj-x8ncFull pull each run (about 240k rows); no backfill.
dob-violationsDOB violations 3h2n-5cm9 plus DOB complaints eabe-havvTwo datasets under one feed id. Violations are incremental; complaints are a full nightly replace at the source and are pulled in full.
ecb-violationsECB violations 6bgk-3dadFull pull each run. Compact-text dates.
oath-hearingsOATH hearings jz4z-kudiFull 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-311311 requests erm2-nwe9Incremental. Resolves by address and borough. Filtered to HPD, DOB, DEP, DSNY, FDNY, DOT. Updates several times a day.
dof-liensDOF tax lien sale list 9rz4-mjekA 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.

ValueMeaningMoves
delta_watermarkEverything changed at or after this moment is loadedForward, nightly. Written only by a run whose fetch was not truncated.
backfill_cursorEverything dated at or after this day is loadedBackward, 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:

  1. Reads the cursor and backfill_window_days.
  2. 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.
  3. 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​

JobRunnerDetails
Nightly deltapg_cron job sync-orchestrator-nightly-delta, 08:10 UTCPosts 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 boxsystemd unit complied-ingest-delta on the ingestion box, 08:10 UTC, running run-delta-daily.shRuns 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 deltaSame 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, identitydb-tests/scripts/ CLIsDirect 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.

VerbEffect
statusPer 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-unitsLinks unit ids on events that landed before unit auto-linking existed.
refreshRuns 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 redrive promotes rows that resolve after buildings has grown (the common case for 311 and OATH).
  • resolve_unlinked_event(p_id, p_building_id) marks one row resolved by hand; it is restricted to platform operators.

Adding a feed​

  1. Create domain/src/ingestion/feeds/<name>.ts exporting a FeedDefinition: a stable id, a sourceIdPrefix unique across feeds, possibleAgencies, fetchSince(sinceIso, deps) returning { rows, truncated, incremental }, and normalize(raw) returning a NormalizedEvent (or null to skip a row). Use fetchAllPages and socrataUrl from feeds/socrata.ts.
  2. For a large dataset, add backfill: socrataBackfill({ datasetUrl, dateField, format, extraWhere }). extraWhere must repeat the feed's permanent filter so a backfill loads exactly the slice the delta keeps current. A feed whose dataset is small omits backfill.
  3. Map status vocabulary through statusNormalization.ts; keep the agency's own wording in statusDetail.
  4. Export it from feeds/index.ts and add it to FEED_REGISTRY in sync-orchestrator/index.ts and to the CLI's registry in db-tests/scripts/ingest.ts.
  5. Add a migration that seeds a feed_sync_state row (not_applicable if it has no backfill) and adds the feed id to the unnest(array[...]) list of the nightly cron job.
  6. Test against real rows saved from the live dataset in domain/tests/ingestion/fixtures (fetch is injected through FeedFetchDeps.fetchJson, so tests never call Socrata), then run redrive-style checks with ingest 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.