---
title: "Adaptive Shard Splitting"
url: "https://nsta1.github.io/Orleans.Lattice/docs/lattice/shard-splitting.html"
source: "https://github.com/NSTA1/Orleans.Lattice/blob/release/9.9/docs/lattice/shard-splitting.md"
package: "Orleans.Lattice"
version: "9.9.0"
documents: "Orleans.Lattice 9.9.0 (release line 9.9)"
built: "2026-10-04"
all-pages: "https://nsta1.github.io/Orleans.Lattice/llms.txt"
bundle: "https://nsta1.github.io/Orleans.Lattice/docs/lattice/llms-full.txt"
---
# Adaptive Shard Splitting

Part of the [Orleans.Lattice documentation](architecture.md).

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](online-reshard.md) 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.

```mermaid
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](consistency.md).

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](consistency.md) 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](snapshot-cursors.md)): 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](state-model.md#activation-replay-rehydrate-and-the-safety-net)).
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](#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](online-reshard.md)). |
| `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](online-reshard.md). |
| `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](consistency.md). |

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](configuration/options-reference-3.md#shardhealingenabled). 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](atomic-writes.md#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](atomic-writes.md#after-the-retention-window-indeterminate-not-inflight).
* **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](online-reshard.md)); autonomic splitting is tuned
exclusively through the `LatticeOptions` listed above.
