Table of Contents

Adaptive Shard Splitting

This page is part of 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 shard-splitting.md, and llms.txt lists every page.

Adaptive shard splitting allows a hot physical shard to split into two at runtime, fully online - no shard is ever taken offline. Splits happen automatically when an autonomic monitor detects a hot shard, and an online reshard drives the same split to grow a tree's shard count. Shard splitting is internal-only: ITreeShardSplitGrain is declared internal and is not reachable from consumer assemblies.

Why

Lattice trees are sharded by hashing keys into a virtual slot space and mapping virtual slots onto physical ShardRootGrain activations. With a fixed shard count, a workload skewed toward a small set of keys will saturate one shard while others sit idle. Adaptive splitting redistributes hot virtual slots to a new physical shard so the load follows the data.

How it works

A split is driven by the internal TreeShardSplitGrain coordinator through five persisted phases, in this order. The source shard S keeps serving reads and writes throughout; the target shard T receives mirrored data and eventually owns the moved slots.

stateDiagram-v2
    [*] --> BeginShadowWrite : split requested for source shard S
    BeginShadowWrite --> Drain : S mirrors moved-slot writes to T, prepared mutations swept across
    Drain --> Swap : every moved-slot entry (live + tombstones) forwarded to T
    Swap --> Reject : S sealed and rejecting, final drain, registry reassigns the moved slots to T
    Reject --> Complete : reject state re-asserted
    Complete --> [*] : final drain pass, then S records the moved slots in its moved-away table
  1. BeginShadowWrite - Coordinator takes the upper half of the virtual slots S owns as the moved set, allocates a fresh target shard index from the registry, persists its intent, and opens S's shadow-write window for the moved slots. From that moment every successful write S applies to a key in a moved virtual slot is also mirrored to T through T's batched merge, preserving the original HLC. Before the drain begins, the coordinator then runs a retroactive prepared-mutation sweep: it walks S's leaf chain, snapshots every in-flight prepared saga mutation whose key hashes into a moved virtual slot, and replays each one into T's pending-transaction buckets under its original transaction id, HLC, origin, vector clock, and expiry, so any prepared write that landed on S before the window opened survives the topology change. A saga the transaction registry already reports as committed or aborted has its terminal applied to T directly instead. The sweep is idempotent per (transaction, key), and a coordinator crash mid-sweep re-runs the whole sweep on recovery. Instrumentation: orleans.lattice.split.retroactive_forward.entries (counter, per replayed mutation) and orleans.lattice.split.retroactive_forward.duration (histogram, total sweep wall-clock). CRDT LWW guarantees correct convergence regardless of how the foreground write and the background drain interleave.
  2. Drain - Coordinator walks S's leaf chain and forwards moved-slot entries (including tombstones) to T with their original HLC timestamps. The drain is chunked and leaf-side filtered: each leaf returns only entries whose virtual slot is in the moved-slot set, and the coordinator flushes to T in batches of SplitDrainBatchSize (default 1024) entries. This bounds peak memory on the coordinator regardless of source shard size, and avoids transferring non-moved entries over the wire. Each pass is bounded by the background-drain budget (BackgroundDrainLeavesPerPass leaves or BackgroundDrainMaxDuration, whichever binds first) and persists a key cursor the next tick resumes from, so the phase advances to Swap only once the whole leaf chain has been swept. Idempotent under retry - re-running merges only converges to the same state.
  3. Swap - Coordinator seals every source leaf for the moved slots and puts S into its reject phase before the registry map flips. From this point any read or write to S for a moved-slot key is refused as stale routing, which freezes the source's committed state for the migrating slots. The coordinator then runs one final authoritative drain pass to T, so the destination is synchronised with the source's now-frozen committed state before any reader can route to T. Only then does it reassign the moved slots to T in the registry's ShardMap, with a single registry call that re-reads the live map, applies the reassignment, and persists it under a fresh Version, so a concurrent split or shard consolidation of the same tree cannot erase either change. Reversing the order - flipping the map before the source rejects - would open a window in which a stale-routing reader could still be served the pre-split value by the source. New router activations immediately route the moved slots to T; stale activations that still cache the old map hit the source's reject gate, invalidate their cached map, fetch the fresh map from the registry, and retry against T - a single transparent retry per call.
  4. Reject - Coordinator re-asserts S's reject phase. S already entered it during Swap and the call is idempotent, so this phase only records the transition; a crash-recovered coordinator that re-enters it changes nothing.
  5. Complete - Coordinator runs one final drain pass to capture any tombstones written during shadow that were not mirrored on the hot path, then tells S to complete the split and clears its own state. Completing also promotes the split's moved slots into S's persisted moved-away slot table, so even after the active reject-phase state is cleared, every subsequent operation on a moved-slot key continues to be refused as stale routing. This guarantees that stale stateless-worker router activations (which may have cached the pre-split shard map) always trigger a map refresh on first use rather than silently returning orphan data. The table is lifted only if a later shard consolidation folds T back into S - and only after T's entries for those slots have been drained onto S - at which point S's leaves re-ship the reclaimed rows to their read caches.

The coordinator state is persisted before any side effect, so a silo crash mid-split is recovered by the keepalive reminder, which resumes from the last persisted phase; every phase is idempotent.

Scan semantics during a split

This section describes the mechanism by which live operations behave during a split. For the consistency contract each ILattice method provides - including under concurrent splits - see Consistency.

Point reads and writes (GetAsync, SetAsync, DeleteAsync, SetIfVersionAsync, GetOrSetAsync, etc.) continue to serve traffic throughout the split: every successful write is mirrored to the new owner during the shadow phase and the reject phase (which precedes the map swap) causes stale activations to transparently retry against the correct shard. The post-Complete MovedAwaySlots rejection extends this for as long as the source does not own those slots again (only a later shard consolidation that folds them back lifts it).

Scans (ScanKeysAsync, ScanEntriesAsync, CountAsync) reconcile against topology changes mid-scan as described below. See Consistency for the guarantee this reconciliation delivers.

How the reconciliation works

Each scan uses a reconciliation algorithm coordinated against the registry's monotonically-incrementing ShardMap.Version, but CountAsync and the ScanKeysAsync / ScanEntriesAsync streams follow two different paths.

CountAsync / CountPerShardAsync - per-slot routing

The orchestrator reads the authoritative ShardMap. While its Version is still 0 - no split has ever been persisted for the tree - it sums each physical shard's own count and accepts the total only if the map is still at version 0 afterwards. Otherwise it partitions the virtual slots by current owner and asks each physical shard to count only the slots it owns (bounded to the requested range), driving each shard's work-bounded count batches to completion so no shard is held for its whole leaf chain. Because each virtual slot is counted exactly once - against whichever shard the map identifies as its current owner - the result is topology-consistent by construction, independent of the source shard's per-split phase. The map version is re-read after the fan-out; if it moved, the count is discarded and retried on the fresh map, bounded by LatticeOptions.MaxScanRetries (default 3). Throws InvalidOperationException on retry exhaustion.

ScanKeysAsync / ScanEntriesAsync - in-line reconciliation

Reconciliation is driven inside the main k-way merge loop rather than as a separate pass. Each shard root reports back:

  • the keys/entries of all keys not in its MovedAwaySlots table (entries it no longer authoritatively owns), and
  • the set of MovedAwaySlots virtual slots it observed during the traversal (used as a topology-stability hint).

Before each priority-queue dequeue, the orchestrator checks whether any live shard cursor has reported new MovedAwaySlots since the last reconciliation step. If so, it queries the current owners of the affected slots via the slot-filtered variants GetSortedKeysBatchForSlotsAsync / GetSortedEntriesBatchForSlotsAsync, loads the reconciled keys into memory, sorts them with the same comparer, and injects them as an additional in-memory cursor into the same priority queue. The merge invariant (global minimum is yielded next) then carries ordering across the topology boundary. A per-call HashSet<string> suppresses duplicates across pre- and post-swap views. A final stability check after the priority queue drains catches the edge case where a split commits after all live cursors finished - reconciled entries from this path are also sorted and injected as a cursor, not appended. Bounded by LatticeOptions.MaxScanRetries.

Snapshot scans - pinned-map slot ownership in the snapshot leaf

A snapshot-isolated scan (OpenSnapshotEntryCursorAsync / OpenSnapshotKeyCursorAsync, the read-only state API, and the Explorer's Data area) cannot use the live in-line reconciliation above, because it must read a single, internally-consistent point in time. Instead the snapshot coordinate pins the registry's ShardMap at open (LatticeSnapshotCoordinate.PinnedShardMap, alongside the pinned ShardMap.Version, the per-shard / per-partition WAL offsets, and the registry HLC). The snapshot fan-out forces a fresh routing read so it opens against the post-split map, then for each fan-out shard it resolves that shard's owned virtual slots under the pinned map and passes them to the shard's snapshot leaf.

The snapshot leaf serves a durable frozen baseline captured at open time (see Snapshot Cursors): the per-shard projection is materialised once by walking the shard's leaf chain and folding each leaf's WAL tail, then persisted and seeded into the leaf with no serve-time WAL replay. Whether a per-key record reaches the served view is resolved by the key's virtual slot under the pinned map - not by the mutation's stamped ShardIndex - applied both when the baseline is folded at capture and again through the leaf's IsKeyOwned filter when the durable rows are seeded. Resolving by slot is what makes the snapshot view correct across a split, because the stamp records the shard that authored a record, which is not the same as the shard that owns the key after the split:

  • A moved key's pre-split copy is physically retained on the donor shard (an orphan). Its stamp still names the donor, but the pinned map now assigns its slot to the target, so the donor's snapshot leaf drops it and only the target surfaces the key - no duplicate.
  • A write routed to the donor for an already-moved slot is shadow-forwarded into the target shard's WAL but keeps the donor's source stamp. Resolving by slot keeps that forwarded record on the target's snapshot leaf, so the leaf's last-writer-wins merge applies it. A stamp-based filter would drop it (its stamp names the donor) and resurrect the pre-forward value - for example a post-split delete forwarded through the donor would be lost and the deleted key would reappear with its stale drained value.

Ownership must be resolved against the pinned map version, not the current one, so the snapshot neither over-excludes (a key not yet moved at the pinned version) nor under-excludes (a key moved after the pinned version). The snapshot k-way merge additionally collapses equal, adjacent keys as a defensive net, but value correctness comes from the leaf-side ownership filter (the merge operates on raw bytes with no HLC and cannot pick the last-writer-wins winner on its own).

Live leaf reactivation - current-map slot ownership in the live leaf

The authoritative live leaf has its own activation-time WAL replay, and the same stamp-versus-slot distinction applies to it. When a live leaf activates (a cold reactivation after deactivation, a silo move, or a crash), it rebuilds its in-memory projection by replaying the WAL through its replay filter - from its snapshot's offset when a persisted snapshot rehydrated it, otherwise across the whole readable window (see State Model). Unlike the snapshot leaf the live leaf has no pinned map - it serves the current point in time - so it resolves per-mutation shard ownership by the key's virtual slot under the current registry ShardMap, fetched once at the start of replay.

This matters only on a cold reactivation that replays a WAL suffix past a checkpoint taken before a shadow-forward. In steady state the live leaf applies the forwarded mutation in real time and folds it into its next checkpoint, so a warm leaf never re-evaluates the record. The trigger is a target leaf that checkpointed before a post-split write was shadow-forwarded through the donor, then reactivated cold:

  • A write routed to the donor for an already-moved slot is shadow-forwarded into the target shard's WAL but keeps the donor's source stamp. Resolving by slot keeps that forwarded record on the target's live leaf (the current map routes its slot here), so the leaf's last-writer-wins merge applies it on replay. A stamp-based filter would drop it - its stamp names the donor - and resurrect the pre-forward value, for example losing a post-split delete and reappearing the deleted key with its stale drained value.
  • Genuine sibling-shard data multiplexed through a shared WAL partition (partitions are keyed by key hash, not by shard) still resolves to another shard under the current map and is dropped, and a donor leaf's own orphan copies of moved slots are dropped because their slot now routes to the target - the live read path already seals those orphans via MovedAwaySlots.

The current map is fetched best-effort and is trusted only when it references the leaf's own shard, so a registry hiccup at activation or a map drawn from a foreign physical shard space falls back to the legacy stamped-ShardIndex axis - a leaf can never reject its own writes. Slot-less legacy leaves (pre-split-feature persisted state with a null ShardIndex) keep their unconditional V1 single-leaf-per-shard semantics and apply on the shard axis regardless.

Trade-offs

  • Order: Keys/Entries are streamed in strict lexicographic (or reverse) order end-to-end, even when splits commit mid-scan. Reconciled entries participate in the same k-way merge as live cursors, so the ordering guarantee is preserved.
  • Memory: scans allocate a HashSet<string> for dedup that grows with the number of distinct keys observed during the scan, plus a per-reconciliation buffer proportional to the number of keys in slots that actually moved during the scan (typically small). For very large trees, prefer the range-bounded overload of ScanKeysAsync / ScanEntriesAsync to bound memory.
  • Latency: when no split has ever occurred, scans take the same fast path as before (one round-trip per shard). The reconciliation passes only run when a shard actually reports moved slots.
  • System trees: a reserved system tree - the lattice registry tree among them - bypasses the reconciliation path. It never participates in adaptive splits and routes by the default shard map without reading one from the registry (for the registry tree that read would deadlock), so its internal scans fan out with no moved-slot tracking. The public count and scan surface refuses system trees outright.

Autonomic detection

Each tree's hot-shard monitor is armed when the tree activates and re-armed from the point, conditional, batched, predicated and CRDT write paths (deletes, atomic batches, bulk loads and merges do not re-arm it) until an arming attempt succeeds, and a keepalive reminder re-activates it after collection. The monitor arms its sampling timer before it registers that keepalive, so a reminder service that is still initialising cannot leave it claiming to run with nothing sampling; a keepalive registration deferred that way is retried after each sampling pass until it succeeds. On each tick (default every 30 s) it:

  1. Polls every physical shard's GetHotnessAsync() in parallel.
  2. Computes ops/sec = (reads + writes) / window.TotalSeconds.
  3. Counts the shard migrations in flight for this tree by asking every physical shard whether it is splitting - the source of an adaptive split reports that it is, and so does the donor of a shard consolidation (fold). If that count is already MaxConcurrentAutoSplits or more, the pass triggers nothing, so a fold in flight takes up one of the tree's autonomic split slots. Because HotShardMonitorGrain is keyed per-tree, the cap is enforced independently per tree - in a multi-tree cluster each tree may have up to MaxConcurrentAutoSplits concurrent splits running simultaneously.
  4. Selects the top-(MaxConcurrentAutoSplits - inFlight) hottest shards whose rate reaches HotShardOpsPerSecondThreshold (default 200 ops/s), skipping any shard already splitting, on cooldown, or owning a single virtual slot. Three shape clauses can refuse a hot shard as well: no shard is admitted once the tree has MaxPhysicalShardsPerTree physical shards (default 256), or while its load is uniform - the hottest shard's rate below HotShardMinSkewRatio (default 1.5) times the median shard rate, the signature of a bulk ingest that a split cannot relieve - and a candidate holding fewer than HotShardMinShardEntries live entries (default 1024) is skipped.
  5. Triggers ITreeShardSplitGrain.SplitAsync on each selected shard in parallel via Task.WhenAll and starts a per-shard cooldown.

Each split runs in its own coordinator activation, keyed {treeId}/{sourceShardIndex}, so independent splits of different source shards within the same tree do not contend on a single coordinator. Concurrent target-index allocation is made collision-free by an atomic registry-side counter, and concurrent shard-map swaps compose because each applies its moved-slot reassignment in a single registry call that re-reads the live map before persisting. Both atomicity guarantees rely on the singleton registry grain running each mutating call to completion before the next (it is not reentrant).

A split is suppressed (whole pass skipped) or a candidate is skipped individually when:

Suppression rule Scope Mechanism
AutoSplitEnabled = false Whole pass Returns early.
Tree younger than AutoSplitMinTreeAge (since monitor activation, default 60 s) Whole pass Returns early.
Resize / reshard / merge / snapshot in progress Whole pass ILattice.IsResize/Reshard/Merge/SnapshotCompleteAsync() returns false.
Any shard has a pending bulk graft Whole pass IShardRootGrain.HasPendingBulkOperationAsync() returns true.
Shard migrations in flight (adaptive splits and fold donors) already at MaxConcurrentAutoSplits Whole pass Count of shards that report they are splitting.
Cluster-wide split ceiling reached (MaxClusterConcurrentAutoSplits set) Per candidate No cluster headroom left in the admission gate; the candidate is deferred to a later tick.
Tree already has MaxPhysicalShardsPerTree physical shards (default 256) Per candidate (every hot shard) Counted on orleans.lattice.split.admission.deferred with reason=shard_ceiling.
Uniform load: hottest shard's rate below HotShardMinSkewRatio (default 1.5) times the median shard rate Per candidate (every hot shard) Counted on orleans.lattice.split.admission.deferred with reason=uniform_load.
Shard holds fewer than HotShardMinShardEntries live entries (default 1024) Per shard One count per candidate that cleared every cheaper clause; counted with reason=low_occupancy.
Shard already splitting Per shard Excluded from candidate set.
Per-shard cooldown active (default 2 min) Per shard In-memory cooldown timestamp.
Shard owns a single virtual slot Per shard Cannot be subdivided further.

Cluster-wide split concurrency (opt-in)

MaxConcurrentAutoSplits is enforced per tree: because the split monitor runs once per tree, each tree counts only its own in-flight migrations. In a multi-tenant or many-tree cluster the summed drain I/O from many trees splitting at once can saturate the storage provider even though no single tree exceeds its own cap.

MaxClusterConcurrentAutoSplits (default null = disabled) opts in to a cluster-wide admission gate - a singleton (well-known integer key 0) - that admits a new autonomic split only while the shard migrations in flight on the trees that set the ceiling, fold donors included, leave headroom under it. The ceiling is enforced in addition to each tree's MaxConcurrentAutoSplits and can only ever lower the number of splits a tree triggers, never raise it. When the option is null the monitor never requests a slot, so nothing is ever denied and the number of splits a tree starts is byte-for-byte identical to running without the option.

A tree with the ceiling unset does still publish an observation-only footprint to that singleton while it has migrations in flight, because the same grain is the cluster's readable split-activity source (see Reading split activity). Those footprints are held in a separate list and never consume admission headroom - a tree that never opted into a shared budget must not be able to throttle one that did - and publication is edge-triggered, so a tree with nothing splitting issues no call at all.

Reading split activity

Split progress is also exposed as a metric (orleans.lattice.split.in_flight), but metrics are write-only in-process: nothing can read one back to decide something. ILatticeAdmin.GetSplitActivityAsync is the readable counterpart, returning a cluster-wide SplitActivityReport (InFlight, ReportingTrees, ObservedAt, AnyInFlight) assembled from the footprints above. Both count shard migrations rather than only splits: a shard donating its slots to a shard consolidation (fold) counts as one in flight, exactly as an adaptive split's source does. It costs a single call to the singleton and never fans out across trees or shards.

Its first consumer is the Orleans.Lattice.Scaling scale-in safety gate, which suppresses scale-in while any split or fold is in flight so a silo is never drained underneath one; it is equally useful to an operator tool or a deployment guard. The count comes from the monitors' sampling passes, so it is not exact. An unsuppressed pass refreshes it once per HotShardSampleInterval; a pass suppressed while the tree is resized, resharded, merged into or snapshotted, during its age grace period, or after autonomic splitting is switched off re-publishes the tree's last count unchanged, so migrations a reshard starts are not counted while it runs; and a tree whose monitor starts with autonomic splitting disabled publishes nothing. Footprints expire, so a silo lost mid-migration cannot pin the count above zero indefinitely.

Per-tree heartbeat footprints (self-healing)

Admission uses a per-tree heartbeat model rather than long-lived permits. On every unsuppressed sampling pass an enabled monitor reports its tree's authoritative in-flight migration count (split sources and fold donors, read from each shard's own record of whether it is splitting) and how many new splits it wants; the gate drops any tree footprint whose time-to-live has lapsed, sums the live in-flight counts of the other trees, and grants new slots only up to the remaining cluster headroom. It then records this tree's footprint (its in-flight count plus any grant) with a fresh expiry of HotShardSampleInterval * 3. Because the count is re-reported from ground truth each pass, there is no permit to leak: a silo that crashes mid-split simply stops refreshing its footprint, so the stale entry lapses at its expiry and the next pass reclaims that share of the ceiling. This self-healing property is what makes the aggregate ceiling safe to enable.

Per-group override

Per-tree options resolve through named IOptionsMonitor<LatticeOptions>.Get(treeName), so a low-traffic tree group can clamp its own MaxConcurrentAutoSplits down (e.g. to 1) while a high-traffic group keeps a higher per-tree cap - all bounded in aggregate by the single global MaxClusterConcurrentAutoSplits ceiling.

Operator questions and the metrics that answer them

Question Metric How to read it
(a) Do I need to enable this? orleans.lattice.split.in_flight (summed across tree) and orleans.lattice.split.candidates_suppressed Both emit even when the gate is disabled. A high steady-state cluster sum with chronically non-zero suppression across many trees means aggregate drain pressure the per-tree cap cannot see.
(b) What ceiling should I pick? orleans.lattice.split.in_flight peak / quantiles Size the ceiling near the aggregate your storage provider absorbs without saturating (correlate with storage-latency panels), leaving headroom above typical peak so the gate only bites during pathological bursts.
(c) How is the enabled gate affecting my system? orleans.lattice.split.admission.deferred (reason=cluster_cap) Flat-zero means the ceiling never binds (raise it or leave the gate off). Sustained non-zero with rising hot-shard latency means the ceiling is too low and is starving legitimate elasticity.

Tunables (LatticeOptions)

Option Default Description
AutoSplitEnabled true Master switch for autonomic splits. When false, HotShardMonitorGrain will not trigger any splits. It does not gate an explicit ReshardAsync, which dispatches splits through the same coordinator to grow the shard count (see Online Reshard).
HotShardOpsPerSecondThreshold 200 Operations/second at or above which a shard is considered hot. Intentionally low so splits occur before throughput degrades.
HotShardSampleInterval 30 s How often the monitor polls hotness counters.
HotShardSplitCooldown 2 min Minimum interval between consecutive splits of the same physical shard.
MaxConcurrentAutoSplits 2 Maximum concurrent splits per tree. The monitor counts the shard migrations already in flight on the tree against it, so a shard consolidation (fold) in flight takes a slot too. Each split runs in its own per-shard coordinator activation; the cap bounds aggregate storage I/O.
MaxClusterConcurrentAutoSplits null Optional cluster-wide ceiling on autonomic split admission across all trees: a new split starts only while the shard migrations in flight on the trees that set it, folds included, leave headroom under it. null disables the gate (per-tree caps only, zero cost); a positive value opts in to a singleton admission gate enforced in addition to each tree's MaxConcurrentAutoSplits.
MaxConcurrentMigrations 4 Maximum concurrent splits (or, for a shrink, shard consolidations) an online reshard (ReshardAsync) keeps in flight. Separate from MaxConcurrentAutoSplits, but not additive with it: the monitor starts no split while a reshard runs, a growing reshard counts the shard migrations already in flight on the tree - autonomic splits and healing folds included - against this cap, and a shrinking reshard waits for any healing fold to finish before it starts its own. See Online Reshard.
SplitDrainBatchSize 1024 Maximum number of moved-slot entries the drain accumulates in memory before flushing to the target shard. Caps coordinator allocation regardless of source shard size.
BackgroundDrainLeavesPerPass 64 Maximum source leaves one Drain-phase pass visits before persisting its key cursor and yielding to the next tick. 0 or less disables the bound. The authoritative drains inside Swap and Complete are never bounded. Shared with the online snapshot copy and the cross-tree merge drain.
BackgroundDrainMaxDuration 10 s Wall-clock net for one Drain-phase pass, for leaves that are individually slow. TimeSpan.Zero disables it and leaves the leaf count as the only bound.
AutoSplitMinTreeAge 60 s Minimum tree age before autonomic splits are allowed; absorbs startup bursts.
HotShardMinSkewRatio 1.5 Minimum ratio of the hottest shard's rate to the tree's median shard rate before any split is admitted, so a uniformly loaded tree (a bulk ingest) is never split. A value at or below 1.0 disables the clause. It is the upper edge of the split/heal dead band whose lower edge is HotShardConsolidationSkewRatio.
HotShardMinShardEntries 1024 Minimum live entries a shard must hold to be split. 0 disables the floor and its per-candidate count probe.
MaxPhysicalShardsPerTree 256 Ceiling on the physical shard count autonomic splits may reach. An explicit ReshardAsync is not gated by it. 0 or less for no ceiling.
MaxScanRetries 3 Maximum bounded retries that a scan (CountAsync, ScanKeysAsync, ScanEntriesAsync) performs when ShardMap.Version keeps moving mid-scan due to concurrent splits. Throws InvalidOperationException on exhaustion. Increase if scans run during very-high split churn. See Consistency.

Automatic over-split healing, which folds shards back together once a tree's load is uniform, is tuned separately - see ShardHealingEnabled, HotShardConsolidationSkewRatio, and MaxConcurrentShardConsolidations in Configuration. While a fold runs, its donor counts as a migration in flight against the split caps above, and healing admits no new fold while any shard of the tree is splitting.

Convergence guarantees

  • No data loss - every write committed to S is either drained, shadow-mirrored, or both, and MergeManyAsync is idempotent under LWW.
  • No prepared-mutation loss - the retroactive sweep at BeginShadowWrite re-stamps every in-flight prepared mutation from S's leaves onto T's _pendingTx buckets, so a SetManyAtomicAsync saga whose Prepare landed on S before the split commits and completes against T with no perceived interruption. Combined with LatticeOptions.TxDecisionRetention (default 60 s), a sweep that installs a pending bucket after the saga's terminal fan-out has already broadcast can still resolve the verdict via the registry tombstone window. See Atomic Writes - Phase 4 Complete.
  • No mixed-round batch across the swap - the coordinator's drain copies a shadowing saga's pre-saga value into T with a migration marker, so between the shard-map swap and the arrival of the saga's backstop terminal on T a reader could otherwise see that key's old value while every sibling key already showed the new one. The coordinator therefore installs a per-key marker naming the shadowing saga, and T's read gate resolves it against the registry: a saga the registry reports as in-flight or aborted is safe (the pre-saga value is the correct answer), and so is any saga whose backstop terminal has already landed on T. Otherwise the read raises StaleShardRoutingException and the deadline-bounded retry loop re-fans once the backstop lands. A saga whose decision has aged out of TxDecisionRetention reads as Indeterminate and takes the same conservative arm as a committed one - serving the migrated pre-saga value there would be an affirmative claim that the saga did not commit, which is exactly what the registry has stopped vouching for. See Atomic Writes - After the retention window.
  • No duplicate authority - after the swap, only T is reachable for moved slots via the public API; orphan entries on S are unreachable and reclaimed when the tree is purged, or when a later shard consolidation retires S and releases its storage, unless a later consolidation folds T back into S, which drains T's entries onto S before lifting S's seal so the survivor's copy is authoritative again, and then releases T's storage.
  • Geometric convergence on a single hot slot - if all heat is in one virtual slot, successive autonomic splits subdivide S's slot set in half each pass, isolating the hot slot in O(log virtualSlotsPerShard) splits.

Scope

Shard splitting is an autonomic concern. ITreeShardSplitGrain is internal infrastructure: once AddLatticeAuth has installed its trust-boundary call filter, starting a split asserts that the call originated inside the cluster, so an external client call to start one is rejected with LatticeAuthorizationDeniedException (a cluster without that filter does not enforce the assertion). There is no public API to trigger or control an individual split. ILattice.ReshardAsync grows a tree's shard count by dispatching splits through the same coordinator (see Online Reshard); autonomic splitting is tuned exclusively through the LatticeOptions listed above.