Governed open lakehouse · banking
A payment made a minute ago, found in half a second, by an agent that sees only what its colleague may see.
An open lakehouse (Iceberg, Polaris, Trino, OPA, Kafka, Spark) built around one hard consumer: an AI assistant on a bank colleague's live customer call. The assistant is the demo. The platform underneath is the point, so every answer shows how the platform produced it.
Unrecognised payment located
£249.99 to QuickTech Electronics Online, 10:26 UTC, status pending. Freeze the card now, then raise a fraud claim.
- Sourcecore banking Postgres commit → Debezium (WAL) → Kafka corebank.core.transactions
- Icebergsnapshot 6170615646954925454, committed by cdc-stream 35 s before the lookup
- Trinoas alice, pinned FOR VERSION AS OF that snapshot · 1 row
- OPArow filter brand IN ('Meridian') · no PII columns read, so nothing to mask
- Auditrow #124, chain hash 9d0053bffa6ba905…, purpose CALL-5251D516
The demo
One platform, three consumers, three acts
A platform is judged by what it lets people do safely. Each act uses the same tables and the same policy, seen from a different consumer.
A colleague on a live fraud call
- The caller says their card was stolen. The assistant hears the intent and looks the caller up with Alice's own token.
- Account cards stay locked until Alice confirms ID&V. The agent can't get round policy: it asks the gateway, and the gateway asks OPA.
- The payment written to core banking a minute ago is already in silver, and the assistant finds it.
- Each card opens to show the snapshot, the Trino query, OPA's decision and the audit row behind it.
- The injection call ("ignore your instructions…") gets nothing. Customer text is data, never instructions.
# open http://localhost:8090, sign in as alice
make live-demo
phone ******4578 · Meridian only · vulnerability shownphone in clear · Meridian + Northgate (complaints team)contact details NULL · all brands · silver deniedSELECT … FROM silver.transactions FOR VERSION AS OF 6170615646954925454gold.customer_360 ← silver.* ← bronze.cdc_events (OpenLineage → Marquez)An analyst on the same tables
- Carol queries gold in Superset SQL Lab, Trino or Grafana, always as herself. She gets aggregates across every brand and no way to identify anyone.
- She can't infer a hidden column by filtering on it: masks apply before predicates.
- "What did the assistant see at 10:26?" is one query: any answer's snapshot id reproduces it exactly.
- Marquez shows where
customer_360came from. Dagster shows the WAP checks that had to pass before it was published.
make sql U=carol Q="SELECT brand, count(*) FROM lakehouse.gold.customer_360 GROUP BY 1"
make sql U=carol Q="SELECT * FROM lakehouse.silver.accounts" # denied by OPA
docker kill cdc-stream (kill -9, mid-call)console keeps serving · reads fine · freshness pausedhealer restarts it · checkpoint resumes at the last Kafka offsetchange made during the outage arrives (no loss) · audit chain intactBreak it on purpose
- Kill the stream in the middle of a call. The colleague keeps working on data that's a few seconds older.
- Kill OPA, Trino or Polaris instead, and the platform fails closed: "data temporarily unavailable", never unmasked data.
- Every one of 11 components recovers by itself. Idempotent bronze writes and MERGEs mean replays don't duplicate.
make chaos # kills each component in turn and measures it
Architecture explorer
Every component, what it does, and how we know it works
Select a component, or trace one of the two paths that matter: how a payment reaches the lakehouse, and how a question on a call is answered.
Tenants
One platform, a stack for every use case
Teams build their own data products and AI tools in their own repos. Each team picks the platform services it needs and runs only those, so a markets demo does not start core banking, and the bank demo does not start markets.
| Blueprint | What it adds | Status |
|---|---|---|
| core | Postgres, object storage, Polaris, Keycloak, OPA, Trino, MCP gateway. Always on, so identity, policy and audit are never optional. | built |
| streaming | Kafka and the tenant reconciler | built |
| orchestration | Dagster, one code server per tenant | built |
| corebank | Change data capture (Debezium and the CDC stream) and the bank's topics. Leave it out and a markets-only stack uses 4.7 GiB instead of 6.4. | built |
| observability · lineage · bi · ai | Grafana and Prometheus, Marquez, the SQL workbench, Live Call Assist | built, opt-in |
# the laptop set: core, streaming, orchestration, tenants make up # everything, for a bigger machine make SCALE=full up # next: a stack per use case, from what each tenant declares make up USE=markets-data
- One pull request onboards a team. A file,
tenants/<team>.yaml, becomes Kafka topics, Polaris namespaces, a Keycloak group, an identity that can write only its own namespaces, and a Dagster code location. Nothing is granted by hand. - The platform tests the platform. A small canary tenant proves reconcile, per-team identity, grants and memory budgets on every change, so no real team's code is needed.
- Tenants test themselves. lakehouse-markets-data owns its streams, data checks and releases. The platform only pins its image tag.
- Memory is a budget. Tenant workloads share 4,608 MB, 4,096 in use. The laptop set measured 6.5 GiB of containers on a 10 GB Docker VM, with no out-of-memory kills.
Coinbase WebSocket → Kafka → bronze, silver, rejectsShadowTraffic → Kafka → bronze, silver, dead-letter topicECB FX rates and the FCDO sanctions list, daily1-minute candles · card auths per day · sanctions hitsevery colleague reads gold · bronze and silver: platform adminsRead the plain-words overview in docs/vision.md, and the decision in ADR 15.
Proof, not claims
What is checked, and where
Each row is an assertion in make verify, make chaos, the eval gate or CI. They run from an empty machine on every push to main.
| Capability | How it is proven | Result |
|---|---|---|
| Fresh data for live calls | Write to core banking, poll as a colleague until visible (freshness_probe.py); the fraud payment written seconds earlier is pinpointed by the agent | ~10 s · < 60 s SLO |
| Row-level security by brand | alice, bob and carol run the same query and get different brands | 3/3 pass |
| Column masking by persona | Phone partial for alice, clear for bob, NULL for carol; carol can't infer a masked column by filtering on it | 4/4 pass |
| Agents inherit the colleague's rights | Token exchange (RFC 8693) so Trino sees the colleague; the agent's answers carry alice's masks; an analyst's agent identifies no one | pass |
| Reproducible answers | Every gateway answer is pinned to an Iceberg snapshot and carries OPA's decision and its audit row | 3/3 pass |
| Tamper-evident audit | Hash chain verified; DELETE on the audit log is refused; denied calls are audited too | pass |
| Bad data never reaches silver | ODCS contracts, quarantine with reasons, WAP publish only after checks; late parents are retried, not dropped | 0 violations |
| Safe AI guidance | Offline eval gate: intent/vulnerability/risk recall, 8 scripted calls, zero ungrounded cards; injection rejected before any SQL | 8/8 · 0 ungrounded |
| Teams onboard without hand grants | The tenant file becomes topics, namespaces, an identity and a code server; the canary tenant writes as itself (200 on its namespace, 403 on silver); a container outside the project joins the platform networks; the tenant memory budget fails lint when exceeded | pass |
| Raw data for platform admins only | Colleagues are refused bronze; the admin reads bronze but every whole-record payload is NULL, including in quarantine; contracts declare the tag | pass |
| AI spend capped and visible | A daily cap on model spend is exported and alerted on; every fallback to rules or the template names its reason; cost and cache hits per call on the dashboard | pass |
| Operable | Every scrape target up (Trino via a machine identity), 20 SLO and burn-rate rules loaded, dashboards from code, Dagster behind SSO | pass |
| Supply chain | Actions pinned by SHA, gitleaks over full history, Trivy, SBOM + provenance, cosign keyless signatures on GHCR images | signed |
Chaos results
Kill each component. Measure what the colleague sees, and time to recover.
Fail closed is the design, not an accident. With OPA, Trino, Polaris or storage down, the assistant says data is temporarily unavailable. It never falls back to unfiltered data, and it tells the colleague not to guess.
Streaming parts degrade gracefully. With Kafka, Debezium or the stream down, reads keep working on slightly older data. A change made during the outage arrives after recovery: the replication slot holds WAL, the checkpoint holds offsets, and writes are idempotent.
Recovery is automatic. Restart policies plus a reconciler that restores desired state, including restarting failed connector tasks. The audit chain is verified intact after every run.
Scale calculator
What 50,000 colleagues means for the architecture
The laptop proves behaviour. Capacity comes from a load model. Move the assumptions and see which design decision the load forces.
Other constants in the model: 500 concurrent analysts at one query per 10 s (about 50/s, cacheable); 64+ transcript partitions for thousands of concurrent calls; OPA as a sidecar with signed bundles so policy never needs a network hop. Details in docs/scale.md.
Run it yourself
From zero to a live call in one command
# Docker with 10 GB (laptop set), uv, make git clone https://github.com/cloudcruncher/open-lakehouse cd open-lakehouse make demo # up, seed, pipelines, then the checks make urls # console, Grafana, Dagster, lineage, SQL workbench make live-demo # play every demo call as alice make chaos # break it on purpose
Limits
What this is not (yet)
- One machine. High availability, multi-AZ Kafka and Trino clusters are designed in scale.md, not run here.
- Synthetic data. 20,000 generated customers and scripted calls. Nothing is real.
- Speech-to-text is simulated. Transcripts arrive on Kafka as a contact-centre platform would send them.
- Rules-first understanding. Claude is used only when an API key is set: to label caller lines and to draft the after-call note from what was said (identity details redacted, never lakehouse records). Its output is bound by the same policy and grounding checks, under a daily spend cap.
- Point lookups go through Trino. At 50k colleagues the recommended design adds a serving store that uses the same OPA policies.