Skip to content

feat(onchain-analytics): correctness and write-safety hardening for the ingestion pipeline - #72

Merged
thalescb merged 12 commits into
masterfrom
remediation/integration
Sep 29, 2026
Merged

thalescb merged 12 commits into
masterfrom
remediation/integration

Conversation

@thalescb

@thalescb thalescb commented Sep 29, 2026 •

Copy link
Copy Markdown
Collaborator

Closes a set of correctness and safety defects in the GoodDollar ingestion pipeline, before any credential able to write the production dataset is issued.


Start here

Review this commit by commit, not file by file. There are 12 commits. Each one is a complete,
self-contained change that builds and passes the full test suite on its own — none of them depends
on a later commit to work. Reading them in order is roughly a tenth of the effort of reading the
combined diff, and it is the order the work was reasoned about.

git log --oneline master..remediation/integration
git show <sha>          # each commit message explains the defect it closes

Every commit message states the defect, the measurement that proved it, and the fix. Those messages
are the real documentation of this change — if you read nothing else, read the twelve messages.

If you have fifteen minutes, read these two commits. They carry the defects with real
consequences and the rest is supporting work:

  • a601672 — concurrent writes duplicated rows while both runs reported success
  • 95eae55 — verification returned "clean" when it had compared nothing

The one-line version of each commit

Commit What it closes
77eec3b Read-only BigQuery metadata manifest tooling
3edd399 Ships the event-surface reference data as a tracked seed, with its era-binding test
7c76a74 Records the release scope, and why Fuse is out of it
0ceb874 Test harness, strict control-plane loader, fail-closed CLI exit codes
4fe4af2 Defines the PipelineRuns outcome columns as an unapplied migration
87b511d PR validation workflow; hardens the daily job
c134a99 Decodes events per contract era, with a confidence grade on each row
93cb0a5 Coverage reports holes it previously hid; sandbox datasets are guarded
95eae55 Reconciles ingested claims against the contracts' own ledgers, on both chains
d1f8535 Checks the reconciliation table's columns at startup, not mid-run
48da430 Batches log capture per chain; reader evidence is measured, not asserted
a601672 Serialises concurrent writes; staging moves out of the production dataset

Where to look, and what to check

pipeline-v5/src/pipeline.ts — the capture path, and the largest single file in the diff.
One batched read per chain now replaces one read per contract, and the result is projected back
onto each contract's own block interval. Check: that a contract can only receive logs its own
address emitted, and nothing below its creation block. The projection is isolated in
src/batch.ts precisely so this is assertable rather than something you have to trust.

pipeline-v5/src/bq.ts — every BigQuery statement passes through here. Two things landed:
staging tables moved to their own dataset, and maximumBytesBilled is attached at a single
chokepoint. Check: that no staging path still composes a name from the production dataset id.
stagingTableName refuses outright if the two datasets are configured the same.

pipeline-v5/src/writelock.ts (new) — a cross-process lease. Two runs of the same block
range previously both staged the same keys, both merged against a snapshot taken before the other
inserted, and both inserted. Check: the lease is released only by the holder that took it, and a
writer refused the lease records the range incomplete rather than complete.

pipeline-v5/src/reconcile.ts — verification against the contracts' own per-day ledgers.
Check: the protocol day is derived from the contract's periodStart, which falls at noon UTC
on both chains. A comparison keyed on the calendar date of a log's timestamp misplaces every claim
made before noon — 28.9% of stored rows. There is a test that fails if that keying ever changes.

gd_dbt/seeds/event_surface.csv (2,825 lines) and contract_deployments.csv (726 lines) —
generated reference data, one row per (contract, era, event). Check the column headers and a
couple of rows, not the whole file.
dbt tests assert their integrity; reading 3,500 CSV lines by
eye is not a good use of review time.

pipeline-v5/test/unit/*.test.ts — roughly half the diff. Skim for shape rather than reading
in full; every one is a new test, none is a modified assertion.

Skip these

  • pipeline-v5/package-lock.json — generated
  • The two large test files (681 and 653 lines) unless a specific behaviour is in question
  • .github/workflows/pipeline-sandbox-integration.yml — defined and disabled by three
    independent guards; it cannot run until a repository variable that does not exist is created

Already verified

  • 543 tests passing across 25 files, 0 skipped.
  • A separate suite of 21 tests pins the specific defects this work targets. All 21 were red when
    it started; all 21 are green now, with none skipped or deferred. Each one fails again if its fix
    is removed — that was checked individually per fix, not assumed.
  • Statement coverage 80.19%, up from 57.55% at the start.
  • tsc --noEmit and eslint both clean. Three CI checks green, including GitGuardian.
  • The concurrency fix is validated by a gate that spawns two genuinely separate OS processes
    against live BigQuery sandbox datasets, run with npm run gate:write-safety.
  • Claim counts and amounts were reconciled against the UBIScheme contracts on both chains. XDC
    matches the contract to the wei on 255 of 255 days.

What this PR does not do

  • No production data was written, and no credential capable of writing it was used or created.
    The identity running this work has tables.create, tables.update, tables.updateData,
    tables.delete and datasets.delete denied on the raw dataset, verified after each change.
  • The declared staging dataset (BlockchainEvents_Staging) is a configuration default and does
    not exist yet
    . Creating it, and scoping a writer identity to it, is the next step and is
    deliberately not in this PR.
  • The scheduled ingestion workflow stays disabled.

Known limitation, stated deliberately

The cross-month duplicate fix bounds the MERGE window by padding it one whole month on each side.
That covers a row displaced by 31 to 62 days; a row displaced 95 days was measured to duplicate
and is not covered. A uniqueness check reports such a duplicate rather than letting it pass
silently. This is a fix within a measured bound, not a general solution, and it is comfortable
inside the twelve-month window this pipeline targets. It is documented at the code that implements
it.


Why one PR rather than several

The twelve commits share files — the seeds, the BigQuery layer and the capture path were each
touched by several of them — so a per-commit split would mean re-committing by hand and would
produce commits that never existed and never passed a test. The per-commit history is the split,
and it is genuinely reviewable in that form.

…ling

Captures a metadata-only manifest of the warehouse and diffs two of them, so a
schema change can be detected without reading any table data. Imports only Node
builtins; bills zero bytes.
…ng test

2,824 rows, one per event per contract era, keyed by the selector the chain puts
in topic0. Every selector was recomputed from the verified interface rather than
recalled. The accompanying test binds each surface row to a real deployment era on
the composite key, which the built-in single-column relationship test cannot do.

Includes the recorded exception for this seed in `.shipgate-allow`: it pairs contract
addresses with era block bounds, which matches the shape of a holder table without
being one. Every address in it is a contract already listed in contract_deployments.
States which chains this release covers, why Fuse was dropped, what retaining its
reference-seed rows means, and how the exclusion is enforced. The control-plane
code cites this file as its decision record, so it lands before the code that
names it.
…nd a fail-closed CLI

Three changes that cannot be separated into commits that each build, because they
share three files:

- A Vitest harness with replaceable BigQuery, RPC, clock and lock seams, so the
  package is testable with no credentials and no network.
- A strict RFC 4180 control plane that rejects a malformed reference file whole
  rather than row by row, and parses every 64-bit bound from its original decimal
  lexeme so the maximum-value sentinel never becomes a different integer.
- A typed run outcome, so a refused run and a run that did no work both exit
  nonzero and the persisted record agrees with the exit code.

Also enforces the declared release scope at the single point every capture,
repair and calibrate command resolves a chain's contracts.
…n unapplied migration

Additive ALTER TABLE only. NOT APPLIED: it carries a DO NOT RUN banner in its first
twenty lines, which the deployment helper honours, because that helper otherwise
executes every .sql in the folder.
…ily job

Runs lint, typecheck and coverage on the pipeline package and parses the dbt project
on every pull request. Depends on the npm scripts and the seed added above, so it
lands last.
…evidence

Binds every captured log to the event definition that was actually live when it was
emitted, and records how confident we are that the binding is right.

- Union ABIs are built newest-implementation-first. Chronological order makes a
  later-era log decode silently wrong rather than raising, which is the failure this
  ordering exists to prevent.
- Any selector carrying two different parameter layouts on one contract is computed
  from the reference data and must be declared. An undeclared one fails the build.
- Contract eras are modelled as validity intervals with three tests: no overlaps, no
  gaps from deployment to head, and exactly one interval covering any queried block.
- Each interval carries a confidence grade and a receipt that resolves inside this
  repository, so a figure can be traced back to how it was established.
- Plausibility bands on every consumed numeric field, because structural tests pass
  on a value that is silently wrong.
…dbox datasets

Two defects where the pipeline reported more confidence than it had, and the control
that keeps the fix for them away from production data.

- A range nobody successfully captured is now reported as an open gap, and the resume
  point no longer reads past a hole it never recorded. Previously an internal hole
  between two clean captures was invisible and ingestion silently skipped it.
- A missing transaction now marks the transaction grain incomplete and the command
  exits nonzero. Previously the grain read complete and the run exited zero, which is
  why the condition went unnoticed.
- Decode coverage is reported beside block coverage, per contract, because a log whose
  event definition is unknown is indistinguishable from an absent one once a query has
  run. Its error count is carried separately: a count of zero means something only
  when no read failed.
- Disposable dataset handling refuses the production datasets by hard-coded name,
  requires a purpose label, and reads that label back off the live dataset before
  deleting anything. The default dataset in configuration is the production one, so a
  refusal that depended on configuration would not have been a refusal.
@thalescb thalescb mentioned this pull request Sep 29, 2026
6 of 9 tasks
… on both chains

The UBIScheme contracts keep their own per-day claimer counts and distributed amounts
in public state, so ingested claims can be checked against the chain rather than
against the warehouse itself. Only the XDC scheme was registered, which left Celo --
half the ingested scope -- with no external check of any kind.

- Celo UBIScheme is registered with the period start the contract returns, 1677672000,
  read at a pinned block with three endpoints agreeing.
- Verification returns what was compared instead of a boolean. "Compared 254 protocol
  days and every one matched" and "compared nothing at all" were previously the same
  value, and both exited zero. Three paths reached it: a chain with no oracle, a
  contract the warehouse holds no rows for, and a window containing no frozen day. A
  result is now clean only when something was compared, and the command exits 2 when
  there was nothing to check.
- An oracle that produced no comparison makes the run a finding rather than a pass.
  Silence from half the checked surface is not evidence that the surface is sound.
- A named day range no longer requires the warehouse to already hold rows for the
  contract. Deriving the window from stored rows is right when no window was given;
  requiring it when the caller supplied one means a contract with no rows can never be
  checked, which is exactly when the answer is worth having.
- An unreadable chain is recorded as unreadable instead of aborting the command, so
  one chain's endpoints having a bad minute no longer discards another chain's
  completed comparison.
- Oracle selection is filtered by release scope, like every other per-chain loop.

Both schemes open their protocol day at noon UTC, so a comparison keyed on the calendar
date of a log's timestamp misplaces every claim made before noon -- 764,706 of the
2,649,450 stored XDC claim rows, 28.9 percent. The comparison already keys on the
contract's own day; there is now a test that fails if that ever changes, and it was
checked by running the same comparison under both keys.
…tartup

`CREATE TABLE IF NOT EXISTS` reaches the right shape in a fresh dataset and silently
skips one that already exists, so a dataset carrying an older copy of
`OracleReconciliation` keeps it indefinitely. The copy in the production dataset
predates the chain dimension: it has `table_id` where this pipeline writes `chain_id`
and `contract_address`, so the rows the reconciliation builds would be rejected by it.

Without the check the run gets as far as the load step -- after reading a year of
per-day contract state -- and fails naming a column rather than the migration to apply.
The two other bookkeeping tables were already asserted this way; this one was not.
…r evidence in measurement

Capture issued one log query per contract while every reader already accepted a list of addresses. Measured against the current registry, that is 146,565 chunk requests on Celo where a batched read needs 3,939, and 175,838 against 7,535 across the three chains in scope.

Capture now issues one batched read per chain and projects the result back onto each contract's own interval. Batching must not cost attribution, so the projection is a separate module with its own tests: a contract receives only the logs its address emitted, only the blocks and transactions those logs point at, and nothing below its own creation block; a failed chunk is recorded against exactly the contracts whose range it overlaps, clipped to that range; an error naming no range is kept rather than dropped.

Reader evidence is now measured rather than asserted. The block pin uses each chain's recorded finality instead of a fixed five blocks, which on Celo was pinning 1,925 blocks inside the reorganisation window. An empty range is confirmed only when two independent endpoints each cover it without error and find nothing, and the rendered reason states the count actually observed. Endpoints are tallied separately for zero-log answers and failures, so an endpoint returning a false zero is visible instead of indistinguishable from one returning rows. Every answer carries an admissibility grade, and the coverage ledger records it.

A capture whose reader reports it is still holding blocks at or below the range start is recorded rollback_eligible rather than complete, so the resume frontier stops there and the range is read again once it settles. Previously the status was decided before the guard was read and the guard only appended a note.

The HyperSync worker call is replaceable, so the index reader's chunk planning, retry bound, short-collection refusal and guard forwarding are covered; that module was at 7 percent because the existing seam replaced the whole reader.
…g from production

Two captures of the same block range running at once each staged the same
merge keys and each merged against a target snapshot taken before the other
inserted, so both inserted. One measured run left RawLogs holding 149 rows
over 102 distinct keys while both processes exited zero and both recorded the
range as completely captured, which is why it went unnoticed.

A cross-process lease now excludes a second writer on the host. It is held by
an exclusive file create, released only by the holder that took it, and
reclaimed at once when its holder is no longer running rather than waiting out
an expiry that has to stay long. A writer that cannot take the lease is
refused, exits nonzero, and records the range incomplete so the next run reads
it again.

Staging tables move to their own dataset. They are created and dropped on every
run, so wherever they live the writer needs permission to delete tables, and
while they lived in the production dataset that meant holding delete permission
on the raw layer. The ordering made it worse: staging cleanup happens after the
MERGE, so a run lacking that permission fails having already written production
rows. Composing a staging name now refuses outright if the two datasets are
configured the same.

Every query and DML statement carries maximumBytesBilled, attached where all of
them pass rather than at each call site, so an over-ceiling job is refused
before it runs instead of billed. Load jobs cannot carry it and are named as
such.

Also: a merge-key uniqueness check, because the padded MERGE window bounds a
cross-month duplicate to a 31-to-62-day displacement and a larger one is not
covered, and an accepted residue is only defensible if it is visible. And the
recent-pin state endpoints are separated from the deep-archive ones, which they
had been conflated with; measured against the reads actually performed, four
free Celo endpoints qualify where the combined list scored one.
@thalescb thalescb changed the title feat(onchain-analytics): test foundation, reference data and safety controls for the ingestion pipeline feat(onchain-analytics): correctness and write-safety hardening for the ingestion pipeline Sep 29, 2026
@thalescb
thalescb merged commit e493247 into master Sep 29, 2026
3 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant