Table of Contents

Replication transport seam (IReplicationTransport)

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 transport.md, and llms.txt lists every page.

IReplicationTransport is the public, pluggable seam over the on-the-wire delivery of replication batches between clusters. It frames the in-process call shape that the outbound shipper uses, decouples that call shape from the bytes-on-the-wire (which is the binary-framing seam's concern), and standardises the receiver-side acknowledgement that drives the sender's per-peer cursor advance.

The contract is intentionally narrow: one method, one batch in, one ack out. There is no per-peer state on the transport, no streaming multiplexing in the public surface, and no transport-specific error vocabulary leaking into the call site - those concerns belong to specific implementations.

API

The interface and value types live in Orleans.Lattice.Replication:

public interface IReplicationTransport
{
    Task<ReplicationAck> SendAsync(
        ReplicationBatch batch,
        CancellationToken cancellationToken);
}

public readonly record struct ReplicationBatch
{
    public string TargetClusterId { get; init; }
    public string TreeName { get; init; }
    public string OriginClusterId { get; init; }
    public ReadOnlyMemory<byte> Payload { get; init; }
    public ReplicationBatchEnvelope? Envelope { get; init; }
    public ReplicationBatchEncodedEnvelope? EncodedEnvelope { get; init; }
}

public readonly record struct ReplicationAck
{
    public bool Accepted { get; init; }
    public HybridLogicalClock HighestAppliedHlc { get; init; }
    public HybridLogicalClock? BlockedAtHlc { get; init; }
    public int? SuggestedBatchSize { get; init; }
    public int? PauseForMs { get; init; }
    public int? SupportedWireVersion { get; init; }
    public uint[]? AdvertisedDictionaryIds { get; init; }
    public AdvertisedCompressionDictionary[]? AdvertisedDictionaries { get; init; }
}
ReplicationBatch member Semantics
TargetClusterId Stable identifier of the destination cluster. Implementations route the call by this value. Required: must be non-null and non-empty.
TreeName Name of the local tree this batch was drawn from. Receivers dispatch their per-tree apply pipeline on this id and keep their per-origin high-water mark per (TreeName, OriginClusterId). Required: must be non-null and non-empty.
OriginClusterId Stable identifier of the local (sending) cluster. Stamped on every captured WalRecord at commit time and surfaced on the batch so transports that frame entries themselves do not need to re-derive the origin from the payload. Required: must be non-null and non-empty.
Payload Opaque, framed batch payload. The byte layout is the responsibility of the binary-framing seam (typically Orleans-serializer-encoded WalRecord records inside a versioned envelope). Implementations treat this as a black box - they do not parse, peek into, or otherwise interpret the bytes. May be empty (heartbeat or keep-alive batch).
Envelope Optional pre-built typed envelope. A transport that frames entries itself sends it verbatim instead of decoding Payload. The anti-entropy repair path (leaf re-replay and the bootstrap fallback) populates it alongside Payload.
EncodedEnvelope Optional pre-encoded framing: a fixed header plus the per-entry segments exactly as the WAL stored them. The shipper populates only this slot - see Framing-only ship path.
ReplicationAck member Semantics
Accepted true when the receiver processed the batch - including when every entry was discarded as a repeat of an already-applied record, dropped because the tree is not enrolled on the receiver, or parked on the dead-letter queue. The canonical gRPC receiver returns false only when it deferred the batch because a coordinated-restore receive fence holds the tree (with a bounded PauseForMs); an apply fault surfaces as a failed call (the send throws), not as Accepted = false. The default no-op transport always returns false.
HighestAppliedHlc The highest HLC the receiver advanced its per-origin high-water-mark to while processing the batch. On an accepted ack the sender advances its per-peer cursor to this value, or to the last shipped entry's HLC when this value is at or below the current cursor (every entry deduplicated or dropped); when Accepted is false the sender does not consume it.
BlockedAtHlc Optional receiver-side blocked-floor pin (lowest HLC across every partially-staged atomic batch). The sender publishes this value to its local IWalCursorRegistry so the producer-side WAL GC AND-s entry.Timestamp < blockedFloor into its trim predicate; null means the receiver has no in-flight admissions for this tree (or is pre-Phase-9 and never stamped the slot). Strictly additive on the wire.
SuggestedBatchSize Optional receiver-side flow-control hint: the largest per-tick batch the receiver would like the sender to ship next, in entries. The sender clamps to [1, options.ShipBatchSize]; null (or any value <= 0) means "no preference" and the sender resumes at its configured ShipBatchSize (the canonical re-acceleration signal). Strictly additive on the wire.
PauseForMs Optional receiver-side flow-control hint: number of milliseconds the sender should pause before its next pump tick. Composes with the shipper's exponential-backoff retry budget via max(currentBackoffDeadline, now + PauseForMs) - a receiver-requested pause never shortens an in-progress backoff. null or <= 0 means "no pause requested". Strictly additive on the wire.
SupportedWireVersion Optional: the highest framing wire version this receiver can decode. An opted-in sender negotiates its target version from it (see Wire Format); null means the receiver did not advertise a capability. Strictly additive on the wire.
AdvertisedDictionaryIds Optional: the shared compression-dictionary ids this receiver can resolve, so an opted-in sender only compresses with a dictionary the peer can decode. null means no capability advertised; an empty array means the capability exists but no dictionary is held. The canonical gRPC receiver sends null in both cases - when its dictionary provider exposes no catalog and when the catalog holds no dictionary - so a sender treats such a peer as capability-unknown. Strictly additive on the wire.
AdvertisedDictionaries Optional fingerprint-bearing successor to AdvertisedDictionaryIds: (id, fingerprint) pairs, so two clusters that map one id to different bytes never negotiate a match. Strictly additive on the wire.

ReplicationBatch is intentionally not Orleans-serialisable: it is the in-process call argument, not the on-the-wire envelope. Wire-format hardening - versioned envelopes, content framing, compression - happens inside Payload and is the binary-framing seam's concern. ReplicationAck is Orleans-serialisable (alias olr.ak) because the receiver returns it to the sender across whatever transport is in use, including in-cluster Orleans RPC bridges.

Send semantics

Three concerns the transport composes for every call:

1. Idempotency at the batch boundary

Receivers de-duplicate re-deliveries by record identity - an exact (origin, hlc, key, op) match in a bounded recent-apply cache, backed by an idempotent leaf-level apply for a repeat that has aged out of the cache - so a transport that retries a batch on transient failure does not cause double-apply. (The per-origin high-water mark is not the drop threshold; only a snapshot-pinned floor from a bootstrap drops entries outright, because everything at or below it is already in the snapshot.) Implementations are free to retry as aggressively as their reliability story requires; the receiver-side dedup is the correctness guarantee.

2. Advance-cursor-on-ack

The sender advances its per-peer cursor only when the ack is accepted: to ReplicationAck.HighestAppliedHlc, or - when that frontier is at or below the current cursor because the receiver deduplicated or dropped every entry - to the last shipped entry's HLC, so the same batch is not re-shipped forever. A rejected ack leaves the cursor in place and the batch is retried after a backoff. This is the canonical at-least-once-delivery, at-most-once-apply contract: a batch may be re-delivered, and the receiver's record-identity dedup is what keeps a re-delivery from applying twice.

There is no partial-apply signal on the ack. The receiver either finishes the batch - an entry that exhausts its MaxApplyRetries budget is parked on the dead-letter queue rather than failing the call - and returns an accepted ack, or an earlier apply fault escapes and the call fails, which leaves the sender's cursor in place so the batch is retried after a backoff. On an accepted ack the sender's per-partition cursors advance past every entry it drained for that batch.

3. Concurrency

Implementations are required to be safe for concurrent invocation across distinct (TargetClusterId, TreeName) pairs - the canonical outbound shipper fans out across peers and trees in parallel. Concurrent invocation against the same (TargetClusterId, TreeName) pair is implementation-defined. With the default ShipMaxInFlight of 1 the canonical shipper serialises calls per pair; raising ShipMaxInFlight makes it issue up to that many concurrent sends per pair (see Receiver Flow Control), so a transport used with pipelining must tolerate concurrent calls for one pair.

Validation

SendAsync throws ArgumentException when:

  • batch.TargetClusterId is null or empty.
  • batch.TreeName is null or empty.
  • batch.OriginClusterId is null or empty.

The interface prescribes no cancellation exception, and the shipped implementations differ: the default no-op transport completes synchronously without observing the CancellationToken, and the gRPC push transport passes the token to the call and lets whatever the gRPC client raises when it fires propagate, rather than translating it to OperationCanceledException.

Registration

AddLatticeReplication registers the default IReplicationTransport implementation as a silo-side singleton:

siloBuilder.AddLatticeReplication(o => o.ClusterId = "site-a");

The default registration is a no-op transport - it validates routing fields, discards the payload, and returns default(ReplicationAck) (i.e. Accepted = false, HighestAppliedHlc = HybridLogicalClock.Zero). The sender's cursor stays put, which is exactly the right behaviour while the rest of the replication pipeline is being wired up but no real transport is configured. Production hosts replace it via standard DI:

sealed class MyTransport : IReplicationTransport
{
    public Task<ReplicationAck> SendAsync(ReplicationBatch batch, CancellationToken cancellationToken)
        => Task.FromResult(default(ReplicationAck));
}

static void Configure(IServiceCollection services)
{
    services.AddSingleton<IReplicationTransport, MyTransport>();
}

Future implementations

The single-method seam is the contract the binary framing and the canonical gRPC transport plug into. The binary framing hardens the byte layout inside Payload; the gRPC transport drives SendAsync as one unary RPC per batch over a cached HTTP/2 channel per peer cluster. Neither changes the call shape established here, and a custom transport plugs into the same seam.

The wire format inside Payload is the concern of IReplicationBatchEncoder. The default registration is the Orleans-serializer-backed binary encoder; hosts swap to a different framing (JSON for HTTP debuggability, content-hash-prefixed for deduplication) by replacing the encoder registration via DI. Transports remain agnostic about which encoder produced the bytes.

The canonical sender + receiver pair ships in the Orleans.Lattice.Replication.Grpc sub-package - see Orleans.Lattice.Replication.Grpc for topology, registration, and operations notes.

Framing-only ship path

The shipper's outbound path is unconditionally framing-only. Every batch the shipper hands to SendAsync carries a populated ReplicationBatch.EncodedEnvelope (a fixed 32-byte header plus length-prefixed pre-encoded entry segments produced by IReplicationBatchEncoder.EncodeFraming). Each entry's bytes are the verbatim segment the WAL stored at append time, read back through the WAL partition's shipping read path - except an entry that pre-ship coalescing folded on a CRDT tree, whose segment is re-encoded once with the combined delta - with no per-tick re-encode through an envelope-level Orleans serializer call, and no producer-side typed-envelope path. ReplicationBatch.Payload and ReplicationBatch.Envelope remain on the contract for receiver-side and test-fixture compatibility, and the anti-entropy repair path still ships them (a typed envelope plus its encoded Payload), but the producer-side shipper writes only EncodedEnvelope.

Custom transports that want to consume the framing bytes directly read them off ReplicationBatch.EncodedEnvelope. There is no separate typed-transport interface or sender-side capability probe - the shipper does not branch on transport type at activation. Bytes-only transports (for example a host-supplied HTTP-framed transport) encode EncodedEnvelope with IReplicationBatchEncoder.EncodeFraming and forward the bytes as-is, as the gRPC transport's marshaller does; the default no-op transport discards the batch.

Caveats

  • Transports do not interpret the payload. A transport that needs to make a routing decision based on payload contents (e.g. shed-load on oversize batches) must do so via batch metadata that the framing seam exposes on the call site, not by parsing Payload itself. Cross-cutting concerns belong on the call envelope; the wire bytes stay opaque.
  • The ack envelope grows additively, never by breaking change. New [Id(n)] slots backed by nullable defaults (the BlockedAtHlc, SuggestedBatchSize, and PauseForMs slots are the existing precedent) are safe to ship on either side of a peering independently because pre-existing receivers and senders decode the slot as null. The single-method SendAsync contract does not change. Receiver-side flow-control hints in particular are wired through a pluggable IReceiverFlowControlPolicy seam.

Metadata pass-through contract

The transport stays dumb about the entries it carries. Specifically, every IReplicationTransport implementation must preserve the causal-plus metadata slots on every WalRecord verbatim across a round-trip:

  • WalRecord.VectorClock - the sparse {originClusterId -> HybridLogicalClock} frontier captured at commit time.
  • WalRecord.DependencySummary - initially aliased one-to-one with VectorClock; reserved as a distinct slot so a future Bloom-filter-shaped summary can ship without re-numbering the wire format.

The transport must not reorder entries, mutate either slot, synthesise an empty frontier when the producer left the slot null (legacy peers and pre-causal-plus entries decode null and the receiver treats that as the empty frontier), or merge the two slots together. Any normalisation, summary derivation, or merge belongs in the producer / receiver, never in the wire layer.

The contract is pinned by TransportMetadataPassthroughContractTests in both Orleans.Lattice.Replication.Tests (loopback transport) and Orleans.Lattice.Replication.Grpc.Tests (the gRPC push transport). A new transport implementation should ship a mirror of that fixture parameterised over its own seam.

Content-hash payload-elision round trip (opt-in, default off)

The push transport above is one-way: one batch in, one ack out. The opt-in content-hash payload-elision feature adds a second, bidirectional exchange in front of the push so the shipper can avoid re-sending payloads a peer already holds. It is built on the same per-(tree, peer) content hashing that drives the default-on re-send-rate measurement (ship.redundant_payloads, enabled by ContentHashDedupEnabled which now defaults to true), but instead of merely counting redundant re-sends it elides them. The elision step itself (ContentHashDedupElisionEnabled) stays opt-in and default-off because it needs the receiver to advertise which hashes it already holds.

The exchange is a default-no-op method on the IReplicationDigestProbeTransport seam (the same bidirectional probe transport the anti-entropy digest probe uses), so existing transports compile and behave unchanged:

public interface IReplicationDigestProbeTransport
{
    Task<ContentManifestResponse> ExchangeContentManifestAsync(
        string targetClusterId,
        ContentManifestRequest request,
        CancellationToken cancellationToken)
        => Task.FromResult(ContentManifestResponse.NotSupported);
}

public readonly record struct ContentManifestEntry
{
    public int EntryIndex { get; init; }
    public string Key { get; init; }
    public ulong ContentHash { get; init; }
    public HybridLogicalClock Hlc { get; init; }
}

public readonly record struct ContentManifestRequest
{
    public string TreeName { get; init; }
    public string OriginClusterId { get; init; }
    public IReadOnlyList<ContentManifestEntry> Entries { get; init; }
}

public readonly record struct ContentManifestResponse
{
    public bool ExchangeSupported { get; init; }
    public IReadOnlyList<int> MissingEntryIndices { get; init; }
    public HybridLogicalClock AdvancedHlc { get; init; }
}

Flow

  1. Sender builds a manifest. When LatticeReplicationOptions.ContentHashDedupEnabled and ContentHashDedupElisionEnabled are both set, the shipper hashes the value-carrying point-Set entries in the drained batch (FNV-1a 64-bit over op + key + range + value, the same digest the measurement uses) and advertises a ContentManifestRequest to the peer. Only eligible entries are manifested - range deletes, saga terminal marks, prepared atomic-batch entries, and zero-HLC entries are never placed in the manifest and always ship verbatim, so atomic-batch boundaries, causal-dependency gating, and per-origin FIFO are preserved.
  2. Receiver answers with the missing set. For each manifest entry the receiver compares the advertised content hash against the content it has already applied for that key. An entry the receiver does not hold (or holds with a different hash) is reported in MissingEntryIndices. An entry the receiver already holds byte-identical is not missing - and if the manifest entry's Hlc is newer than the receiver's recorded per-origin high-water mark (the idempotent re-set of an identical value; the index holds only each key's content hash, not a per-key clock), the receiver advances its per-origin high-water-mark via a metadata-only apply and reports the advanced clock in AdvancedHlc, all without the payload travelling.
  3. Sender ships only the missing payloads. The shipper drops every elided entry from the outbound batch and ships the remainder through the ordinary IReplicationTransport.SendAsync push, then advances its per-peer cursor past the whole originally-drained range (the receiver advanced its high-water-mark for the elided entries during the exchange). When every entry is elided no batch is shipped at all. The shipper records the sender side of the exchange on three counters tagged tree, peer, and tenant: ship.manifest_exchanges (one per outbound batch that advertised a manifest and received a pull-missing reply), ship.elided_payloads (set-entry payloads dropped from the batch), and ship.elided_payload_bytes (their summed pre-encoded wire-segment length).

Default-off and rolling-upgrade safety

The default ExchangeContentManifestAsync returns ContentManifestResponse.NotSupported (ExchangeSupported = false), so a transport (or peer) that has not implemented the pull-missing RPC reports "not supported" and the shipper permanently falls back - for the rest of the activation - to shipping the full batch verbatim, byte-identical to today. Capability is learned lazily per shipper activation: the first eligible batch attempts the exchange, and a "not supported" reply latches elision off until the grain re-activates. With elision disabled (the default) the exchange is never attempted and the wire bytes are identical to a build without the feature.

The elision composes with sender-side multi-batch ship pipelining. A configured pipelining window (ShipMaxInFlight > 1) is preserved while elision is enabled - the per-batch manifest exchange runs inline in the bounded-pipelining drain loop, so it no longer collapses the window to one. A batch every entry of which the receiver already holds (a fully-elided batch) ships no envelope; it advances the durable per-peer cursor strictly in FIFO order through the same in-flight queue via a synthetic already-completed ack, so per-origin FIFO, causal-dependency gating, atomic-batch boundaries, and advance-strictly-on-ack cursor semantics hold across the whole window exactly as on the serial path. The full-range cursor-advance inputs are captured before the exchange, so the cursor still advances past every originally-drained entry regardless of how many were elided, and the synthetic zero-latency ack is excluded from the adaptive batch-size controller so it cannot skew adaptive sizing.

gRPC binding

The Orleans.Lattice.Replication.Grpc sub-package binds ExchangeContentManifestAsync to a real unary RPC, ExchangeContentManifest, mirroring the existing ProbeDigest binding: the client invoker reuses the same long-lived, HTTP/2-multiplexed per-peer GrpcChannel cache and the same shared-secret auth interceptor the push and probe RPCs use. The marshaller wraps the already-aliased ContentManifestRequest / ContentManifestResponse value types in reference-typed boxes (the gRPC Method<,> class constraint) and writes their Orleans-serialized bytes straight into the gRPC stream's buffer writer.

A peer that has not bound the method answers Unimplemented, and a peer that is momentarily unreachable answers Unavailable; the client invoker catches both and returns ContentManifestResponse.NotSupported, so the sender's existing capability-latch falls back to shipping the full batch verbatim with no per-hop wire-version pre-check. This makes enabling elision on one side of a peering rolling-upgrade safe.

On the receiver, the gRPC service handler first refuses the call (PermissionDenied) when the tree is not enrolled for replication there, when the call carries no stamped origin, or when the request's OriginClusterId names a different origin than the one stamped on the call, then resolves the durable per-origin high-water mark for the tree, projects the receiver's applied-content index onto the manifest's keys, and computes the missing set with the same pure planner the in-process path uses. For an entry the receiver already holds whose Hlc is newer than the recorded high-water-mark, the handler performs a durable metadata-only TryAdvanceAsync on the high-water-mark grain (no payload travels) and reports the advanced clock in AdvancedHlc. The handler increments three receiver-side counters tagged tree, the origin peer, and tenant: receiver.content_manifest_exchanges (one per exchange answered), receiver.content_entries_elided (entries the receiver reported it already holds), and receiver.content_hwm_advances (one per exchange whose durable high-water-mark advance succeeded).

Receiver applied-content index

The receiver answers "which hashes do I already hold?" from a bounded, in-process, best-effort per-tree index mapping Key -> ContentHash (the same FNV-1a digest the sender manifests). It is populated as point-Set writes apply through the receiver-side applier, removed on a point Delete, and cleared for the whole tree on a DeleteRange (the index has no range query, so the coarse clear avoids a stale "already holds" answer for a removed key). The index is never serialized and never travels on the wire.

The index is gated behind ContentHashDedupEnabled: when the master switch is off it stays empty and the populate path is off-path-free. It is a best-effort cache - a cold, never-populated, or evicted entry simply omits the key, so the planner reports the entry as missing and the sender ships it, which is always safe. Only last-writer-wins point mutations that are not part of a not-yet-visible atomic-batch prepare phase are recorded; CRDT-mode entries are skipped because the receiver merges rather than overwrites them and so they are never elision-eligible.

Scope. The manifest engine, the shipper-side elision wiring, the capability gating, the gRPC binding (client invoker + server handler), the receiver applied-content index, and the options/metrics/dashboard surface ship in this package family and are covered by unit, in-process loopback, and gRPC-binding tests. The cross-cluster round trip advances the remote receiver's high-water-mark for the identical-content-newer-clock case via a durable metadata-only apply on the high-water-mark grain. Transports that do not implement ExchangeContentManifestAsync keep the default no-op, which leaves every peering wire-identical to today.