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.
IMutationObserveris 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 callGetAsyncitself 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 freshGuidper public call. - Atomic batches (
SetManyAtomicAsync) share a singleGuidacross 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 withIsPrepared = truebefore 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 anIMutationObserver, so an observer cannot tell from this seam alone whether a prepared emit committed. - Per-shard fan-out (
DeleteRangeAsync) shares a singleGuidacross every per-shard emit. - Multi-key writes (
SetManyAsync) share a singleGuidacross 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. |