Orleans.Lattice.Replication
This page documents Orleans.Lattice.Replication 9.9.0, in the documentation for Orleans.Lattice 9.9.0 (release line 9.9), built 2026-10-04. It is also published as markdown, with every table and list, at README.md, and llms.txt lists every page.Cross-cluster replication for Orleans.Lattice - captures every mutation at commit time, ships it between Orleans clusters under the source cluster's HLC, and applies it on the receiver with CRDT-aware merges, causal delivery, snapshot bootstrap, and dead-letter quarantine.
What is it?
Orleans.Lattice.Replication is the end-to-end cross-cluster replication subsystem that layers on top of Orleans.Lattice. It is more than a wire format - it covers the full producer/transport/receiver pipeline plus the operational surface around it:
- Capture. Mutations are intercepted at commit time on the producing cluster and written to a per-tree WAL via the pluggable
IWalStorageProviderseam (in-memory default; optional Azure Table Storage and file-system backends). - Ship. A per-
(tree, peer)shipper tails the WAL and ships batches to each peer over a push transport (IReplicationTransport, with gRPC as the canonical binding - one unary call per batch over a long-lived channel). - Apply. Inbound entries flow through
IReplicationApplier, which gates each entry against this receiver's own per-tree enrollment and merge mode, drops entries a pinned bootstrap snapshot already covers, suppresses repeated records by exact(origin, hlc, key, op)identity (shadow-forward de-duplication), parks entries whose causal dependencies have not arrived in the causal-apply buffer, and performs CRDT-aware merges for the configured merge mode. - Bootstrap. New or fallen-off-the-log peers seed via the snapshot subsystem (
ISnapshotProvider,IRemoteSnapshotTransport,LatticeRemoteSnapshotService,RemoteSnapshotProvider) coordinated by the per-tree bootstrap coordinator (ILatticeBootstrapCoordinator) with crash-resumable state. The fall-off detector starts a bootstrap automatically; a new peer is seeded by an explicit re-seed request (ILatticeReplicationAdmin.RequestSnapshotAsyncorILatticeBootstrapCoordinator.BootstrapAsync). - Operate. Dead-letter quarantine for poison entries, per-tree merge-mode resolution, operator-driven re-seed, fall-off-log detection, admin introspection (
ILatticeReplicationAdmin,ILatticeWalIntrospection), shared-secret-based mutual auth between clusters, and first-class metrics (apply duration, lag, FIFO violations, bootstrap retries, dead-letter rates) are all in-scope.
No external broker, no shared database, no host-level outgoing-call filter.
It supports:
- Per-tree opt-in with a declared
LatticeMergeMode. - Origin-stamped HLC on every record, with at-most-once apply per exact
(origin, hlc, key, op)record identity. - Causal+ delivery - vector-clock-stamped entries with receiver-side dependency satisfaction across point writes, atomic multi-key writes, and structural shadow-forwards; maintenance writes (tombstone compaction) are never replicated, so they add no edges to the causal graph.
- Active-active topology: any peer can write to any tree; conflicting updates converge deterministically.
- Atomic batch delivery - replicated
SetManyAtomicAsyncarrives on every peer as a single visible unit. - gRPC push transport - one unary call per batch over a cached HTTP/2 channel per peer cluster.
- Pluggable
IReplicationTransportseam - gRPC is the canonical implementation; in-process and custom transports plug into the same contract. - Snapshot bootstrap for new and re-seeded peers; auto-bootstrap on fall-off-the-log.
- Per-tree dead-letter queue for poison entries; replication continues past them.
- First-class per-peer metrics and lag observability.
- Pluggable
IWalStorageProviderdurability seam - in-memory default, optional Azure Table Storage and file-system backends.
Core Properties
- Convergent under concurrent writes. Two clusters writing to the same key arrive at the same final state, deterministically, without coordination.
- Causally consistent. A receiver never observes a write before the writes it causally depends on - across point writes, atomic multi-key writes, and structural shadow-forwards.
- Cycle-safe. Origin attribution is durable metadata on every record, not ambient context - replicating into and back out of a peer cluster never loops a mutation back to its source.
- At-most-once apply. Re-delivery of the same
(origin, hlc, key, op)is idempotent. Counters do not double-increment, sets do not re-add. - No host-level coupling. Replication is produced by the silo at commit time. Hosts neither install outgoing-call filters nor route mutations through their own pipeline.
Behaviour is validated end-to-end by active-active convergence chaos tests across the replicated CRDT catalogue running through real AddLatticeReplication silos.
Features
| Feature | What it gives you | Docs |
|---|---|---|
| Active-active topology | Any peer can write to any tree. Multi-cluster concurrent updates converge to the same state by CRDT mode, not by post-merge LWW-on-bytes. | Replication Modes |
| Anti-entropy (opt-in) | A default-off safety net for silent divergence: a digest probe flags a shard whose content digest differs from a peer's (a conservative trigger, not proof of divergence), a read-only Merkle walk localises it to leaves, and targeted WAL re-replay or a scoped snapshot fallback repairs it behind an opt-in gate, a rate cap, and a circuit breaker. | Automatic drift remediation |
| At-most-once apply | Re-delivery of the same record is idempotent: an exact-identity recent-apply cache suppresses a repeated (origin, hlc, key, op) record, and the typed merges are themselves idempotent, so counters, sets, and registers are never double-applied. |
Replication Apply |
| Atomic batch delivery | Replicated SetManyAtomicAsync arrives on every receiver as a single visible unit. No reader observes a partial-set state across clusters. |
Replication Apply |
| Auto-bootstrap on fall-off-log | When the fall-off detector finds a peer's per-origin high-water mark behind the oldest WAL entry still retained for that peer, the receiver re-seeds from a fresh snapshot automatically (AutoBootstrapOnFallOffLog, on by default). The built-in periodic check compares against the local WAL; a sender only trims entries its shipper has not yet acknowledged when WalRetention is set - see WalRetention. |
Auto-Bootstrap |
| Causal+ ordering | A receiver never observes a write before its causal dependencies - point writes, atomic multi-key writes, and structural shadow-forwards all preserve causal order. | WAL |
| Change feed | IChangeFeed is a public, in-process pull feed over a tree's locally-authored WAL entries for bridges, tests, and in-process projections; the shipper tails the WAL directly and does not use it. |
Change Feed |
| Coordinated multi-cluster restore | Restoring a backup into a replicated tree runs as an all-or-nothing cross-cluster saga: every cluster cuts over together or rolls back together, so no peer re-advances the restored cut and no reader observes a torn restore. | Coordinated Restore |
| Dead-letter queue | An entry whose apply keeps failing is quarantined per tree after a configurable retry budget (MaxApplyRetries); entries the receiver's merge-mode or tenant-isolation gate refuses, the causal-apply buffer evicts or fails to apply when it drains, or the sender cannot encode are parked at once. Replication continues past them. |
Dead-Letter Queue |
| gRPC push transport | One unary gRPC call per batch over a long-lived, HTTP/2-multiplexed channel per peer. Push latency is sub-second, well below reminder-cadence pull. | Orleans.Lattice.Replication.Grpc |
| Health check | ASP.NET Core / Kubernetes IHealthCheck reporting Degraded when entries-behind, last-contact age, or consecutive-error streak crosses a soft bound, and Unhealthy when a hard bound is crossed or the degraded state outlasts the configured grace window. An opt-in inbound-silence signal covers receive-side liveness. |
Health Check |
| Observability | Per-peer entries-behind, bytes-behind, ship-in-flight, consecutive-errors, and last-contact gauges, plus receiver-side apply lag and duration histograms, on LatticeReplicationMetrics. |
Observability |
| Origin-stamped HLC | Every replicated record carries (originClusterId, hlc). Cycles break naturally, transitive topologies preserve causality, and applies are idempotent by record identity: a repeated (originClusterId, hlc, key, op) record is suppressed or re-applies as a no-op. |
Replication Apply |
| Per-tree opt-in + per-key filter | Declare which trees replicate and (optionally) which keys within a tree the shipper sends. Granular enough to ship operator-visible labels while keeping per-shift counters out of the incremental stream; the key filter is not applied to snapshot exports or the opt-in anti-entropy repair paths, so a peer that bootstraps from this cluster still receives every key. | Replication Modes |
| Pluggable transport | IReplicationTransport is the public seam. gRPC is the canonical implementation; in-process and custom transports plug into the same contract. |
Transport |
| Receiver-side flow control | The receiver stamps optional SuggestedBatchSize / PauseForMs hints onto every ack; the sender clamps its per-tick batch cap and pauses on request. A struggling receiver throttles in-band without timing out RPCs; a recovered receiver re-accelerates by lifting the hints. |
Receiver Flow Control |
| Runtime per-tree replication config | Enable or disable replication for a tree at runtime under a fixed merge mode, distributed as the converging sys-replication-config system tree. Flip it once on any cluster and every peer converges; concurrent divergent modes are detected rather than one silently overwriting the other, and the tree resolves to no mode until an operator settles it. |
Runtime Replication Config |
| Snapshot bootstrap | New or re-seeded peers receive a point-in-time snapshot, then switch to incremental shipping at the snapshot's HLC. | Snapshot Bootstrap |
| System-tree replication | Enrol the reserved membership + auth-policy trees so identity and authorization converge across sites. Replication-applied writes bypass the access gate under a system-origin scope; an optional strict policy-epoch fence closes the revoke window per tree. | System-Tree Replication |
| Transport security | Shared-secret authentication between clusters, required by default, with the secret sourced from environment variables, configuration, or a custom ILatticeReplicationSecretSource; by default a caller must also present the secret configured for the cluster it claims to be, and the gRPC binding refuses plaintext peer endpoints unless explicitly allowed. |
Transport Security |
| Typed CRDT deltas | The wire carries typed deltas for OR-Set, PN-Counter, VersionVector, MV-Register, OR-Map, RGA sequence, OR / RW flags, G-Counter, G-Set, RW-Set, and max / min bounded registers; last-writer-wins trees ship the opaque value bytes. A receiver merges a CRDT-mode tree by mode, not by opaque-byte LWW. | Deltas |
Quick Start
Add replication on top of an existing Orleans.Lattice silo. The minimum end-to-end setup is AddLattice + AddLatticeReplication + a transport. The snippet below keeps everything in memory: AddLattice registers the in-memory WAL provider, and the storage callback given to it here registers in-memory grain storage. For state that survives restarts, add AddWalStorage (or a storage package such as the Azure Table WAL) and give AddLattice a durable grain-storage provider as well - leaf state and snapshots, and the replication shipper cursors, high-water marks, and bootstrap state, persist there. On the producer/sender silo:
var builder = WebApplication.CreateBuilder();
builder.Host.UseOrleans(silo =>
{
silo
.AddLattice((s, storageName) => s.AddMemoryGrainStorage(storageName))
.AddLatticeReplication(opts =>
{
opts.ClusterId = "site-a";
opts.ReplicatedTrees = new Dictionary<string, LatticeMergeMode>(StringComparer.Ordinal)
{
["operator-labels"] = LatticeMergeMode.OrSet,
["machine-counters"] = LatticeMergeMode.PnCounter,
};
opts.ReplicationPeers = new[] { "site-b" };
});
});
// Cross-cluster gRPC binding - one entry per peer cluster wires both
// the live-push transport and the bootstrap snapshot transport.
builder.Services.AddLatticeReplicationGrpc(grpc =>
{
grpc.Peers["site-b"] = new Uri("https://site-b.example:5001");
});
On the receiver silo, register the binding (same single helper) and map the endpoint routes:
var builder = WebApplication.CreateBuilder();
builder.Services.AddLatticeReplicationGrpc();
var app = builder.Build();
app.MapLatticeReplicationGrpc();
Both clusters also need a shared secret, because the receiver refuses unauthenticated calls by default. Set LATTICE_REPLICATION_SECRET on every silo of both clusters; with the default origin binding a single cluster-wide secret must be the same value on both. See Transport Security for per-peer secrets, custom secret sources, and rotation.
For a working multi-cluster example exercising HLC-ordered facts, typed OR-Set replication, and gRPC push, see the MultiSiteManufacturing project under samples/.
Default efficiency posture (versions greater than v7.1.0)
A stock AddLatticeReplication deployment ships with the safe efficiency bundle on out of the box, so the minimal setup above is already coalesced, compressed, and measured:
- Pre-ship coalescing (
PreShipCoalescingEnabled, defaulttrue) collapses redundant per-key versions off the cross-cluster wire before they ship. It is convergent by construction - a last-writer-wins tree keeps the highest-HLC version per key, recognised CRDT shapes (registered OR-Map shapes included) delta-merge, and an unregistered OR-Map shape or an opaque payload ships verbatim. - Content-hash dedup measurement (
ContentHashDedupEnabled, defaulttrue) records the payload re-send rate (ship.redundant_payloads). It is observability-only and never alters the shipped bytes; the actual payload elision (ContentHashDedupElisionEnabled) stays opt-in. - Dict-less Zstandard framing compression (
FramingCompression, defaultLatticeCompression.Zstd) compresses the framing tail of every batch whose encoded entries reach theFramingCompressionMinBatchBytes(512-byte) threshold. Every current-wire-version peer decodes it with no extra wiring; shared dictionaries (ZstdDictionary) stay opt-in.
The wire bytes change shape but remain backward-compatible to decode by any current-wire-version peer, so no coordinated rollout is required. Each knob is individually overridable back to the prior behaviour:
siloBuilder.AddLatticeReplication(opts =>
{
opts.ClusterId = "site-a";
// Opt out of any part of the safe efficiency bundle individually.
opts.PreShipCoalescingEnabled = false; // ship every version verbatim
opts.ContentHashDedupEnabled = false; // stop the re-send-rate measurement
opts.FramingCompression = LatticeCompression.None; // ship the framing tail uncompressed
});
Reference
For day-to-day use and operations:
- Architecture - the producer-to-receiver pipeline, the seams it attaches to, and the invariants it preserves end to end.
- API Reference - the public types, seams, registration helpers, and extension points.
- Configuration - every
LatticeReplicationOptionsknob, its default, and per-tree scope. - Chaos Tests - the cross-cluster, gRPC transport, and Azure Table WAL chaos suites.
- Coordinated Restore - the all-or-nothing cross-cluster restore saga, its dispatch rule, the public saga-participant SPI, and reliability under duress.
- Replication Modes - per-tree opt-in,
LatticeMergeModeselection, per-key filter. - Observability -
LatticeReplicationMetricsinstruments, per-peer lag, error counters. - Dead-Letter Queue - quarantine model, operator surface, replay.
- Snapshot Bootstrap - point-in-time bootstrap, snapshot HLC, incremental cutover.
- Auto-Bootstrap - fall-off-the-log detection and automatic re-seed.
- Automatic drift remediation - operator playbook for the opt-in anti-entropy stack: default-off posture, opt-in path, metrics surface, and the version-skew / WAL-trimmed / circuit-breaker failure-mode matrix.
- Transport Security - shared-secret authentication, HTTPS-by-default, custom secret sources, env-var convention.
- Health Check - the ASP.NET Core
IHealthCheckthat turns per-peer lag, contact age, and error streaks into aHealthy/Degraded/Unhealthyverdict. - System-Tree Replication - enrolling the reserved membership and auth-policy trees so identity and authorization converge across sites.
- Runtime Replication Config - runtime per-tree enable / disable distributed as the
sys-replication-configtree: the static anchor, the compiled snapshot, the dynamic membership / merge-mode seams, and fail-closed ambiguity resolution.
For internals (the "how"):
- Change Feed -
IChangeFeedseam, per-partition offset cursor, async enumerable shape. - Anti-entropy digest probe - the detection stage: a low-frequency, read-only pass comparing the content digest of each shard the tree's live shard map routes to against every peer's.
- Anti-entropy Merkle walk - the localisation stage: a read-only top-down descent that narrows a shard mismatch to the diverged leaves and their covering ranges.
- Anti-entropy leaf re-replay - the repair stage: re-ships the retained WAL entries covering those ranges down the ordinary TX-aware apply path, de-duplicated at the receiver.
- Anti-entropy bootstrap fallback - the repair path taken when re-replay cannot reach the divergence, such as a WAL trimmed past the divergence point.
- Anti-entropy remediation guards - the opt-in switch, rate cap, and circuit breaker wrapping the repair stages; detection is never gated by them.
- Replication Apply - receiver-side applier, per-origin high-water-mark, recent-apply cache, atomic batch buffering.
- Replication Drivers - production drivers that turn the dormant seams into a running pipeline.
- Transport -
IReplicationTransportseam, batch shape, acks. - Orleans.Lattice.Replication.Grpc - canonical transport: unary push RPC (plus a server-streaming snapshot RPC), channel reuse, custom marshallers.
- Receiver Flow Control -
IReceiverFlowControlPolicyseam, ack-stamped hints, sender clamping / pause composition. - Wire Format -
ReplicationBatchEnvelope,IReplicationBatchEncoder, wire version negotiation. - Deltas - typed CRDT delta records on the wire.
- WAL - partitioned replication write-ahead log, turn-safe batching, causal+ entry schema.