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 isHybridLogicalClock.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 = falsewithHighWaterMark = HybridLogicalClock.Zero, records the apply-duration outcomerejected-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 thesys--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 = falsewithHighWaterMark = HybridLogicalClock.Zero, the entry is enqueued to the tree's dead-letter queue taggedmode_mismatch, and the apply-duration outcomerejected-mode-mismatchis 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, ortenant_suspendedreason, records the matchingrejected-foreign-tenant/rejected-tenant-offline/rejected-tenant-suspendedapply-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 = falsewithDeferred = true, and the sender keeps its cursor and re-ships once the fence lifts. A single-entry apply records the deferral under thededupapply-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.TreeIdis null or empty.entry.OriginClusterIdis null or empty.entry.Op == Seton anLwwRegisterentry andentry.Valueis null, orentry.Op == Seton a CRDT-mode entry that carries neither a typedDeltanor a full-stateValue.entry.Op == DeleteRangeandentry.EndExclusiveKeyis null, or the range delete carries atomic-batch metadata (AtomicBatchSize > 0).- A prepared entry (
IsPrepared == true) carries an emptyTransactionId, or a preparedSetcarries a nullValue. - A saga terminal mark (
TxCommit/TxAbort) carries no usable shard index or an emptyTransactionId.
InvalidOperationException is thrown when:
entry.Modehas no apply rule, orentry.Opis not a point-apply kind. The merge-mode gate normally intercepts an unknown wire mode first, dead-lettering it asmode_mismatch, because it cannot match the locally resolved mode.- An
OrMaptree 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.Zeroand 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
Xagainst 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
GetAsyncreads the persisted per-origin HWM and a singleGetPinnedFloorAsyncreads 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
localVcDirtyflag forces a re-fetch on the next causal-dep check. - A single
TryAdvanceAsyncadvances 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
LwwRegisterSet/Deleteentries that pass classification are deferred and merged in one batched grain call per run, and non-prepared typed-CRDTSetentries that carry a delta (every CRDT mode exceptOrMap) fold in one batched delta call. Range deletes, saga terminal marks, prepared entries, CRDT-mode deletes, andOrMapor 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
DedupwithHighWaterMark = HybridLogicalClock.Zeroand 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
PinSnapshotAsyncand 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
ApplyAsyncso behaviour is bit-identical with the legacy receiver for the trivial case. - Tombstone-reap envelopes are not acknowledged as
dedupon 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, andIsPreparedare preserved on the wire end-to-end. The receiver consumesTransactionIdandIsPreparedto drive the prepared / terminal staging path, andAtomicShardCountto drive the per-source-shard arrival tally on terminal records. It also readsAtomicBatchSizeto 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 forwardsAtomicBatchSizeandAtomicBatchIndexwith every prepared write into the per-transaction staging hop.- The receiver-side multi-key atomic apply seam (and its associated
Atomic+Applyvalue 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 localSetManyAtomicAsyncsaga insideOrleans.Latticeuses the same point-apply seam as a non-saga write, with theIsPreparedflag 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, orKeyPrefixesonLatticeReplicationOptions- seewal.md.TxCommitandTxAbortrecords are exempt from the per-key filter as described above.