Packages

Framework-agnostic Elixir CDC consumer for Postgres logical replication (pgoutput) with zero-loss delivery to a pluggable sink — exactly-once for transactional sinks, at-least-once (duplicate-bounded) for non-transactional sinks.

Retired package: Removed QueryBuilder.table_columns/0 used by ash_replicant; use 1.5.1

Current section

Files

Jump to
replicant README.md
Raw

README.md

# Replicant

A framework-agnostic Elixir CDC consumer for Postgres logical replication
(`pgoutput`), delivering committed row changes to a pluggable **sink** with
**zero data loss**: the replication slot advances only after the sink has
durably persisted the transaction.

Replicant is **tenant-blind and classification-blind** — the reliable CDC
consumer sibling to [`arcadic`](https://github.com/baselabs/arcadic).
Multitenancy, classification, and Ash resources live one layer up, in the
[`ash_replicant`](https://hex.pm/packages/ash_replicant) sink adapter.

> **Status:** 1.4.1 is the current release (see CHANGELOG `[1.4.1]`). It corrects
> the records and test geometry behind the plugin decoders' parity claim, adds the
> 15 plugin lane, and fixes the third parallel RI-FULL key-flag path — the plugin
> decoders themselves arrived in 1.4.0: read a
> PostgreSQL 9.6 to 15 server through its existing `pglogical` or `wal2json`
> output plugin with the same sink contract, watermark, and halts
> (ADR-0009; see CHANGELOG `[1.4.0]`). 1.3.0 hardened the value layer every sink
> receives — multidimensional arrays of every casted type, locale-honest `money`,
> lossless `timetz`, signed `type_modifier`, a raising-free `lsn_from_string/1`,
> and stricter malformed-frame decoding (ADR-0008).
> Replicant owns
> the replication slot via `Postgrex.ReplicationConnection`, acks only after the
> sink durably commits (ack-after-checkpoint), halts fail-closed on slot
> invalidation, and is proven by a real-PG16 crash-injection suite
> (loss = 0, effect-dup = 0). The 1.0 contract includes initial snapshot/backfill (incl. a
> resumable incremental mode), a lib-owned checkpoint
> store for non-transactional sinks, batched checkpointing, sink-owned atomic
> batch delivery, in-progress-transaction streaming, consumer-side disk
> spill for oversized transactions, multi-publication per pipeline, and
> logical-decoding messages have all shipped. See "How it streams" below.

## Highlights

- **Sink-owned, transaction-granularity exactly-once** — the unit of delivery
  and of the watermark is the *transaction*, keyed by its single `commit_lsn`
  (every row in a pgoutput proto-v1 transaction shares one commit LSN). A sink
  skips any transaction whose `commit_lsn <= checkpoint` and upserts rows by
  table PK; that is at-least-once plus an idempotent sink, which is the only
  honest way to reach exactly-once without two-phase commit.
- **Value-free errors, logs, and telemetry** — every row value is assumed to
  be PII or a secret. Decode failures are caught and scrubbed into a
  `Replicant.Error` that never carries raw WAL bytes; telemetry metadata is
  allowlisted to LSNs, table/slot names, counts, durations, booleans, and error
  classes (the full event and halt-reason reference lives in
  [`usage-rules.md`](usage-rules.md)).
- **Identifier-validated SQL** — slot and publication names pass through
  `Replicant.Identifier.validate/1` (a strict Postgres-identifier allowlist)
  before they reach SQL, closing the raw-interpolation surface in the
  upstream parser this library vendors from.
- **TOAST-sentinel aware** — an UPDATE that doesn't touch a TOASTed column
  sends a sentinel, not the value. Replicant surfaces it as a first-class
  `unchanged: [col]` list on `Replicant.Change`, so a sink knows exactly which
  columns to leave untouched on upsert, instead of overwriting them with a
  placeholder.
- **Fail-closed on destructive schema drift** — a replica-identity change or a
  dropped column is classified `:destructive` and halts, rather than silently
  emitting incomplete or misattributed rows. A sink that can adapt its own
  schema implements the optional `handle_schema_change/2` callback to accept or
  veto destructive changes at the migration window instead of halting.
- **Actual replication-session identity** — `IDENTIFY_SYSTEM` runs on the exact
  replication connection before checkpoint lookup. A source-aware sink can
  accept or reject `%Replicant.SessionIdentity{}` synchronously on every connect
  and reconnect; no separate-connection preflight is treated as authoritative.
- **Column names stay strings** — never `String.to_atom`, so a wide or
  attacker-influenced schema cannot exhaust the atom table.

## LSN representation

A Postgres LSN is exposed as a single `non_neg_integer` — the 64-bit value
`(xlog_file <<< 32) ||| xlog_offset` — so that ordinary integer comparison is
correct WAL ordering, and the same value feeds the wire-level standby status
update:

```elixir
Replicant.lsn_to_string(0x16E3778)
#=> "0/16E3778"
```

Use `Replicant.lsn_to_string/1` for display; LSNs are WAL positions, not row
data, so they are permitted in telemetry metadata. The exactly-once watermark
check is plain integer comparison: `txn.commit_lsn <= checkpoint`. The inverse
`Replicant.lsn_from_string/1` returns `{:ok, lsn} | {:error, :invalid_lsn}`
(since 1.3.0 — it never raises, so a caller's input can never reach a crash
report).

## What your sink receives (casting)

`%Change{}.record` values are cast from Postgres's text output. The full
contract is [ADR-0008](docs/adr/0008-casting-lenient-value-preserving-fallback.md);
the three rules that surprise people:

- **Multidimensional arrays nest** — `numeric[][]`, `timestamptz[][]`,
  `jsonb[][]` and every other casted array type deliver nested lists with the
  same per-element semantics as their scalar clause (since 1.3.0; `NULL`
  elements are `nil` at any depth). `interval[]`/`timetz[]` deliver raw-string
  elements, mirroring their scalar clauses.
- **`money` is locale-honest** — money output follows the server's
  `lc_monetary`. The strict C/en-US shapes deliver a `Decimal`; anything else
  (e.g. `de_DE` `"1.234,56"`) delivers the original string — never a
  silently-wrong Decimal (since 1.3.0).
- **`timetz` delivers the raw string** — fractional seconds and offset
  preserved; there is no Elixir type for time-with-offset (same as `interval`).

Everything else is lenient-by-default: a value that fails its type's parse is
delivered as the original string; only genuinely-malformed input (a malformed
`numeric`, non-hex `bytea`) raises into the value-free decode boundary and
halts `:decode_failure`.

## How it streams

A running pipeline is two processes under a `:one_for_all` supervisor (three in
lib/checkpoint-store mode, which adds `Replicant.CheckpointStore` as the first child):

- **`Replicant.Connection`** (`Postgrex.ReplicationConnection`) owns the
  replication slot and the socket. It answers a keepalive with the **last
  durably-checkpointed LSN** while a published transaction is in flight (never
  advancing the slot past un-persisted data). An idle **state-mirror** pipeline
  advances the slot to the server WAL position so a quiet-but-filtered
  publication does not pin WAL; an `:append_log` pipeline deliberately retains
  the durable checkpoint so a reused origin ahead of it remains a detectable
  gap. It decodes each WAL message behind the
  value-free boundary and forwards it to the assembler — it never runs the sink,
  so it is always free to answer keepalives. It advances the ack asynchronously
  when the sink signals a durable commit, and halts fail-closed on slot
  invalidation (`wal_status = 'lost'` / `conflicting`), a decode failure, or a
  sustained sink-lag backlog (the bounded in-flight window).
- **`Replicant.AssemblerServer`** applies the sink synchronously, off the
  keepalive path, and halts fail-closed on a destructive schema change or a
  sink write fault.

Because the ack reports the durable checkpoint while a transaction is in flight
(advancing over filtered WAL only for an idle state mirror, where it carries no
publication data),
a crash between dispatch and persist re-delivers from the durable `confirmed_flush`
and the idempotent sink dedups — the exactly-once seam that `walex`'s
fire-and-forget `wal_end + 1` ack does not have.

## PostgreSQL version support

Replicant is **tested on PostgreSQL 9.6, 12, 15, 16, 17, and 18** — the CI matrix runs
the full suite against all six majors (Docker-only, `wal_level=logical`; the 9.6 and 12
rows carry the `pglogical` and `wal2json` output plugins for the decoder feature below).
Capabilities are gated by the server's `server_version_num`, so a single build runs
correctly across the range:

| Capability | PG9.6 | PG12 | PG15 | PG16 | PG17 | PG18 |
|---|:---:|:---:|:---:|:---:|:---:|:---:|
| Logical streaming, snapshot, checkpoint, exactly-once (decoder plugins) | ✅ `pglogical`/`wal2json` | ✅ all three decoders | ✅ | ✅ | ✅ | ✅ |
| pgoutput (`decoder: :pgoutput`, the default) | ❌ refused at start | ✅ | ✅ | ✅ | ✅ | ✅ |
| Slot-invalidation columns queried | none (pre-13 tier) | none (pre-13 tier) | `wal_status` | `+ conflicting` | `+ invalidation_reason, synced` | same as 17 |
| Failover slots (`failover: true`) | ❌ rejected | ❌ rejected | ❌ rejected | ❌ rejected | ✅ | ✅ |

### Reading pre-15 servers through pglogical or wal2json (ADR-0009)

A server that has no `pgoutput` (9.6) or sits outside the tested 15-18 range (10-14) can
be read through its existing logical-decoding output plugin instead:

    Replicant.start_link(
      connection: [hostname: "legacy.internal", ...],
      slot_name: "replicant_orders",
      decoder: :wal2json,                      # or :pglogical
      tables: [{"public", "orders"}],          # wal2json: the table list
      # replication_sets: ["default"],         # pglogical: the replication sets
      sink: MySink,
      go_forward_only: true
    )

The same `Replicant.Sink` contract, `commit_lsn` watermark, checkpoint modes and halt
semantics apply; a fixture transaction delivers byte-identically (after LSN, xid and
timestamp normalization) across pgoutput and both plugins — compared on one server by
CI's 12 plugin row (pgoutput vs `pglogical` vs `wal2json` on the 12 primary, plus both
plugins again on a real 9.6 secondary started beside it) in
`test/integration/decoder_parity_test.exs`; the truncate + transactional-message
variant runs the same way on any 15+ server carrying the plugins (ADR-0009 records the
one such run and why CI wires none). The same 9.6 secondary also carries a connected
snapshot leg (`test/integration/plugin_snapshot_pg96_test.exs`): PostgreSQL 9.6 exports
its snapshot name in a two-part form (`00004E57-1`; 10+ export three parts), and a
`snapshot: true` back-fill adopts it, handing off to streaming gap-free and dup-free
while a writer commits across the handoff.

**Keyless tables never stream silently missing updates.** wal2json drops an
update/delete on a table with no replica-identity index and `REPLICA IDENTITY ≠ FULL`
with only a server-side warning — no wire signal exists. Replicant therefore halts
fail-closed at start (`{:decoder, :table_keyless}`) for such a configured table. A
genuinely insert-only table opts in explicitly:

      decoder: :wal2json,
      tables: [{"public", "audit_events"}],
      allow_keyless_tables: true   # insert-only semantics: U/D on keyless tables are
                                   # invisible to the stream, by your own declaration

**Dropped columns halt fail-closed on every decoder.** pgoutput and pglogical re-emit
relation metadata after DDL, so a dropped column classifies `:destructive` immediately.
wal2json cannot re-emit (format 2 carries no relation messages), so Replicant splits the
ambiguous wire absence two ways: an INSERT missing a cached column, or an UPDATE missing
a fixed-width column (int/float/bool/date/time/timestamp/interval/uuid — never stored
out-of-line), is a drop and halts immediately; the remaining case (a dropped TOASTable
column on an update-only table, wire-identical to the unchanged-TOAST sentinel) is
bounded by a periodic catalog re-read — `schema_check_interval` (default `30_000` ms,
wal2json-only) — which routes the subset relation through the same destructive
classification. See `usage-rules.md` for the full halt table.

The slot-invalidation query selects only the columns that exist on the connected major
(`conflicting` was added in PG16, `invalidation_reason`/`synced` in PG17), so it never errors
on an older server. On PG17+ Replicant reads the authoritative `invalidation_reason` column (a
superset of the PG15/16 signals) and supports **failover slots** for HA.

### Failover slots (PG17+)

Pass `failover: true` to `Replicant.start_link/1` to create the replication slot with the
`FAILOVER` option, so PostgreSQL syncs it to physical standbys:

    Replicant.start_link(
      connection: [hostname: "primary.internal", ...],
      slot_name: "replicant_orders",
      publication: "orders_pub",
      sink: MyApp.OrdersSink,
      failover: true            # PG17+ only; on PG15/16 the pipeline halts {:config, :failover_unsupported}
    )

After a failover, repoint the connection at the promoted primary — the slot already exists
there with its `confirmed_flush` position, so replication resumes with zero loss. Slot sync
(`sync_replication_slots`), promotion, and DNS/endpoint changes are operator concerns.
**Do not point Replicant at an unpromoted standby's synced slot** — it cannot be consumed
until promotion, and Replicant halts fail-closed (`{:slot_synced_unpromoted}`) rather than
looping.

## The 5 critical rules (see [`docs/INVARIANTS.md`](docs/INVARIANTS.md) for the full text)

1. **No row value in an error, log, or telemetry event.**
2. **Validate identifiers** before they reach SQL.
3. **Exactly-once is at-least-once + a transaction-watermark-idempotent
   sink** — never claim a naked exactly-once.
4. **Unchanged TOAST is a sentinel, not a value** — never overwrite it.
5. **Stay tenant-blind** — multitenancy and classification live in
   `ash_replicant`, never here.

## Installation

```elixir
def deps do
  [
    {:replicant, "~> 1.0"}
  ]
end
```

Contributors working from a source checkout can point Mix at the local checkout
via a path override, as shown in the Livebook.

## Interactive tour (Livebook)

[![Run in Livebook](https://livebook.dev/badge/v1/blue.svg)](https://livebook.dev/run?url=https%3A%2F%2Fraw.githubusercontent.com%2Fbaselabs%2Freplicant%2Fmain%2Fnotebooks%2Fgetting_started.livemd)

[`notebooks/getting_started.livemd`](notebooks/getting_started.livemd) is a
runnable, self-verifying tour: it starts a live pipeline against a local Postgres,
streams `INSERT`/`UPDATE`/`DELETE` through a tiny sink you can watch, then
demonstrates the unchanged-TOAST sentinel, transaction-granularity exactly-once,
snapshot/backfill, and logical-decoding messages. Click the badge to open it in
[Livebook](https://livebook.dev), or read it rendered on
[HexDocs](https://hexdocs.pm/replicant/getting_started.html). The notebook's code
is executed against live PostgreSQL 15/16/17/18 on every CI run
(`test/integration/livebook_getting_started_test.exs`), so it never drifts from the
library.

## Reference example (docker)

[`examples/replication_pipeline`](examples/replication_pipeline/README.md) is the deployment
shape consumers actually build, as one `docker compose up`: a Postgres source
(`wal_level=logical` + publication) → replicant as an OTP release in one container
→ a Postgres destination, with an idempotent upsert-by-PK replica (unchanged-TOAST
aware), a value-free receipts ledger, session-identity binding, and a durable
commit-LSN checkpoint. It is a **go-forward change replicator** — pre-existing
source rows are not backfilled; `snapshot: true` is the one-flag alternative
(demonstrated interactively in the Livebook tour above). The `reference-example`
CI job runs the example's own static gates (format, compile-warnings, credo,
dialyzer), then builds the stack every push and proves delivery, TOAST-sentinel
survival, and restart-resume with zero duplicate effects — the public sink
API's canary.

## Usage

Start a pipeline against a standby with `Replicant.start_link/1`, pointing it at
a sink that implements `checkpoint/0` + `handle_transaction/1`:

```elixir
Replicant.start_link(
  connection: [hostname: "standby.internal", port: 5432, username: "u",
               password: "p", database: "orders", ssl: true],
  slot_name: "replicant_orders",
  publication: "orders_pub",
  sink: MyApp.OrdersSink,
  go_forward_only: false
)

defmodule MyApp.OrdersSink do
  @behaviour Replicant.Sink

  @impl true
  def handle_session_identity(identity, %{slot_name: slot, publication: publications}) do
    # Before checkpoint/0 runs: bind or compare {identity.system_identifier,
    # identity.database, slot, publications}. Return :ok only when safe to resume.
    :ok
  end

  @impl true
  def checkpoint, do: {:ok, MyApp.Repo.last_committed_lsn()}

  @impl true
  def handle_transaction(%Replicant.Transaction{commit_lsn: lsn} = txn) do
    # In ONE DB transaction: skip if lsn <= checkpoint, else upsert txn.changes
    # by table PK and persist lsn as the new checkpoint. Then:
    {:ok, lsn}
  end
end
```

**Start modes.** A `:state_mirror` sink starting from an empty checkpoint must declare
its intent — `go_forward_only: true` (stream only new changes), or `snapshot: true`
(**backfill** the current state, then hand off to streaming at the snapshot LSN with zero
gap and zero duplication). A non-empty checkpoint simply resumes. `snapshot: true`
requires the sink to also implement `handle_snapshot/2` (batch upsert; clear the table on
`first_for_table?`) and `handle_snapshot_complete/1` (durably persist the handoff
checkpoint); a mid-snapshot crash halts fail-closed (`:snapshot_incomplete`) for an
operator to drop the slot and retry.

**Snapshot columns.** Both snapshot modes read only the published column set, in
physical table order. On PostgreSQL 15+, pgoutput uses each requested publication's
effective column list; older servers have no column lists. Identical lists for a
shared table are accepted; differing lists, including an unrestricted publication
combined with a partial list, halt `:snapshot_column_list_conflict` before any
snapshot rows are read or delivered. pglogical applies `set_att_list` from the
requested local replication sets and uses the same conflict refusal. wal2json has
table filters but no publication column lists, so it retains all configured-table
columns. A snapshot selects the same columns and rows as the change stream.
Snapshots are INSERT-equivalent. Both checkpoint modes and all incremental paths
apply server row filters: PG15+ publication filters combine with OR only among
publications that publish INSERT; an unfiltered INSERT publication permits all
rows. With no INSERT publication, the snapshot delivers no rows. pglogical 2.4.8
combines non-NULL set filters with AND, even beside an unfiltered set.
A table with columns but no publishable projection halts with
`:snapshot_no_publishable_columns` and the table name. A genuinely zero-column
table, including one whose columns were all dropped, keeps the empty-reset
behavior of 1.4.2. PG12–15 discovery excludes generated columns; PG16+ uses the
publication's effective column list, including PG18's opt-in stored generated
columns. These repairs were verified live on PostgreSQL 18.6; PG12–15 selection
was verified offline, with a version-gated live regression for those servers.
[ADR-0002](docs/adr/0002-multi-publication-per-pipeline.md) documents PG15's
additional stream refusal when dropped columns prevent full-list collapse.

Key columns are never added to that projection. PostgreSQL refuses UPDATE/DELETE
when the published list omits replica-identity columns; INSERT-only lists can omit
them. If the complete primary key is unavailable, incremental snapshots use the
existing whole-table window instead of keyset pagination. Sinks must support that
keyless record shape; Replicant does not manufacture a missing key.

**Resumable incremental snapshot.** For large tables where an all-or-nothing `snapshot: true`
(a crash at 99% of a multi-hour COPY redoes everything) is the adoption blocker, opt into the
**incremental** mode instead:

```elixir
snapshot: [mode: :incremental, chunk_rows: 1000, max_pending_chunks: 4]
```

It creates a **durable** slot and streams immediately, while a linked reader backfills the
publication's tables in PK-ordered keyset **chunks** on its own connection — each chunk
consistency-bracketed by read-only LSN watermarks and collision-corrected against the live
stream (a concurrent write to a backfilling row wins; the stale chunk row is dropped). Progress
is durable per chunk, so a crash, halt, or reconnect **resumes from the last applied chunk**
rather than restarting. Chunks arrive through the same `handle_snapshot/2` callback; sink-owned
mode gives effect-once chunks, lib mode dup ≤ 1 chunk (never loss). Keyed drop-cap contention and
the PK-less whole-table fallback both halt `:snapshot_table_contended` after three
contention-discarded attempts; reconnects do not consume that reader-local budget.
`snapshot: true` remains the point-in-time option and is untouched.

A sink-owned adapter that durably arms the attempt before the first chunk returns
`{:ok, :backfill_pending}` from `snapshot_progress/0`. If the process dies after slot creation but
before a chunk token commits, the restarted pipeline reads the live slot origin and starts a fresh
reader at the safe floor instead of abandoning the backfill for stream-only delivery. Lib mode
persists and consumes the equivalent marker internally.

**Lib-owned checkpoint (non-transactional sinks).** Pass a `:checkpoint_store`
(`[connection: <postgrex opts>, table: "replicant_checkpoints"]`) to flip the pipeline
into **lib mode**: the library writes the checkpoint to a durable Postgres table **after**
the sink persists (checkpoint-after-persist), so a non-transactional sink (files, S3,
Kafka, external APIs) needs to implement only `handle_transaction/1`. The guarantee is
**at-least-once, duplicate bounded to one transaction, never loss** — not effect-once (a
non-transactional sink cannot dedup). A store outage (connect-read or mid-stream write) is
bounded: the pipeline retries `max_retries` times (default 5) `retry_backoff_ms` apart
(default 1000 ms — ~5s of outage tolerated) then halts fail-closed; a permanent fault
(schema mismatch / invalid config) halts immediately. With `snapshot: [mode:
:incremental]` the store additionally owns a `progress_table` (default
`"replicant_snapshot_progress"`) carrying the opaque backfill progress token.
PK bounds are row data; `row_filter` text is publication configuration and can
include literal constants. Persist the token without logging it. Both table names
are configurable.

**Persistent replication-command errors (`max_command_retries`, default 5).** A replication
command that fails *before the stream starts* — e.g. `CREATE_REPLICATION_SLOT` when the
server's `max_replication_slots` are exhausted, a second consumer already holding the slot,
or a forward-incompatible result shape — otherwise reconnects forever (`auto_reconnect`).
Replicant bounds this: after `max_command_retries` failed connect cycles without the stream
establishing, the pipeline **halts fail-closed and stays idle** instead of livelocking, and
emits `[:replicant, :connection, :command_error_halt]` (metadata: `attempt`, `max_retries`,
`slot_name` — value-free, no error content, Critical Rule 1). Set `max_command_retries: 0`
to halt on the first such fault. The bound is a *cycle count*, not a wall-clock time (the
failing reconnect is a hot loop). Only errors *before the first replication frame* are
bounded — once the stream is flowing, a later transient outage resets the counter and
still self-heals, and a server that is simply down (connection refused) keeps retrying
untouched, exactly as before.

**Sink-owned atomic batch delivery (transactional sinks).** For a **transactional sink** that
can persist multiple rows + checkpoint in one database transaction, pass a top-level
`batch_delivery: [max_transactions: 100, max_delay_ms: 1000]` to accumulate committed
transactions and deliver them as a batch:

```elixir
Replicant.start_link(
  connection: [hostname: "standby.internal", port: 5432, username: "u",
               password: "p", database: "orders", ssl: true],
  slot_name: "replicant_orders",
  publication: "orders_pub",
  sink: MyApp.OrdersSink,
  batch_delivery: [max_transactions: 100, max_delay_ms: 1000],
  go_forward_only: false
)

defmodule MyApp.OrdersSink do
  @behaviour Replicant.Sink

  @impl true
  def checkpoint, do: {:ok, MyApp.Repo.last_committed_lsn()}

  @impl true
  def handle_batch(transactions) do
    # In ONE DB transaction: skip any commit_lsn <= checkpoint, else upsert all
    # rows from all transactions by table PK, and persist the batch's highest
    # commit_lsn as the new checkpoint. Then:
    {:ok, List.last(transactions).commit_lsn}
  end
end
```

The batch flushes when it reaches `max_transactions` transactions, after `max_delay_ms`
milliseconds idle, or when the batch's WAL span (LSN-span lag) hits an auto-derived
safety cap (derived from `:max_inflight_lag`). Because the rows + checkpoint write is
atomic, effect-once is **preserved** (dup=0, loss=0) — stronger than lib-mode's
`checkpoint_store: [batch: …]` which is per-transaction delivery with batched checkpointing.
The sink must implement both `checkpoint/0` (resume) and `handle_batch/1`; it cannot
use `checkpoint_store`, and any `handle_transaction/1` implementation is ignored.
Emits `[:replicant, :sink, :batch_committed]` telemetry once per flush.

**Consumer-side disk spill (oversized transactions).** By default a single in-progress streamed
transaction is bounded by the in-flight window: one larger than `max_inflight_lag` (default
64 MiB) halts fail-closed. Opt into **disk spill** to reassemble such a transaction partly on
disk and still deliver it effect-once:

```elixir
Replicant.start_link(
  connection: [...], slot_name: "replicant_orders", publication: "orders_pub",
  sink: MyApp.OrdersSink, go_forward_only: false,
  max_inflight_lag: 64 * 1024 * 1024,
  streaming: [
    max_concurrent_txns: 64,
    spill: [dir: "/var/lib/replicant/spill", max_spill_bytes: 1024 * 1024 * 1024]
  ]
)
```

A transaction whose resident bytes cross `max_inflight_lag` spills its oldest changes to a per-txn file
under `dir`; at commit it is delivered as a **lazy, single-pass, disk-backed** `%Transaction{}` whose
`changes` streams the spilled frames + the resident tail. There are **two ceilings**: the resident RAM
bound `max_inflight_lag` (the spill trigger) and the disk bound `max_spill_bytes` (a transaction
exceeding it halts `:spill_exhausted`). Defaults: `max_spill_bytes` is `16 × max_inflight_lag`; `dir`
is a `0700` subdir of the OS temp dir.

**Delivery obligation.** A spilled transaction's `changes` is a **single-pass** `Enumerable` valid only
*during* the `handle_transaction/1` (or `handle_batch/1`) call — iterate it with `Enum`/`Stream` and do
**not** call `length/1`, `Enum.to_list/1`, or re-iterate it (any of which forces the whole transaction
back into RAM, defeating spill), and do not retain it past the call. The usual List-backed `changes`
still works exactly as before; only an oversized spilled transaction delivers the lazy form.

**Operator guidance.** Spill files are ephemeral non-fsync'd scratch (`0600`, value-free on fault),
deleted on commit/abort/reset/halt and swept per-slot on (re)connect. Replicant does **not** encrypt
them — if the source rows are sensitive, point `dir` at an encrypted/secure volume; a custom persistent
`dir` is yours to clean on decommission (the default OS temp dir is cleared by the OS).

**Multi-publication per pipeline.** `publication:` accepts a single name (the default) **or a list**.
Pass a list to stream the union of several publications through one slot — every name is validated
and deduplicated, and overlapping tables are de-duped by pgoutput on the wire:

```elixir
Replicant.start_link(
  connection: [...],
  slot_name: "replicant_orders",
  publication: ["orders_pub", "audit_pub"],   # was: publication: "orders_pub"
  sink: MyApp.OrdersSink,
  go_forward_only: false
)
```

Every requested publication must exist; a missing one **halts fail-closed** at connect time rather
than silently streaming the subset (a `START_REPLICATION` that names a missing publication would
otherwise stream only the found subset).

**Logical-decoding messages (`pg_logical_emit_message`).** Opt into Postgres logical-decoding
**messages** with `messages: true` to receive outbox-pattern / heartbeat rows emitted via
`pg_logical_emit_message`. The durability guarantee depends on the message's transactional flag
(stated honestly in the `handle_message/2` docs):

- **Transactional** (`transactional => true` to `pg_logical_emit_message`) — rides
  `%Transaction.messages` and inherits the transaction path's `commit_lsn` effect-once dedup
  (effect-once, just like the row changes in the same txn).
- **Non-transactional** (`transactional => false`) — arrives standalone via the
  `handle_message/2` sink callback and is **at-least-once** (no dedup key; duplicates are possible
  on reconnect). This is the honest guarantee — do not claim effect-once for it.

```elixir
Replicant.start_link(
  connection: [...],
  slot_name: "replicant_orders",
  publication: "orders_pub",
  sink: MyApp.OrdersSink,
  messages: true,                            # opt-in; sink MUST implement handle_message/2
  go_forward_only: false
)

defmodule MyApp.OrdersSink do
  @behaviour Replicant.Sink

  @impl true
  def checkpoint, do: {:ok, MyApp.Repo.last_committed_lsn()}

  @impl true
  def handle_transaction(%Replicant.Transaction{commit_lsn: lsn, messages: msgs} = txn) do
    # msgs :: [Replicant.Decoder.Messages.Message.t()] — transactional messages ride the
    # txn and inherit its commit_lsn effect-once dedup. Persist them with the txn's rows.
    {:ok, lsn}
  end

  @impl true
  # context :: %{lsn: lsn}; non-transactional messages ONLY (transactional ride handle_transaction/1).
  # At-least-once: a reconnect can re-deliver — your effect must be idempotent. The message's
  # `content`/`prefix` are USER BYTES (Critical Rule 1: never log or surface them in telemetry).
  def handle_message(%Replicant.Decoder.Messages.Message{} = msg, %{lsn: lsn}) do
    :ok
  end
end
```

`messages: true` requires the sink to implement `handle_message/2`; a sink missing it is rejected at
start (`:messages_unsupported`) rather than silently dropping non-transactional messages later.

### Slot origin for go-forward append consumers

A go-forward **append-log** sink can learn the LSN its slot streams from — the *consistent-point
origin* — by implementing the optional `handle_slot_origin/2` callback. It fires on every connect and
reconnect, before `START_REPLICATION`, for both a freshly-created and a reused slot:

```elixir
@impl true
def handle_slot_origin(origin, %{slot_name: slot, reused?: reused?}) do
  # reused? == false → origin is the CREATE_REPLICATION_SLOT consistent_point (a new slot).
  # reused? == true  → origin is max(durable checkpoint, confirmed_flush_lsn), the
  #                    effective START_REPLICATION origin for a resumed slot.
  # Return :ok to proceed; return {:error, _}/raise to halt fail-closed (e.g. on a gap
  # past your last appended LSN) instead of silently skipping WAL.
  :ok
end
```

`context` is value-free (the slot name and a boolean — never row bytes). A sink that does not
implement the callback is unaffected: no extra query, unchanged streaming. If PostgreSQL cannot
supply a valid logical-slot origin, Replicant halts before the callback and streaming with
`:slot_origin_unavailable`; it never reports a fabricated zero origin.

An `:append_log` sink — declared by implementing the optional `sink_kind/0`
callback — never uses the filtered-WAL idle advance. This is what
makes `origin > durable checkpoint` unambiguously an out-of-band gap instead of
a legitimate keepalive side effect. On a quiet append publication in a busy
cluster, publish a normal heartbeat transaction (a row or admitted logical
message) to advance the durable checkpoint and release retained WAL.

## Development

```bash
mix deps.get
mix test
mix quality   # format --check-formatted + credo --strict + dialyzer
```

The published redaction, identifier-validation, and tenant-blind invariants live
in [`docs/INVARIANTS.md`](docs/INVARIANTS.md). Repository contributors also follow
the checkout's repository-specific contributor contract.

## Roadmap

The v1 CDC core and every delivery slice have shipped and are closeout-reviewed
against a real-PG16 crash-injection suite:

- **Offline core** — decode / assemble / validate / redact behind the value-free boundary.
- **Live streaming + exactly-once** — the `Postgrex.ReplicationConnection` that owns the slot with ack-after-checkpoint, slot-invalidation fail-closed halt, and the bounded in-flight window (loss = 0, effect-dup = 0).
- **Initial snapshot / backfill** — `EXPORT_SNAPSHOT` → `COPY` → stream-at-snapshot-LSN, gap-free and dup-free.
- **Resumable incremental snapshot** — `snapshot: [mode: :incremental]`: a durable-slot, PK-ordered keyset-chunk backfill interleaved with the live stream and collision-corrected against it, resuming from the last applied chunk after a crash/halt/reconnect (sink-owned effect-once chunks / lib dup ≤ 1 chunk; PK-less whole-table fallback).
- **Lib-owned checkpoint store** — a durable Postgres checkpoint written *after* persist for **non-transactional** sinks (at-least-once, dup bounded to one transaction, never loss), with bounded retry-then-halt on store faults.
- **Batched checkpointing (lib mode)** and **sink-owned atomic batch delivery** — amortize the checkpoint write / the sink's own commit across a batch of transactions.
- **In-progress-transaction streaming** (`pgoutput` v2) and **consumer-side disk spill** — reassemble and deliver a transaction larger than memory, effect-once, instead of halting.
- **Multi-publication per pipeline** — `publication: [p1, p2]` streams the union of several publications through one slot, with fail-closed missing-publication detection.
- **Logical-decoding messages** — opt-in `messages: true` delivers `pg_logical_emit_message` payloads: transactional messages ride `%Transaction.messages` (effect-once); non-transactional via `handle_message/2` (at-least-once).

The sibling libraries live one layer up from this tenant-blind core:

- **[`ash_replicant`](https://github.com/baselabs/ash_replicant)** — the Ash /
  multitenancy / classification sink adapter. Its coordinated 1.0 release will
  require Replicant 1.x so the actual-session identity check cannot be resolved away.
- **[`ash_onetime`](https://hex.pm/packages/ash_onetime)** — the authoritative
  idempotency-key / one-time-nonce admission layer for Ash/Postgres actions.
  AshReplicant uses idempotency, not nonce rejection, for retryable
  non-transactional logical-message effects. Replicant transaction replay remains
  governed by the durable commit-LSN watermark; it is not replaced by AshOnetime.

## Credits

- [**walex**](https://github.com/cpursley/walex) — the `pgoutput` byte parser,
  OID-to-type database, type caster, and array parser this library vendors
  from (MIT). See `NOTICE` for the full attribution chain (cainophile,
  Supabase Realtime, epgsql).
- The `postgrex`/`ash_postgres` split that inspired `arcadic` and
  `ash_arcadic` also shapes the `replicant`/`ash_replicant` layering.

## License

MIT — see [LICENSE](LICENSE). Third-party attributions in [NOTICE](NOTICE).