Skip to content

Commit 4dd4369

Browse files
committed
Add committed projection outbox and protected live query feeds
1 parent 76a99d3 commit 4dd4369

29 files changed

Lines changed: 1312 additions & 88 deletions

‎README.md‎

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,9 +8,9 @@ The original load-testing prototype has been replaced. This repository implement
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, 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.
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, a committed projection outbox with generation pins, protected resumable document changes and scalar live queries, 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, 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.
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, automatic retention and generation cleanup, 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

@@ -151,6 +151,26 @@ Queries require a matching point/equality index or explicit `AllowFullScan`. Sca
151151

152152
Search accepts typed vector spaces and explicit text/vector fields. Both branches use one authorized read cut. Exact vector scores and BM25 ranks are combined with weighted RRF using one-based ranks. The managed ANN and graph retrieval extensions are tracked separately.
153153

154+
## Change feeds and live queries
155+
156+
Every successful mutation appends to a private, per-partition system outbox in the same transaction as its canonical effects and outcome. This is separate from business event streams and retained topics. Public `ReadChangesAsync` returns projected document before/after images and deletion metadata; it requires `ChangesRead` and `DocumentsRead`. Current principal, row visibility and sensitive-field policies are checked on every page.
157+
158+
```csharp
159+
var liveQuery = KeyLoadQuery<Order>.From(partition, "orders")
160+
.Where(order => order.Status == "new").Take(100).ToRequest(allowFullScan: true);
161+
var initial = await client.StartLiveQueryAsync(new(liveQuery));
162+
initial.ThrowIfFail();
163+
var next = await client.ReadLiveQueryAsync(new(liveQuery, initial.Value.Cursor));
164+
next.ThrowIfFail();
165+
// Apply Upsert/Remove by canonical document ID; retain next.Value.Cursor after applying its page.
166+
```
167+
168+
The initial complete snapshot and change cursor share one read gate. Live predicates use Q1's same validator, field-use binder and scalar evaluator. Deltas are unordered result-set upserts/removals; ranking, top-k, graph and ANN subscriptions are separate profiles. The initial result must fit its row and byte budgets. Delta pages have examined-position and output-byte budgets; the caller bounds its maintained result set. Re-reading a cursor can repeat changes, so apply them by ID/revision. Continue while `HasMore` even if a filtered page is empty.
169+
170+
Cursors bind incarnation, tenant/partition, collection, principal/policy and schema; live cursors also bind the normalized query. A row ACL change requires a fresh snapshot. Lost retention returns `HistoryUnavailable`; invalid policy/scope returns `TokenInvalidated`. Clear the old result set and take a new snapshot in either case. These cursors can continue on another caught-up RF3 voter.
171+
172+
System projection APIs require cluster administration. A consumer defines its index generation and resource/mutation filter. `ReadProjectionAsync` issues a signed contiguous batch; `CommitProjectionAsync` atomically stores same-partition effects, a replay receipt and its checkpoint. Failed effects leave the checkpoint unchanged. A released generation rejects old batch tokens. Active consumer/rebuild checkpoints pin history against `PurgeOutboxAsync`; public stateless cursors do not pin it. Quotas stop producers before partial publication. Projection effects also consume that quota, so capacity must be available before processing; reserved progress capacity and automatic cleanup are still part of governor qualification. See the [change-feed and outbox contract](docs/design/change-feeds.md).
173+
154174
## Backups and Cartograph
155175

156176
The canonical backup includes the checksummed redo journal, database identity and a SHA-256 manifest. Domain data, schemas, credentials, outcomes and inbox receipts are journaled together. The journal can begin with a verified checkpoint followed by newer transaction frames. ZoneTree files can be rebuilt from that canonical history. Restore validates every manifest file, creates a new incarnation, resets consensus routing metadata and leaves queue dispatch paused.
@@ -174,12 +194,12 @@ dotnet run --project src/KeyLoad.Cli -- copy-artifact backups/snapshot.ctg archi
174194
| Project | Responsibility |
175195
| --- | --- |
176196
| `KeyLoad.Abstractions` | Typed protocol, identities, limits, errors and ordered key codec |
177-
| `KeyLoad.Core` | Catalog, atomic mutation compilation, documents, streams, queues, graph and samples |
197+
| `KeyLoad.Core` | Catalog, atomic mutation compilation, documents, outbox/change feeds, streams, queues, graph and samples |
178198
| `KeyLoad.Storage.ZoneTree` | File ownership, redo recovery, consistent apply gate and verified backups |
179199
| `KeyLoad.Replication` | Durable Raft append barrier, state-machine apply and bounded writer coordination |
180200
| `KeyLoad.Orleans` | Physical-shard command facade and consensus-backed membership |
181201
| `KeyLoad.Security` | Scope, row and sensitive-field policies |
182-
| `KeyLoad.Query` | Bounded SQL, exact vectors, BM25 and fusion |
202+
| `KeyLoad.Query` | Shared SQL/JSON/C# queries, scalar live queries, exact vectors, BM25 and fusion |
183203
| `KeyLoad.Server` | HTTP API, authenticated peer transport and bootstrap |
184204
| `KeyLoad.Client` / `KeyLoad.Cli` | .NET SDK and administrative commands |
185205
| `KeyLoad.Artifacts` | Optional Cartograph archive and ManagedCode storage transport |
@@ -195,7 +215,7 @@ dotnet test --project tests/KeyLoad.RecoveryTests --no-build --no-restore
195215
dotnet test --project tests/KeyLoad.IntegrationTests --no-build --no-restore
196216
```
197217

198-
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.
218+
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, subscription effects/inbox/checkpoint publication and projection effects/outbox/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, live-query continuation and topic delivery on surviving voters, reject a minority write and erase/restart one replica to require snapshot catch-up of documents, outbox entries, generation checkpoints and inbox outcomes. No manually running AppHost is needed for tests.
199219

200220
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.
201221

‎docs/design/change-feeds.md‎

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
# Committed outbox, document changes and scalar live queries
2+
3+
The system outbox is a logical mutation stream inside one atomic partition. It is separate from the private redo journal, business event streams, retained topics and queue delivery. Each entry stores a monotonically increasing sequence, its command-local ordinal, commit token, evaluation time, accepted mutation and receipt. Document mutations also retain canonical before/after images. The outbox entry, strict indexes, domain effect, command outcome and apply watermark share one compiled redo frame. Failed batches publish no entries; a same-command retry does not append again.
4+
5+
Old stores begin collecting entries when running this kernel; existing documents are included by a live-query snapshot rather than synthesized into past history. Entries and consumer checkpoints are canonical backup/snapshot state. Reopening, compaction and RF3 installation preserve them; administrative restore changes the incarnation and invalidates old tokens.
6+
7+
## System projection lifecycle
8+
9+
These APIs require a current cluster administrator. Raw mutations and historical images are confined to this boundary; public feeds expose only authorized projections.
10+
11+
| SDK/API | Contract |
12+
| --- | --- |
13+
| `ConfigureProjectionAsync` / `admin/projections/configure` | Create an immutable named consumer/index generation with resource and mutation-kind filters. An empty filter accepts every resource/kind. `StartAfter` must be a retained cut; the default starts at the retained beginning. Register a rebuild consumer before reclaiming its required tail. |
14+
| `ReadProjectionAsync` / `admin/projections/read` | Read from its persisted checkpoint and issue a signed contiguous batch. Filtered positions advance the batch's through-position. No checkpoint changes until commit. |
15+
| `CommitProjectionAsync` / `admin/projections/commit` | Apply canonical effects within the same atomic partition and persist the checkpoint and replay receipt together. The batch token binds incarnation, consumer, generation and exact start/end positions. An empty tail batch cannot produce effects. |
16+
| `ReleaseProjectionAsync` / `admin/projections/release` | Release that generation's retention pin and fence old tokens. Use a new consumer identity for a new filter, generation or rebuild. |
17+
| `OutboxStatusAsync` / `admin/outbox/status` | Inspect the retained range, byte/record counters and consumer checkpoints. |
18+
| `PurgeOutboxAsync` / `admin/outbox/purge` | Reclaim a bounded contiguous prefix, at most through the lowest active checkpoint. Released consumers no longer pin history. |
19+
20+
Effects failing a CAS, permission or quota check roll back both effects and checkpoint. A different command retrying a completed batch with the same effects returns its original receipt and `AlreadyProcessed`; changed effects conflict. Current generation and effect permissions are checked before returning a cached result. A concurrent batch based on a checkpoint which has advanced must reread. Checkpoints cannot be supplied as arbitrary unauthenticated offsets.
21+
22+
The default outbox quota is 100,000 entries and 1 GiB per atomic partition; 64 named projection consumers are allowed. Projection read budgets default to 100 examined entries and 4 MiB of selected entry payload, with maxima of the configured result count and 16 MiB. Each immutable entry is loaded separately and an undelivered entry remains after the returned through-position. A first entry larger than the caller's byte budget returns `BudgetExceeded`. Canonical projection effects append their own outbox entries; filters should exclude derived targets to avoid feedback. Effects require available outbox quota. Reserved progress capacity, automatic retention, expired receipts/generation cleanup and external-file atomic manifests remain qualification work.
23+
24+
External files must use a replay-safe journal or atomic generation manifest before advancing this checkpoint. An HTTP checkpoint commit alone does not establish external-file durability. The current atomic effects path demonstrates canonical document projections, rather than claiming a qualified ANN/FTS provider lifecycle.
25+
26+
## Protected document feed
27+
28+
`POST /v1/changes/read` (`ReadChangesAsync`) requires `ChangesRead | DocumentsRead` on its collection. `Beginning` starts at the currently retained first position; `Now` captures the current tail. Every page, including an empty tail, returns a signed cursor. This profile polls bounded pages and creates no server subscription, lease or public retention pin.
29+
30+
A returned change contains the sequence/commit, entity reference, revision, deletion flag and safely projected before/after images. Before images require their historical row ACL; all images also require current row visibility and the after-image ACL. Deleted after images carry metadata and no document payload. Sensitive field classification uses the same secure projector as ordinary document reads. No raw outbox payload or unrelated resource identity is returned.
31+
32+
Positions are shared by mutations in the atomic partition, so filtered sequences can contain gaps in the public output. `Limit` bounds examined positions, rather than promising that many visible changes. Byte limits bound the selected change payloads. Empty filtered pages with `HasMore` must still be continued. Saving a cursor after applying its page supports at-least-once reconnect within the retained window; saving it first can lose application delivery. Consumers apply ID/revision updates idempotently.
33+
34+
The cursor binds incarnation, partition, collection, principal, policy epoch, schema and row-visibility epoch. A document ACL change advances that epoch in the same atomic write and rejects old cursors with `TokenInvalidated`, requiring a new authorized snapshot. Reauthorization precedes reading. Principal revocation rejects the call; a policy update fences old cursors. Prefix reclamation past the cursor yields `HistoryUnavailable`. Both conditions require clearing the old projection and resynchronizing. Cursors expire after 24 hours and can move between caught-up voters with the same cluster identity.
35+
36+
## Scalar live-query profile
37+
38+
`POST /v1/query/live/start` captures the complete initial Q1 result and outbox tail under one storage read gate. It requires query/document/change-read grants and ordinary predicate field-use grants. The query must be unordered, scalar and read-only, with no query continuation or `EXPLAIN`. A result larger than `Limit` or its byte budget fails rather than returning an incomplete subscription baseline.
39+
40+
`POST /v1/query/live/read` binds the same normalized AST and evaluates each authorized before/after image with the ordinary Q1 evaluator. A matching after image emits `Upsert` using the ordinary query projector; an image leaving the predicate or being deleted emits `Remove` without a payload. Nonmatching changes still advance the cursor. Protected values can be used by a granted predicate while remaining omitted in returned rows.
41+
42+
The profile exposes a complete bounded initial snapshot and bounded delta pages. It retains no server result-set state; the application bounds its maintained set and polling rate. `ORDER BY`, top-k, aggregates, joins, graph and ANN subscriptions need separate cost/freshness profiles. Live-query calls share ordinary query admission, validation, deadline, candidate-byte and field-lineage controls. Starting at a different predicate/projection or changing policy/row visibility invalidates the cursor. Query snapshot plus tail removes the initial subscription race; it does not provide indefinite history retention or cross-partition coverage.
43+
44+
Unit tests compare seeded live deltas with ordinary Q1 results, exercise concurrent snapshot/write, protected input grants, row/policy fencing, byte boundaries and replay. Real process-kill tests cover projection effects/outbox/receipt/outcome/checkpoint publication. RF3 scenarios cover continuation after leader loss and checkpoint/receipt/outbox recovery after native snapshot catch-up. Power-loss, external manifests, multi-shard coverage and endurance remain distinct qualification gates.

‎docs/design/query-q1.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ Fields use bounded JSON pointers; `/@id` and `/@revision` address canonical meta
4646
| `LIMIT` | Positive result count within configured bounds. | Ordinary query permission. |
4747
| `EXPLAIN` | Reports the authorized access path, atomic partition and scan budget. | The same binding permissions as execution. |
4848

49-
Full scans require explicit opt-in. Candidate count, result bytes, request bytes, predicate depth/nodes, parameter count, projection count, sort keys, execution time and simultaneous queries are bounded. This profile still requires the broader tenant memory/CPU governor and batch executor qualification recorded in the implementation tracker.
49+
Full scans require explicit opt-in. Candidate count and loaded record bytes, serialized result bytes (including identity/redaction metadata), request bytes, predicate depth/nodes, parameter count, projection count, sort keys, execution time and simultaneous queries are bounded. Index document dereferences charge their bytes before the complete candidate set is materialized. The manifest reports `MaxCandidateBytes` and the `Q1`, `documentChangeFeed` and `scalarLiveQuery` read profiles. This profile still requires the broader tenant memory/CPU governor and batch executor qualification recorded in the implementation tracker.
5050

5151
Each page uses one consistent storage read gate. Its cursor is signed and binds the normalized query, principal/policy epoch, schema, node/read generation, original cut and collection data version. Document changes atomically advance that version, including ACL changes. Unrelated catalog or collection writes preserve the cursor; a change to its source rejects continuation with `CursorExpired`. Principal revocation/policy change is rechecked before every page. Replica installation changes the read generation; compaction preserves it. A cursor is a short-lived continuation contract, not an indefinitely retained MVCC snapshot.
5252

0 commit comments

Comments
 (0)