Coordinated multi-cluster restore
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 coordinated-restore.md, and llms.txt lists every page.Restoring a backup into a tree that is replicated across clusters is not a local operation. If one cluster swapped in the restored data on its own, its peers would still be shipping their pre-restore writes, and the restored cut would be re-advanced (partially overwritten) the moment replication resumed. Worse, a reader on another cluster could observe a torn mix of restored and pre-restore state while the swap propagated. This is the failure mode tracked by issue #1169.
Coordinated restore closes that gap. When the restore target is currently replicated, the restore is promoted into an all-or-nothing cross-cluster saga: every participating cluster prepares the restored data, then a single global decision either cuts every cluster over together or rolls every cluster back (subject to the fence-timer bound described under The saga phases). No peer re-advances the restored cut, and no reader ever observes a torn or half-restored tree.
When a restore becomes a saga
The decision is a function of the target tree's current replication membership, never of where the backup was originally captured:
- Target tree is replicated - the restore runs as a coordinated saga across the local cluster and every current replication peer. A backup captured on a single cluster and restored into a replicated tree still runs the saga, because it is the target that must stay consistent across clusters.
- Target tree is not replicated (or the backup package is deployed single-cluster) - the restore runs as a plain local restore with no saga. There is nothing to coordinate.
A backup set (multiple trees captured together) restores as one unit: if any member tree is replicated, the whole set restores under a single saga spanning the local cluster and every current replication peer, so every member tree flips together on every participating cluster or none does.
The initiating cluster authorizes the restore with the caller's identity before
it decides whether to dispatch a saga, so an unauthorized caller never reaches a
peer: it checks the Restore capability over the target tree (over every member
tree for a set) and, when a single-tree restore retargets onto a different tree,
the Backup capability over the tree the backup was captured from. A refused
capability check throws LatticeAuthorizationDeniedException. A single-tree
restore into a replicated target then runs the capacity check on the initiating
cluster and refuses an infeasible target with a
LatticeRestoreValidationException that names the refusal but not the target's
stored size or shard count.
Before any cluster prepares, the coordinator probes every current peer over the
saga control channel and refuses to start - with a
LatticeRestoreValidationException naming the unreachable peers - if any peer
cannot be reached, because a partial coordinated restore is never allowed. An
aborted saga likewise surfaces to the caller as a
LatticeRestoreValidationException, after every cluster has been compensated.
The saga phases
The coordinator drives every enlisted participant on every cluster through three phases:
- Prepare - each participant builds the restored data into a shadow alongside the live tree. Prepare is unfenced and resumable: the live tree keeps serving reads and accepting writes throughout the build, and a participant that restarts mid-build resumes from its checkpoint rather than restarting from zero. Each participant returns a vote.
- Commit - reached only if every participant on every cluster voted to
commit. Each participant engages a short per-tree write fence, atomically
swaps the tree's alias to the restored shadow, then unblocks local writes.
The fence covers every shard the tree's routing reaches when it engages -
including a shard an adaptive split added, on the physical copy an earlier
resize or restore put behind the alias - and every release lifts exactly
that set.
The alias swap is bounded by tree ownership like every alias change: a
registered
ITreeOwnershipGuard(theOrleans.Lattice.Appspackage registers one) is consulted before the alias is written, and a denial throwsLatticeTreeOwnershipDeniedExceptionwithout swapping. The restored shadow records the target tree as its origin, so the guard admits a restore into the tree it was built for. Local writes resume as soon as the cutover completes, but cross-cluster shipping and receiving stay paused until the saga completes globally, so an early-flipping cluster cannot re-advance the restored cut. The write fence is held only for the cutover, not for the whole build, so healthy clusters are not write-starved while a large tree builds. The write fence also self-lifts on a bounded cutover deadline (five minutes after it engages), so a stalled cutover never fences local writes indefinitely; that release leaves shipping and receiving paused until the saga completes globally. - Abort - reached if any participant voted to abort. Every participant that prepared is compensated: its shadow is reverted and garbage collected and the pre-restore tree is left untouched.
Two guarantees make this safe under failure:
- Single global decision. The coordinator reaches exactly one commit-or-abort decision after collecting every vote, and delivers that one decision to every cluster that voted to commit; a cluster that voted to abort has already compensated its own prepared work. A participant never observes a mixed outcome.
- Bounded fence-timer auto-compensation. A prepared participant waits for the decision under a bounded cutover-fence timer (five minutes; distinct from the per-tree write fence, which engages only at commit). If the coordinator is lost before it delivers a decision, the timer expires and the participant auto-compensates (aborts), so a prepared single-tree restore cannot leak after a coordinator loss. A backup-set restore is not covered in full: the auto-compensation lifts that cluster's fence but does not garbage collect the member shadows it built. The timer starts when that cluster finishes its own prepare and fires whether or not the coordinator is still alive, so a commit decision that reaches it more than five minutes after it prepared finds it already compensated: the participant does not apply the late commit, and the coordinator does not treat that refusal as a failure. The coordinator separately bounds the prepare phase: it aborts a saga whose prepare is still being retried an hour after the saga started.
Reliability under duress
Every participant runs an admission pre-flight before it builds a shadow: it probes the backup's self-describing size and topology and votes to abort if that probe fails or its capacity check refuses the target, so the saga fails fast with a clear vote rather than failing mid-build. The shipped capacity check admits every target, so today the pre-flight catches an unprobeable backup rather than a tree too large for the cluster. A participant whose build fails permanently (a restore validation failure, such as an artifact absent from the sink or failing its content-digest check) or exhausts its bounded retry budget (a set restore does not retry: the first member build that fails aborts the whole set) votes to abort and garbage collects its partial shadow, leaving no orphaned shadow state; the whole saga then rolls back all-or-nothing.
The sink must be shared, and that is checked at capture time
Every cluster in the saga resolves the backup's manifest chain from its own
configured ILatticeBackupSink. A coordinated restore therefore only works if
every cluster's sink is the same storage. Point each region at an isolated
sink and the misconfiguration is invisible: each capture succeeds locally, each
local health check passes, and the fault only surfaces as an all-or-nothing saga
abort at restore time - after the operator has spent weeks relying on backups
that were never restorable.
That check now runs at capture/startup time instead. The replication package registers a real cross-cluster sink-sharing probe over the backup package's no-op default (the same layering trick the saga dispatcher itself uses - see Enlisting your own resource in the saga), so the backup package never takes a dependency on replication. When a tree is replicated and the deployment has peers, each cluster writes a tiny self-naming marker into its own sink and reads every peer's marker back out of that same sink. A marker that is missing while its peer answers the saga control channel proves the sinks are separate; a marker missing from an unreachable peer is merely undecided and is re-probed on the next backup-health sweep. Over the shipped gRPC saga control channel, though, the peer refuses that reachability call, because it carries an empty saga id and the receiving service rejects one as an invalid argument, so a reachable peer still counts as unreachable and a missing marker leaves the verdict undecided rather than refuting the sink.
The verdict is logged at start, annotated onto every affected backup's health
report (so it shows as a Warning in the Health column of the backup catalogue
in the Explorer's Backups area), and can be made to block silo start outright.
Nothing is probed at all - no sink write and no network call - when no tree is
replicated or the deployment has no peers.
See backup configuration for the enforcement modes and their defaults, and disaster recovery for how the verdict reaches the health surface.
Triggering a restore
Use the ordinary backup restore surface. Restoring into a replicated target transparently runs the saga; the caller does not opt in.
using Orleans.Lattice.Backup;
// Restore an entire captured backup set as one coordinated unit. When any member
// tree is replicated this runs as a single all-or-nothing cross-cluster saga
// across the local cluster and every current replication peer; otherwise it
// runs as a plain local per-member restore.
var restoreService = client.ServiceProvider.GetRequiredService<ILatticeBackupRestoreService>();
IReadOnlyList<LatticeRestoreResult> results =
await restoreService.RestoreSetAsync("your-backup-set-id", cancellationToken);
A single-tree restore uses RestoreAsync(LatticeRestoreRequest, CancellationToken)
in the same way; see the backup package restore docs
for the request shape and result fields.
Enlisting your own resource in the saga
The built-in restore participant is one participant among many. An application that holds a resource which must flip atomically alongside a replicated restore (for example an external projection or a downstream index) can enlist its own participant so it runs in the same saga, under the same unanimous prepare and single global decision.
Implement the public ISagaParticipant interface (in Orleans.Lattice.Replication)
and register it with AddLatticeSagaParticipant<TParticipant>(name) on the silo
builder. The interface has four methods:
PrepareAsync- prepare the resource set this participant hosts for the saga and return aSagaParticipantPrepareResultcarrying the vote. The work may be long-running and must be idempotent and resumable.CommitAsync- make the prepared mutation durable. Idempotent.AbortAsync- compensate (roll back) the prepared resource set. Compensation must be total: once a participant votes to commit it must always be able to undo that prepare.GetStatusAsync- report the phase the participant currently holds, without changing state.
The optional name argument passed to AddLatticeSagaParticipant is used for
diagnostics and logging only; it never affects the saga wire contract. A
participant that hosts nothing for a given saga prepares vacuously (votes to
commit) rather than blocking the saga. Registration is idempotent per participant
type and form: repeated unnamed calls, or repeated named calls (the first name is
kept), enlist it once, but mixing an unnamed and a named call for the same type
enlists it twice, so every saga drives it through each phase twice - use one form
per participant type.
Guardrails. Every method must be idempotent, and a participant that cannot
guarantee total compensation must vote to abort from PrepareAsync rather than
preparing. These match the intra-cluster cross-tree saga contract.
Observability
The saga emits OpenTelemetry instruments on the replication meter
(orleans.lattice.replication): saga phase durations, participant vote / commit
/ abort counts, per-tree write-fence window durations, and compensation counts by
cause. See observability for the
full instrument list and tags.