Snapshot / bootstrap export
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 snapshot-bootstrap.md, and llms.txt lists every page.Orleans.Lattice.Replication ships an ISnapshotProvider seam used by
the snapshot/bootstrap protocol to seed a newly-joining peer (or a
peer that has fallen off the WAL) before switching it to incremental
replication.
The seam is registered by AddLatticeReplication and resolved per
host via TryAddSingleton, so a host that needs a more efficient
storage-specific export can pre-register its own implementation
before calling AddLatticeReplication.
Public surface
| Type | Shape | Purpose |
|---|---|---|
ISnapshotProvider |
Task<SnapshotStream> ExportAsync(string treeName, HybridLogicalClock asOfHlc, CancellationToken ct) + Task<SnapshotStream> ExportAsync(string treeName, string sourceClusterId, HybridLogicalClock asOfHlc, CancellationToken ct) + Task<SnapshotStream> ExportAsync(string treeName, IReadOnlyList<LeafReReplayRange> ranges, HybridLogicalClock asOfHlc, CancellationToken ct) |
Streaming as-of-HLC export of a tree's primary state. The three-arg overload carries the sender-cluster identifier and is the one the bootstrap coordinator invokes; intra-cluster implementations inherit a default interface method that delegates to the two-arg overload after validating sourceClusterId. The range-scoped overload - used by the bootstrap fallback - yields only entries inside the half-open [StartKey, EndKey) ranges; its default interface method filters the whole-tree export client-side. |
SnapshotStream |
sealed class with TreeName, AsOfHlc, CausalStableFrontier (VersionVector), Entries (IAsyncEnumerable<SnapshotEntry>) |
Carries the export metadata + entry stream produced by ExportAsync. |
SnapshotEntry |
readonly record struct with Key, Value, Timestamp, IsPrepared, IsTombstone, TransactionId, SourceShardIndex, AtomicBatchSize, AtomicBatchIndex, ExpiresAtTicks, Delta, Mode |
A single exported record stamped with its commit-time HLC so the receiver can pin the value at exactly that timestamp. A committed-projection row sets Key, Value, Timestamp, and ExpiresAtTicks; a prepared saga row additionally sets IsPrepared, IsTombstone, TransactionId, and the typed CRDT Delta / Mode (see Snapshot and in-flight atomic visibility). SourceShardIndex is reserved and always 0. |
SnapshotEntry is alias olr.se.
Semantics
asOfHlc = HybridLogicalClock.Zerodisables the upper-bound filter and includes every live entry in the tree. This is the common case when seeding a fresh peer that has no incremental cursor yet.asOfHlc > Zerofilters out entries whose stamped commit-time HLC is strictly greater thanasOfHlc. After the drain the receiver's snapshot pin replaces both its per-origin high-water-mark vector and its pinned causal floor with the snapshot's causal-stable frontier (the source cluster's own coordinate sealed at or above every entry the drain applied), and the steady-state floor gate inIReplicationApplierthen drops any incremental entry at or below that frontier, which keeps the handoff exactly-once across the snapshot/incremental boundary. See "Bootstrap drain bypasses the pinned-floor gate and the high-water-mark advance" below for the receiver-side state machine that keeps the in-drain apply idempotent without relying on that gate.CausalStableFrontieris the producer's causal-stable frontier at snapshot time - the pointwise minimumVersionVectoracross every consumer that has reported a vector throughIWalCursorRegistry.GetCausalStableAsync. When no consumer has reported a VC-shaped cursor (single-peer cluster, fresh deployment, host using the legacy HLC-only overload), the provider falls back to the producer's per-tree local vector clock from the per-tree high-water-mark store's current vector - a strict superset of the meet that is safe as a snapshot cut-point. The receiver records this(asOfHlc, frontier)cut-point when the export opens and pins the frontier on the per-tree high-water-mark store only after every snapshot entry has been applied, so the causal dependency check on the first incremental entry after the pin runs from a non-empty frontier.- Tombstoned and expired keys are not emitted. Only live entries
reach the receiver through the committed projection; the tombstone
state is reconstructed from the incremental WAL after the snapshot
completes. The one exception is an in-flight saga's prepared delete,
which ships as a prepared row with
IsTombstoneset (see Snapshot and in-flight atomic visibility). - A live key's TTL is carried. Every exported row - a
committed-projection row as well as a prepared saga row - carries the
source entry's absolute
ExpiresAtTicks(0for a durable key), so on a last-writer-wins tree a key that has a TTL on the source expires at the same instant on the bootstrapped peer. On a typed CRDT tree the receiver folds a committed row's full state through a state-based merge that does not apply the carried expiry, so the key is written there as a durable entry.
Default implementation
The default snapshot provider enumerates the local tree
via the public ILattice.EntriesAsync surface and stamps each entry
with its commit-time HLC via ILattice.GetWithVersionAsync. It is
correct for intra-cluster seeding (snapshot-as-a-tool: an operator
snapshots a tree and restores it later in the same cluster, where the
local tree is the authoritative source) but pays a per-key version
round-trip on top of the leaf-chain enumeration. A future revision
will swap to a single-pass streaming HLC-threshold scan once the core
library exposes a version-bearing leaf-scan primitive; hosts that need a
faster export today can register their own ISnapshotProvider via DI.
Cross-cluster bootstrap uses a separate receiver-side seam.
On a receiver whose local tree is empty (e.g. a fresh cluster joining
an existing federation), the default snapshot provider would
yield zero entries because it reads the receiver's own tree rather
than the sender's. The "Cross-cluster transport contract" section
below documents the IRemoteSnapshotTransport abstraction and the
sender-side LatticeRemoteSnapshotService; the "Receiver-side
adapter" section documents RemoteSnapshotProvider, the
IBootstrapSnapshotSource implementation the bootstrap state
machine drains from when an IRemoteSnapshotTransport is registered.
The seam is split from ISnapshotProvider so a single silo can
simultaneously act as snapshot sender (its local ISnapshotProvider
is exported to peer receivers via LatticeRemoteSnapshotService) and
snapshot receiver (its IBootstrapSnapshotSource drains from an
upstream peer through the registered transport). Registering an
IRemoteSnapshotTransport alongside AddLatticeReplication is the
active-active signal: the receiver-side seam auto-flips to
RemoteSnapshotProvider, no separate opt-in call required.
Cross-cluster transport contract
The first step of the cross-cluster bootstrap pipeline is the
transport-shaped seam that delivers a snapshot stream from a sender
cluster to a receiver cluster. It is a separate abstraction from the
live-incremental IReplicationTransport so a host can plug a different
binding for the bulk snapshot path (HTTP, blob-store, gRPC) without
disturbing the live tail pipeline.
| Type | Shape | Purpose |
|---|---|---|
IRemoteSnapshotTransport |
Task<RemoteSnapshotMetadata> GetMetadataAsync(string treeName, string sourceClusterId, HybridLogicalClock fromAsOfHlc, CancellationToken ct) + IAsyncEnumerable<SnapshotEntry> RequestSnapshotAsync(string treeName, string sourceClusterId, HybridLogicalClock fromAsOfHlc, CancellationToken ct) |
Transport-shaped sub-interface used by a cross-cluster ISnapshotProvider adapter to fetch a snapshot from a sender cluster. |
RemoteSnapshotMetadata |
readonly record struct with TreeName, SourceClusterId, AsOfHlc, CausalStableFrontier |
Snapshot cut-point captured atomically with the start of the entry stream; alias olr.sm. |
Semantics
- Two RPCs, one cut-point. The receiver invokes
GetMetadataAsyncfirst to capture the sender's cut-point, then invokesRequestSnapshotAsyncwith the sametreeName/sourceClusterId/fromAsOfHlctuple to drain the stream. The metadata RPC returns the(AsOfHlc, CausalStableFrontier)pair the receiver records before the drain and pins on the per-tree high-water-mark store once the drain completes, so the snapshot/incremental handoff stays exactly-once even though metadata and stream travel on separate calls. - Point-in-time view. Implementations MUST guarantee that entries
committed on the sender after the metadata cut-point do not leak
into the corresponding stream call. Receivers treat the stream as a
point-in-time view at
metadata.AsOfHlc; a moving-target stream would violate the cut-point pin and break the causal-stable handoff of the first incremental entry. - Concurrency. Implementations are safe to invoke concurrently
across distinct
(treeName, sourceClusterId)pairs. Concurrent invocation against the same pair is implementation-defined; receivers serialise per pair through the bootstrap coordinator. - Argument validation. Both methods throw
ArgumentNullExceptionwhentreeNameorsourceClusterIdisnullandArgumentExceptionwhen either argument is empty or whitespace-only. - Prepared-transaction state travels in the entry stream. The
metadata DTO carries only the cut-point. In-flight saga state rides
the entry stream as prepared rows (see
Snapshot and in-flight atomic visibility),
so a transport must stream each
SnapshotEntryverbatim -IsPrepared,TransactionId,Delta, andModeincluded - for the bootstrapping peer to keep every saga all-or-nothing across the bootstrap.
Contract test fixture
Implementations import RemoteSnapshotTransportContractTests from
the replication test project and derive a concrete fixture overriding
CreateTransportAsync to plug the transport in front of a
sender-side StubSenderSnapshotProvider. The inherited acceptance
suite pins:
GetMetadataAsyncreturnstreeName/sourceClusterId/AsOfHlc/CausalStableFrontiermatching the staged sender snapshot.RequestSnapshotAsyncstreams every staged entry verbatim.RequestSnapshotAsyncyields an empty stream when the sender has no entries.- Metadata-then-stream is consistent under concurrent sender writes: entries staged after the metadata cut-point do not leak.
ArgumentNullException/ArgumentExceptioninvariants hold for both RPCs.- The stream observes cancellation tokens during enumeration.
The exemplar InMemoryRemoteSnapshotTransport in the replication
test project wraps a local ISnapshotProvider and is the smallest
reference shape for what a wire-bound implementation must preserve.
Sender-side handler
LatticeRemoteSnapshotService is the canonical sender-side
implementation of IRemoteSnapshotTransport. It is registered as a
singleton by AddLatticeReplication and is the seam concrete bindings
(gRPC, in-process loopback, custom HTTP) delegate to when an inbound
metadata/stream RPC arrives on a producer silo. The handler is
deliberately binding-agnostic: the same instance is shared across every
concrete binding the host registers.
| Type | Purpose |
|---|---|
LatticeRemoteSnapshotService |
Sender-side IRemoteSnapshotTransport handler. Validates routing arguments, refuses a tree that is not enrolled for replication on this cluster, invokes the local ISnapshotProvider, and returns the resulting cut-point metadata or streams the resulting entries. Stateless and safe for concurrent invocation across distinct (treeName, sourceClusterId) pairs. |
Semantics
- Sender-side enrollment gate. The requested tree name comes from
the peer, so it is re-resolved against this cluster's own
replication enrollment before anything is read. A tree that is not
enrolled here - or a handler with no enrollment source to decide -
is refused with
UnauthorizedAccessException, so a peer that holds the mesh secret cannot stream out a tree this cluster keeps local, such as thesys-authorization and identity trees. The check sits on the handler, so every binding (gRPC, loopback, custom) inherits it. - Delegation to the local provider.
GetMetadataAsynccallsISnapshotProvider.ExportAsync(treeName, fromAsOfHlc, ct)and returns aRemoteSnapshotMetadatacarrying the resultingSnapshotStream.AsOfHlcandCausalStableFrontier. A pairedRequestSnapshotAsynccall invokesExportAsyncagain with the samefromAsOfHlcfilter and drainsSnapshotStream.Entries. - Point-in-time view. The canonical default snapshot provider
reads the producer's causal-stable frontier once at the start of
the export and applies the receiver-supplied
fromAsOfHlcfilter at entry-emission time. WhenfromAsOfHlcis greater thanZero, writes committed on the sender after the metadata cut-point are excluded from the matching stream call by the per-entry HLC filter, so the stream remains a point-in-time view atmetadata.AsOfHlceven though the metadata RPC and the stream RPC are separate calls. AZerofromAsOfHlc- what a fresh bootstrap passes - disables the filter, so that stream can include writes committed after the metadata call. - Host-replaceable provider. Hosts that register a custom
ISnapshotProvider(for example, a storage-backend-aware export) drive the handler through the standard DI seam: replacing the provider replaces whatLatticeRemoteSnapshotServiceexports, no per-binding configuration change required. A custom provider must preserve the point-in-time semantics above; otherwise the cross-cluster bootstrap may lose entries committed during the drain window.
Receiver-side adapter
RemoteSnapshotProvider is the canonical receiver-side
IBootstrapSnapshotSource implementation. The bootstrap state
machine drains from IBootstrapSnapshotSource; AddLatticeReplication
registers a factory that resolves the seam to RemoteSnapshotProvider
when an IRemoteSnapshotTransport is registered alongside it (the
active-active default), or to a local snapshot-source wrapper
over the silo's local ISnapshotProvider otherwise (the
single-cluster recovery path). The silo's local ISnapshotProvider
is untouched in either case, so a silo that is also a snapshot
sender keeps serving outbound requests from peer receivers via
LatticeRemoteSnapshotService while its own bootstrap drains from
the registered transport.
| Type | Purpose |
|---|---|
IBootstrapSnapshotSource |
Receiver-side seam consumed by the bootstrap state machine. Split from ISnapshotProvider so sender and receiver roles can coexist on a single silo without one DI slot overwriting the other. |
RemoteSnapshotProvider |
Cross-cluster IBootstrapSnapshotSource. Calls IRemoteSnapshotTransport.GetMetadataAsync once to capture the sender-side cut-point, then drains RequestSnapshotAsync and yields each entry through the existing SnapshotStream shape. Stateless and safe for concurrent invocation across distinct (treeName, sourceClusterId) pairs. |
| Local snapshot-source wrapper | Single-cluster default IBootstrapSnapshotSource. Forwards both ExportAsync overloads to the silo's local ISnapshotProvider. |
Semantics
- Three-arg overload only. The cross-cluster adapter implements
only the three-arg
ExportAsync(treeName, sourceClusterId, asOfHlc, ct)overload; the legacy two-arg overload throwsInvalidOperationExceptionbecause the adapter cannot address a sender peer without the sender cluster id. The bootstrap coordinator always invokes the three-arg overload with the value read from the coordinator's persisted source-cluster id, so this branch is unreachable in the normal flow and indicates an integration bug when it fires. - Metadata-then-stream consistency. The adapter calls
GetMetadataAsynconce and pins the returnedAsOfHlcandCausalStableFrontieron theSnapshotStreamit returns; the entries on that stream come from a pairedRequestSnapshotAsynccall against the same(treeName, sourceClusterId, fromAsOfHlc)tuple. The transport contract guarantees the two calls describe the same snapshot, so the receiver-side state machine pins a point-in-time frontier even though the metadata and stream RPCs are separate. - Transport-agnostic. The adapter knows nothing about the
concrete binding (gRPC, in-process loopback, custom HTTP); the
host wires the binding by registering an
IRemoteSnapshotTransportsingleton in DI alongsideAddLatticeReplication. - Active-active by default. A silo that registers an
IRemoteSnapshotTransportbecomes both a sender (serving outbound snapshot requests viaLatticeRemoteSnapshotService, which drives the localISnapshotProvider) and a receiver (draining inbound snapshots throughRemoteSnapshotProvider, which drives the registered transport). The two seams do not collide: each grain consumes its own DI slot. Hosts that want to force the local-only bootstrap path even with a transport present can pre-register a customIBootstrapSnapshotSourcebeforeAddLatticeReplicationand the default factory becomes a no-op.
Sample wiring on a silo that bootstraps from a peer cluster (registering the transport is sufficient - the bootstrap seam auto-flips to the cross-cluster adapter):
using Microsoft.Extensions.DependencyInjection;
using Orleans.Hosting;
siloBuilder.AddLatticeReplication(opts => opts.ClusterId = "site-b");
// The concrete transport binding the host plugs in (gRPC, custom
// HTTP, or a test-only loopback) implements IRemoteSnapshotTransport
// against the sender cluster. Registered as a singleton so the
// adapter resolves a stable instance per silo. A real implementation
// connects to the sender's binding endpoint; the no-op shown here
// stands in for the host-supplied implementation in this snippet.
// Registering the transport is sufficient: AddLatticeReplication's
// IBootstrapSnapshotSource factory observes the registration and
// resolves the seam to RemoteSnapshotProvider automatically.
siloBuilder.Services.AddSingleton<IRemoteSnapshotTransport>(_ =>
throw new NotImplementedException("Plug in your transport binding."));
A reference for the receiver-side wiring lives in the replication
test project's RemoteSnapshotProviderIntegrationTests, which uses
an in-process transport stub to round-trip the receiver-side adapter
against the real sender-side LatticeRemoteSnapshotService running
on a peer cluster.
gRPC binding
The Orleans.Lattice.Replication.Grpc package ships the canonical
IRemoteSnapshotTransport binding alongside the live-push
IReplicationTransport. Both transports share a single options
type, one peer-URI map, and the same TLS, shared-secret, and
channel-reuse conventions.
| Type | Role |
|---|---|
| Client-side gRPC snapshot transport | Client-side IRemoteSnapshotTransport that hosts the cross-cluster GetMetadataAsync unary call and the RequestSnapshotAsync server-streaming call. One gRPC channel is cached per sourceClusterId for the transport's lifetime, built through the same hardened channel pipeline (TLS gate and shared-secret call credentials) as the live-push transport, which keeps its own channel cache. |
| Sender-side gRPC snapshot service | Sender-side ASP.NET Core gRPC service that delegates each call to the local LatticeRemoteSnapshotService (which in turn drives the host-registered ISnapshotProvider). |
LatticeReplicationGrpcOptions |
Per-peer endpoint map (Peers), TLS-required-by-default gate (AllowPlaintextEndpoints), optional per-channel configuration hook (ConfigureChannel), and an override for the local cluster id used in outbound headers. The same options instance drives both the live-push transport and the snapshot transport. |
A silo can simultaneously serve outbound snapshot requests and bootstrap its own tree from a peer; the default composition is active-active. A single helper pair wires both directions:
using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.DependencyInjection;
using Orleans.Lattice.Replication;
using Orleans.Lattice.Replication.Grpc;
siloBuilder.AddLatticeReplication(opts => opts.ClusterId = "site-b");
// Cross-cluster gRPC binding: registers both the live-push
// IReplicationTransport (outbound batches) and the snapshot
// IRemoteSnapshotTransport (bootstrap pulls). Registering the
// snapshot transport auto-flips the receiver-side
// IBootstrapSnapshotSource seam to the cross-cluster adapter so
// the bootstrap state machine drains through gRPC; the sender-side
// ISnapshotProvider remains the local tree.
siloBuilder.Services.AddLatticeReplicationGrpc(opts =>
{
opts.Peers["site-a"] = new Uri("https://snap.site-a.example/");
});
// In the ASP.NET Core endpoint composition:
// app.MapLatticeReplicationGrpc();
A silo that only sends snapshots (it never bootstraps from a peer)
leaves Peers empty; a silo that only receives snapshots omits the
MapLatticeReplicationGrpc endpoint mapping. The active-active
default is opt-in by registration, not by a separate role flag.
The shared-secret authentication interceptor that gates the existing
orleans.lattice.replication.LatticeReplication push service also
recognises the orleans.lattice.replication.LatticeRemoteSnapshot
service (and the orleans.lattice.replication.LatticeSaga control
channel) and enforces on every RPC shape - unary, server-streaming,
client-streaming, and duplex - so the same shared-secret credential
(LATTICE_REPLICATION_SECRET on the sender, checked against the
receiver's accepted set, LATTICE_REPLICATION_ACCEPTED_SECRETS, or a
custom ILatticeReplicationSecretSource) and the same
LatticeReplicationSecurityOptions.RequireAuthentication switch cover
inbound snapshot calls without additional wiring. With
LatticeReplicationSecurityOptions.BindCredentialToOriginCluster on
(the default), the interceptor also refuses a snapshot call that
carries no x-lattice-replication-origin header, or whose presented
secret is not the one the exporting cluster would itself send to the
named origin - see
Configuration.
The client translates RpcException(StatusCode.Cancelled) raised while
the caller's own cancellation token is cancelled into the canonical
OperationCanceledException, so receivers can rely on the
same cancellation contract whether the transport is the gRPC binding,
the in-process loopback, or a host-supplied custom binding.
Sample usage
using Orleans.Lattice;
ISnapshotProvider provider = client.ServiceProvider.GetRequiredService<ISnapshotProvider>();
SnapshotStream snapshot = await provider.ExportAsync("orders", HybridLogicalClock.Zero, cancellationToken);
await foreach (SnapshotEntry entry in snapshot.Entries.WithCancellation(cancellationToken))
{
// Apply each entry on the receiver. Use the entry's commit-time
// Timestamp so transitive replication paths (A -> B -> C) preserve
// the originating HLC.
_ = entry.Key;
_ = entry.Value;
_ = entry.Timestamp;
}
VersionVector frontier = snapshot.CausalStableFrontier;
HybridLogicalClock asOf = snapshot.AsOfHlc;
_ = (frontier, asOf);
In a host the ISnapshotProvider is resolved from DI on the sender
side; the default snapshot provider is shown above for illustration. The
receiver pins the snapshot's CausalStableFrontier on its per-tree
high-water-mark store after draining the entry stream, so the causal
dependency check on the first incremental entry after the pin runs
from a non-empty frontier.
Receiver-side bootstrap state machine
The bootstrap state machine that drains an ISnapshotProvider export
on the receiver, applies every entry through the local apply seam
preserving the source HLC, and pins the snapshot's causal-stable
frontier on the per-tree high-water-mark grain ships as the public
ILatticeBootstrapCoordinator seam. Triggered by the fall-off
detector (when the per-tree maintenance pass finds a peer's per-origin
high-water mark behind the oldest entry that peer authored in the head
window of the local WAL partitions - see Auto-Bootstrap) and by operator-driven
re-seed flows.
| Type | Shape | Purpose |
|---|---|---|
LatticeBootstrapState |
enum with members Idle, RequestingSnapshot, ApplyingSnapshot, IncrementalHandoff, LiveIncremental, Failed |
The state machine's observable position for a single tree. |
BootstrapCoordinatorStatus |
readonly record struct (LatticeBootstrapState Phase, string? SourceClusterId) |
Observable status snapshot returned by GetStatusAsync; carries the phase plus the in-flight source cluster id (or null when no bootstrap is in flight). |
ILatticeBootstrapCoordinator |
Task<LatticeBootstrapState> GetStateAsync(string treeName, CancellationToken ct) + Task<BootstrapCoordinatorStatus> GetStatusAsync(string treeName, CancellationToken ct) + Task BootstrapAsync(string treeName, string sourceClusterId, CancellationToken ct) |
Public façade over the per-tree bootstrap coordinator grain. Registered as a singleton by AddLatticeReplication; the state machine itself lives in a per-tree internal grain whose cluster-wide single activation provides cross-silo mutual exclusion. |
LatticeBootstrapTransientFaultClassifier |
public static class exposing bool IsTransient(Exception) |
Default classifier consumed by the bootstrap drain's bounded-retry seam. Returns true for TimeoutException, HttpRequestException, SocketException, IOException, Orleans' EnumerationAbortedException (an expired cross-grain enumeration session), aggregate wrappers of those, and gRPC RpcException carrying Unavailable, DeadlineExceeded, or Aborted. Hosts can compose this with a custom predicate via LatticeReplicationOptions.BootstrapTransientRetry.RetryableExceptionClassifier. |
State transitions
Idle
└─► RequestingSnapshot (BootstrapAsync invoked; ExportAsync issued)
└─► ApplyingSnapshot (snapshot stream open; draining Entries)
└─► IncrementalHandoff (entries drained; pinning AsOfHlc + CausalStableFrontier)
└─► LiveIncremental (terminal - incremental replication is live)
Any state ──► Failed (any thrown exception; restart is a fresh BootstrapAsync call)
Semantics
Kickoff-and-poll API.
BootstrapAsyncis an idempotent kickoff: it persists the bootstrap intent, schedules the background phase pump on a 2-second grain timer plus a 1-minute keepalive reminder, and returns. Callers pollGetStateAsyncfor progress. This avoids the 30-second Orleans RPC timeout for long-running snapshot drains and decouples caller liveness from the bootstrap workflow.One bootstrap per tree at a time, cluster-wide. The state machine is hosted in an internal per-tree Orleans grain, so every silo's
BootstrapAsynccall for a given tree id routes to the same activation. The grain reads the persistedInProgressflag on entry: a concurrent call from the same source cluster is a no-op (idempotent retry); a concurrent call from a different source cluster throwsInvalidOperationException. No distributed lock or external coordination is required - Orleans' single-activation invariant plus the durable in-progress flag is the synchronisation primitive. Concurrent bootstraps of different trees route to different activations and run in parallel.Durable, crash-resumable state. The grain uses the same keepalive-reminder plus phase-timer coordinator pattern as the core tree-resize coordinator. Phase, source cluster id, and a
LastAppliedHlccursor are persisted to theLatticeOptions.StorageProviderNamestorage provider. After a silo crash, Orleans reactivates the grain on a surviving silo within the keepalive reminder period and the phase pump resumes from the persisted phase. DuringApplyingSnapshotthe cursor (the highest source HLC applied so far) is persisted every 100 entries. On resume the coordinator re-opens the export with no upper bound (HybridLogicalClock.Zero), exactly as a fresh drain does. The cursor is never passed as the export'sasOfHlc:ISnapshotProvider.ExportAsynctreats that argument as a strict upper bound, not as a resume point, and emits entries in leaf-chain order rather than HLC order, so a resume bounded at the cursor would silently drop every not-yet-applied entry stamped above it. The cursor feeds only the source-origin seal pinned at the handoff. Re-applying the entries the interrupted drain already applied is a correctness no-op, because the receiver-side LWW reconciliation on each leaf grain (plus the per-leaf recently-terminal short-circuit and the per-tx registry no-op described under "Bootstrap drain bypasses the pinned-floor gate and the high-water-mark advance" below) absorbs it. A freshBootstrapAsynckickoff resets the cursor toZero.Failedis restartable. On any thrown exception inside the phase pump the state transitions toFailed(persisted) and the pump tears down. A subsequentBootstrapAsynccall restarts the cycle fromRequestingSnapshot.Bounded retry on transient transport faults. When the snapshot drain throws an exception classified as transient by
LatticeReplicationOptions.BootstrapTransientRetry.RetryableExceptionClassifier(default:LatticeBootstrapTransientFaultClassifier.IsTransient-TimeoutException,HttpRequestException,SocketException,IOException,EnumerationAbortedException, aggregate wrappers, and gRPCRpcExceptioncarryingUnavailable,DeadlineExceeded, orAborted), the coordinator retries the drain in-place using a bounded exponential backoff (default:DefaultBootstrapMaxAttempts = 4attempts, initial delay500 ms, capped at30 s). Each retry re-opens the export with no upper bound (the resume rule above); the drain applies without the pinned-floor gate, and receiver-side LWW reconciliation is what makes re-applying the entries the failed attempt already applied a correctness no-op. Every retry increments theorleans.lattice.replication.bootstrap.transient_retriescounter (LatticeReplicationMetrics.BootstrapTransientRetries) so operators can dashboard the rate. Non-transient faults still pivot toFailedon the first failure exactly as they did before this seam landed; budget exhaustion re-throws the captured transient and pivots toFailedas the terminal outcome. SetBootstrapTransientRetry.MaxAttempts = 1to disable retries entirely (fail-fast).Source HLC + origin preservation. Every snapshot entry is applied through
IReplicationApplier.ApplyAsync, the same canonical inbound apply seam used by live-incremental replication, carrying the entry's commit-timeTimestampand the suppliedsourceClusterId. On a last-writer-wins tree the receiver keeps thatTimestamp, so transitive replication paths (A -> B -> C) preserve the originating HLC; a typed-CRDT row folds through a state-based merge that the receiver writes at a fresh local HLC.Bootstrap and live-incremental share the apply seam. Routing the snapshot drain through
IReplicationAppliermeans every decorator stacked on the applier - dead-letter tracking and any host-supplied per-key change observer - fires identically for bootstrap-arrived entries and live-incremental entries. (Bootstrap entries carry no vector clock, so the applier's causal-dependency gate never parks them in the causal-apply buffer.) A receiver that catches up via bootstrap therefore raises the same observable side-effects as a receiver that catches up via the WAL tail, so UI live-update hooks and audit observers see the bootstrap window rather than missing it.Bootstrap drain bypasses the pinned-floor gate and the high-water-mark advance. The applier reads an ambient bootstrap-apply flag on every inbound call; when the flag is set (the bootstrap coordinator opens one scope around the entire drain) the pinned-causal-floor check and the post-apply high-water-mark advance are skipped, and the steady-state
orleans.lattice.replication.apply.fifo_violationstracker is not fed. This is required because the snapshot exporter enumerates shards/leaves in arbitrary order rather than HLC order: per-shard HLCs are not globally monotonic across a single bootstrap stream, so advancing the high-water-mark mid-drain can suppress a still-pending saga key with a strictly-earlier source HLC and break per-saga all-or-nothing visibility on the bootstrapped peer; and feeding those out-of-order HLCs to the FIFO diagnostic would register every out-of-order shard arrival as a violation. The drain is still idempotent end-to-end because:- Receiver-side LWW on each leaf grain reconciles concurrent arrivals of the same key by HLC, with a replica-invariant tie-break (tombstone, expiry, then value bytes) ahead of the origin id, so a re-applied snapshot entry that has already been delivered is a no-op rather than an over-write.
- The per-leaf recently-terminal short-circuit on the leaf grain suppresses a re-arriving saga terminal whose bucket has already drained, so saga-terminal re-delivery is correctness-preserving.
- The per-tree transaction registry's "repeat-same-outcome no-op"
drops a commit/abort mark that matches a transaction id already
in the requested terminal state. Note the bound: that guarantee
holds only while the decision is still recorded. Once the tree has
forgotten the saga and its tombstone has been physically pruned,
the registry has nothing left to recognise the repeat against and a
late redelivery records the verdict afresh. The safe redelivery
window is
LatticeOptions.TxDecisionRetention, not "forever".
The post-drain snapshot pin atomically replaces the per-origin high-water-mark vector and the pinned causal floor with the snapshot's causal-stable frontier, so the bootstrap-to-incremental handoff retains exactly-once semantics on the live tail. Range deletes, saga terminal records (which carry the saga's own terminal HLC), and tombstone-reap envelopes are routed before the pinned-floor gate, so the bootstrap scope does not change how they apply.
Live-incremental dedup is unchanged. The snapshot-pinned causal floor in the applier suppresses any re-delivery of bootstrap-arrived entries through the live-incremental path - the pin at the end of the drain sets each origin's floor to its coordinate in the snapshot's causal-stable frontier, so a live entry at or below that coordinate is deduped canonically.
Snapshot/incremental handoff is exactly-once. The coordinator pins the snapshot's causal-stable frontier on the per-tree high-water-mark store after every snapshot entry has been applied, first sealing the source cluster's own coordinate at or above the highest HLC the drain applied and the oldest source-authored entry the local WAL still retains, so the fall-off detector cannot read the retained baselines as a trim gap. The pin replaces both the high-water-mark vector and the pinned floor (the
AsOfHlcpassed alongside it is currently ignored). The applier's floor gate then makes any incremental entry whose timestamp is at or below the pinned frontier a no-op, so the snapshot/incremental boundary is exactly-once regardless of overlap.Tombstones in custom providers are skipped. Committed (non-prepared) snapshot entries whose
Valueisnull(not emitted by the default provider, but permissible from a host-suppliedISnapshotProvider) are skipped rather than applied as deletes, as are prepared rows with an emptyTransactionId. A prepared row withIsTombstoneset is applied as a prepared delete.Per-tree merge mode is honoured on bootstrap. Every
WalRecordemitted by the bootstrap drain is stamped with the merge modeILatticeMergeModeResolverresolves for the tree - theLatticeReplicationOptions.ReplicatedTreesdeclaration, or the runtime configuration on a host that enables it. Trees declared asLatticeMergeMode.OrSetorLatticeMergeMode.PnCountertherefore merge bootstrap-arrived entries under the tree's CRDT semantics: a committed-projection row carries the primitive's full state, which the receiver folds through the primitive's state-based merge rather than the per-entry delta fold live-incremental entries use. A tree the resolver returns no mode for (and a tree declared asLwwRegister) is stampedLatticeMergeMode.LwwRegister; for a tree the resolver returns no mode for, the receiver's own enrollment gate then drops every stamped entry as not enrolled, so such a drain applies nothing. The mode is resolved once at the start of the drain, not per entry: the resolver is on the hot path's allocation budget but is invariant for the lifetime of a single drain.
Sample usage
ILatticeBootstrapCoordinator coordinator = client.ServiceProvider
.GetRequiredService<ILatticeBootstrapCoordinator>();
await coordinator.BootstrapAsync("orders", sourceClusterId: "site-a", cancellationToken);
LatticeBootstrapState state = await coordinator.GetStateAsync("orders", cancellationToken);
_ = state; // LatticeBootstrapState.LiveIncremental once the bootstrap completes
Operator-driven re-seed
Beyond the receiver-driven auto-bootstrap path (ILatticeFallOffLogDetector), the package exposes an explicit operator-facing entry point for scheduled bootstraps - a new peer joining, a bandwidth-constrained initial sync, or a post-disaster re-bootstrap. The seam is ILatticeReplicationAdmin.RequestSnapshotAsync; honoured requests delegate to the same ILatticeBootstrapCoordinator.BootstrapAsync driving the auto-bootstrap path, so every re-seed - operator-driven or detector-driven - flows through one state machine.
| Type | Shape | Purpose |
|---|---|---|
ILatticeReplicationAdmin |
Task<OperatorReseedDecision> RequestSnapshotAsync(string treeName, string sourceClusterId, CancellationToken ct) |
Public façade that gates the request behind a per-(tree, sourceClusterId) rate limit before delegating to the bootstrap coordinator. |
ILatticeReplicationAdmin |
Task<OperatorReseedDecision> ForceRequestSnapshotAsync(string treeName, string sourceClusterId, CancellationToken ct) |
Opt-in bypass that skips the rate-limit check entirely. Intended for disaster-recovery and scheduled re-seed scenarios where a real cross-cluster drain may exceed the configured window. Every call is audit-logged at Information. |
OperatorReseedDecision |
readonly record struct with Triggered, LastRequestedAt, RetryAfter |
Diagnostic return value indicating whether the call invoked the coordinator and, when denied, how long the operator should wait before retrying. Both overloads share this return shape. |
Semantics
- Per-
(tree, sourceClusterId)rate limit.LatticeReplicationOptions.OperatorReseedMinInterval(default1 minute) bounds the minimum gap between honoured requests for the same pair. A second request inside the window returnsTriggered = falsewithRetryAfterset to the remaining time; the coordinator is not invoked and no exception is thrown.TimeSpan.Zerodisables the rate limit entirely (every request reaches the coordinator). - Process-local rate-limit table. The default implementation tracks honoured requests in process memory only; a silo restart resets the rate-limit window for every pair. Cross-silo coordination is not required because
ILatticeBootstrapCoordinatoris itself idempotent under concurrent invocations against the same tree from the same source cluster (the per-tree internal grain absorbs the second call as a no-op) and rejects mismatched-source concurrent kickoffs asInvalidOperationException. The rate limit is therefore a fairness mechanism, not a correctness one. - Timestamp updates only on success. The dictionary timestamp is stamped only after the coordinator call returns successfully, so a thrown coordinator exception (transport failure, conflicting in-flight bootstrap from a different source) does not consume the rate-limit budget against the operator. The same rule applies to both the rate-limited overload and the bypass overload.
- Per-tree options resolution. The minimum interval is resolved per-tree via
IOptionsMonitor<LatticeReplicationOptions>.Get(treeName), so different replicated trees can run different re-seed cadences without separate seam instances. - Argument validation.
treeNameandsourceClusterIdmust be non-null and non-empty (ArgumentNullExceptionwhennull,ArgumentExceptionwhen empty); the cancellation token is observed before the rate-limit check and propagated to the underlying coordinator. The bypass overload validates identically.
Force-bypass semantics
The rate limit was originally sized assuming the underlying snapshot drain is intra-cluster and fast. Once cross-cluster transport is in play, a real re-seed against a large tree may exceed the configured window, and a routine retry would be denied even though the previous call's drain has not completed. ForceRequestSnapshotAsync is the escape hatch for that case.
- Always reaches the coordinator (on success). When the coordinator call returns,
Triggered = trueunconditionally; the bypass overload never returns a denied decision.RetryAfteris alwaysnull. - Audit-logged on every call. A
LogLevel.Informationline taggedFORCEand carrying bothtreeandsourceClusterIdis emitted before the coordinator dispatch, so a log-tailing operator sees the bypass even if the coordinator call subsequently throws. - Stamps the dictionary on success. A successful bypass updates the rate-limit dictionary timestamp so a follow-up routine
RequestSnapshotAsyncinside the window correctly observes the bypass as the last honoured request. This preserves the operator's mental model that the limiter knows about every actual re-seed, not just the rate-limited ones. - Failure does not consume the rate-limit budget. A coordinator exception leaves the dictionary unchanged so a follow-up routine call is still honourable.
Sample usage
ILatticeReplicationAdmin admin = client.ServiceProvider
.GetRequiredService<ILatticeReplicationAdmin>();
OperatorReseedDecision decision = await admin.RequestSnapshotAsync(
"orders", sourceClusterId: "site-a", cancellationToken);
if (!decision.Triggered)
{
// Rate-limited: the operator should wait `decision.RetryAfter` before
// retrying. The previously honoured request is still driving the
// bootstrap coordinator if one was kicked off recently.
_ = decision.RetryAfter;
_ = decision.LastRequestedAt;
return;
}
// Triggered: poll the coordinator for state-machine progress.
LatticeBootstrapState state = await client.ServiceProvider
.GetRequiredService<ILatticeBootstrapCoordinator>()
.GetStateAsync("orders", cancellationToken);
_ = state;
Sample usage (force bypass)
ILatticeReplicationAdmin admin = client.ServiceProvider
.GetRequiredService<ILatticeReplicationAdmin>();
// Disaster recovery: the previous routine re-seed against a large
// cross-cluster tree is still draining and we need to retry under
// a different sourceClusterId without waiting for the configured
// OperatorReseedMinInterval window. ForceRequestSnapshotAsync
// always reaches the coordinator and audit-logs at Information.
OperatorReseedDecision decision = await admin.ForceRequestSnapshotAsync(
"orders", sourceClusterId: "site-b", cancellationToken);
// On success, Triggered is always true and RetryAfter is null.
// A coordinator exception (e.g. a conflicting in-flight bootstrap
// against a different source) propagates verbatim; catch and
// inspect at the caller.
_ = decision.Triggered;
_ = decision.LastRequestedAt;
Snapshot and in-flight atomic visibility
The default snapshot provider preserves cross-cluster saga atomic visibility across the bootstrap boundary: a saga whose prepare-commit pair straddles the producer's snapshot cut is observed by the bootstrapped peer either at every key or at none, never at a strict subset.
The export operates in two passes against a single frozen view of the producer's tree-wide transaction-registry decisions, unioned across every registry shard of the tree:
Prepared rows pass (runs first). Walks every shard's leaf chain and emits a
SnapshotEntrywithIsPrepared = truefor every(transactionId, key)pair in any leaf's per-tx pending bucket whose registry status in the captured snapshot isInFlight,Indeterminate, or absent. The emitted row carries the source-stamped prepare-time HLC verbatim, plusIsTombstone,TransactionId,ExpiresAtTicks, and the typed CRDTDelta/Modeso the receiver can route it identically to a steady-state prepared WAL record and a prepared CRDT entry folds its delta on the terminal commit.Committed projection pass. Drains the source tree's entries via the resilient
ScanEntriesAsyncwrapper overILattice.EntriesAsync, under the same frozen registry scope. Sagas the snapshot recorded asCommittedsurface their prepared value as the live one; sagas recorded asAbortedare dropped; sagas stillInFlight(orIndeterminate) against the snapshot are hidden from the committed scan because the prepared rows pass above has already shipped them. Each emitted row carries the key's value, commit-time HLC, and absoluteExpiresAtTicksfrom the same per-key version read. The wrapper matters here because the export is long-running and latency-prone: a per-key version read and a (potentially proxied, cross-cluster) stream write interleave between pulls, so the source grain enumerator can idle-expire or be reclaimed by a deactivation mid-stream.ScanEntriesAsyncrecovers from the resultingEnumerationAbortedExceptionand deterministically resumes from the last yielded key, so a reclaimed enumerator is re-opened instead of aborting the whole snapshot stream (which would fail the receiver's bootstrap drain and leave its high-water-mark pinned, re-triggering the fall-off detector on its next tick).
The receiver replays prepared rows through the receiver-side prepared-set and prepared-delete apply paths into its per-tx pending bucket. The matching terminal record - delivered subsequently by the post-snapshot incremental WAL stream - flips visibility atomically via the transaction-terminal apply path and the receiver's local transaction-registry linearization point.
Ordering matters. The prepared rows pass runs before the
committed projection pass because a source-side terminal that drains
a pending bucket between the two would otherwise erase the saga from
the export entirely: the committed pass under the frozen snapshot
would hide the saga (the snapshot still says InFlight), the
prepared pass would find the bucket already drained, and the
prepare-time WAL records - stamped at HLC <= asOfHlc - would never
re-arrive through the post-snapshot incremental stream (which starts
at asOfHlc). Capturing prepared rows first guarantees every
InFlight saga's per-key state is shipped to the receiver's pending
bucket, and the incremental stream delivers the terminal record to
flip visibility.
An aged-out decision is shipped, not dropped. A decision whose
tombstone has outlived LatticeOptions.TxDecisionRetention is carried
in the frozen snapshot as Indeterminate rather than omitted from it.
Omission used to fold two different facts onto one reading - "no
decision was ever recorded" and "a decision was recorded and the source
is no longer entitled to report it". Locally that collapse was
invisible, but in this export it was load-bearing in the wrong
direction: an aged-out Committed saga reached the receiver as
absence, was read as still preparing, and could never be corrected,
because the terminal that would have flipped it was already behind the
incremental stream the receiver drains after the snapshot. The source
held "committed", the receiver held "preparing", permanently, with no
repair path on either side. Carrying the row explicitly means absence
in this payload once again means only what it says, and an
Indeterminate saga's prepared rows ship (pass 1 above) so the
receiver holds exactly what the source holds.
A saga the producer's registry recorded as Committed before the
snapshot is naturally folded into the committed projection by the
leaf scan's pending-transaction read step - which honors
the frozen registry scope - so the receiver observes the post-saga
value directly without a separate prepared/terminal round trip. The
same applies in reverse for Aborted: the prepared mutation is
correctly dropped from the committed pass and not shipped as a
prepared row.