Skip to content

Commit 4d0c3a7

Browse files
committed
Bound command admission by bytes and scope with reserved control work
1 parent c226ee0 commit 4d0c3a7

18 files changed

Lines changed: 442 additions & 24 deletions

File tree

‎README.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,8 @@ The standalone server requires explicit cluster identity, voter endpoints, share
3535

3636
The default peer connection/RPC/request timeouts are 500/1500/2000 milliseconds, below the 4000–8000 millisecond election range. These values are configurable through the corresponding `KeyLoad` options; validation requires that ordering. Native snapshots are produced every `KeyLoad:SnapshotThreshold` committed entries (default 1024). Small snapshot catch-up is tested; large transfers still need qualification against the configured RPC deadline.
3737

38+
Command admission bounds queued and active commands by node count/bytes and verified tenant/principal counts. The default data lane has 256 slots and 128 MiB of retained payload accounting; ACK/renew, membership and dispatch commands have a separate bounded control reserve. A full lane returns `ResourceExhausted` before this attempt's Raft acceptance. Configure `KeyLoad:CommandAdmission` and inspect `AdmissionStatusAsync` as a cluster administrator. See the [admission contract](docs/design/command-admission.md) for defaults, scheduling and the remaining memory/disk qualification work.
39+
3840
## .NET client
3941

4042
The SDK uses `ManagedCode.Communication.Result<T>` and typed protocol records. Supply an API key in application configuration. Keep command IDs stable across retries: an interrupted write response has an unknown outcome.

‎docs/design/command-admission.md‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
# Command admission and reserved control work
2+
3+
Each RF3 node admits commands before enqueueing or forwarding them. `CommandAdmissionGovernor` bounds the count and retained payload bytes of queued **and active** commands, with separate tenant and principal counts. The scope comes from the node's verified principal catalog, rather than tenant or role claims in the submitted JSON. Multiple API keys for the same principal share its budget.
4+
5+
The default data lane permits 256 commands and 128 MiB of retained payload accounting, with caps of 128 per tenant and 64 per principal. Accounting includes the UTF-16 payload, twice its UTF-8 byte count for serialization/escaping in the Raft envelope and 4 KiB of envelope overhead. This measures the admitted command payload model, not total process RSS or native storage allocations.
6+
7+
Queue ACK/NACK/renew, subscription delivery commands, Orleans membership and dispatch control use an independent reserve: 16 commands and 2 MiB, with eight commands per tenant/principal and a 64 KiB maximum UTF-8 control payload. Processing effects and projection effects use the data lane. Both lanes remain bounded. A single apply worker selects control work first, then selects a waiting data command after at most eight consecutive control commands. Each lane preserves FIFO. An active replication cannot be preempted; the existing replication deadline still applies. Native Raft RPCs use their own transport and do not enter this command queue.
8+
9+
Admission failure returns `ResourceExhausted` before this attempt takes the command ID or appends a Raft entry. An earlier interrupted attempt with that ID may already have committed; continue using the same ID until its durable outcome is resolved. A follower preserves an explicit authenticated leader admission rejection; an interrupted or malformed leader response remains `UnknownWriteOutcome`. Once admitted, cancelling the caller's response does not release the reservation or imply that the command was cancelled. The worker releases it only when replication/forwarding completes or fails. Shutdown rejects queued responses with `UnknownWriteOutcome` and releases their reservations; an active command retains its reservation until its worker exits. Scope counters are removed when their last reservation is released.
10+
11+
Configure a standalone node through `KeyLoad:CommandAdmission:<property>` or equivalent environment variables, such as `KeyLoad__CommandAdmission__MaxRetainedBytes`. The AppHost forwards configured values to its three voters. These **local** admission limits may differ between nodes. Canonical mutation/outbox/transaction limits must remain identical between voters; local throttling never changes the apply decision for an already committed entry.
12+
13+
Cluster administrators can inspect `GET /v1/admin/admission` or `KeyLoadClient.AdmissionStatusAsync`, which returns the configured limits and current data/control usage. The endpoint uses the normal authenticated quorum-read path.
14+
15+
Unit tests cover concurrent count/byte limits, tenant/principal isolation, control reservation, bounded scheduling, cancellation and shutdown. A separate Aspire RF3 test exhausts the entire data-byte budget, verifies that catalog writes remain unpublished, and commits dispatch control under the rejected command IDs. Continued routing and that ID reuse prove rejection before canonical acceptance. The normal RF3 suite independently exercises queue and subscription ACK/processing, leader loss, minority rejection and snapshot catch-up.
16+
17+
Shared query/search working-set admission, public request-body admission before deserialization, native cache/memtable limits, disk-floor throttling, rebuild/compaction reserves, automatic retention and sustained RSS/endurance qualification remain tracked work. This queue governor alone does not qualify those resources.

‎docs/implementation/durability-audit.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,8 @@ The committed system outbox stores accepted mutations and document before/after
2828

2929
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.
3030

31+
The cluster coordinator reserves bounded command count and retained payload bytes before enqueueing, scoped by verified tenant/principal. Short delivery, membership and dispatch commands have a separate control reserve and bounded priority. Accepted commands retain their reservation across caller cancellation until the worker finishes; shutdown releases queued reservations with an unknown-outcome response. Local admission configuration does not alter canonical apply decisions. An RF3 test exhausts data admission and verifies unpublished catalog writes, unclaimed rejected IDs, continued dispatch commits and ready Orleans routing on all voters. Request-body allocations before deserialization and total native/RSS/disk budgets remain separate pending qualification.
32+
3133
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.
3234

3335
## Verified and outstanding qualification

‎docs/implementation/kernel-qualification.json‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
"restore": "locked-mode passed",
88
"build": "passed with zero warnings",
99
"unitTests": {
10-
"passed": 94,
10+
"passed": 104,
1111
"failed": 0
1212
},
1313
"recoveryTests": {
@@ -31,7 +31,7 @@
3131
"evidence": "CI retains artifacts/qualification/crash-trials-*.jsonl"
3232
},
3333
"rf3Integration": {
34-
"passed": 4,
34+
"passed": 5,
3535
"failed": 0,
3636
"scenarios": [
3737
"leader process kill",
@@ -50,7 +50,9 @@
5050
"cross-voter live-query continuation after leader loss",
5151
"generation-pinned outbox effects and checkpoint replay after leader loss",
5252
"outbox entries, projection checkpoints and receipts preserved across snapshot installation",
53-
"protected document change-feed and live-query HTTP projection"
53+
"protected document change-feed and live-query HTTP projection",
54+
"data admission exhaustion rejects before canonical acceptance",
55+
"reserved dispatch control and Orleans routing with exhausted data admission"
5456
]
5557
},
5658
"durability": [
@@ -60,9 +62,9 @@
6062
"powerLoss": "not qualified",
6163
"endurance72Hours": "pending",
6264
"crossPlatformCI": {
63-
"qualifiedCommit": "4dd4369b8cde399dfcbbcc3fc538aec1768dc0d9",
64-
"run": "https://github.com/managedcode/KeyLoad/actions/runs/36911698298",
65+
"qualifiedCommit": "c226ee0c0659be1045b6b151d796cd4c2ae5cf93",
66+
"run": "https://github.com/managedcode/KeyLoad/actions/runs/36915868306",
6567
"platforms": ["Linux", "macOS", "Windows"],
66-
"currentProjectionProgressChanges": "pending"
68+
"currentCommandAdmissionChanges": "pending"
6769
}
6870
}

‎docs/implementation/status.json‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -310,7 +310,13 @@
310310
"src/KeyLoad.Query/QueryEngine.cs",
311311
"src/KeyLoad.Core/ProjectionOutbox.cs",
312312
"tests/KeyLoad.UnitTests/QueryAdapterTests.cs",
313-
"tests/KeyLoad.UnitTests/ProjectionProgressTests.cs"
313+
"tests/KeyLoad.UnitTests/ProjectionProgressTests.cs",
314+
"src/KeyLoad.Core/CommandAdmissionGovernor.cs",
315+
"src/KeyLoad.Core/AdmittedCommandQueue.cs",
316+
"tests/KeyLoad.UnitTests/CommandAdmissionTests.cs",
317+
"tests/KeyLoad.UnitTests/CommandQueueTests.cs",
318+
"tests/KeyLoad.IntegrationTests/AdmissionClusterTests.cs",
319+
"docs/design/command-admission.md"
314320
]
315321
},
316322
"KL-041": {
Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
namespace KeyLoad;
2+
3+
// Node-local admission settings may differ between voters; they never change canonical apply decisions.
4+
public sealed record CommandAdmissionLimits
5+
{
6+
public int MaxCommands { get; init; } = 256;
7+
public long MaxRetainedBytes { get; init; } = 134_217_728;
8+
public int MaxTenantCommands { get; init; } = 128;
9+
public int MaxPrincipalCommands { get; init; } = 64;
10+
public int ReservedControlCommands { get; init; } = 16;
11+
public long ReservedControlBytes { get; init; } = 2_097_152;
12+
public int MaxControlPayloadBytes { get; init; } = 65_536;
13+
public int MaxTenantControlCommands { get; init; } = 8;
14+
public int MaxPrincipalControlCommands { get; init; } = 8;
15+
public void Validate()
16+
{
17+
if (MaxCommands < 1 || MaxRetainedBytes < 1 || MaxTenantCommands < 1 || MaxPrincipalCommands < 1
18+
|| ReservedControlCommands < 1 || ReservedControlBytes < 1 || MaxControlPayloadBytes < 1
19+
|| MaxTenantControlCommands < 1 || MaxPrincipalControlCommands < 1)
20+
throw new ArgumentException("Command admission limits must be positive.");
21+
}
22+
}
23+
public sealed record CommandAdmissionSnapshot(int Commands, long RetainedBytes, int ControlCommands, long ControlRetainedBytes,
24+
int ActiveTenantScopes, int ActivePrincipalScopes);
25+
public sealed record NodeAdmissionStatus(CommandAdmissionLimits Limits, CommandAdmissionSnapshot Usage);

‎src/KeyLoad.AppHost/Program.cs‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,10 @@
4141
foreach (var resource in nodes)
4242
{
4343
resource.WithEnvironment("KeyLoad__PublicEndpoint", resource.GetEndpoint("http"));
44+
foreach (var option in new[] { "MaxCommands", "MaxRetainedBytes", "MaxTenantCommands", "MaxPrincipalCommands", "ReservedControlCommands",
45+
"ReservedControlBytes", "MaxControlPayloadBytes", "MaxTenantControlCommands", "MaxPrincipalControlCommands" })
46+
if (builder.Configuration[$"KeyLoad:CommandAdmission:{option}"] is { } value)
47+
resource.WithEnvironment($"KeyLoad__CommandAdmission__{option}", value);
4448
for (var index = 0; index < nodes.Length; index++) resource.WithEnvironment($"KeyLoad__Peers__{index}", nodes[index].GetEndpoint("http"));
4549
}
4650
if (benchmarkMode) BenchmarkResources.Add(builder, nodes, admin, benchmarkRoot);

‎src/KeyLoad.Client/KeyLoadClient.cs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,4 +102,6 @@ public Task<Result<bool>> ConfigureApiKeyAsync(Guid commandId, ApiKeyRecord key,
102102
=> Send<bool>("/v1/admin/api-keys", new ConfigureApiKeyRequest(key), true, commandId, cancellationToken);
103103
public Task<Result<NodeStatus>> StatusAsync(CancellationToken cancellationToken = default)
104104
=> Send<NodeStatus>("/v1/status", null, false, null, cancellationToken);
105+
public Task<Result<NodeAdmissionStatus>> AdmissionStatusAsync(CancellationToken cancellationToken = default)
106+
=> Send<NodeAdmissionStatus>("/v1/admin/admission", null, false, null, cancellationToken);
105107
}
Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,62 @@
1+
namespace KeyLoad.Core;
2+
3+
// One reader, bounded admission, FIFO within each lane and at most eight control commands before a waiting data command.
4+
public sealed class AdmittedCommandQueue(CommandAdmissionGovernor governor)
5+
{
6+
private readonly object gate = new();
7+
private readonly Queue<Pending> commands = [], controls = [];
8+
private readonly SemaphoreSlim available = new(0);
9+
private int controlBurst;
10+
private bool stopped;
11+
public Pending Enqueue(ReplicatedOperation operation, PrincipalRecord principal, int payloadBytes,
12+
CancellationToken cancellationToken = default)
13+
{
14+
if (principal.Id != operation.PrincipalId)
15+
throw Errors.Fail(ErrorCode.PermissionDenied, "Command admission requires the verified principal identity.");
16+
var lease = governor.Reserve(operation.Kind, principal, payloadBytes, operation.PayloadJson.Length, cancellationToken);
17+
try
18+
{
19+
lock (gate)
20+
{
21+
cancellationToken.ThrowIfCancellationRequested();
22+
if (stopped) throw Errors.Fail(ErrorCode.ResourceExhausted, "The node command queue has stopped accepting operations.");
23+
var pending = new Pending(operation, lease);
24+
(CommandAdmissionGovernor.IsControl(operation.Kind) ? controls : commands).Enqueue(pending);
25+
available.Release(); return pending;
26+
}
27+
}
28+
catch { lease.Dispose(); throw; }
29+
}
30+
public async ValueTask<Pending?> ReadAsync(CancellationToken cancellationToken = default)
31+
{
32+
await available.WaitAsync(cancellationToken).ConfigureAwait(false);
33+
lock (gate)
34+
{
35+
if (controls.Count > 0 && (commands.Count == 0 || controlBurst < 8))
36+
{ controlBurst = Math.Min(controlBurst + 1, 8); return controls.Dequeue(); }
37+
if (commands.Count > 0) { controlBurst = 0; return commands.Dequeue(); }
38+
if (stopped) return null;
39+
throw new InvalidOperationException("The command queue signal is inconsistent.");
40+
}
41+
}
42+
public void Stop()
43+
{
44+
lock (gate)
45+
{
46+
if (stopped) return;
47+
stopped = true;
48+
foreach (var queue in new[] { controls, commands })
49+
while (queue.TryDequeue(out var pending)) pending.Fail(Errors.Fail(ErrorCode.UnknownWriteOutcome,
50+
"The node stopped before returning a command outcome. Query or retry the same command ID."));
51+
available.Release();
52+
}
53+
}
54+
public sealed class Pending(ReplicatedOperation operation, CommandAdmissionGovernor.Lease lease)
55+
{
56+
private readonly TaskCompletionSource<OperationResult> completion = new(TaskCreationOptions.RunContinuationsAsynchronously);
57+
public ReplicatedOperation Operation { get; } = operation;
58+
public Task<OperationResult> Completion => completion.Task;
59+
public void Complete(OperationResult result) { lease.Dispose(); completion.TrySetResult(result); }
60+
public void Fail(Exception exception) { lease.Dispose(); completion.TrySetException(exception); }
61+
}
62+
}
Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
namespace KeyLoad.Core;
2+
3+
public sealed class CommandAdmissionGovernor
4+
{
5+
private readonly object gate = new();
6+
private readonly Dictionary<(bool Control, string Id), int> tenants = [];
7+
private readonly Dictionary<(bool Control, string Id), int> principals = [];
8+
private int commands, controlCommands;
9+
private long retainedBytes, controlRetainedBytes;
10+
public CommandAdmissionLimits Limits { get; }
11+
public CommandAdmissionGovernor(CommandAdmissionLimits? limits = null)
12+
{
13+
Limits = limits ?? new(); Limits.Validate();
14+
}
15+
public static bool IsControl(OperationKind kind)
16+
=> kind is OperationKind.Delivery or OperationKind.SubscriptionDelivery or OperationKind.Membership or OperationKind.SetDispatch;
17+
public Lease Reserve(OperationKind kind, PrincipalRecord principal, int payloadBytes, int payloadCharacters,
18+
CancellationToken cancellationToken = default)
19+
{
20+
cancellationToken.ThrowIfCancellationRequested();
21+
if (payloadBytes < 0 || payloadCharacters < 0)
22+
throw new ArgumentOutOfRangeException(nameof(payloadBytes));
23+
// JSON embedded in the Raft envelope escapes quotes/backslashes again; reserve two UTF-8 copies.
24+
var bytes = checked(payloadCharacters * 2L + payloadBytes * 2L + 4_096);
25+
var control = IsControl(kind);
26+
if (control && payloadBytes > Limits.MaxControlPayloadBytes)
27+
throw Errors.Fail(ErrorCode.ResourceExhausted, "The control command exceeds its reserved byte budget.");
28+
var tenantKey = (control, principal.TenantId); var principalKey = (control, principal.Id);
29+
lock (gate)
30+
{
31+
cancellationToken.ThrowIfCancellationRequested();
32+
var count = control ? controlCommands : commands;
33+
var used = control ? controlRetainedBytes : retainedBytes;
34+
var maxCount = control ? Limits.ReservedControlCommands : Limits.MaxCommands;
35+
var maxBytes = control ? Limits.ReservedControlBytes : Limits.MaxRetainedBytes;
36+
var maxTenant = control ? Limits.MaxTenantControlCommands : Limits.MaxTenantCommands;
37+
var maxPrincipal = control ? Limits.MaxPrincipalControlCommands : Limits.MaxPrincipalCommands;
38+
var tenantCount = tenants.GetValueOrDefault(tenantKey); var principalCount = principals.GetValueOrDefault(principalKey);
39+
if (count >= maxCount || bytes > maxBytes - used || tenantCount >= maxTenant || principalCount >= maxPrincipal)
40+
throw Errors.Fail(ErrorCode.ResourceExhausted, "The node, tenant or principal command admission budget is exhausted.");
41+
if (control) { controlCommands++; controlRetainedBytes += bytes; }
42+
else { commands++; retainedBytes += bytes; }
43+
tenants[tenantKey] = tenantCount + 1; principals[principalKey] = principalCount + 1;
44+
return new(this, control, principal.TenantId, principal.Id, bytes);
45+
}
46+
}
47+
public CommandAdmissionSnapshot Snapshot()
48+
{
49+
lock (gate) return new(commands, retainedBytes, controlCommands, controlRetainedBytes, tenants.Count, principals.Count);
50+
}
51+
private void Release(Lease lease)
52+
{
53+
lock (gate)
54+
{
55+
if (lease.Released) return;
56+
lease.Released = true;
57+
if (lease.Control) { controlCommands--; controlRetainedBytes -= lease.Bytes; }
58+
else { commands--; retainedBytes -= lease.Bytes; }
59+
Decrement(tenants, (lease.Control, lease.Tenant)); Decrement(principals, (lease.Control, lease.Principal));
60+
}
61+
}
62+
private static void Decrement(Dictionary<(bool Control, string Id), int> scopes, (bool Control, string Id) key)
63+
{
64+
if (scopes[key] == 1) scopes.Remove(key); else scopes[key]--;
65+
}
66+
public sealed class Lease : IDisposable
67+
{
68+
private readonly CommandAdmissionGovernor owner;
69+
internal bool Control { get; }
70+
internal string Tenant { get; }
71+
internal string Principal { get; }
72+
internal long Bytes { get; }
73+
internal bool Released { get; set; }
74+
internal Lease(CommandAdmissionGovernor owner, bool control, string tenant, string principal, long bytes)
75+
{ this.owner = owner; Control = control; Tenant = tenant; Principal = principal; Bytes = bytes; }
76+
public void Dispose() => owner.Release(this);
77+
}
78+
}

0 commit comments

Comments
 (0)