Table of Contents

Replication apply seam (IReplicationApplier)

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

IReplicationApplier is the public, in-process inbound seam over the per-tree apply pipeline. It installs a single WalRecord authored on a remote cluster onto the local tree while preserving the remote cluster's origin id end-to-end and, except where Source-HLC and origin preservation notes otherwise, its HybridLogicalClock, and it filters re-delivery via a snapshot-pinned causal floor plus a shadow-forward identity cache and per-key last-writer-wins idempotence, so at-least-once transports become at-most-once apply.

The contract is deliberately neutral: there is no transport binding, no per-peer state, no ack envelope. It is the seam custom transports and integration tests plug into.

API

The interface and result type live in Orleans.Lattice.Replication:

public interface IReplicationApplier
{
    Task<ApplyResult> ApplyAsync(
        WalRecord entry,
        CancellationToken cancellationToken = default);

    Task<ApplyResult> ApplyBatchAsync(
        IReadOnlyList<WalRecord> entries,
        CancellationToken cancellationToken = default);
}

public readonly record struct ApplyResult
{
    public bool Applied { get; init; }
    public HybridLogicalClock HighWaterMark { get; init; }
    public bool Deferred { get; init; }
}
ApplyResult member Semantics
Applied true when the entry was merged onto the local tree; false when the entry was filtered out as a re-delivery (its Timestamp was at or below the origin's pinned snapshot floor, or its identity tuple hit the shadow-forward cache), parked on the causal-apply buffer to await its dependencies, deferred by a restore saga's receive fence, acknowledged as a no-op because it is a tombstone-reap envelope, or rejected as inapplicable (its OriginClusterId matched the local cluster id and would have looped, its tree is not enrolled for replication on this receiver, its wire merge mode disagreed with the locally-resolved mode, or the tenant-isolation gate refused it). For batch calls, true if any entry in the batch was newly merged.
HighWaterMark For point applies (Set / Delete) this is the per-origin HWM after the call - equal to entry.Timestamp when entry.Timestamp advanced the frontier, or the current HWM otherwise (including when Applied is false). For range deletes, saga terminal marks, tombstone-reap envelopes, local-origin no-op rejections, receive-fence deferrals, and receiver-side enrollment / merge-mode / tenant-isolation rejections - none of which reads or advances the HWM - this is HybridLogicalClock.Zero. For batch calls, the pointwise maximum HWM across every distinct origin in the batch.
Deferred true only when a cross-cluster restore saga's durable receive fence has paused inbound apply for the tree, so the entry (or run) was not applied and must be re-shipped once the fence lifts. For batch calls, true if any entry or run in the batch was deferred. Receive paths turn a deferred result into a not-accepted, cursor-preserving ack; every other Applied == false outcome is terminal and lets the sender advance past the entry.

Apply semantics

The applier composes these concerns for every call:

1. Source-HLC and origin preservation

For LwwRegister mode point applies route through the core library's apply seam, which persists the entry's value with the supplied Timestamp and OriginClusterId verbatim - no fresh local HLC is stamped. This is what unlocks transitive replication (A -> B -> C with A's HLC intact) and deterministic LWW resolution against concurrent local writes.

For typed CRDT modes, a steady-state entry carries the producer's typed delta in WalRecord.Delta (authored via the accessor at commit time). The applier forwards that delta verbatim through the same CrdtDelta-recording grain seam used by both the batch path and a locally-authored CRDT write, preserving the source OriginClusterId and folding the delta into the visible state in one grain turn. The source Timestamp is preserved only when the entry is folded as part of a multi-entry batched run, where the batch path re-stamps each item with its source HLC; a per-entry apply - a single-entry batch, an OrMap entry, a causal-buffer drain, or a dead-letter replay - takes a fresh local HLC at the merge point. The receiver therefore records a CrdtDelta revision with per-member ADDED/REMOVED changes - identical history fidelity to a local write - rather than a flattened full-value Set. The fold is wrapped in a LatticeOriginContext.With(originClusterId) scope so the receiver's commit-time observer publishes the foreign origin and the producer-side ship loop filters the resulting entry out. A bootstrap committed-projection row carries the full state in WalRecord.Value with no delta; it has no per-delta shape, so it folds via state-based merge under optimistic concurrency and stays a full-state set, written at a fresh local HLC and without the row's ExpiresAtTicks, so the key is stored as a durable entry. For OrMap mode the receiver resolves the concrete (TKey, TValue) shape through CrdtShapeRegistry (populated by ISiloBuilder.AddOrMapShape<TKey, TValue>); an OrMap-mode apply against an unregistered tree faults with a clear configuration-error message.

Range deletes carry the producer's issue HLC: the authoring cluster pins one HLC for the whole range-delete fan-out, stamps every tombstone it writes with it, and publishes it on the WalRecord. The receiver walks the leaf chain locally and pins every tombstone it writes to that same HLC, so a range delete authored at T cannot overwrite a foreign-origin write whose HLC is strictly greater than T; the remote OriginClusterId rides through an ambient LatticeOriginContext scope so the receiver-side change-feed observer publishes it on every emitted LatticeMutation. A predicate-filtered range delete also ships the explicit set of keys it matched, and the receiver tombstones exactly those keys instead of re-deriving membership from the range bounds. Only an entry persisted by an older producer carries HybridLogicalClock.Zero; for such an entry the receiver falls back to stamping each tombstone with a freshly-ticked local HLC.

2. Snapshot-pinned causal floor (and per-origin high-water-mark)

The applier tracks two distinct per-(TreeId, OriginClusterId) quantities on the high-water-mark grain:

  • The per-origin high-water-mark is the max-applied source HLC. It advances monotonically after every successful point apply and drives FIFO / causal ordering, observability, and the bootstrap handoff. It is not a drop criterion for steady-state point writes.
  • The pinned causal floor is written only by PinSnapshotAsync (the bootstrap handoff) and records the snapshot frontier below which every mutation is already contained in the pinned snapshot. When no snapshot has been pinned the floor is HybridLogicalClock.Zero.

Before applying a point entry the applier reads the pinned floor; if entry.Timestamp <= floor the call is a no-op (Applied = false), because everything at or below a pinned snapshot frontier is already durably present. Otherwise the entry is admitted and its at-most-once guarantee rests on per-key last-writer-wins idempotence at the leaf (a superseded or duplicate write merges to the same state) plus the shadow-forward identity cache below. After a successful point apply the HWM advances monotonically; a laggard's lower advance becomes a no-op.

This is deliberately not a entry.Timestamp <= hwm drop. The per-origin HLC stream is not monotonic in write-ahead-log / ship order: HLCs are stamped per leaf (each leaf carries its own clock) and the write-ahead-log partitions by key hash, so many leaves interleave in one partition and a genuinely-new point write can arrive with a source HLC below the running max-applied HLC. Dropping such an entry on a scalar hwm comparison silently strands it - the cross-cluster data-loss regime this seam is designed to avoid. A legitimately out-of-order-but-new entry that lands below the current HWM is applied and increments apply.fifo_violations (observability only). Typed CRDT modes consult the floor and advance the HWM the same way; state-based merge is naturally idempotent, so the floor gate just short-circuits redundant grain calls for the below-snapshot backlog.

Range deletes bypass the floor by design. Range applies are naturally idempotent at the leaf layer: re-running a range delete on already-tombstoned keys merges to the same state, so dedupe is unnecessary.

3. Local-origin rejection

A WalRecord whose OriginClusterId matches the local cluster id is rejected as a no-op (Applied = false). This check is the receiving cluster's only enforcement that a local-origin entry is never applied back onto its authoring cluster - it is not defence-in-depth behind the sender's outbound origin filter. That filter runs on the sending cluster and decides what the sender ships, not what this receiver accepts, and a hand-built apply pipeline or a test can hand the applier such an entry directly.

4. Receiver-side enrollment and merge-mode gate

Before the concerns above run, the applier gates every inbound entry against this receiver's own per-tree replication configuration, re-resolving the tree's enrollment and merge mode locally instead of trusting the wire. The peer-supplied OriginClusterId is unverified and WalRecord.Mode is a peer-controlled header field, so neither is taken on faith. The gate yields two rejections:

  • Not enrolled here (dropped). A tree that is not enrolled for replication on this receiver is dropped: the call returns Applied = false with HighWaterMark = HybridLogicalClock.Zero, records the apply-duration outcome rejected-not-replicated, and is not dead-lettered. A non-enrolled tree id is peer-controlled, so parking it in a dead-letter queue would let a hostile peer spawn unbounded dead-letter-queue activations; dropping keeps the rejection cheap and bounded. This closes the gap where a peer that holds the mesh secret could otherwise write a tree the cluster deliberately kept cluster-local by not enrolling it - the reserved-prefix core-tree guard covers only the _lattice_ core trees, not the sys--prefixed authorization and identity trees.
  • Enrolled but wire mode mismatched (dead-lettered). A tree that is enrolled but whose peer-supplied wire mode disagrees with the locally resolved merge mode is dead-lettered: the call returns Applied = false with HighWaterMark = HybridLogicalClock.Zero, the entry is enqueued to the tree's dead-letter queue tagged mode_mismatch, and the apply-duration outcome rejected-mode-mismatch is recorded. The tree is enrolled and therefore a bounded id, so parking the entry cannot be abused to spawn unbounded activations. Re-resolving the mode locally rather than trusting the wire field stops a peer from overriding the local merge algebra by shipping a different mode.

The merge mode is always re-resolved locally through the receiver's per-tree resolver (ILatticeReplicationContext.ResolveMergeMode, falling back to the raw LatticeReplicationOptions.ReplicatedTrees map); the wire Mode field is only ever compared against that resolution, never adopted. An applier with neither an injected replication context nor a ReplicatedTrees map has no enrollment signal, so the gate fails closed: every inbound entry is dropped as rejected-not-replicated (not dead-lettered) and a one-time warning is logged. Production registers the replication context, so the gate is always evaluable there.

A run the gate rejects - either way, or for want of an enrollment source - is also never recorded as inbound contact in ReplicationPeerStats, on any receive path, so a peer cannot plant a tree id of its choosing in the peer statistics or the peer-status report. The inbound half of that state is additionally capped, because the origin id of an admitted run is still the peer's own claim; see observability.

Two further receiver-side gates run after an entry clears enrollment, before any high-water-mark read:

  • Tenant isolation (dead-lettered). When tenancy is on, the owning tenant is derived from the tree id alone - never from a wire field - and an entry whose tenant does not exist here, is not resident in the region serving this receiver, or has been suspended or disabled is refused. It is dead-lettered with the foreign_tenant, tenant_offline, or tenant_suspended reason, records the matching rejected-foreign-tenant / rejected-tenant-offline / rejected-tenant-suspended apply-duration outcome, and leaves the high-water-mark unchanged. The rejection is not a deferral, so the receive path still acknowledges the batch and the sender advances past the entry rather than re-shipping it; once the tenant exists, becomes resident, or is reinstated, the parked entry can be replayed from the dead-letter queue. With tenancy off the gate is inactive and costs nothing.
  • Restore receive fence (deferred). While a cross-cluster restore saga has paused inbound apply for the tree, the entry is not applied: the call returns Applied = false with Deferred = true, and the sender keeps its cursor and re-ships once the fence lifts. A single-entry apply records the deferral under the dedup apply-duration outcome; the batch path defers a multi-entry run whole and records no apply-duration sample for it.

The batch path applies the same classification once per run, except that it checks the restore receive fence first: a run for a fenced tree is deferred whole before it is classified. The wire mode is part of the run key - a run is a contiguous (TreeId, OriginClusterId, Mode) segment, so a mode change starts a new run that is classified on its own - and the representative first entry therefore classifies the whole run. A rejected run neither merges nor advances the per-origin high-water-mark; every entry still records its matching apply-duration outcome so per-entry receiver observability is preserved, while a single warning is logged per run rather than per entry to avoid a log-flood amplification from a hostile peer.

5. Shadow-forward dedupe cache

A structural rewrite that shadow-forwards a user write into a different shard - a shard split, or a shard consolidation (merge), which reuses the split's shadow-write window - generates a duplicate-emit pair: one entry from the originating shard's commit, one from the destination shard's commit, both carrying identical (originClusterId, timestamp, key, op) identity tuples, because the forward carries the original last-writer-wins value and its HLC. An atomic-write abort is not a source of such pairs: it issues no per-key rollback writes, only an abort recorded in the tree's transaction registry and TxAbort terminals that discard the prepared writes. Because the pinned-floor gate does not drop above-floor point writes, both deliveries reach the apply path; the identity cache is what collapses the redundant second grain hop before it happens.

The applier holds a per-tree bounded FIFO cache of recently-applied identity tuples (LatticeReplicationOptions.ShadowForwardDedupeCacheSize, default 4096, validator floor 64). The cache is consulted after the pinned-floor gate so floor-deduped entries do not pollute it (which preserves operator-driven re-pin semantics where lowering the pinned floor must re-admit previously-deduped identity tuples). On cache hit the apply is suppressed with Applied = false and the apply-duration histogram is tagged outcome=shadow-forward-dedup. Range deletes bypass the cache - they are applied before it is consulted - and the leaf layer is naturally idempotent for range applies.

The cache is a fast-path optimisation, not the correctness backstop. It suppresses the duplicate-emit pair before the apply grain hop; if an entry it would have caught has been evicted under sustained churn, the duplicate still re-applies to the same state under per-key last-writer-wins idempotence at the leaf, so an eviction can never cause a divergent re-merge - it only costs one redundant grain hop.

6. Causal-dependency gate

Entries authored with causal-plus tracking carry a VectorClock frontier. Before applying such an entry the receiver fetches its local vector clock and checks that every component of the entry's frontier is dominated-or-equal locally; an entry with an unsatisfied dependency is parked in the per-tree bounded causal-apply buffer and retried each time a later apply advances the local clock. The buffer is bounded by CausalBufferMaxEntries (default 1024) and CausalBufferMaxBytes (default 16 MiB); when a new park would exceed either bound, the oldest parked entries are evicted to the tree's dead-letter queue with reason hlc_skew, and each evicted entry's shadow-forward identity reservation is released, so a dead-letter replay or a peer re-delivery of it is applied rather than suppressed as a duplicate. Two frontier components are exempt from the check:

  • The entry's own origin diagonal. The per-origin high-water-mark tracks that origin's own FIFO progression, so requiring the local clock to dominate the diagonal would deadlock the very entry being applied.
  • The receiver's own cluster id. The receiver-side local vector clock tracks only foreign-applied frontiers - it never advances its own diagonal - but the receiver durably holds every write it authored itself, so any dependency on one of the receiver's own writes is trivially satisfied. Without this exemption a peer entry whose frontier references a write the receiver originated (for example, site C's post-partition write that causally follows site A's pre-partition write, once an A-C partition heals) would park forever against a perpetually-zero self-component and stall convergence.

Validation

ApplyAsync throws ArgumentException when:

  • entry.TreeId is null or empty.
  • entry.OriginClusterId is null or empty.
  • entry.Op == Set on an LwwRegister entry and entry.Value is null, or entry.Op == Set on a CRDT-mode entry that carries neither a typed Delta nor a full-state Value.
  • entry.Op == DeleteRange and entry.EndExclusiveKey is null, or the range delete carries atomic-batch metadata (AtomicBatchSize > 0).
  • A prepared entry (IsPrepared == true) carries an empty TransactionId, or a prepared Set carries a null Value.
  • A saga terminal mark (TxCommit / TxAbort) carries no usable shard index or an empty TransactionId.

InvalidOperationException is thrown when:

  • entry.Mode has no apply rule, or entry.Op is not a point-apply kind. The merge-mode gate normally intercepts an unknown wire mode first, dead-lettering it as mode_mismatch, because it cannot match the locally resolved mode.
  • An OrMap tree has no registered (TKey, TValue) shape.
  • A typed CRDT state-merge exhausts its CAS retry budget under sustained contention on the target key.

OperationCanceledException is thrown when the supplied CancellationToken is already cancelled or fires during a grain call.

Registration

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

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

Resolve it from inside a silo-side service (typically a transport adapter or a hosted-service inbound pipeline) via constructor injection on IReplicationApplier. The applier is not exposed on the cluster client - it is a silo-local seam by design, because the apply path must run inside the cluster that owns the receiving tree.

Threading and concurrency

The applier is a silo-wide singleton. It holds no per-call state, only process-local per-tree structures - the shadow-forward dedupe cache, the causal-apply buffer, and the FIFO-diagnostic tracker - and all durable coordination flows through the per-origin high-water-mark grain (single-threaded under Orleans turn semantics) and the per-tree apply grain (StatelessWorker). The HWM grain serialises each individual read and advance, but it does not serialise whole ApplyAsync calls: two concurrent deliveries for the same (tree, origin) pair can both pass the floor gate before either advances the HWM, which is why the shadow-forward identity cache and per-key last-writer-wins idempotence at the leaf, not the HWM, absorb a concurrent duplicate. Calls for different pairs are independent.

Bootstrap handoff

The pinned causal floor is the explicit handoff contract for the bootstrap protocol: a newly-bootstrapped peer calls PinSnapshotAsync, which atomically installs both the per-origin HWM and the pinned floor at the snapshot's authoring frontier, then resumes incremental replication from that pinned frontier with exactly-once apply guarantees across the snapshot / incremental boundary. PinSnapshotAsync replaces the floor rather than max-merging it, so a peer that rewinds to an older snapshot (a pinned value lower than the receiver's prior frontier) correctly lowers the floor and re-admits entries above the new, lower cut. The grain therefore exposes this unconditional pin alongside the monotonic HWM advance and a GetPinnedFloorAsync read.

Caveats

  • Range deletes preserve the producer's issue HLC, not per-leaf HLCs. The wire carries the single HLC the producer pinned for the whole range delete and the receiver stamps every tombstone with it, so LWW resolution against a concurrent write compares against the producer's authoring HLC rather than the receiver's local clock. An entry persisted by an older producer carries HybridLogicalClock.Zero and falls back to fresh local HLCs, where LWW resolution depends on the local clock at apply time. Idempotence at the leaf layer is what makes a re-applied range delete safe.
  • The HWM is per-origin, not per-shard. A receiver applying entries from origin X against a tree split into N shards advances a single HWM row keyed (tree, X) regardless of which shard the entry targets. This is intentional: the HWM contract is the bootstrap-handoff seam, and bootstrap operates per-origin not per-shard.

Batch apply path

Inbound transports deliver batches of WalRecord records, not single entries: a 256-entry gRPC push from a single producer is one network round-trip carrying 256 mutations. ApplyBatchAsync is the seam that lets the receiver process such a batch as one logical operation rather than 256 independent ApplyAsync calls - for each contiguous run of entries sharing a tree, origin, and merge mode it reads the high-water-mark and the pinned causal floor once, merges the run's plain point writes in one batched grain call instead of one apply call per entry, advances the high-water-mark once, and drains the causal-apply buffer once at the end of the run instead of after every successful apply.

The default-interface-method body provides backward-compatible semantics: it loops over ApplyAsync and aggregates the per-entry results - Applied if any entry was newly merged, the pointwise-maximum HighWaterMark, and Deferred if any entry was deferred - so any custom IReplicationApplier written before the batch seam existed continues to work without changes, and a restore saga's receive-fence deferral returned by its ApplyAsync still reaches the receive path as a deferred, cursor-preserving result. The shipped applier overrides the batch path with the optimised implementation described below.

Run grouping

The optimised batch path walks the inbound list and identifies maximal contiguous runs of entries that share the same (TreeId, OriginClusterId, Mode) tuple (a well-formed batch carries one merge mode per run, so including the mode never splits a legitimate run). For a 256-entry batch shipped by a single producer the entire batch is one run; for an interleaved batch (e.g. a snapshot recovery merge that intersperses entries from two origins) the path emits one run per contiguous group, each amortised independently. Within a run:

  • A single GetAsync reads the persisted per-origin HWM and a single GetPinnedFloorAsync reads the pinned causal floor at the start of the run.
  • The pinned floor is constant for the run, so every entry in the run is dedup-tested against the same floor with no further round trips; there is no in-batch running-HWM accumulator (a below-max-applied-HLC entry is a genuine write under non-monotonic per-origin HLC, not a duplicate, so it must not be dropped mid-run).
  • Causal-dependency entries fetch the local vector clock lazily on first use and reuse it until an apply has occurred, at which point a localVcDirty flag forces a re-fetch on the next causal-dep check.
  • A single TryAdvanceAsync advances the persisted HWM to the highest applied HLC at the end of the run.
  • The causal-apply buffer is drained once, if the run advanced the persisted HWM.
  • Non-prepared LwwRegister Set / Delete entries that pass classification are deferred and merged in one batched grain call per run, and non-prepared typed-CRDT Set entries that carry a delta (every CRDT mode except OrMap) fold in one batched delta call. Range deletes, saga terminal marks, prepared entries, CRDT-mode deletes, and OrMap or delta-less CRDT entries flush the pending batch first and take their own per-entry apply hop.

For a 256-entry single-origin LwwRegister batch this collapses roughly 4 x 256 = 1,024 grain round-trips on the per-entry path (a high-water-mark read, a pinned-floor read, a point apply, and a high-water-mark advance per entry) to about four (one of each, with the 256 point applies folded into one batched merge) - the dominant receiver-side cost on every inbound push.

Preserved per-entry semantics

Every classification the per-entry path produces survives the batch path, except the tombstone-reap no-op (the last bullet):

  • Range-delete entries bypass the pinned-floor gate and apply unconditionally (a range apply is naturally idempotent at the leaf layer).
  • Local-origin runs classify every entry as Dedup with HighWaterMark = HybridLogicalClock.Zero and emit no grain calls.
  • Below-floor dedup against the run's pinned causal floor is bit-equivalent to per-entry dedup against a freshly-read floor, because the floor is written only by PinSnapshotAsync and is therefore constant across a run.
  • Causal-park is exercised per-entry; only the local-vector-clock fetch is lazy.
  • Per-entry instrumentation (ApplyDuration, ApplyLag, ApplyFifoViolations) is recorded inside the per-entry loop so observability is preserved verbatim.
  • Single-entry batches defer to ApplyAsync so behaviour is bit-identical with the legacy receiver for the trivial case.
  • Tombstone-reap envelopes are not acknowledged as dedup on the batch path: it has no no-op branch for them, so a multi-entry batch carrying one faults on that entry, because the point-apply step has no rule for it. The dead-letter-tracking decorator then falls back to per-entry apply, which does acknowledge the envelope as a no-op. The sender never ships these envelopes, so only an older shipper or a hand-built caller can deliver one.

Failure model

Per-entry failures inside the batch surface as ApplyAsync-equivalent exceptions. The gRPC receiver endpoint wraps the batch call in a transport-level exception so the sender's backoff/retry loop kicks in for the whole batch - partial-batch acceptance is not a guarantee the seam offers. The dead-letter-tracking applier decorator falls back to per-entry routing when any entry in the batch already has retry history, or when the batch call throws part-way, so its DLQ accounting is exact. On that per-entry fallback (and for a single-entry batch) the decorator itself records the inbound per-peer contact, which the batch path otherwise records once per run; entries from runs the failed batch attempt already recorded can therefore record contact a second time for the same push (see Observability).

Parallel apply across independent runs

Under multi-tree load the per-run walk can serialise otherwise-independent work: a batch that interleaves runs from several trees applies them one after another even though they share no per-tree state, inflating apply latency and apply.lag (which now also drives receiver back-pressure, so slow applies translate directly into sender throttling).

LatticeReplicationOptions.ApplyMaxParallelRuns bounds how many independent runs the batch path may apply concurrently. Independence is defined at the tree granularity:

  • Runs targeting distinct trees may apply in parallel. Distinct trees share no per-tree state - separate causal-apply buffers, shadow-forward dedupe caches, high-water-mark grains, and apply grains - so concurrent apply cannot reorder or interleave their work.
  • Runs that share a tree (different origins of the same tree) stay in one ordered group and apply strictly sequentially in write-ahead-log order. The per-tree causal-apply buffer and shadow-forward dedupe cache are shared across a tree's origins, so keeping same-tree runs serialised guarantees those structures observe the exact access order the fully-sequential path produces.

Parallelism is therefore only ever introduced across independent runs, never within one. Every within-run ordering invariant holds unchanged regardless of the configured degree of parallelism: per-origin FIFO, the causal dependency gate and its bounded buffer, per-origin high-water-mark monotonicity, and atomic-batch (saga) apply boundaries. A multi-entry run still collapses to a single batched merge; an atomic batch still applies as a unit on its owning run.

The effective degree of parallelism for a given batch is the largest ApplyMaxParallelRuns configured for any tree in that batch (the option resolves per tree), clamped to the number of distinct trees present in it, and is bounded by a per-batch semaphore so concurrency can never amplify local WAL saturation beyond the configured cap. It is surfaced on the apply.parallel_runs histogram (see observability).

Default posture: fully sequential. ApplyMaxParallelRuns defaults to 1, which is exactly the historical behaviour - the batch path walks every run in order and awaits each before the next. The single-tree batch (the overwhelmingly common inbound shape, since the transport ships per-(tree, peer)) always takes the allocation-free sequential walk regardless of the configured value, because cross-tree parallelism is moot when there is only one tree. Raise the value conservatively, per workload, only after validating parallel apply for that topology.

Cross-cluster atomic visibility - receiver seam

SetManyAtomicAsync sagas authored on the source cluster ride the standard WAL replication transport: every prepared per-key write emits a Set / Delete WalRecord with IsPrepared = true and a non-empty TransactionId, and the saga's terminal phase emits one TxCommit (or TxAbort) WalRecord per touched shard. The shipper preserves these records verbatim; the receiver seam interprets them through three additional internal apply hops:

Apply hop Wire trigger Receiver behaviour
Prepared set Op == Set && IsPrepared == true Stages the write under the saga's TransactionId in the destination leaf's per-tx pending bucket. The visible projection is unchanged - public readers (GetAsync, KeysAsync, etc.) do not observe the prepared entry.
Prepared delete Op == Delete && IsPrepared == true Stages a tombstone under the saga's TransactionId in the same pending bucket. The pre-saga value remains visible to public readers until the terminal arrives.
Transaction terminal Op == TxCommit or Op == TxAbort Records this per-source-shard terminal arrival in the per-tree transaction registry (keyed by txid, source shard index, commit/abort outcome, and atomic shard count) to tally arrivals. While the tally is not final the registry mark stays unset and the receiver leaves' pending buckets stay in place so reads remain all-or-nothing. Only on the final arrival does the receiver mark the per-tree transaction-registry entry and pre-fan the terminal across the transitive split-forward closure of every observed source-shard in a single parallel hop. On commit every pending entry under the TransactionId flips into the visible projection; on abort the pending entries are dropped.

The batch-apply classifier excludes any entry with IsPrepared == true from the batched LWW fast-path so prepared Set / Delete records are always routed through the per-entry prepared-set / prepared-delete apply hops. Without this exclusion the prepared writes would commit directly into the receiver leaf's visible projection and the saga's terminal mark would find no matching pending entries to flip - so the cross-cluster reader would observe the prepared write as visible before the registry gate flipped, purely as a function of whether the inbound run happened to be batched or single-entry. Unprepared writes continue to consume the batched merge path.

The per-source-shard arrival tally is the receiver-side multi-shard atomic-visibility gate. A saga that touched N source shards emits N independent terminal records, one per source shard, that ship through the change feed under independent backpressure / batching cadences. Each terminal carries the saga's authoritative touched-shard count in the additive WalRecord.AtomicShardCount slot, which the receiver feeds into the terminal-arrival tally to compute finality. A receiver running a pre-gate producer sees atomicShardCount == 0 on every terminal, which the gate treats as "no expected-total information" and falls back to first-terminal-wins semantics - equivalent to the pre-gate behaviour and wire-compatible across mixed-version deployments.

The producer-side ship filter explicitly bypasses per-tree KeyFilter and KeyPrefixes for TxCommit and TxAbort records: a saga whose prepared keys passed the filter must have its terminal delivered or the receiver-side pending bucket leaks. The empty-origin guard and the cycle-break filter still run before the bypass, so a malformed or self-loopback terminal is still rejected.

Cross-tree terminals (receiver barrier)

A terminal that belongs to a cross-tree atomic write (IGrainFactory.SetManyAtomicAsync) carries two additional slots - WalRecord.CrossTreeOperationId and WalRecord.CrossTreeParticipants (the canonical participant tree-id set). The applier threads these into the transaction-terminal apply hop as a crossTreeOperationId plus a receiver-scoped wait set. The wait set is the participant set intersected with the trees this receiver actually replicates (its per-tree enrollment - the LatticeReplicationOptions.ReplicatedTrees declaration or a runtime-enabled tree); the tree that received the terminal is always included. A participant tree not replicated here is excluded, so a cross-tree batch spanning a mix of replicated and non-replicated trees stays valid - the barrier completes on the present subset rather than waiting forever on a tree that never ships here.

Once a tree's per-source-shard gate is final, a cross-tree terminal does not flip that tree's registry directly. Instead the receiver durably registers the tree's local txid as delegated to a receiver coordinator grain (keyed by (originClusterId, operationId)) and notifies it of this tree's arrival and commit/abort vote. The coordinator decides only once a terminal has arrived for every tree in the wait set, committing iff every arrived tree voted commit. Before the decision, a delegated read on any participating tree's registry resolves InFlight against the coordinator, so every tree stays invisible (an unreachable coordinator resolves Indeterminate, which likewise keeps the keys invisible); after it, the receiver flips every participating tree together. The coordinator only ever returns the decision (it never calls back into a tree grain); the calling tree grain performs the per-tree finalizes - itself inline, siblings via their apply grains - so there is no circular wait. A null/empty crossTreeOperationId routes the terminal through the legacy single-tree gate unchanged.

Public readers therefore observe the receiver-side same-cluster atomic-visibility property end-to-end: at every point in time, either every key the saga prepared on the receiver is at its post-saga value (after the commit terminal applies) or none of them is (during the prepare window or after an abort). The HLC the visible value carries is the source cluster's HLC verbatim - the receiver's wall-clock progression does not bump it - so transitive LWW resolution (A -> B -> C with A's HLC intact) holds across saga output identically to single-key cross-cluster writes.

What ships today:

  • WalRecord.AtomicBatchSize, AtomicBatchIndex, AtomicShardCount, TransactionId, and IsPrepared are preserved on the wire end-to-end. The receiver consumes TransactionId and IsPrepared to drive the prepared / terminal staging path, and AtomicShardCount to drive the per-source-shard arrival tally on terminal records. It also reads AtomicBatchSize to recognise saga prepare-phase entries (which bypass the pinned-floor and causal-dependency gates) and to reject a range delete that carries atomic-batch metadata, and it forwards AtomicBatchSize and AtomicBatchIndex with every prepared write into the per-transaction staging hop.
  • The receiver-side multi-key atomic apply seam (and its associated Atomic + Apply value types) was deleted by the universal- visibility ship. Cross-cluster atomic visibility is provided exclusively by the per-key prepared / per-shard terminal-mark apply hops described above; the local SetManyAtomicAsync saga inside Orleans.Lattice uses the same point-apply seam as a non-saga write, with the IsPrepared flag selecting the staging behaviour.
  • Local single-tree atomic visibility (within one cluster) is shipped end-to-end via the per-tree transaction registry linearization point; see Atomic Writes for the protocol and Consistency for the read-path dial-back. The cross-cluster receiver seam reuses the same registry grain.
  • The producer-side per-key WAL filter shipped earlier. Hosts that need to bound the change feed at commit time configure ReplicatedTrees, KeyFilter, or KeyPrefixes on LatticeReplicationOptions - see wal.md. TxCommit and TxAbort records are exempt from the per-key filter as described above.