Table of Contents

Mutation observers

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 mutation-observers.md, and llms.txt lists every page.

Part of Lattice Public API Reference.

IMutationObserver is a grain-side extensibility hook invoked synchronously after a mutation is durably committed, before the grain method returns to the caller. It is the primary seam for replication write-ahead logs, change-feed producers, and external audit consumers. It does not see every commit: bulk loads, a leaf split's key moves and the other paths listed under Emission points and shape commit without publishing.

Mutation observers vs. tree events. Both surface "something changed" notifications.

  • IMutationObserver is in-process, synchronous, and carries the full value bytes. It runs on the grain's scheduler before the write returns, so the caller's latency includes the observer's latency. Use it when a downstream component (replication WAL, outbox) must see the value at commit time and must be on the write path - typically another library, not application code.
  • Tree events are out-of-process, asynchronous, and metadata-only (kind, tree, key, shard, operation id, and a wall-clock AtUtc - no value bytes and no HLC). They ride Orleans Streams. Use them for UI updates, cache invalidation, dashboards, audit projections - anything that can tolerate at-most-once delivery and is willing to call GetAsync itself when it needs the value.

A single write typically fires both: the observer first (inline, with value), then an event (post-commit, metadata-only).

Register one or more observers in the silo DI container. They are resolved as IEnumerable<IMutationObserver>, so multiple can coexist. When none is registered the hook is zero-cost.

LatticeMutation.TreeId is always the logical tree id the caller addressed, and it stays stable across an alias swap: after a resize, a shadow-cutover restore (and its revert) or a schema remediation routes a tree's writes to a physical copy, observers keep receiving the logical id, so an observer keyed by tree id keeps working without reacting to the swap. The logical id is accepted only when it was routed to the exact physical tree that commits the write; a write that reaches a physical tree directly, without going through its logical tree, reports that physical tree's own id. The routing tier overwrites the routed identity on every routed write, and under AddLatticeAuth the capability-stripping call filter also strips it from external clients, so a caller cannot choose the id observers see. Aliasing changes no callback coverage: the paths that publish callbacks, and the ones that deliberately do not (bulk load and bulk append, and saga terminal records, which are WAL-only), are the same as on an unaliased tree. The WAL records remain keyed by the physical tree; the replication shipper decodes them back to the logical tree it ships.

Because the callback is on the write path, its cost is measured and attributed: every invocation is timed onto the orleans.lattice.observer.duration histogram, tagged observer (the observer's full CLR type name), tree, and tenant. That turns "one of our observers is slow" into a named series on the same OpenTelemetry pipeline as the traffic it slows down. The measurement is taken on the faulting path too, so an observer that throws slowly is just as visible as one that returns slowly, and spans only your callback - the warning Lattice logs for a faulting observer is not billed to it. The instrument costs nothing when no observer is registered, and nothing beyond a boolean read when no metrics listener is attached.

public sealed class MyReplicationObserver : IMutationObserver
{
    public Task OnMutationAsync(LatticeMutation mutation, CancellationToken ct)
    {
        // Inspect mutation.TreeId, mutation.Kind, mutation.Key, mutation.Value,
        // mutation.Timestamp, mutation.IsTombstone, mutation.ExpiresAtTicks,
        // mutation.OriginClusterId, mutation.VectorClock,
        // mutation.TransactionId, mutation.Category,
        // mutation.AtomicBatchSize, mutation.AtomicBatchIndex,
        // mutation.Delta, and for DeleteRange
        // also mutation.EndExclusiveKey.
        return Task.CompletedTask;
    }
}
siloBuilder.ConfigureServices(services =>
    services.AddSingleton<IMutationObserver, MyReplicationObserver>());

Emission points and shape

Mutation Kind Shape
SetAsync (all overloads), SetIfVersionAsync and GetOrSetAsync (when they write), SetManyAsync, SetManyWherePredicateAsync (the keys it writes), SetManyAtomicAsync, SetManyAtomicWhereAsync and each tree of a cross-tree atomic write (their prepare-phase writes, flagged IsPrepared), and every CRDT write (ApplyCrdtDeltaAsync, ApplyCrdtDeltaManyAsync, the typed accessors) Set One event per key. Value holds the committed bytes; Timestamp is the stamped HLC; ExpiresAtTicks carries the TTL deadline (or 0 for no-expiry).
DeleteAsync, and the staged deletes of a mixed set-and-delete SetManyAtomicAsync batch or a cross-tree atomic write (their prepare-phase writes, flagged IsPrepared) Delete One event per tombstoned key. IsTombstone is true, Value is null. A DeleteAsync of an absent key publishes nothing.
DeleteRangeAsync and DeleteRangeWherePredicateAsync (including each step of a delete-range cursor, which issues one over its sub-range) DeleteRange One event per bounded page of each shard's walk (not per key and not per user call): a shard whose part of the range fits in one page - bounded by LatticeOptions.MaxLeavesPerScanPage and MaxScanPageDuration - emits exactly one, and a longer walk emits one per page. Emitted even when the page matched zero live keys so replication consumers propagate the range unconditionally; only a DeleteRangeAsync on a tree that was never created returns 0 without publishing. Key carries the page's start - startInclusive on a shard's first page, the resume key on each later one - and EndExclusiveKey carries endExclusive, so the key pair differs between pages: consumers that need exactly-once per user call group on TransactionId, which every emit of one call shares. Timestamp is the single issue HLC the call stamped on every tombstone it wrote across the fan-out, so every emit of one call carries the same value. A predicate-filtered range delete also carries the page's exact matched keys in MatchedKeys (null for an unconditional range delete or a page that matched nothing).
Tree merge (MergeAsync), the copy an Online SnapshotAsync - or a resize, which runs one - drains into its destination, a backup restore that merges into its target, and replication apply of plain values Set / Delete Re-published so downstream consumers react exactly as they do to a foreground commit: every non-empty merged batch publishes one event per key in the batch, carrying the key's committed row (the pre-existing value where the incoming one lost last-writer-wins), whether or not the merge changed anything. A snapshot or resize copy is published under the destination's physical tree id, and while an online copy runs, each write the source forwards to the destination is published there again. Other replicated writes follow the rows above: a range delete as DeleteRange (one that carried matched keys arrives as per-key Deletes through this row), and CRDT deltas and prepared writes as Set / Delete.

These paths commit without publishing anything: the bulk loads (BulkLoadAsync, BulkAppendChunkAsync, and the streaming extension), which seed leaves directly; the copy of an Offline SnapshotAsync and a backup restore that bulk-loads a single full backup (into a shadow tree, or over the whole of a tree that did not exist), which use the same bulk seeding; a leaf split moving keys to its new sibling; a shard split's drain and forwarding and a shard consolidation, which move already-authored values between shards; replaying the write-ahead log or loading a snapshot when a leaf activates; a prepared write becoming visible when its saga commits (it was published once, flagged IsPrepared); and the saga terminal and compaction records described below. A write that commits while its leaf is splitting is published by that leaf and again by the leaf that receives the key, so an observer can see it twice.

MutationKind also defines TxCommit, TxAbort, and Tombstone. Those are write-ahead-log record kinds (WalRecord.Op) - an atomic-write saga's per-shard commit and abort terminal marks, and a compaction pass's tombstone reap marks - and are never delivered to an IMutationObserver.

Transaction correlation (TransactionId)

Every emit carries a Guid TransactionId that lets observers detect when several payloads belong to the same enclosing user call or saga:

  • Single-key writes (SetAsync, SetIfVersionAsync, GetOrSetAsync, DeleteAsync) get a fresh Guid per public call.
  • Atomic batches (SetManyAtomicAsync) share a single Guid across every per-key emit - including crash-recovery replays. An aborted batch issues no per-key rollback writes (its prepared writes never became visible; the abort is recorded once and broadcast to the affected shards), so it produces no compensating emits. Those per-key emits are the saga's prepare-phase writes, published with IsPrepared = true before the outcome is known; the commit or abort mark that resolves them is written to the write-ahead log only and is never delivered to an IMutationObserver, so an observer cannot tell from this seam alone whether a prepared emit committed.
  • Per-shard fan-out (DeleteRangeAsync) shares a single Guid across every per-shard emit.
  • Multi-key writes (SetManyAsync) share a single Guid across every per-key emit.
  • Convergence paths that re-publish an already-authored value (tree merge, an online snapshot's copy) emit Guid.Empty; cross-shard migration traffic (split shadow-forward and drain, shard consolidation) is not published at all.

Observers that batch by transaction group on TransactionId; observers that do not care simply ignore the field.

Pre-merge author's delta (Delta)

Every emit may carry an optional byte[]? Delta slot that lets a producer attach the author's pre-merge delta alongside the post-merge committed value. Lattice never opens the payload itself; consumers decode it based on WalRecord.Mode (the declared LatticeMergeMode of the tree the entry belongs to). The delta lets deterministic replay - and active-active CRDT receivers - reach the same convergence the original writer reached, which the post-merge Value bytes alone cannot guarantee for non-LWW CRDTs.

Stamp the slot via the LatticeDeltaContext ambient helper:

using (LatticeDeltaContext.With(new byte[] { 1, 2, 3 }))
{
    await tree.SetAsync("k", new byte[] { 4, 5, 6 }, cancellationToken);
}

You rarely stamp it yourself: every CRDT write (ApplyCrdtDeltaAsync, ApplyCrdtDeltaManyAsync, and the typed accessors built on them) stamps its typed delta on the emit automatically, and the atomic-write saga supplies each staged CRDT entry's delta per entry - the replication package's typed-delta receiver dispatch reads the stamped slot and applies via MergeDelta automatically.

Atomic-batch metadata (AtomicBatchSize / AtomicBatchIndex)

Every emit produced by an in-flight SetManyAtomicAsync saga carries int AtomicBatchSize (total entry count) and int AtomicBatchIndex (zero-based position). Single-key writes outside a saga emit 0 / 0.

using (LatticeAtomicBatchContext.With((5, 2)))
{
    // A forwarder for a remote saga can stamp the ambient directly;
    // the typical caller does not need to - the saga stamps it
    // on its own per-key writes.
}

The slots are independent of OriginClusterId, VectorClock, and Category. Wire-compatible: missing slots on legacy persisted state decode to 0. LatticeAtomicBatchContext also exposes CurrentIndexMap, CurrentDeltaMap, and CurrentDeleteSet, with With overloads that stamp them, so a saga can give each entry its own batch index, typed CRDT delta, and delete marker.

Remaining LatticeMutation slots

LatticeMutation carries further wire-compatible slots, each decoding to its default (false, 0, null, User, or LwwRegister) on state persisted before it existed. Most exist for the write-ahead log and replication rather than for a local observer: a per-key emit populates Category, IsPrepared and ShardIndex, a range-delete emit populates Category and MatchedKeys, and a local observer sees the others at their defaults.

Slot Type Meaning
Category MutationCategory User (the default) for a caller-driven write, or Maintenance for a library-internal structural write; replication-aware observers do not ship Maintenance emits across clusters. Independent of OriginClusterId.
IsPrepared bool true for an atomic-write saga's prepare-phase write, which becomes visible only when a later TxCommit terminal under the same TransactionId flips it (a TxAbort drops it).
ShardIndex int The logical chain-shard index that authored the mutation; activation-time replay uses it to ignore records a sibling shard sharing the same WAL partition wrote.
IsBackstop bool true for a cross-migration last-writer-wins backstop write, authored for a saga key whose prepare-phase shadow-forward was lost to a concurrent split or drain.
AtomicShardCount int On saga terminal marks, the number of shards the saga touched, so a cross-cluster receiver can hold visibility until every per-shard terminal has arrived; 0 on every other mutation.
IsMerge bool true for a merge write (replication apply, tree merge, snapshot restore, split redistribution, or migration import) or a compaction reap, so a consumer can tell it from a foreground write.
MatchedKeys IReadOnlyList<string>? The exact keys a predicate-filtered DeleteRange matched, so replay and replication tombstone exactly that set.
CrossTreeOperationId string? On a cross-tree atomic write's sub-saga terminal, the caller's operationId.
CrossTreeParticipants IReadOnlyList<string>? On the same terminals, the ordinal-sorted participant tree-id set.
Mode LatticeMergeMode The declared convergence rule, carried so WAL replay can re-fold a prepared CRDT mutation's typed delta. The observer publish does not populate it, so a local observer sees the LwwRegister default; a consumer that needs the mode reads the WAL record's Mode.