Open Lakehouse
ci passinge2e 69/69release signed

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.

~10 ssource commit to visible for a colleague (CDC)
73end-to-end checks on the full stack for every merge, from an empty machine
11/11components killed and recovered with no data lost
Alice · contact centre · MeridianLive call

Unrecognised payment located

£249.99 to QuickTech Electronics Online, 10:26 UTC, status pending. Freeze the card now, then raise a fraud claim.

⚡ 483 msCARD-FRAUD-01ID&V passed
Below the waterlinehow the platform answered
  1. Sourcecore banking Postgres commit → Debezium (WAL) → Kafka corebank.core.transactions
  2. Icebergsnapshot 6170615646954925454, committed by cdc-stream 35 s before the lookup
  3. Trinoas alice, pinned FOR VERSION AS OF that snapshot · 1 row
  4. OPArow filter brand IN ('Meridian') · no PII columns read, so nothing to mask
  5. 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.

The Live Call Assist console during a card-fraud call: transcript, guidance cards, customer panel, and the platform x-ray of the card that found the payment
Recorded from the running stack: a scripted caller, real CDC, real policy. It ends on the platform x-ray.

A colleague on a live fraud call

  1. The caller says their card was stolen. The assistant hears the intent and looks the caller up with Alice's own token.
  2. 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.
  3. The payment written to core banking a minute ago is already in silver, and the assistant finds it.
  4. Each card opens to show the snapshot, the Trino query, OPA's decision and the audit row behind it.
  5. The injection call ("ignore your instructions…") gets nothing. Customer text is data, never instructions.
CDC freshnessrow + column securityon-behalf-oftime travelaudit chain
# open http://localhost:8090, sign in as alice
make live-demo

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.

Architecture: core banking to Kafka and Spark into Iceberg; Trino with OPA serves the MCP gateway, analysts and dashboards Sources Streaming Lakehouse Serve · govern · consume

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.

BlueprintWhat it addsStatus
corePostgres, object storage, Polaris, Keycloak, OPA, Trino, MCP gateway. Always on, so identity, policy and audit are never optional.built
streamingKafka and the tenant reconcilerbuilt
orchestrationDagster, one code server per tenantbuilt
corebankChange 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 · aiGrafana and Prometheus, Marquez, the SQL workbench, Live Call Assistbuilt, 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.
tradesCoinbase WebSocket → Kafka → bronze, silver, rejects
card authsShadowTraffic → Kafka → bronze, silver, dead-letter topic
referenceECB FX rates and the FCDO sanctions list, daily
gold1-minute candles · card auths per day · sanctions hits
accessevery colleague reads gold · bronze and silver: platform admins

Read 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.

CapabilityHow it is provenResult
Fresh data for live callsWrite 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 brandalice, bob and carol run the same query and get different brands3/3 pass
Column masking by personaPhone partial for alice, clear for bob, NULL for carol; carol can't infer a masked column by filtering on it4/4 pass
Agents inherit the colleague's rightsToken exchange (RFC 8693) so Trino sees the colleague; the agent's answers carry alice's masks; an analyst's agent identifies no onepass
Reproducible answersEvery gateway answer is pinned to an Iceberg snapshot and carries OPA's decision and its audit row3/3 pass
Tamper-evident auditHash chain verified; DELETE on the audit log is refused; denied calls are audited toopass
Bad data never reaches silverODCS contracts, quarantine with reasons, WAP publish only after checks; late parents are retried, not dropped0 violations
Safe AI guidanceOffline eval gate: intent/vulnerability/risk recall, 8 scripted calls, zero ungrounded cards; injection rejected before any SQL8/8 · 0 ungrounded
Teams onboard without hand grantsThe 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 exceededpass
Raw data for platform admins onlyColleagues are refused bronze; the admin reads bronze but every whole-record payload is NULL, including in quarantine; contracts declare the tagpass
AI spend capped and visibleA 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 dashboardpass
OperableEvery scrape target up (Trino via a machine identity), 20 SLO and burn-rate rules loaded, dashboards from code, Dagster behind SSOpass
Supply chainActions pinned by SHA, gitleaks over full history, Trivy, SBOM + provenance, cosign keyless signatures on GHCR imagessigned

Chaos results

Kill each component. Measure what the colleague sees, and time to recover.

Time to recover in seconds for each of 11 components after kill -9
fails closed (no data rather than wrong data) keeps serving, freshness pauses isolated (console only)

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.

About half the contact centre is on shift at peak.
About 6 in the first minute, then about 1 a minute.
–simultaneous live calls
–lookups per second, steady
–lookups per second, call-start bursts

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.