Skip to content

Repository files navigation

UPSTREAM

Desired-state intelligence for DataHub

Apache-2.0 Python 3.11 DataHub v1.7.0 tests

DataHub stores what is; Upstream decides what should be

DataHub is the system of record. Upstream is the system of intelligence.

You declare what the estate should look like. Upstream perceives what it actually looks like through DataHub's own retrieval surface, measures the gap as a single number, explains it with cited evidence, and proposes priced changes for a human to approve — then writes the ratified ones back into DataHub.

DataHub stores the metadata. Upstream decides what it means and what to do about it.

Status. Both controllers run end to end against a live DataHub: the estate is measured, the gaps are closed, and a failure is diagnosed, fixed, verified and written back. There is an operator console over a live HTTP API, with a human gate in front of every write. If something is described here, it runs — what is not yet proven is in LIMITATIONS.md rather than implied to work.

        ┌───────────────────────────────────────────┐
        │                 DATAHUB                   │
        │            SYSTEM OF RECORD               │
        │  metadata · lineage · assertions          │
        │  incidents · timeline · properties        │
        └───────────────────┬───────────────────────┘
                            │  read through DataHub's own Skills
                            ▼
        ┌───────────────────────────────────────────┐
        │                UPSTREAM                   │
        │          SYSTEM OF INTELLIGENCE           │
        │  desired-state reconciliation · memory    │
        │  simulation · RCA · risk · confidence     │
        └───────────────────┬───────────────────────┘
                            │  proposals + evidence
                            ▼
                 ┌──────────────────────┐
                 │    HUMAN APPROVAL    │
                 └──────────┬───────────┘
                            │  ratified changes only
                            ▼
              WRITE BACK INTO DATAHUB
   incidents · runbooks · assertions · properties · tags

Read the arrows literally. Metadata flows up into intelligence and conclusions flow back down into the system of record. Upstream is a client of DataHub in both directions and a store of record for nothing.

DataHub is responsible for Upstream is responsible for
metadata storage · lineage · search reasoning · reconciliation
governance · the entity graph memory · simulation
timeline · APIs · ingestion proposals · orchestration · approval

Three consequences, each enforced in code rather than asserted here. Upstream has no metadata store — no catalog table, no lineage table, no entity registry. If a fact is not in DataHub, Upstream does not know it, and the fix is to write it to DataHub rather than cache it. Upstream has no ingestion — where this repository runs datahub ingest, it is running DataHub's. Every conclusion is a DataHub write: RCAs become incidents, runbooks become documentation, learned checks become assertions, risk becomes structured properties.

Quickstart

cp .env.example .env                 # add a DataHub token and one provider key
pip install -e ".[dbt]"              # Python 3.11 — acryl-datahub does not support 3.12+

make up                              # DataHub OSS quickstart, version-pinned
make seed                            # DataHub's own sample estate, mirrored, then decayed
make score                           # the gap between estate.yaml and the graph
make reconcile                       # close it, behind a human gate
make break scenario=vendor-drift     # a silent failure
make investigate urn="urn:li:dataset:(urn:li:dataPlatform:dbt,demo.main.fct_revenue_daily,PROD)"

Prefer a screen to a terminal? make api and make console-live put all of the above behind five screens at http://localhost:5273. Want neither? make eval-replay re-scores the committed run records with no API key, no network and no DataHub.

make seed loads DataHub's showcase-ecommerce and bootstrap sample packs, so the estate a judge gets is the same one every number below was measured against. Reproducibility here is a property of the design, not of luck.

The operator console

Everything above is also a surface. Five screens, served by a FastAPI backend over the same projectors that generate the committed artifacts — so the live payload and the recorded one cannot drift.

make api             # the read, write and event API on :8000
make console-live    # the console against it on :5273
make console         # or against the committed recordings, with no backend at all
Screen Answers
Mission Control What is broken, and what is waiting on a person
Investigation Detail Why the system believes it — the evidence ledger, cited
Simulation Workspace Which remediation to approve, and what each one is worth
Approval Center Every controller's pending work, in one queue
Proposal Review Exactly what one approval writes, aspect by aspect

The frontend computes nothing. Filtering, ordering, paging, counting, pricing and the empty-state copy are all backend projections; React renders what it is handed. A payload the projector cannot fill arrives with the reason it is empty rather than a zero.

read     route → dependency → projector → response
write    route → dependency → controller → writer → receipt → projector → response
live     GET /api/events — the stored receipt, streamed verbatim

Every mutation answers with a MutationReceipt: what was written, to which aspect, and the urn DataHub returned — or the reason nothing was. That one object is also the line appended to the journal and the frame the event stream broadcasts, so there is no second schema and no translation layer between the write path and the live one. A receipt names the screens it invalidated; each store re-asks its projector rather than patching state by hand.

Approving is not the same as writing, and the receipt distinguishes them. A batch that half-wrote answers 200 partial, never a 500 — discarding the writes that landed is the one outcome an operator cannot recover from.

Full contract, endpoint table and design notes: console/README.md.

The spec, and the dial

estate.yaml declares what the estate should look like: tiers, and the requirements each tier must meet. Selectors resolve to URNs through DataHub search, so the spec follows the estate as it grows rather than naming assets one by one.

make score      # conformance of the graph against estate.yaml
tier-1     16 assets    128 checks     75 passed     59%
tier-2      2 assets     10 checks      4 passed     40%
raw         8 assets     24 checks      2 passed      8%
overall    26 assets    162 checks     81 passed     50%

Scoring is deterministic and involves no model: one pure function per requirement, scored by checks rather than by assets, so the number moves smoothly as gaps are closed. Every read reaches it through DataHub's own retrieval tools — 137 invocations across six skills, with a logged SDK fallback where a tool returns the wrong answer. The full report, with all 81 violations, is in examples/conformance/before.md.

The raw tier scores 8% because the vendor feed boundary carries no schema contract, no owner and no test. That is not a scoring artefact — it is the reason the drift scenario propagates silently.

Closing the gap

Measuring drift is a report. Closing it is the product.

make reconcile-dry    # propose and price, write nothing
make reconcile        # the same run, with a human gate before any write
make reset            # undo every write, back to the seeded estate
observe ──> plan ──> ratify ──> actuate ──> verify
   │          │         │          │           │
  50%     45 priced   human      45 writes    83%
        proposals     gate

Every proposal is grounded and priced before anyone sees it: an owner comes from git history, the lineage graph, a sibling asset or the domain — never from a guess; a description is written from the columns and upstreams the read path already returned; a check is a deterministic template. Proposals with no simulated benefit are dropped rather than shown.

Projected +32.7%, measured +32.7%. The what-if and the what-is run through the same checks.py over the same context object, so the number on the ratify screen is the number you get. A disagreement here is a bug, not a rounding difference.

The gate is a real pause — LangGraph interrupt() with a checkpointer — and there is a test that fails if it ever becomes a pass-through.

28 violations survive on purpose. Freshness requirements need an evaluation, not a declaration; some assets have no owner anywhere in the graph to infer from; others ask for a judgement (a business term, a retention period) that belongs to a human. Each one is listed with its reason in examples/proposals/reconcile.md. An agent that filled those in would be guessing into a catalog people trust.

What is verified today

Every DataHub surface Upstream depends on has been executed against a live OSS quickstart and read back. The working shapes, the deviations from the documentation, and the parts that are not proven are recorded in spike/DATAHUB-USAGE.md.

make up      # start DataHub (version-pinned, host-port-safe)
make seed    # load DataHub's sample estate, mirror it, add run history, decay it
make spike   # execute every read and write, and read each one back

make spike reports PASS only for writes it could read back afterwards; anything it could not confirm comes back as PARTIAL or FAIL. A write that emits cleanly but never lands is a failure — that distinction matters, because a malformed assertion write passed silently here until the read-back was added.

The demo estate

Upstream reasons over a real metadata graph, so the demo uses DataHub's own sample data — loaded, mirrored into DuckDB for forensic SQL, then broken on purpose by a deterministic scenario library. See demo-stack/README.md.

make scenarios                      # list the failure scenarios
make break scenario=vendor-drift    # inject one, then re-ingest so DataHub sees it

The headline scenario is a silent one: a vendor renames a column, a staging cast turns the downstream mart 85% NULL, and no job fails. DataHub's Timeline API reports the rename as a single MODIFY event — that is the signal the correlator is built on.

Diagnosing a failure

make investigate urn="urn:li:dataset:(urn:li:dataPlatform:dbt,demo.main.fct_revenue_daily,PROD)"
make investigate-dry urn="..."   # diagnose and price the fix, write nothing
make context urn="..."           # what the agent perceives, and which skill produced each fact
triage ──> lineage ──> correlate ──> forensics ──> critic ──┬──> report ──> …
   │          │            │             │            │     └──> abstain
 8 skills   6 hops    timeline +     bounded      confidence
 + memory   upstream  run history    read-only    from 5 signals
 + risk               MODIFY, not    SQL
                      REMOVE+ADD

The confidence number is computed, never reported by the model. Five orthogonal signals, each visible in the trace and in the write-up:

Signal Weight On the golden run
evidence — cited support × temporal fit × diversity 0.60 0.83
memory — a recorded conclusion on this subgraph 0.10 0.00 cold · 0.60 warm
timeline — the change precedes the symptom 0.10 1.00
forensic — a query independently confirms it 0.10 1.00
simulation — the fix removes the symptom 0.10 0.50 drafted · 1.00 verified

A hypothesis that cites no evidence is dropped in code. A hypothesis with no timeline event, query result or run record is multiplied by 0.6, because a story without a hard artifact is a story.

On the drift scenario the run concludes at 0.75 — the vendor rename, citing the MODIFY event and a query showing the mart's discount column is 100% NULL while the source is 82%. On a healthy asset the same command abstains at 0.70 and names the one probe that would settle it: see examples/traces/ for both.

Fixing it, and making it not happen again

Diagnosis is where most agents stop. The rest of the loop is the product.

critic ──> simulate ──> [ human gate ] ──> remediate ──> verify ──> scribe
             │                                 │            │          │
          price the                        the diff     rebuild in   4 artifacts
          fix first                                      a sandbox   in DataHub

Nothing is written before the gate, and the gate shows numbers rather than a recommendation: the conformance the fix would move, the blast radius, how much earlier a boundary check would have caught this, and the confidence decomposition behind the claim.

The patch is constructed, not generated. The remediator reads the live schema, finds the transformation still naming the column the vendor removed, and rewrites it to prefer the new name and fall back to the old — then adds a contract test at the staging boundary, one hop from the vendor instead of three transformations downstream at the mart. A model writes the pull-request prose, where being wrong is cheap. Nothing is applied: the diff is written to examples/pull-requests/ and checked with git apply --check, because a branch is a human's decision.

Then it is verified before it is believed. The transformation project and the warehouse are copied to a temporary directory, the patch is written there, the affected models are rebuilt and the failing check is re-evaluated — on the copy. 13/13 tests pass after the patch is what moves the fifth confidence signal from 0.5 to 1.0 and the run from 0.75 to 0.85. A fix nobody confirmed leaves the incident open, and there is a test that fails if that ever stops being true.

Four artifacts land in DataHub, and they are what make the next run cheaper:

Artifact Where it goes Who reads it
Incident + RCA the symptom asset, resolved only if verified a human, in the UI
Runbook the cause asset, where the next engineer looks a human
fragilityScore · lastRootCause · mttrMinutes · knownFailureModes both assets the next run
Immunization assertions the ingestion boundary DataHub, continuously

The citations in the RCA are written by the code, not the model. It compounds in DataHub rather than in Upstream: delete this container and the learning survives.

The second time is cheaper

Re-break the same scenario and run the same command:

cold warm
plan walk lineage → correlate → forensics → critic confirm the remembered cause → forensics → critic
tool calls 48 8
forensic queries 2 1
evidence collected 18 11
confidence 0.75 → 0.85 0.82 → 0.92

The trace names the incident it matched. Recognition comes from the structured properties this system wrote for itself to read — not from prose similarity, which cannot work: a symptom line is a dozen words and an RCA is three hundred, so a union denominator scores a perfect match at 0.007. That bug was invisible until there was something to recall.

A warm start reorders the investigation; it does not skip the evidence or the critic. If the memory is wrong, the confidence does not materialise and the run abstains.

Every model call is recorded to a cassette, so LLM_MODE=replay reproduces a run with no API key and no network. Replay reproduces a run, not a question: prompts carry the evidence that run collected, and the correlator reads a window relative to now, so a re-seeded estate asks something different and says so rather than pretending.

The numbers

Ten seeded cases from the scenario library, scored by evals/runner.py against the committed run records in evals/runs/.

make eval           # inject every scenario and run the loop, live
make eval-replay    # re-score the committed records — no API key, no network, no DataHub
Scenario Cases Root cause @1 Correct abstention Median wall-clock
vendor schema drift 4 1/1 3/3 109s
silent skipped job 2 0/1 1/1 107s
boundary assertion failure 2 1/2 n/a 64s
unrelated noise (no cause) 1 n/a 1/1 160s
recurrence (warm memory) 1 1/1 n/a 71s
overall 10 3/5 5/5 95s

Confidence separates right answers from wrong ones by 0.256 — mean 0.850 on correct conclusions against 0.594 on everything else. That row is the one that decides whether the five-signal model is load-bearing or decorative, which is why it is published rather than described.

Cold against warm, on the same failure: the warm run confirms a remembered cause instead of re-deriving it, and finishes in a fraction of the tool calls — the full comparison is in evals/results.md.

Both misses are explained in evals/results.md, and neither is a wrong answer — the system abstained both times. One is a gap in the demo estate: the freshness injector ages rows in the warehouse but publishes no staleness signal to DataHub, so there is nothing for the read path to perceive. The other is the confidence model working as specified: a real failing assertion with no timeline evidence behind it reaches 0.60 on a single hard artifact and declines, because requiring several independent signals to move together is what a 0.85 is meant to mean.

Prioritising: eight signals, one number

Risk decides which investigation runs first, how much it may spend, and how proposals are ordered at the gate. One heuristic is a number a reviewer can argue with; eight independent graph facts, each separately inspectable, is a model.

risk = 0.20 · business_criticality      tier · domain · exec-facing consumers
     + 0.15 · historical_fragility      decay-weighted incident history
     + 0.15 · assertion_health_gap      1 − passing / required
     + 0.10 · lineage_depth             normalised hops × fan-out
     + 0.10 · governance_violations     open estate.yaml violations on the asset
     + 0.10 · incident_similarity       resemblance to past incidents here
     + 0.10 · business_impact           glossary terms, revenue domains, dashboards
     + 0.10 · (1 − confidence)          the controller's own uncertainty is a risk
                                        ─────
                                         1.00

The eighth term is the one people miss: acting on a shaky conclusion is riskier than acting on a solid one, so confidence feeds risk and risk gates actuation. risk.explain() prints the arithmetic next to the priority in every trace.

What this is, in nine rows

Capability What it means
Desired-state reconciliation A declared intent and a control loop that closes the gap. Not a report — a reconciler.
Institutional memory The agent reads its own past conclusions out of DataHub before it thinks. The organisation's experience is an input, not a wiki nobody opens.
Evidence-based reasoning Every hypothesis cites indices into a ledger. Uncited claims are dropped in code, not discouraged in a prompt.
Multi-step investigation Triage → lineage → timeline → forensics → critic, with budgets, a bounded second pass, and a real abstention path.
Simulation Blast radius, governance impact, conformance and MTTR change — computed before a proposal is shown, deterministically, by the same code as the live dial.
Proposal generation Conclusions become typed, cited, risk-ranked, priced proposals. Never silent mutations.
Continuous learning Each resolution registers the check that would have caught it, so the same failure is cheaper — or prevented — next time.
Explainability Risk, confidence and conformance are each arithmetic you can read in one function. The RCA is the ledger, rendered.
Human-in-the-loop governance Nothing is written without ratification. Read-only SQL, drafted patches only, per-run budgets. Autonomy bounded by design rather than by hope.

Documents

ARCHITECTURE.md The layering, the pipeline, the read path, and every component contract
SAFETY.md The gates, the thresholds, the guardrails, and what this system refuses to do
LIMITATIONS.md What is not proven, and what the demo leans on
spike/DATAHUB-USAGE.md The verified DataHub surface, with the read-path matrix and three platform findings
evals/results.md The grid, the confidence calibration, and cold against warm
examples/README.md Every committed artifact, and the command that regenerates it
console/README.md The operator console, the read API, the write path and the event stream
PROJECT.md The full write-up: the problem, the architecture, every mechanism and what each one is for
DEMO.md The shot list: what the demo shows, in order
contrib/CONTRIBUTIONS.md What went back to DataHub: a merged-ready patch, three findings, and the reconcile-estate Skill

Layout

estate.yaml    the desired state: tiers, requirements, and how controllers may act
upstream/      the package
  ├ spec/        the declared estate, and the loader that resolves selectors through search
  ├ datahub/     the read path (8 Skills) and every writer, one per aspect
  ├ agents/      triage, lineage, correlator, forensics, critic, simulator, remediator, scribe
  ├ controllers/ the two control loops, the shared write dispatch, and the ratify controller
  ├ conformance/ the dial: one pure function per requirement, no model
  ├ console/     the read model — one projector per screen, plus the receipt journal
  └ api/         FastAPI routes: read, mutate, stream. Orchestration only
console/       the operator console (Vite + React), and its wire contracts
contrib/       the Skill contributed back to DataHub
spike/         the verified DataHub read/write surface, and the doc that records it
demo-stack/    the estate: sample-pack loaders, DuckDB mirror, chaos + governance decay
evals/         the ten-case suite, its runner, and the committed run records
examples/      generated reports and LLM cassettes, committed as evidence
landing-page/  the public page, built from those reports — `make landing`
tests/         fixtures pinned from live responses, so parsing is testable offline

Requirements

Python 3.11 (acryl-datahub does not support 3.12+), Docker with at least 8 GB, and one LLM API key — the router is provider-agnostic and .env.example documents four options, with NVIDIA NIM's free tier as the default. DataHub image and dependency versions are pinned deliberately; spike/requirements.txt records what floating them cost.

Install with the [dbt] extra. The transformation layer is what makes the drift scenario propagate, so make seed needs it; a bare pip install -e . gets as far as the DuckDB mirror and then stops at dbt: No such file or directory.

make eval-replay needs none of it: no key, no network, no DataHub.

Licence

Apache-2.0. See LICENSE.

About

Desired-state reconciliation for the data platform, built on DataHub

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages