---
title: "Mutation observers - Lattice Public API Reference"
url: "https://nsta1.github.io/Orleans.Lattice/docs/lattice/api/mutation-observers.html"
source: "https://github.com/NSTA1/Orleans.Lattice/blob/release/9.9/docs/lattice/api.md?plain=1#L702-L906"
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"
---
# Mutation observers

Part of [Lattice Public API Reference](../api.md).

`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](#emission-points-and-shape) commit
without publishing.

> **Mutation observers vs. [tree events](../events.md).** 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`](../metrics/instrument-catalog-7.md#mutation-observers)
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.

```csharp verify
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;
    }
}
```

```csharp verify
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 `Delete`s 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:

```csharp verify
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`.

```csharp verify
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`. |

Previous: [ILatticeAdmin](ilatticeadmin.md). Next: [Caller-credential propagation (LatticeCredentialContext)](caller-credential-propagation-latticecredentialcontext.md). Contents: [Lattice Public API Reference](../api.md).
