ADR 0009: Data movement is a mirror of declared facts, not an embedded engine
- Status: accepted
- Date: 2026-09-28
Context
A tessera app increasingly needs to move data between systems: replicate a Postgres table into an analytics warehouse, stream change-data into a second database, or derive external rows into the app's own CQRS model. The obvious implementation is to embed a movement engine (a Temporal-driven orchestrator like PeerDB, a Kafka/Debezium pipeline, or a CDC client library with its own state machine) alongside the generated app.
That conflicts with the laws:
- Law 2 (Zero Tax) / ADR 0001 (closed vocabulary): an embedded engine is an escape hatch — an open, Turing-complete data pipeline with its own configuration, not a declared fact the plan owns.
- Law 3 (One Model) / ADR 0002 (One Resolve): if a runtime engine decides routes, checkpoints, and schemas, the "model" of the data flow lives in two places (the DSL and the engine), and they can drift.
- Law 4 (Thin Renderers) / ADR 0004: the flow's facts (which table, which keys, which columns map to which, the mode, the batch size) must be resolved exactly once, in the plan, and the renderer only prints.
The question is how to get PeerDB-grade sync (backfill + CDC, schema mapping, resumable checkpoints) while keeping the decision in the plan and the runtime a deterministic projection of it.
Decision
Two declarations, one resolved plan, a verbatim runtime.
Data movement is two new top-level tessera declarations — peer and
mirror — that resolve (fail-closed) in mosaic-core into PeerPlan and
MirrorPlan. The renderer emits a small, deterministic sync engine
(app/src/sync.rs) plus a verbatim runtime core
(app/src/syncrt.rs, copied byte-for-byte from crates/mosaic-sync). No
engine is embedded; the app is the mover.
1. peer — a named, kinded connection (the secret is an env-var name)
peers:
- name: legacy
kind: postgres # postgres | mysql | clickhouse | file
conn: MOAIC_PEER_LEGACY # an ENV VAR NAME, never a literal
A peer is a fact: name, kind, and conn. conn is the name of an
environment variable that holds the connection string; the value is read at
runtime. No credential ever lands in the model or the generated source (the
same law as resource.auth).
2. mirror — source → target with per-table mode
mirrors:
- name: orders_to_warehouse
source: legacy # a peer
target: warehouse # a peer, or the literal `model`
batch_rows: 500
tables:
- name: orders
keys: [id]
mode: cdc # cdc | snapshot | poll
where: "status != 'draft'"
map: { customer: customer_id } # target_col: source_col (rename)
- name: order_items
keys: [id]
mode: snapshot
A mirror names a source peer and a target (a peer or the reserved
model), a batch_rows (default 1000), an optional poll interval (for
mode: poll), and a set of tables. Each table declares its keys (the
primary key, required for cdc), its mode, an optional where row filter,
an optional map (a rename: target_col: source_col; unmapped columns
pass through), and — for target: model — into / into_delete (the CQRS
command the insert/update and the delete dispatch).
3. Modes are resolved to a concrete strategy in the plan
cdc(Postgres source): a keyset backfill to a cursor, then apgoutputlogical-replication stream (viapgwire-replication, the one external crate reused) that decodesR/I/U/Dmessages, normalizes each change, and checkpoints the commit LSN.snapshot: a one-shot keyset backfill (no stream).poll(any source): a watermark/keyset poll everypollinterval, for sources without a logical decoder (MySQL, or a chosen Postgres table).
snapshot/poll targets accept any peer; cdc requires a Postgres source
and a Postgres/model target. The plan rejects a cdc table without keys,
an unknown source/target, a model target without into, and a
non-Postgres source with mode: cdc — all fail-closed, like every other
resolver (ADR 0002).
4. Sinks are projections of the resolved change
Each decoded change is a normalized, named row (Op + RowMap + key
RowMap). The sink is a pure function of that row + the MirrorPlan:
- Postgres sink — writes to a
_mosaic_raw_<mirror>staging table, then aMERGEnormalizes it into the destination table (TOAST-safe: an unchanged cell is absent from the payload and theMERGEkeeps the stored value viaCASE WHEN _row ? col). Deletes set_mosaic_deleted = TRUE(soft delete). - ClickHouse sink —
INSERTinto aReplacingMergeTree(the_mosaic_synced_atcolumn is the version). - File sink — one JSON object per line (append-only archive).
- Model sink — the mosaic differentiator: the row is dispatched as a
command onto the app's own aggregate (
into/into_delete), so an external table becomes state the app already serves, guards, and projects. This is the "derive" that no external engine offers.
5. One external crate, reused byte-exactly
pgwire-replication (0.4, Apache-2.0/MIT, pure Rust, no libpq) is the single
reused dependency for the cdc path. The rest of the runtime is in-house
(crates/mosaic-sync), emitted verbatim into the app so the golden output
pins it byte-for-byte and there is no version drift between the crate and the
generated app. Candidates evaluated and rejected: deltaforge (a standalone
service, heavy rdkafka/deno_core tree — not an embeddable crate), d-engine
(not CDC — an embeddable Raft KV store / etcd alternative), dbmazz
(ELv2-licensed, so its code cannot be reused; its Sink trait and LSN
checkpoint concepts are adopted here instead).
6. State is a file the app owns
The sync engine persists its cursors, LSNs, watermarks, and per-mirror
last_error to a JSON file (MOAIC_SYNC_STATE, default sync-state.json)
and is resumable: on restart it resumes the backfill cursor / replication LSN
rather than re-reading. GET /api/sync reports per-mirror state;
POST /api/sync/{mirror}/once triggers a manual pass.
Consequences
- The data flow is a declared fact:
peer/mirrorresolve once in the plan (fail-closed), and the renderer only prints the engine + the verbatim runtime. No second source of truth, no embedded orchestrator (Laws 2–4). - Deriving external rows into the model (
target: model) is a first-class capability no external engine has: a legacy table becomes CQRS state with the app's guards, projections, REST, web, and site — for the price of amirror. - The runtime is a projection, so it is golden-pinned: a byte change to
crates/mosaic-syncor thesync.rscodegen fails the conformance suite until the golden is re-pinned (Law 6). - Secrets stay env-var names (Law 2 / the
resource.authprecedent); a credential can never be committed. - Cost: one new crate (
mosaic-sync) and two new declarations. Thecdcpath addspgwire-replication+tokio-postgresto the generated app'sCargo.tomlonly when a Postgres peer /cdcmode is declared (Zero Tax). MySQL addsmysql_asynconly when a MySQL peer is declared; ClickHouse uses HTTP/JSON (no new crate); the file sink usesstd.
Proof
mosaic-conformance::tests::examples_match_goldenspinsexamples/data-sync/golden/byte-for-byte (the fullsync.rs+ verbatimsyncrt.rs+ the gatedCargo.toml). A change to the codegen or the runtime that is not re-pinned fails CI.mosaic-syncunit tests (17) prove the pgoutput decoder byte-exactly:decode_real_capture_frames,decode_real_capture_delete_frame, anddecode_real_capture_update_framedecode real frames captured from a live PostgreSQL 15 stream (including the full-width deleteKtuple and the keyless update), plus chunk-spanning truncation and theto_changesmapping.mosaic-coreplan tests (data_mesh_tests, 15) prove the resolver is fail-closed: unknown peer/kind/mode,cdcwithoutkeys, amodeltarget withoutinto, and a non-Postgres source withmode: cdcall reject.- The dependency gate is proven by the golden
Cargo.toml(deps appear only for the declared kinds) and byrenderemitting them conditionally.