Replication Fan-Out
Replication Fan-Out
Multiplex one PostgreSQL replication slot to N downstream clients over gRPC. Clients connect, receive a consistent snapshot, then stream live changes — no need for each service instance to hold its own slot.
Architecture
Server configuration
The replication protocol is served by a single engine-global listener on its
own port (default 4002), separate from the OAM/Query port (4001). That one
listener serves every replication-fanout target on the server, routing each
request by table — see ADR-007.
The stock laredo-server binary starts it automatically whenever at least one
fan-out target is configured.
# Engine-global replication listener. Optional — when omitted, the port
# defaults to 4002 as soon as a replication-fanout target exists.
fanout {
grpc { port = 4002 }
}
tables = [{
source = pg
schema = public
table = config_document
targets = [{
type = replication-fanout
journal {
max_entries = 1000000
max_age = 24h
}
# In-memory snapshots: how often a base snapshot is cut and how many to keep.
snapshot {
interval = 5m
retention { keep_count = 5, max_age = 1h }
}
max_clients = 500
client_buffer {
max_size = 50000
policy = drop_disconnect # or slow_down
}
heartbeat_interval = 5s
}]
}]
Each target field maps onto a target/fanout.New
option; omit any to keep its default. The snapshot block governs the target's
in-memory snapshots (interval + retention) — durable on-disk snapshot stores
are not yet wired into the config. To let a client that has fallen behind the
hot journal catch up from cold object storage, register a read-only archive on
the replication service as shown under Cold-tier replay.
Client protocol
The client connects via a single server-streaming gRPC call:
- Handshake — client declares its state (fresh or has a local snapshot)
- Catch-up — server sends full snapshot or journal delta
- Live streaming — changes in real time
Fresh client: → FULL_SNAPSHOT → journal catch-up → live
Resuming client: → DELTA (journal entries from last sequence) → live
Failing-over client: → DELTA (journal entries from last source position) → live
Stale client: → FULL_SNAPSHOT (too far behind journal) → live
When the server is draining it sends a GoAway message; the client re-dials and
resumes on another instance by source position. See
Failover & zero-downtime deploys.
Go client
client, err := fanout.New(
fanout.ServerAddress("laredo-server:4002"),
fanout.Table("public", "config_document"),
fanout.LocalSnapshotPath("/var/lib/myapp/laredo-cache"),
fanout.WithIndexedState(
memory.LookupFields("instance_id", "key"),
memory.AddIndex("by_instance", []string{"instance_id"}, false),
),
fanout.ClientID("myapp-instance-abc123"),
)
client.Start()
client.AwaitReady(30 * time.Second)
row, ok := client.Lookup("inst_abc", "rulesets/default")
client.Listen(func(old, new laredo.Row) {
// react to changes
})
client.Stop() // saves local snapshot for fast restart
Subscription filtering
A client can ask for only the rows it cares about by attaching one or more
column predicates to its Sync request. The server applies them uniformly
across the snapshot, journal catch-up, and live phases, so the subscriber gets a
consistent slice of the table — and rows it filters out never cross the wire.
The common use is partition scoping: a single equality predicate on a
partition column gives each subscriber its own slice of a shared table without a
separate pipeline. This is the right tool for multi-tenant fan-out — one fan-out
target for an events table, each tenant subscribing with tenant_id equal to
its own id — with hard isolation, because other tenants' rows are never sent.
From the Go client
client := fanout.New(
fanout.ServerAddress("laredo-server:4002"),
fanout.Table("public", "events"),
fanout.WithFilterEquals("tenant_id", "acme"), // partition scope
// ...AND-combined with any further predicates:
// fanout.WithFilterPrefix("key", "rulesets/"),
// fanout.WithFilterIn("region", "eu", "us"),
)
The client's local replica then contains only matching rows; Count(), Get,
Lookup, All, Listen, and any secondary indexes operate over the filtered
slice.
Predicate types
| Client option / wire field | Matches |
|---|---|
WithFilterEquals(field, v) / equals |
the column equals v. Numbers compare numerically (a predicate of 42 matches an integer or float 42); strings and bools by value. |
WithFilterPrefix(field, p) / prefix |
the (string) column starts with p. |
WithFilterIn(field, v...) / in |
the column equals one of the values. |
Multiple predicates are AND-combined. A predicate is evaluated against the
post-change row for inserts and updates, and the pre-change row for deletes;
TRUNCATE always passes (it is structural). A missing or null column never
matches.
Caveats
- Filter columns must be stable for a row's lifetime. Subscription filtering
assumes a row never changes the partition it belongs to. If an update moves a
row from one partition value to another, a subscriber on the old value is not
told the row left, and a subscriber on the new value sees it as an update. For
partition keys like
tenant_idthis never happens; design filter columns to be immutable. - Filtered resume. A filtered client resumes from the last sequence it
received. If no matching change occurs for a long time while the journal
prunes past that point, a reconnect falls back to a (filtered) full snapshot —
correct, just without the delta optimization. Size
journal.max_entries/journal.max_agewith your sparsest partition in mind, the same way you size for your slowest consumer. - Filtering is about delivery, not memory. The fan-out target still holds the whole table in memory; filtering controls only what is sent to each subscriber (and the isolation between them).
Cold-tier replay
When a client's resume position predates the bounded change journal, the server
would normally fall back to a full snapshot of live state — losing the delta
and the history the client missed. If you also run
laredo-snapshotter to archive the table to object
storage (base snapshot + diffs + manifest), the server can instead replay from
that cold archive: it reconstructs catch-up from the newest base snapshot and
the diffs after it (or, when the client's position aligns to a diff boundary,
only the diffs), then hands off gaplessly to the hot journal and the live
stream. The client receives the ordinary wire messages, so no client change is
needed; the handshake mode is SYNC_MODE_REPLAY_ARCHIVE.
This is the tool for consumers that fall far behind the journal — an onboard or edge node reconnecting after a long outage, a regional replica catching up, a service restarting after downtime. They resume with the history they missed from cheap object storage, instead of re-downloading the whole live table.
Enabling it
From HOCON (stock laredo-server)
Add an archive block to the replication-fanout target. laredo-server builds
the reader and registers it for that table's (schema, table) automatically
(ADR-007,
EDR-0005):
tables = [{
source = pg
schema = public
table = events
targets = [{
type = replication-fanout
journal { max_entries = 1000000, max_age = 24h }
# Local archive — the offline-first case (an onboard node archives to local
# disk and resumes from it).
archive {
store = local
store_config { path = "/var/lib/laredo/archive/events" }
format = jsonl # jsonl | protobuf; default jsonl
key_prefix = "laredo/public.events/" # MUST match the snapshotter's prefix
}
}]
}]
For an S3 archive, set store = s3 and the bucket/region instead:
archive {
store = s3
store_config { bucket = laredo-archive, prefix = "laredo/", region = us-east-1 }
format = jsonl
key_prefix = "laredo/public.events/"
}
key_prefix must match the prefix laredo-snapshotter
wrote this table under — otherwise the reader finds no manifest and cold replay
declines to a live snapshot. Both local and s3 stores are built through the
shared snapshotter/destwire package, the same path the snapshotter writes
through. S3 uses the ambient AWS credential chain (environment, IRSA,
instance/task roles); named profiles and assume-role are not wired for the
archive — embed the engine-level API below if you need them.
From the library (custom embedding)
Register a read-only archive reader on the replication service, pointed at the same destination/prefix the snapshotter writes:
import (
"github.com/zourzouvillys/laredo/service/replication"
"github.com/zourzouvillys/laredo/snapshotter"
"github.com/zourzouvillys/laredo/snapshotter/format/jsonl"
)
reader, _ := snapshotter.NewReader(archiveDest, "laredo/public.events/", jsonl.New())
svc := replication.New(engine,
replication.WithArchive("public", "events", reader))
Without a registered archive, behaviour is unchanged (too-stale clients get a full live snapshot).
How the handoff stays gapless
The archive and the hot journal advance independently and share only the
source position (WAL LSN). Cold replay streams the archive up to its head
position, then resumes the hot journal from the first entry after that
position — the same ResumeSequenceForPosition bound that failover uses. Resume
is at-least-once across the boundary; clients apply by primary key, so a
brief overlap is safe.
Caveats
- Correctness never depends on the archive. If the archive is missing, unreadable, an unknown manifest version, or has fallen so far behind that the journal pruned past its head (a genuine gap), the server logs it and falls back to a full live snapshot before sending any data. Cold replay is purely an optimization on the stale path.
- Object-storage latency. Cold catch-up reads object storage and is slower than a hot delta. It runs off the live path and under the usual per-client backpressure, but a fleet reconnecting at once reads the archive hard — size and monitor accordingly.
- Read-only.
laredo-snapshotterremains the sole writer of the archive; the replication service only reads it.
Cascading: a fan-out as a source
A downstream engine can treat an upstream fan-out as a source, so fan-outs
chain: one upstream slot fans out to regional engines, each of which
re-serves its own fan-out, in-memory targets, or HTTP sync to local consumers —
no PostgreSQL slot per leaf. source/fanout is a SyncSource over
client/fanout, so it inherits snapshot, resume, cold-tier replay, and failover.
For the design rationale, see
EDR-0004 — Cascading replication.
import (
"github.com/zourzouvillys/laredo"
srcfanout "github.com/zourzouvillys/laredo/source/fanout"
"github.com/zourzouvillys/laredo/target/memory"
)
src := srcfanout.New(
srcfanout.ServerAddress("upstream-laredo:4002"),
srcfanout.Table("public", "events"),
srcfanout.ClientID("region-eu"),
)
engine, _ := laredo.NewEngine(
laredo.WithSource("upstream", src),
laredo.WithPipeline("upstream", laredo.Table("public", "events"), memory.NewIndexedTarget(memory.LookupFields("id"))),
)
The downstream engine receives the upstream's baseline and live changes and applies them to its targets like any other source.
Positions across the cascade
A fan-out's changes carry the upstream's source position (a PostgreSQL WAL
LSN). source/fanout orders them with a pluggable comparator that defaults to
PostgreSQL-LSN order — so a cascade over a PostgreSQL-backed fan-out is
zero-config. For a non-PostgreSQL upstream, supply your own:
srcfanout.WithPositionComparator(func(a, b string) int { /* order a vs b */ })
Caveats
- Lag compounds. Each hop adds propagation delay and a re-materialized in-memory copy. Size chains accordingly.
- Loops are yours to avoid. Nothing stops A→B→A if mis-wired; topology hygiene is the operator's responsibility.
- Pause is best-effort.
Pause/Resumeflip state but keep the upstream connection.
Failover & zero-downtime deploys
A client can hand off from one laredo-server instance to another — during a
rolling deploy or when an instance is being replaced — and resume without a
full re-sync.
Why it works: resume by source position
Each fan-out instance assigns its own journal sequence numbers, so a
sequence is meaningless on a different instance. The coordinate that is stable
across instances is the source position — the PostgreSQL WAL LSN. Every
journal entry carries its source_position, and the client persists the LSN of
the last change it applied. On failover it resumes from that LSN.
This requires that both instances tail the same PostgreSQL (each with its own replication slot, both reading the same publication). The same row change then yields the same LSN on every instance, so any instance can answer "give me everything after LSN X".
The handoff
- The draining instance sends a
GoAwaycontrol message on each client's stream (it keeps the stream open so the client can overlap the cutover). - The client re-dials its configured address — a load balancer or headless
DNS routes the new connection to a healthy instance — and resumes with
last_known_source_positionset to its last applied LSN. - The new instance replies with a
DELTA(no snapshot) containing every change after that LSN, then goes live. - Once the new stream has caught up, the client closes the old stream cleanly, letting the draining instance finish shutting down.
If the new instance has already pruned the journal back past the client's LSN, it falls back to a full snapshot automatically — correctness is preserved, only the optimization is lost.
Resume is at-least-once: during the brief overlap (and when an instance is momentarily behind), the client may re-apply changes it already has. Targets are keyed by primary key and apply idempotently, so this is safe.
Triggering a drain
On shutdown (rolling deploys). Start laredo-server with --drain-grace:
laredo-server --config laredo.conf --drain-grace 15s
On SIGTERM the server marks /health/ready unready (so the load balancer
deregisters it), tells fan-out clients to hand off, waits up to the grace period
for them to leave, then stops.
On demand (operator). Use the OAM admin RPC to drain a running instance — e.g. to cordon it before maintenance:
LaredoOAMService.DrainReplication { schema: "public", table: "config_document" }
Omit schema/table to drain every fan-out target. drain_deadline_seconds
optionally caps how long the server keeps serving drained streams.
Operations
- Put fan-out instances behind an L4/L7 load balancer or headless service and point clients at that address — not at individual pods.
- Size
journal.max_entries/journal.max_agefor your slowest consumer and your deploy cadence: an instance must retain entries back to the LSN of any client that might fail over to it, or that client re-snapshots. - Set
--drain-gracecomfortably above your LB's deregistration delay plus a typical client catch-up.
Troubleshooting
- Clients full-resync on every deploy — the new instance's journal doesn't
reach back to their LSN. Increase
journal.max_entries/max_age, or shorten the deploy gap. - New connections land on the draining instance — the LB hasn't
deregistered it yet. Ensure it scrapes
/health/readyand increase--drain-grace. GoAwaynever arrives — the fan-out target isn't being drained. Confirm--drain-grace > 0(SIGTERM path) or thatDrainReplicationmatched a target.
Consistency guarantees
- Snapshot + journal = complete state: no gaps, no duplicates
- Strict ordering: journal entries delivered in sequence order
- Atomic handoff: journal is pinned during snapshot send — no gap between snapshot and stream
- Portable resume: clients resume across instances by source position (WAL LSN), not per-instance sequence
- At-least-once on reconnect/failover: clients must be idempotent (they apply by primary key)
- At-least-once on reconnect: clients must be idempotent for the entry at their declared sequence