Mosaic one model, many lenses

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 a pgoutput logical-replication stream (via pgwire-replication, the one external crate reused) that decodes R/I/U/D messages, normalizes each change, and checkpoints the commit LSN.
  • snapshot: a one-shot keyset backfill (no stream).
  • poll (any source): a watermark/keyset poll every poll interval, 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 a MERGE normalizes it into the destination table (TOAST-safe: an unchanged cell is absent from the payload and the MERGE keeps the stored value via CASE WHEN _row ? col). Deletes set _mosaic_deleted = TRUE (soft delete).
  • ClickHouse sink — INSERT into a ReplacingMergeTree (the _mosaic_synced_at column 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/mirror resolve 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 a mirror.
  • The runtime is a projection, so it is golden-pinned: a byte change to crates/mosaic-sync or the sync.rs codegen fails the conformance suite until the golden is re-pinned (Law 6).
  • Secrets stay env-var names (Law 2 / the resource.auth precedent); a credential can never be committed.
  • Cost: one new crate (mosaic-sync) and two new declarations. The cdc path adds pgwire-replication + tokio-postgres to the generated app's Cargo.toml only when a Postgres peer / cdc mode is declared (Zero Tax). MySQL adds mysql_async only when a MySQL peer is declared; ClickHouse uses HTTP/JSON (no new crate); the file sink uses std.

Proof

  • mosaic-conformance::tests::examples_match_goldens pins examples/data-sync/golden/ byte-for-byte (the full sync.rs + verbatim syncrt.rs + the gated Cargo.toml). A change to the codegen or the runtime that is not re-pinned fails CI.
  • mosaic-sync unit tests (17) prove the pgoutput decoder byte-exactly: decode_real_capture_frames, decode_real_capture_delete_frame, and decode_real_capture_update_frame decode real frames captured from a live PostgreSQL 15 stream (including the full-width delete K tuple and the keyless update), plus chunk-spanning truncation and the to_changes mapping.
  • mosaic-core plan tests (data_mesh_tests, 15) prove the resolver is fail-closed: unknown peer/kind/mode, cdc without keys, a model target without into, and a non-Postgres source with mode: cdc all reject.
  • The dependency gate is proven by the golden Cargo.toml (deps appear only for the declared kinds) and by render emitting them conditionally.