Durable Streams

Read with: Durable Streams Are Mesh Nodes (why the mesh, not a stream provider, is the durable channel) · Hub Disposal Model (why an item in flight survives a recycle) · Asynchronous Calls (everything below is IObservable, nothing is async).

A durable stream is an append-only, per-stream ordered sequence of items that survives a process restart, a silo roll and a hub recycle, and is consumed inside a hub: a hub declares in its configuration that it consumes a stream, every activation of that hub resumes where the last one stopped, and each item is handed to the handler on the hub's own message loop.

A durable stream IS saved mesh nodes. Nothing about it lives anywhere else:

{Key}/_DurableStream/{Namespace}                                     DurableStream      owner's mesh_nodes
{Key}/_DurableStream/{Namespace}/_DurableStreamItem/000000000001     DurableStreamItem  owner's durable_stream_items
{Key}/_DurableStream/{Namespace}/_DurableStreamItem/000000000002     DurableStreamItem
…

Producing

hub.PublishDurable(new DurableStreamId("DocumentLog", documentPath), new DocumentLogAppend(text, offset))
   .Subscribe(sequence => …, ex => logger.LogWarning(ex, "append failed"));

PublishDurable is cold. It sends AppendDurableStreamItemRequest to the stream node's hub, which runs every append, claim, acknowledgement and release through one serial chain: it creates the item node {seq:D12}, records LastSequence, and only then answers with the sequence. When the publish emits, the item is a saved mesh node. The stream node is created on first use, once per stream per process (DurableStreamDirectory).

🚨 Ensure before you route. A request routed to a stream node that does not exist yet is answered NotFound, and on the Orleans route that refusal outlived the node's creation: a resend right after a successful create was refused the same way (measured on the Orleans TestCluster). That is why the directory creates the stream node before the first request instead of discovering it by routing to it.

Consuming inside a hub

// in the NodeType's (or any hub's) configuration
config.WithDurableStreamConsumer<DocumentLogAppend>(
    "DocumentLog",
    hub => hub.Address.Path,                                // the owner whose stream this hub consumes
    (hub, item) => hub.GetMeshNodeStream(hub.Address.Path)  // runs ON THE HUB'S LOOP
        .Update(node => Append(node, item.Payload, item.Sequence))
        .Select(_ => Unit.Default));                        // complete = done; then it is acknowledged

On every activation the consumer:

  1. Stays dormant until there is a stream. A NodeType that consumes a family is carried by every hub of that type, and most of them never get a stream (every uploaded document carries the DocumentLog consumer). Activation therefore costs one existence query, a listing read and never a point read of a node that may be absent. If the stream node does not exist, the consumer creates nothing, claims nothing and arms no timer until a producer's append wakes the hub (DurableStreamWake). Consumer requests (claim, ack, release) never create a stream node; only a producer does.
  2. Claims the stream (ClaimDurableStreamRequest). The stream node's hub grants the lease when the stream is free or the claimant already holds it. Every grant starts a new epoch.
  3. Reads the items strictly after the acknowledged Checkpoint. One synced query over the items' namespace provides both halves: its initial result is the backlog, its later Added changes are the live tail (the mesh change feed, Postgres LISTEN across pods). Items are released strictly in sequence. One that arrives ahead of a missing predecessor waits for it; one at or below the cursor is a duplicate and is dropped.
  4. Hands each item to the hub as a request to itself (DurableStreamTurn), so the handler runs on the hub's action block, serially with its messages. Once the handler's observable completes, the turn acknowledges the sequence (AckDurableStreamRequest, carrying the epoch) and replies. Only then is the next item handed over.
  5. Releases the lease when the hub shuts down, once the item in flight is done. That is the orphan event: the stream node's Subscriber clears, OrphanedAt is stamped and LastSubscriber is named.

A document's existing _DocumentPart/{index:D6} children can be consumed as a stream without a second copy, by passing a DurableStreamSource<T> (items namespace, id → sequence, node → payload):

config.WithDurableStreamConsumer<DocumentPart>(
    "DocumentLog",
    hub => hub.Address.Path,
    hub => new DurableStreamSource<DocumentPart>(
        DocumentPartPaths.PartNamespace(hub.Address.Path),
        node => DocumentPartPaths.TryParsePartIndex(node.Id, out var i) ? i + 1 : null,
        (node, options) => node.ContentAs<DocumentPart>(options)!),
    (hub, part) => …);

The lease and the checkpoint still live on {document}/_DurableStream/DocumentLog. Only the items are read from where they already are.

A plain, unleased reader is hub.ReadDurable<T>(stream, afterSequence): same ordering, no acknowledgement.

The lease IS the subscription

question where the answer is
does the stream have a live subscriber, and who? DurableStreamState.Subscriber (+ SubscribedAt, LeaseEpoch)
how far has it got, and when did it last advance? Checkpoint, CheckpointAt
did the last subscriber leave? Subscriber == null, OrphanedAt, LastSubscriber. Watch the stream node: its change is the event
who decides between two would-be subscribers? the stream node's own hub. Claims are serialized and exactly one is granted
can a subscriber that lost the lease still advance the checkpoint? no. An acknowledgement carrying an old epoch is fenced

Takeover. A claim against a held lease makes the stream's hub ask the holder for a sign of life (PingRequest, DurableStreamHub.HolderProbeBudget = 5 s). It grants the lease to the new claimant (new epoch, logged as TOOK OVER) only when the holder cannot be reached. Reaching a node address activates its hub, so a node that owns a lease always answers. Ownership stays with the address, and only an address no silo can serve any more (a gone portal or worker) loses it. A refused consumer claims again on the orphan event, on a wake, or every ReclaimInterval (30 s). That interval is the upper bound on taking over from a holder that died without releasing.

Wake. An append to an orphaned stream, or a release that leaves unacknowledged items, makes the stream's hub wake the owner (Key). The wake is a read of the owner node, not a message: a recycled owner is still quiescing for a moment after it released, and refuses every new message then. The mesh read re-probes that refusal until a fresh activation answers, and the activation starts the consumer. A DurableStreamWake then nudges a consumer that was waiting for the lease.

Guarantees

Limits

Why nodes and not a stream provider

Every alternative needs a second store with its own sequence, retention, access model and recovery path. Those are exactly the properties the mesh already has for nodes. The durable part is the item node; ordering and ownership are decided by one hub. Delivery rides the change feed every pod already runs. An Orleans persistent-stream provider over a separate queue table was built first and dropped in favour of this: it duplicated the store, needed its own offsets per cluster and queue, and still needed the store for catch-up. On Orleans the stream node and the consuming node are grain-hosted hubs like any other, so the same hub code runs unchanged on the Monolith and the Distributed portal.

How to check a live cluster

# lease grants, takeovers, orphans — the stream hub's own sentences
{namespace="<namespace>"} |= "[DurableStreams]"

# a consumer that stopped on an item (it retries; a repeating line is a poison item or a broken handler)
{namespace="<namespace>"} |= "Durable stream consumer" |= "was NOT acknowledged"

A stream's state is its node: get @{Key}/_DurableStream/{Namespace}. Its items are search namespace:{Key}/_DurableStream/{Namespace}/_DurableStreamItem.

Pinned by

DurableStreamConsumerTest (Monolith, seven cases, including dormancy: five consumers without a stream create nothing until a publish wakes exactly one) and OrleansDurableStreamConsumerTest (Orleans TestCluster: grain-hosted consumer and stream hub, client producer). Both were falsified when they landed. Without the acknowledgement the recycle case re-processed items, and with live-only delivery (the memory-stream shape) every case lost items. When the consumer was made to claim at every activation, the dormancy case found five stream nodes it should never have created.