Table of Contents

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 IWalStorageProvider seam (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.RequestSnapshotAsync or ILatticeBootstrapCoordinator.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 SetManyAtomicAsync arrives 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 IReplicationTransport seam - 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 IWalStorageProvider durability 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, default true) 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, default true) 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, default LatticeCompression.Zstd) compresses the framing tail of every batch whose encoded entries reach the FramingCompressionMinBatchBytes (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 LatticeReplicationOptions knob, 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, LatticeMergeMode selection, per-key filter.
  • Observability - LatticeReplicationMetrics instruments, 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 IHealthCheck that turns per-peer lag, contact age, and error streaks into a Healthy / Degraded / Unhealthy verdict.
  • 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-config tree: the static anchor, the compiled snapshot, the dynamic membership / merge-mode seams, and fail-closed ambiguity resolution.

For internals (the "how"):

  • Change Feed - IChangeFeed seam, 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 - IReplicationTransport seam, 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 - IReceiverFlowControlPolicy seam, 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.