Skip to content

Commit cabf439

Browse files
committed
Add retained topics and durable subscription processing
1 parent 2bc7ccc commit cabf439

18 files changed

Lines changed: 1072 additions & 27 deletions

File tree

‎README.md‎

Lines changed: 40 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,16 @@
11
# KeyLoad
22

3-
KeyLoad is a .NET database built around atomic transaction domains, durable ordered storage, and an RF3 Raft cluster. Documents, event streams, work queues, graph edges, samples and vector sidecars share the same transactional command path. Orleans provides the routing layer; each node owns its ZoneTree materialization and journals.
3+
KeyLoad is a .NET database built around atomic transaction domains, durable ordered storage, and an RF3 Raft cluster. Documents, event streams, retained topics, subscription groups, work queues, graph edges, samples and vector sidecars share the same transactional command path. Orleans provides the routing layer; each node owns its ZoneTree materialization and journals.
44

55
The original load-testing prototype has been replaced. This repository implements the new [architecture and development plan](docs/design/architecture-v0.3.uk.md), with the [original HTML edition](docs/design/architecture-v0.3.uk.html) preserved alongside it.
66

77
## Development status
88

99
This is an early implementation of the clustered kernel. The default server topology has three persistent voting nodes and requires a majority for writes and strong reads. It needs no external database, Redis or message broker.
1010

11-
Implemented surfaces include document CRUD and field patches, partition-scoped unique/composite equality indexes, command deduplication, stream expected-revision append/read, scheduled work queues with fenced leases and retry/DLQ, atomic inbox completion, bounded property-graph traversal, ordered samples, exact vector search, BM25 and weighted reciprocal-rank fusion, a bounded read-only Q1 SQL dialect, API keys, field omission and field-use policies, verified backups, native Raft snapshots and empty-replica catch-up, offline journal compaction, a .NET SDK and CLI.
11+
Implemented surfaces include document CRUD and field patches, partition-scoped unique/composite equality indexes, command deduplication, stream expected-revision append/read, retained topics, per-source durable subscriptions with bounded delivery windows and contiguous checkpoints, scheduled work queues with fenced leases and retry/DLQ, atomic inbox completion, bounded property-graph traversal, ordered samples, exact vector search, BM25 and weighted reciprocal-rank fusion, a bounded read-only Q1 SQL dialect, API keys, field omission and field-use policies, verified backups, native Raft snapshots and empty-replica catch-up, offline journal compaction, a .NET SDK and CLI.
1212

13-
The [104-task implementation tracker](docs/implementation/status.json) records the remaining work. Automatic canonical journal maintenance, large and interrupted snapshot transfer qualification, shard movement, distributed multi-shard query planning, managed HNSW, topics and subscription groups, CDC, retention, schema migrations, external backup stores and full release qualification remain under development. The current server uses one replicated physical shard which contains many independent atomic partitions.
13+
The [104-task implementation tracker](docs/implementation/status.json) records the remaining work. Automatic canonical journal maintenance, large and interrupted snapshot transfer qualification, shard movement, distributed multi-shard query planning, managed HNSW, subscription coverage discovery across partitions and rebalance, CDC, retention, schema migrations, external backup stores and full release qualification remain under development. The current server uses one replicated physical shard which contains many independent atomic partitions.
1414

1515
The advertised profiles are `ProcessDurable` for embedded storage and `QuorumProcessDurable` for the cluster. The kernel has process-kill recovery tests and the cluster has real leader-loss and minority tests. Power-loss qualification, broader platform qualification and the 72-hour endurance gate are still required before advertising `LocalDurable`, `QuorumDurable` or production readiness.
1616

@@ -85,6 +85,42 @@ foreach (var delivery in received.Value.Deliveries)
8585

8686
`CommitProcessing` atomically stores the same-partition effects, inbox receipt and ACK. A stale lease cannot acknowledge a later delivery. External side effects need their own idempotency or outbox integration.
8787

88+
## Retained topics and subscriptions
89+
90+
Topics retain events once per source. Each durable group owns its delivery window, attempts, leases and checkpoint. Groups can also subscribe to one stream generation using `EventSourceKind.Stream` and a stream ID. Source identity includes the atomic partition; cross-partition coverage is still being developed.
91+
92+
```csharp
93+
await client.ConfigureResourceAsync(Guid.NewGuid(), new("acme", "shop",
94+
new ResourceDefinition("activity", ResourceKind.Topic, "order-processing")));
95+
var source = new EventSourceRef(partition, "activity", EventSourceKind.Topic);
96+
var group = new SubscriptionRef(source, "projection");
97+
await client.ConfigureSubscriptionAsync(new(Guid.NewGuid(), group,
98+
new SubscriptionDefinition("root"), SubscriptionStart.FromBeginning));
99+
100+
var publish = new CommandRequest(Guid.NewGuid(), partition,
101+
[new PublishTopic("activity", [new("activity-1", "OrderCreated", "{\"number\":1}")])]);
102+
(await client.CommitAsync(publish)).ThrowIfFail();
103+
104+
var groupReceived = await client.ReceiveSubscriptionAsync(new(Guid.NewGuid(), group, MaxEvents: 10));
105+
groupReceived.ThrowIfFail();
106+
foreach (var delivery in groupReceived.Value.Deliveries)
107+
{
108+
var processing = new SubscriptionProcessingRequest(Guid.NewGuid(), group, delivery.Token,
109+
"order-projection", ExecutionGeneration: 1,
110+
[new PutDocument("orders", "projection-" + delivery.Event.Data.EventId,
111+
delivery.Event.Data.PayloadJson, ExpectedRevision: 0)]);
112+
(await client.CommitSubscriptionProcessingAsync(processing)).ThrowIfFail();
113+
}
114+
```
115+
116+
Use the application's data principal for a deployed group; `root` above is the local development principal. Delivery projects the intersection of that data principal's and the worker's current field permissions. A missing required input grant stops delivery. Revocation or a policy epoch change invalidates old payload retries and acknowledgement tokens.
117+
118+
Keep publish, receive and processing IDs stable while resolving an uncertain outcome. Event IDs are unique within a retained topic or stream generation: identical reuse returns `DuplicateEventId`, and changed content returns `Conflict`. Subscription processing stores effects, inbox and ACK in one commit. Replay of the same handler/input returns `AlreadyProcessed` and `OriginalEffectsToken`; a newly leased replay is acknowledged without applying those effects again. A different authorized worker can recover the same inbox outcome after current effect permissions are checked.
119+
120+
ACK only advances a contiguous prefix. ACKs for positions 1 and 3 leave checkpoint 1 until position 2 finishes. `MaxWindow` bounds gaps and active deliveries. NACK uses bounded exponential backoff; exhausted attempts park the unresolved position and pause the group. `SeekSubscriptionAsync` advances its generation, fences previous leases and pauses the group until an explicit resume through `SetSubscriptionPausedAsync`.
121+
122+
`FromNow` records the current source tail inside the creation commit. `FromCursor` uses a signed, principal-bound cursor from `ReadEventSourceAsync`; even an empty page at the tail returns a cursor for later catch-up. Topics enforce retained event/byte quotas and reject the entire producer batch when full. Retention reclamation, group deletion/filter migration and completion-audit compaction remain tracked work.
123+
88124
## Queries and search
89125

90126
```csharp
@@ -143,7 +179,7 @@ dotnet test --project tests/KeyLoad.RecoveryTests --no-build --no-restore
143179
dotnet test --project tests/KeyLoad.IntegrationTests --no-build --no-restore
144180
```
145181

146-
Tests use xUnit and Microsoft.Testing.Platform. Recovery qualification runs 1000 seeded real-process kills, checks complete-frame corruption and verifies a clean backup restore. Additional process kills cover checkpoint publication and all native Raft append paths immediately after acknowledgement. Integration tests own the Aspire lifecycle and use three independent server processes with isolated persistent directories. They kill the elected leader, retry the same command, verify atomic effects on surviving voters, reject a minority write and erase/restart one replica to require snapshot catch-up. No manually running AppHost is needed for tests.
182+
Tests use xUnit and Microsoft.Testing.Platform. Recovery qualification runs 1000 seeded real-process kills, checks complete-frame corruption and verifies a clean backup restore. Additional process kills cover checkpoint publication, native Raft append acknowledgement and subscription effects/inbox/checkpoint publication. Integration tests own the Aspire lifecycle and use three independent server processes with isolated persistent directories. They kill the elected leader, retry the same command, verify atomic effects and topic delivery on surviving voters, reject a minority write and erase/restart one replica to require snapshot catch-up of documents, group checkpoints and inbox outcomes. No manually running AppHost is needed for tests.
147183

148184
Benchmarks can be started with `dotnet run -c Release --project benchmarks/KeyLoad.Benchmarks`. Comparative PostgreSQL/Marten/Wolverine and multi-node scaling qualification are part of the development plan; no comparative performance claims are made yet.
149185

‎docs/implementation/durability-audit.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,12 @@ Format 2 adds a checkpoint prefix with ordered live records, per-frame SHA-256 c
1616

1717
The .NEXT 6.8.1 persistent Raft log owns terms, votes, log matching and consensus commit. `DurableRaftLog` requires native manual checkpoints (`FlushInterval = Timeout.InfiniteTimeSpan`), serializes its public append/commit paths and awaits the native foreground `FlushAsync` before acknowledgement. Background flush notifications do not establish this boundary for an isolated first append or snapshot installation. Four regressions invoke append through the actual `IPersistentState` interface, kill the process immediately after ACK and reopen the tail using private page memory; they do not rely on graceful disposal or shared memory pages surviving exit.
1818

19-
The state machine produces native Raft snapshots from a verified canonical store cut. A snapshot-only regression verifies that installation acknowledges without a subsequent append and reopens at that cut. The RF3 test erases one stopped follower's data directory, requires native snapshot installation, checks the fresh node identity and read generation and verifies documents and command outcomes through the HTTP API. Incoming bodies are authenticated and rewound for both `Body` and `BodyReader`; outgoing one-shot payloads are serialized once into a bounded disk spool. Raft term, snapshot metadata, request ID and content type are included in the peer signature.
19+
The state machine produces native Raft snapshots from a verified canonical store cut. A snapshot-only regression verifies that installation acknowledges without a subsequent append and reopens at that cut. The RF3 test erases one stopped follower's data directory, requires native snapshot installation, checks the fresh node identity and read generation and verifies documents, retained events, subscription checkpoints, inbox receipts and command outcomes through the HTTP API. Incoming bodies are authenticated and rewound for both `Body` and `BodyReader`; outgoing one-shot payloads are serialized once into a bounded disk spool. Raft term, snapshot metadata, request ID and content type are included in the peer signature.
2020

2121
The bounded leader writer replicates a trusted operation with one leader-chosen evaluation time. It returns success after consensus commit and local state-machine apply, then resolves the persisted outcome under current authorization. Cancellation or response loss after admission has an unknown write outcome; retries must retain the command ID. A minority cannot acknowledge writes.
2222

23+
Per-source subscriptions journal bounded delivery references, contiguous checkpoints and signed lease generations. Processing effects, inbox receipts, ACK and the command outcome share one compiled redo frame. Seven process-kill scenarios reopen the real store before retrying a processing command and require the effect and checkpoint to agree; a second inbox retry must not apply the effect twice. Seek is an explicit generation fence. Retained topic quotas fail the producer batch before publication; reclamation and cross-partition coverage are separate pending work.
24+
2325
Strong reads on a leader force a quorum barrier tied to that leadership term. Followers request this barrier through an authenticated peer endpoint and wait for local application of the returned committed position. The provider's follower synchronization API alone is not used as KeyLoad's strong-read acknowledgement. Text and vector search, SQL predicates, row authorization and returned field projection each stay within one storage read gate.
2426

2527
Orleans membership is a replicated catalog record, bootstrapped through consensus before Orleans starts. Grains route commands; they do not own files or durability. The first topology contains one physical shard and many separately scoped atomic partitions.

‎docs/implementation/kernel-qualification.json‎

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,14 +7,15 @@
77
"restore": "locked-mode passed",
88
"build": "passed with zero warnings",
99
"unitTests": {
10-
"passed": 26,
10+
"passed": 42,
1111
"failed": 0
1212
},
1313
"recoveryTests": {
14-
"passed": 43,
14+
"passed": 50,
1515
"failed": 0,
1616
"seededProcessCrashes": 1000,
1717
"checkpointPublicationCrashes": 8,
18+
"subscriptionProcessingCrashes": 7,
1819
"nativeRaftAcknowledgementCrashes": 4,
1920
"nativeSnapshotTests": 2,
2021
"peerSecurityTests": 6,
@@ -34,11 +35,14 @@
3435
"leader process kill",
3536
"same-command outcome retry",
3637
"atomic inbox/effects/ACK",
38+
"retained topics and contiguous subscription ACK after leader loss",
39+
"subscription inbox replay without duplicate effects",
3740
"minority read/write rejection",
3841
"voter restart",
3942
"unsigned peer RPC rejection",
4043
"empty-replica native snapshot catch-up",
41-
"command outcome preserved across snapshot installation"
44+
"command outcome preserved across snapshot installation",
45+
"topic events, group checkpoints and inbox receipts preserved across snapshot installation"
4246
]
4347
},
4448
"durability": [
@@ -47,5 +51,10 @@
4751
],
4852
"powerLoss": "not qualified",
4953
"endurance72Hours": "pending",
50-
"crossPlatformCI": "pending"
54+
"crossPlatformCI": {
55+
"qualifiedCommit": "2bc7ccc93997a36555c71553403d5539d64b90d7",
56+
"run": "https://github.com/managedcode/KeyLoad/actions/runs/36885624099",
57+
"platforms": ["Linux", "macOS", "Windows"],
58+
"currentTopicSubscriptionChanges": "pending"
59+
}
5160
}

‎docs/implementation/status.json‎

Lines changed: 24 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -595,6 +595,7 @@
595595
"status": "in_progress",
596596
"evidence": [
597597
"src/KeyLoad.Abstractions/Contracts.cs",
598+
"src/KeyLoad.Abstractions/Subscriptions.cs",
598599
"docs/implementation/durability-audit.md"
599600
]
600601
},
@@ -603,6 +604,8 @@
603604
"status": "in_progress",
604605
"evidence": [
605606
"src/KeyLoad.Core/Events.cs",
607+
"src/KeyLoad.Core/EventSources.cs",
608+
"src/KeyLoad.Core/SubscriptionGroups.cs",
606609
"tests/KeyLoad.UnitTests/TransactionTests.cs"
607610
]
608611
},
@@ -611,15 +614,18 @@
611614
"status": "in_progress",
612615
"evidence": [
613616
"src/KeyLoad.Core/Events.cs",
614-
"tests/KeyLoad.UnitTests/TransactionTests.cs"
617+
"tests/KeyLoad.UnitTests/TransactionTests.cs",
618+
"tests/KeyLoad.UnitTests/SubscriptionTests.cs"
615619
]
616620
},
617621
"KL-084": {
618622
"title": "ReadStream, feed cursors і catch-up",
619623
"status": "in_progress",
620624
"evidence": [
621625
"src/KeyLoad.Core/Events.cs",
622-
"tests/KeyLoad.UnitTests/TransactionTests.cs"
626+
"src/KeyLoad.Core/EventSources.cs",
627+
"tests/KeyLoad.UnitTests/TransactionTests.cs",
628+
"tests/KeyLoad.UnitTests/SubscriptionTests.cs"
623629
]
624630
},
625631
"KL-085": {
@@ -656,16 +662,24 @@
656662
},
657663
"KL-089": {
658664
"title": "Durable groups і contiguous checkpoints",
659-
"status": "pending",
660-
"evidence": []
665+
"status": "in_progress",
666+
"evidence": [
667+
"src/KeyLoad.Core/SubscriptionGroups.cs",
668+
"tests/KeyLoad.UnitTests/SubscriptionTests.cs",
669+
"tests/KeyLoad.RecoveryTests/SubscriptionRecoveryTests.cs",
670+
"tests/KeyLoad.IntegrationTests/ClusterTests.cs"
671+
]
661672
},
662673
"KL-090": {
663674
"title": "Inbox і CommitProcessing",
664675
"status": "in_progress",
665676
"evidence": [
666677
"src/KeyLoad.Core/Messaging.cs",
678+
"src/KeyLoad.Core/SubscriptionGroups.cs",
679+
"src/KeyLoad.Core/MutationAuthorization.cs",
667680
"tests/KeyLoad.UnitTests/MessagingTests.cs",
668-
"tests/KeyLoad.UnitTests/SecurityAndQueryTests.cs"
681+
"tests/KeyLoad.UnitTests/SecurityAndQueryTests.cs",
682+
"tests/KeyLoad.RecoveryTests/SubscriptionRecoveryTests.cs"
669683
]
670684
},
671685
"KL-091": {
@@ -708,7 +722,9 @@
708722
"status": "in_progress",
709723
"evidence": [
710724
"src/KeyLoad.Security/AuthorizationPolicy.cs",
725+
"src/KeyLoad.Core/MutationAuthorization.cs",
711726
"tests/KeyLoad.UnitTests/SecurityAndQueryTests.cs",
727+
"tests/KeyLoad.UnitTests/SubscriptionTests.cs",
712728
"tests/KeyLoad.RecoveryTests/PeerSecurityTests.cs"
713729
]
714730
},
@@ -722,6 +738,7 @@
722738
"status": "in_progress",
723739
"evidence": [
724740
"src/KeyLoad.Storage.ZoneTree/ZoneTreeStore.cs",
741+
"src/KeyLoad.Core/EventSources.cs",
725742
"src/KeyLoad.Artifacts/BackupArtifact.cs",
726743
"tests/KeyLoad.UnitTests/ArtifactTests.cs"
727744
]
@@ -732,7 +749,8 @@
732749
"evidence": [
733750
"src/KeyLoad.Replication/ReplicatedStateMachine.cs",
734751
"tests/KeyLoad.IntegrationTests/ClusterTests.cs",
735-
"tests/KeyLoad.RecoveryTests/RecoveryTests.cs"
752+
"tests/KeyLoad.RecoveryTests/RecoveryTests.cs",
753+
"tests/KeyLoad.RecoveryTests/SubscriptionRecoveryTests.cs"
736754
]
737755
},
738756
"KL-100": {

0 commit comments

Comments
 (0)