Repository navigation
feat(mesh): durable streams are saved mesh nodes, consumed inside hubs - #6018
Conversation
A durable stream is a stream node {Key}/_DurableStream/{Namespace} whose own hub
orders and leases it, plus one write-once item node per appended item in the
owner partition's new durable_stream_items satellite table. Nothing about a
stream lives outside the mesh.
- hub.PublishDurable(stream, payload): the stream hub creates the item node,
records LastSequence, then answers with the sequence (durable on emit).
- config.WithDurableStreamConsumer<T>(namespace, streamKey, handler): claims
the lease, reads items after the acknowledged checkpoint (one synced query:
backlog + live Added), hands each item to the hub as a request to itself so
the handler runs on the hub's loop and a recycle's quiesce waits for it,
acknowledges after the handler completes, releases on shutdown.
- The lease IS the subscription: epochs fence stale holders, a claim takes
over only from an unreachable holder, a release orphans the stream (the
stream node's change is the event), an orphaned stream with pending items
wakes its owner by a re-probing read of the owner node.
- DurableStreamSource<T> lets a consumer read an existing ordered child set
(e.g. _DocumentPart children) without a second copy.
Tests: DurableStreamConsumerTest (Monolith, 6) and
OrleansDurableStreamConsumerTest (Orleans TestCluster, grain-hosted consumer
and stream hub). Falsified: without the acknowledgement 3/6 fail (recycle
re-processes items); with live-only delivery 6/6 fail.
Doc: Doc/Architecture/DurableStreams.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Test Results (shard 4)660 tests - 179 660 ✅ + 14 2m 37s ⏱️ +8s Results for commit c692bb5. ± Comparison against base commit 174f317. This pull request removes 216 and adds 37 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
Test Results (shard 2) 2 files - 2 2 suites - 2 4m 32s ⏱️ +20s Results for commit c692bb5. ± Comparison against base commit 174f317. This pull request removes 77 and adds 665 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
Test Results (shard 1) 2 files - 1 2 suites - 1 4m 27s ⏱️ -49s Results for commit c692bb5. ± Comparison against base commit 174f317. This pull request removes 669 and adds 2 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
…stream A NodeType that consumes a stream family is carried by every hub of that type, most of which never get a stream (every uploaded document will carry the DocumentLog consumer). Activation now costs ONE existence query (a listing read, never a point read of a possibly absent node): a stream that exists is claimed; one that does not is claimed only when a producer's append wakes the hub (DurableStreamWake). Consumer requests (claim, ack, release) never create a stream node — only a producer does. Test: five active consumers without a stream create no stream node over a window; a publish to one wakes exactly that one. Falsified: claiming at every activation (the previous behaviour) creates 5 stream nodes and the test fails. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Test Results (shard 3) 3 files - 1 3 suites - 1 15m 59s ⏱️ + 2m 54s Results for commit c692bb5. ± Comparison against base commit 174f317. This pull request removes 37 and adds 1187 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
Test Results (shard 5) 4 files ± 0 4 suites ±0 9m 33s ⏱️ - 3m 12s Results for commit c692bb5. ± Comparison against base commit 174f317. This pull request removes 1205 and adds 14 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
Test Results 15 files - 5 15 suites - 5 40m 29s ⏱️ -34s Results for commit c692bb5. ± Comparison against base commit 174f317. This pull request removes 317 and adds 19 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
… period A stream with a producer and no consumer (or an owner whose NodeType carries no consumer) read its owner on every append. The wake now happens once per orphan period; granting or releasing a lease starts a new one. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
# Conflicts: # src/MeshWeaver.Mesh.Contract/SatelliteTableMapping.cs
|
Shard 5 red is an unrelated flake, not this diff. It is
I have re-run the failed jobs. |
# Conflicts: # src/MeshWeaver.Mesh.Contract/SatelliteTableMapping.cs
|
🩹 PR babysitter (build instance) is re-running the failed jobs of run(s) 37139221887 on head Why: inherited from its base 'main': 4 pull requests of Systemorph/MeshWeaver fail identically — 'lane / Automatic review answered' concluded failure: Process completed with exit code 1. — the base 'main' moved from 174f317 (what the red run tested) to 1bafcb2. Validated: Base-red, not this diff: 4 pull requests of Systemorph/MeshWeaver fail identically on 'lane / Automatic review answered' concluded failure: Process completed with exit code 1 (same check, same words), and the base 'main' the red run tested has since moved from 174f317 to 1bafcb2 — both commits in the proposal's evidence. Approving the update-branch: it merges the moved base into the branch and re-runs run 37139221887 against it, which drains the base's red rather than re-testing the old merge commit. No prior heals on this head. It does not merge, push or dequeue. A red after this re-run is left for the owner (rbuergi). |
Durable streams: saved mesh nodes, consumed inside hubs
A durable stream is a set of saved mesh nodes. It is consumed inside a hub, and it survives a recycle, a pod roll and a restart. Design of record:
Doc/Architecture/DurableStreams(new page, in this PR).Shape
Subject+Concatchain. Sequences are therefore contiguous, and the lease has exactly one owner.Subscriber=nullandOrphanedAtare set, and that change to the stream node is the orphan event.Checkpoint,CheckpointAt).SatelliteTableMapping("_DurableStreamItem", "durable_stream_items", "DurableStreamItem"). The segment is 18 characters, longer than_ThreadMessageand_DocumentPart, so it wins the longest-segment placement rule.ensure_partition_schemagenerates the table from the mapping list, so the Postgres side needs no code change.Added), and items are released strictly in order.DurableStreamSource<T>reads an existing ordered child set, such as a document's_DocumentPartchildren, as a stream without a second copy.API
IObservable<long> hub.PublishDurable<T>(DurableStreamId stream, T payload)— cold; emits the sequence once the item node is saved.config.WithDurableStreamConsumer<T>(string ns, Func<IMessageHub,string> streamKey, Func<IMessageHub, DurableStreamItem<T>, IObservable<Unit>> handler)(plus an overload that takes aDurableStreamSource<T>).hub.ReadDurable<T>(stream[, source], afterSequence)— an unleased reader.public sealed record DurableStreamId(string Namespace, string Key),public sealed record DurableStreamItem<T>(DurableStreamId Stream, long Sequence, T Payload).Three behaviours found while building this
cannot process AckDurableStreamRequest … RunLevel=Quiescing). For that reason, ack, release and claim are issued fromhub.GetMeshHub().ReadIssuingHub(), never from the consumer.GetMeshNodeOutcome), followed by aDurableStreamWakenudge.DurableStreamDirectory), never discovered by routing to them.Tests (executed)
DurableStreamConsumerTest(Monolith): 6/6. It covers:OrleansDurableStreamConsumerTest(Orleans TestCluster, grain-hosted consumer and stream hub, client producer): 1/1.MeshWeaver.Graph.Testfull 2460/2460,MeshWeaver.Hosting.Testfull 1094 passed, 1 skipped,MeshWeaver.Documentation.Testguards 662/662.DurableStreamPostgresTests). It checks that items are rows of the owner partition'sdurable_stream_itemstable and that the checkpoint is on the stream node's row. It passes against pgvector, and it fails with42P01 … durable_stream_items does not existwhen the mapping is removed.Limits (stated in the doc)
Deploy
No migration and no DbVersion bump. Existing partitions get
durable_stream_itemsfromensure_partition_schema. Nothing is served until a hub opts in withWithDurableStreamConsumer, so there is no address to recycle.Pairs-with: none — this PR removes no public type or member.
Implementers: none — no member is added to an existing public interface.
Mirror-sync: none — no i18n catalog key is added or changed.
🤖 Generated with Claude Code