From bfc23f93248dfc7a7d9ddf529db6f05cc382a91a Mon Sep 17 00:00:00 2001
From: Zacgoose <107489668+Zacgoose@users.noreply.github.com>
Date: Sun, 13 Sep 2026 00:47:19 +0800
Subject: [PATCH 1/2] feat(hosting): compress /api responses and bill
compressed egress bytes
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Make dynamic /api response compression a first-class, independently-toggled
capability, and fix the egress cap to bill the compressed on-the-wire size
instead of the pre-compression body.
Compression
- App:Api:Compression / CRAFT_API_COMPRESSION (default on), resolved in
CraftRoles.ApiCompressionEnabled, independent of the static Frontend.Compression
toggle — an origin behind a CDN that compresses static assets still compresses
its API JSON, and vice versa.
- App:Api:CompressionLevel / CRAFT_API_COMPRESSION_LEVEL (Fastest | Optimal |
SmallestSize | NoCompression) applied to both providers, surfaced on the startup
"[System] Compression" line. Default Optimal: measured on 2 vCPU, Brotli Optimal
~8.4x vs Fastest ~4.9x at the same ~14% CPU and p95, so it is a near-free win.
SmallestSize stays off by default (Brotli q11 pegged both cores, p95 in the tens
of seconds); drop to Fastest/NoCompression on a very small SKU.
- Program.cs splits the response pipeline by path (UseWhen): /api and static
compression are governed by separate toggles, so an /api body is never
double-compressed.
Egress accounting (policy + counting, no compression logic)
- ApiEgressLimiterMiddleware now only classifies and sheds (429), flagging greenlit
API requests via ApiEgressWireCounterMiddleware.ChargeItemKey.
- New ApiEgressWireCounterMiddleware runs outside the response compressor (outermost
/api link) and records the flagged response's wire bytes, so the cap charges what
actually leaves the box. It must sit outside compression because the compressor
only flushes its trailing block on unwind — a counter inside undercounts.
Tests
- ApiCompressionPipelineTests: real Kestrel (dynamic port), six pipeline shapes
(bare, wire counter, UseWhen, routed endpoint, charset, full realistic chain) all
confirm /api responses come back compressed and the counter does not suppress it.
Adds Microsoft.AspNetCore.TestHost as a permanent test dependency.
- Egress limiter/counter unit tests (shed vs flag, never touches UI, counts the
compressed size when compression runs inside it, records only flagged requests,
restores the body on throw). CraftRoles tests for the two toggles' independence
and the level resolver (env override, unrecognised falls back to Fastest).
perf-harness
- run-compression.ps1: correctness + accounting proof (identity uncompressed vs gzip
compressed, ledger delta == bytes received in both) then CPU-vs-bandwidth per
encoding under k6 + docker stats.
- run-compression-levels.ps1: sweeps the level (recreating per level) reporting
ratio, CPU and p95 for each level x encoding.
- api_load.js gains an ENC knob (Accept-Encoding request header).
---
Services/Configuration/ApiSettings.cs | 46 +++
Services/Configuration/CraftSettings.cs | 4 +
.../Hosting/ApiEgressLimiterMiddleware.cs | 38 +-
.../Hosting/ApiEgressWireCounterMiddleware.cs | 67 ++++
.../Hosting/CraftHostBuilderExtensions.cs | 38 +-
Services/Hosting/CraftRoles.cs | 19 +-
Services/Program.cs | 54 ++-
appsettings.example.jsonc | 19 +
perf-harness/README.md | 40 ++-
perf-harness/docker-compose.api.yml | 3 +
perf-harness/k6/api_load.js | 14 +-
.../scripts/run-compression-levels.ps1 | 183 ++++++++++
perf-harness/scripts/run-compression.ps1 | 326 ++++++++++++++++++
.../ApiCompressionPipelineTests.cs | 185 ++++++++++
.../ApiEgressLimiterMiddlewareTests.cs | 132 ++-----
.../ApiEgressWireCounterMiddlewareTests.cs | 208 +++++++++++
tests/Craft.Tests/Craft.Tests.csproj | 1 +
tests/Craft.Tests/CraftRolesTests.cs | 87 +++++
18 files changed, 1333 insertions(+), 131 deletions(-)
create mode 100644 Services/Configuration/ApiSettings.cs
create mode 100644 Services/Hosting/ApiEgressWireCounterMiddleware.cs
create mode 100644 perf-harness/scripts/run-compression-levels.ps1
create mode 100644 perf-harness/scripts/run-compression.ps1
create mode 100644 tests/Craft.Tests/ApiCompressionPipelineTests.cs
create mode 100644 tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs
diff --git a/Services/Configuration/ApiSettings.cs b/Services/Configuration/ApiSettings.cs
new file mode 100644
index 0000000..7db5f9e
--- /dev/null
+++ b/Services/Configuration/ApiSettings.cs
@@ -0,0 +1,46 @@
+namespace Craft.Configuration;
+
+///
+/// Policy for dynamic /api responses (the PowerShell / native dispatch pipeline), kept separate
+/// from so the API surface can be governed on its own terms rather than
+/// riding on the static-frontend toggles.
+///
+public class ApiSettings
+{
+ ///
+ /// Whether the host compresses dynamic /api responses on the fly (Brotli preferred, gzip
+ /// fallback, negotiated from the caller's Accept-Encoding). Default true.
+ ///
+ /// Independent of , which governs static assets:
+ /// turning static compression off (for example because an upstream CDN already compresses the
+ /// static bundle) does NOT turn this off, and vice versa. That separation is the whole point — a
+ /// CDN in front of the origin typically does not re-compress API JSON, so the origin should keep
+ /// doing it even when static compression is delegated to the edge.
+ ///
+ ///
+ /// This is a server capability, not a mandate: a response is only compressed when the caller
+ /// advertises an accepted encoding; a client that sends none is served identity. Overridable via
+ /// the CRAFT_API_COMPRESSION environment variable (true/false), which wins over this setting.
+ ///
+ ///
+ public bool Compression { get; set; } = true;
+
+ ///
+ /// On-the-fly compression level for the response compressors (Brotli and gzip): one of
+ /// Fastest, Optimal, SmallestSize, or NoCompression (a
+ /// name). Default Optimal.
+ ///
+ /// Optimal is the default because it is a near-free win over Fastest: measured on a 2-vCPU container
+ /// (PerfJson payload), Brotli Optimal compressed ~8.4x versus Fastest's ~4.9x at the same CPU
+ /// (~14%) and latency, and gzip likewise improved at the same cost. SmallestSize is
+ /// deliberately NOT the default: Brotli's SmallestSize is quality 11, which on a small shared-core
+ /// container pegged both cores (~180% CPU) and drove p95 into the tens of seconds — never use it for
+ /// dynamic /api. Drop to Fastest (or NoCompression) on a very small SKU if the
+ /// compressor is seen competing with the PowerShell worker pool. Applies to on-the-fly compression
+ /// generally (dynamic /api and the static fallback for assets without a precompressed sibling;
+ /// precompressed .br/.gz siblings are built ahead of time and unaffected). Overridable
+ /// via CRAFT_API_COMPRESSION_LEVEL. An unrecognised value falls back to Fastest.
+ ///
+ ///
+ public string CompressionLevel { get; set; } = "Optimal";
+}
diff --git a/Services/Configuration/CraftSettings.cs b/Services/Configuration/CraftSettings.cs
index 897d0ae..dd93f3a 100644
--- a/Services/Configuration/CraftSettings.cs
+++ b/Services/Configuration/CraftSettings.cs
@@ -79,6 +79,10 @@ public class CraftSettings
/// Frontend serving policy — CSP header injection (EasyAuth handles auth/redirects).
public FrontendSettings Frontend { get; set; } = new();
+ /// Dynamic /api response policy (compression), independent of the static frontend.
+ /// See .
+ public ApiSettings Api { get; set; } = new();
+
///
/// Deployment roles (capabilities) — which parts of the host this process serves. One image, three
/// independent switches. See . Also settable via the CRAFT_SERVE_FRONTEND /
diff --git a/Services/Hosting/ApiEgressLimiterMiddleware.cs b/Services/Hosting/ApiEgressLimiterMiddleware.cs
index 4f37b8d..ca66a21 100644
--- a/Services/Hosting/ApiEgressLimiterMiddleware.cs
+++ b/Services/Hosting/ApiEgressLimiterMiddleware.cs
@@ -3,15 +3,24 @@
namespace Craft.Hosting;
///
-/// Sheds app-only API traffic once the instance has served its daily egress budget, and records the
-/// outbound bytes of the responses it lets through. Runs only for API clients (client-credentials
-/// callers); every interactive request passes straight through after a single header check.
+/// Sheds app-only API traffic once the instance has served its daily egress budget, and marks the
+/// requests it lets through for billing. Runs only for API clients (client-credentials callers); every
+/// interactive request passes straight through after a single header check.
+///
+/// This is the policy half of the egress feature: it classifies the caller and decides whether
+/// to reject. It no longer counts bytes itself — the outbound size is measured post-compression by
+/// , which runs outside the response compressor. When this
+/// middleware lets an API request through it sets
+/// on the request, and the wire counter records the response's on-the-wire bytes for exactly those
+/// requests. Splitting it this way keeps the cap billing what actually transits the network (the
+/// compressed body) while the shed decision stays here, after auth, where the caller is known.
+///
///
/// Registered only when egress accounting is enabled (hosted env, or forced) — see
/// CraftHostBuilderExtensions.AddCraftEgressLimiter. Placement mirrors the rate limiter: after
/// the auth middleware (so an app-only caller's AppId is resolved) and after static file serving (so a
/// page load's assets are never charged). Enforcement (the 429) only bites once a budget is configured;
-/// with no budget it counts silently, which is the accounting-only rollout phase.
+/// with no budget it flags silently, which is the accounting-only rollout phase.
///
///
public sealed class ApiEgressLimiterMiddleware
@@ -50,21 +59,12 @@ public async Task InvokeAsync(HttpContext context)
return;
}
- // Count the body this request writes. Swapping Response.Body captures both a direct
- // Body.WriteAsync and Response.WriteAsync(string), since the response writer is re-adapted onto
- // our stream — see CountingStream.
- var original = context.Response.Body;
- var counting = new CountingStream(original);
- context.Response.Body = counting;
- try
- {
- await _next(context);
- }
- finally
- {
- context.Response.Body = original;
- _ledger.Record(counting.BytesWritten);
- }
+ // Greenlit: mark the request so ApiEgressWireCounterMiddleware (running outside the response
+ // compressor) bills its on-the-wire bytes. We don't count here — a counter at this position
+ // would see the pre-compression body and miss the compressor's final flush, which unwinds
+ // further out. The shed body above is deliberately left unflagged, so it is never charged.
+ context.Items[ApiEgressWireCounterMiddleware.ChargeItemKey] = true;
+ await _next(context);
}
private async Task RejectAsync(HttpContext context)
diff --git a/Services/Hosting/ApiEgressWireCounterMiddleware.cs b/Services/Hosting/ApiEgressWireCounterMiddleware.cs
new file mode 100644
index 0000000..41f857b
--- /dev/null
+++ b/Services/Hosting/ApiEgressWireCounterMiddleware.cs
@@ -0,0 +1,67 @@
+namespace Craft.Hosting;
+
+///
+/// Measures the on-the-wire size of each /api response and, for the requests the egress limiter
+/// greenlit, records it into the . It is the accounting half of the egress
+/// feature and holds no compression logic of its own — it simply counts whatever bytes leave the
+/// box, whether that is identity or a Brotli/gzip body produced by the generic /api response
+/// compression that runs just inside it.
+///
+/// Why a second middleware, and why here. The egress cap is meant to bill the bytes that
+/// actually transit the network, so the counter has to sit outside compression: response
+/// compression only emits its trailing block when its middleware unwinds, so a counter placed inside it
+/// would both count the pre-compression body and miss the final flush. This middleware is therefore
+/// registered as the outermost link of the /api pipeline — ahead of
+/// UseResponseCompression — so its wraps the socket and sees the
+/// compressed output, and its finally runs only after compression has fully flushed.
+///
+///
+/// The decision of which requests to bill stays entirely in ,
+/// which runs later (after auth, so it can classify the caller and shed when over budget). When it lets
+/// an API request through it sets on the request; this middleware records
+/// the wire bytes only when that flag is present, so UI traffic, anonymous requests, static assets and
+/// shed 429s are never charged. Registered only when egress accounting is enabled — see
+/// CraftHostBuilderExtensions.AddCraftEgressLimiter and the /api pipeline in Program.cs.
+///
+///
+public sealed class ApiEgressWireCounterMiddleware
+{
+ ///
+ /// Request item set by on an API request it lets through,
+ /// signalling this middleware to bill the response's wire bytes. Absent for UI, anonymous, static
+ /// and shed requests, which are never charged.
+ ///
+ public const string ChargeItemKey = "Craft.Egress.Charge";
+
+ private readonly RequestDelegate _next;
+ private readonly EgressLedger _ledger;
+
+ public ApiEgressWireCounterMiddleware(RequestDelegate next, EgressLedger ledger)
+ {
+ _next = next ?? throw new ArgumentNullException(nameof(next));
+ _ledger = ledger ?? throw new ArgumentNullException(nameof(ledger));
+ }
+
+ public async Task InvokeAsync(HttpContext context)
+ {
+ ArgumentNullException.ThrowIfNull(context);
+
+ // Swapping Response.Body captures every write path (Body.WriteAsync, Response.WriteAsync, the
+ // BodyWriter pipe) — see CountingStream. Restored in the finally, whether the pipeline completed
+ // or threw. Only requests the limiter flagged are recorded; the flag is set downstream (inner)
+ // and read here after the whole pipeline — including compression's flush — has unwound.
+ var original = context.Response.Body;
+ var counting = new CountingStream(original);
+ context.Response.Body = counting;
+ try
+ {
+ await _next(context);
+ }
+ finally
+ {
+ context.Response.Body = original;
+ if (context.Items.TryGetValue(ChargeItemKey, out var charge) && charge is true)
+ _ledger.Record(counting.BytesWritten);
+ }
+ }
+}
diff --git a/Services/Hosting/CraftHostBuilderExtensions.cs b/Services/Hosting/CraftHostBuilderExtensions.cs
index 28222ff..48531b6 100644
--- a/Services/Hosting/CraftHostBuilderExtensions.cs
+++ b/Services/Hosting/CraftHostBuilderExtensions.cs
@@ -185,8 +185,38 @@ public static LogLevel AddCraftLogging(this WebApplicationBuilder builder)
private static readonly string[] second = new[] { "application/json", "text/json", "application/javascript", "text/javascript" };
+ ///
+ /// Resolves the on-the-fly compression level from the CRAFT_API_COMPRESSION_LEVEL env var
+ /// (which wins) or App:Api:CompressionLevel, parsed as a
+ /// name (case-insensitive). An unset or unrecognised value is
+ /// — the safe default on a small, shared-core container.
+ ///
+ public static CompressionLevel ResolveCompressionLevel(CraftSettings settings, Func env)
+ {
+ ArgumentNullException.ThrowIfNull(settings);
+ ArgumentNullException.ThrowIfNull(env);
+
+ var raw = env("CRAFT_API_COMPRESSION_LEVEL");
+ if (string.IsNullOrWhiteSpace(raw)) raw = settings.Api.CompressionLevel;
+
+ return Enum.TryParse(raw, ignoreCase: true, out var level)
+ ? level
+ : CompressionLevel.Fastest;
+ }
+
+ /// Convenience overload resolving against the real process environment.
+ public static CompressionLevel ResolveCompressionLevel(CraftSettings settings) =>
+ ResolveCompressionLevel(settings, Environment.GetEnvironmentVariable);
+
/// Response compression, matching Azure Static Web Apps behaviour.
- public static IServiceCollection AddCraftResponseCompression(this IServiceCollection services)
+ ///
+ /// Compression level applied to both the Brotli and gzip providers. Defaults to
+ /// — deliberate on a small container, where the request-path
+ /// CPU of Optimal/SmallestSize usually costs more than the bytes it saves. Resolve it from config
+ /// with .
+ ///
+ public static IServiceCollection AddCraftResponseCompression(
+ this IServiceCollection services, CompressionLevel level = CompressionLevel.Fastest)
{
ArgumentNullException.ThrowIfNull(services);
@@ -199,10 +229,8 @@ public static IServiceCollection AddCraftResponseCompression(this IServiceCollec
second);
});
- // Fastest, not Optimal: these run on the request path on a small container, where the extra
- // CPU costs more than the bytes saved. Precompressed .br/.gz siblings cover the static assets.
- services.Configure(o => o.Level = CompressionLevel.Fastest);
- services.Configure(o => o.Level = CompressionLevel.Fastest);
+ services.Configure(o => o.Level = level);
+ services.Configure(o => o.Level = level);
return services;
}
diff --git a/Services/Hosting/CraftRoles.cs b/Services/Hosting/CraftRoles.cs
index c44f186..03c8049 100644
--- a/Services/Hosting/CraftRoles.cs
+++ b/Services/Hosting/CraftRoles.cs
@@ -24,7 +24,7 @@ public sealed class CraftRoles
{
private CraftRoles(bool frontend, bool http, bool background,
bool responseCacheEnabled, bool healthEnabled, string healthPath,
- bool compressionEnabled)
+ bool compressionEnabled, bool apiCompressionEnabled)
{
Frontend = frontend;
Http = http;
@@ -33,6 +33,7 @@ private CraftRoles(bool frontend, bool http, bool background,
HealthEnabled = healthEnabled;
HealthPath = healthPath;
CompressionEnabled = compressionEnabled;
+ ApiCompressionEnabled = apiCompressionEnabled;
}
/// Serve static web content from Frontend/.
@@ -67,10 +68,20 @@ private CraftRoles(bool frontend, bool http, bool background,
///
/// When false the host serves all static content raw/identity: precompressed .br/.gz
- /// siblings are not served and on-the-fly compression is not applied.
+ /// siblings are not served and on-the-fly compression is not applied. Governs static content
+ /// only — dynamic /api compression is .
///
public bool CompressionEnabled { get; }
+ ///
+ /// Whether dynamic /api responses are compressed on the fly (Brotli/gzip, negotiated from the
+ /// caller's Accept-Encoding). Default true, and deliberately independent of
+ /// : an origin behind a CDN that compresses the static bundle should
+ /// still compress its API JSON, so the two toggles do not share a switch. Env override
+ /// CRAFT_API_COMPRESSION wins over App:Api:Compression.
+ ///
+ public bool ApiCompressionEnabled { get; }
+
///
/// Resolves roles and derived toggles from configuration and the environment.
///
@@ -115,9 +126,11 @@ public static CraftRoles Resolve(CraftSettings settings, Func e
if (!healthPath.StartsWith('/')) healthPath = "/" + healthPath;
var compressionEnabled = EnvFlag.Read(env, "CRAFT_COMPRESSION") ?? settings.Frontend.Compression;
+ var apiCompressionEnabled = EnvFlag.Read(env, "CRAFT_API_COMPRESSION") ?? settings.Api.Compression;
return new CraftRoles(frontend, http, background,
- cacheEnabled, healthEnabled, healthPath, compressionEnabled);
+ cacheEnabled, healthEnabled, healthPath,
+ compressionEnabled, apiCompressionEnabled);
}
/// Convenience overload reading the real process environment.
diff --git a/Services/Program.cs b/Services/Program.cs
index 2f2c7fb..0263819 100644
--- a/Services/Program.cs
+++ b/Services/Program.cs
@@ -56,6 +56,7 @@
var healthEnabled = roles.HealthEnabled;
var healthPath = roles.HealthPath;
var compressionEnabled = roles.CompressionEnabled;
+var apiCompressionEnabled = roles.ApiCompressionEnabled;
// Kestrel limits, logging sinks, compression, the service graph and the rate limiter all live in
// Services/Hosting/CraftHostBuilderExtensions.cs.
@@ -64,7 +65,8 @@
// The resolved level is logged at startup and also gates PowerShell stream capture.
var configuredLogLevel = builder.AddCraftLogging();
-builder.Services.AddCraftResponseCompression();
+var compressionLevel = CraftHostBuilderExtensions.ResolveCompressionLevel(craftSettings);
+builder.Services.AddCraftResponseCompression(compressionLevel);
// ── Native C# endpoints and scheduled tasks ────────────────────────────────────────────────────────
// Discovered before the container is built so the endpoint/task types, the central handler and any
// application service module they ship can be registered into it. Costs nothing when no assemblies
@@ -269,11 +271,47 @@ void RunInitialization()
string.Join(", ", CraftSettings.Auth.DevRoles));
}
-// Response compression must be before static files. Skipped entirely when compression is disabled
-// (App:Frontend:Compression=false / CRAFT_COMPRESSION=false) so everything is served raw/identity.
+// Response body pipeline, split by path so /api and static compression are governed independently and
+// the egress cap can bill the size that actually leaves the box. Both must sit before static file
+// serving (a compressor has to wrap the static middleware to compress its output).
+//
+// /api : [egress wire counter?] -> [/api response compression?]
+// The counter is OUTSIDE compression on purpose — it measures the compressed, on-the-wire
+// bytes, and its finally runs only after the compressor has flushed. See
+// ApiEgressWireCounterMiddleware. Compression here is generic; it knows nothing of accounting.
+// other : [static/frontend on-the-fly compression?]
+// Precompressed .br/.gz siblings are served by the static pipeline itself; this is the
+// on-the-fly fallback for static assets without a sibling and for the SPA fallback HTML.
+//
+// A request takes exactly one branch, so an /api body is never re-compressed by the static compressor.
+// /api compression (App:Api:Compression / CRAFT_API_COMPRESSION, default on) is independent of static
+// compression (App:Frontend:Compression / CRAFT_COMPRESSION): a CDN that compresses the static bundle
+// does not re-compress API JSON, so the origin keeps doing it even with static compression delegated.
+var egressAccountingEnabled = CraftSettings.RateLimit.Egress.ResolvedEnabled;
+
+if (egressAccountingEnabled || apiCompressionEnabled)
+{
+ app.UseWhen(IsApiPath, api =>
+ {
+ if (egressAccountingEnabled) api.UseMiddleware();
+ if (apiCompressionEnabled) api.UseResponseCompression();
+ });
+}
+
if (compressionEnabled)
- app.UseResponseCompression();
-logger.LogInformation("[System] Static compression: {State}", compressionEnabled ? "enabled (precompressed .br/.gz + on-the-fly fallback)" : "DISABLED (raw/identity)");
+ app.UseWhen(ctx => !IsApiPath(ctx), sf => sf.UseResponseCompression());
+
+logger.LogInformation(
+ "[System] Compression — static: {Static} /api: {Api} level: {Level} (egress accounting: {Egress})",
+ compressionEnabled ? "on (precompressed .br/.gz + on-the-fly)" : "off (raw/identity)",
+ apiCompressionEnabled ? "on (on-the-fly br/gzip)" : "off (identity)",
+ compressionLevel,
+ egressAccountingEnabled ? "on — billing compressed wire bytes" : "off");
+
+// True for /api (and /api/...), case-insensitive — the only paths the egress counter and /api
+// compressor act on. Static, /.auth, /healthz and SPA routes take the other branch.
+static bool IsApiPath(HttpContext c) =>
+ c.Request.Path.StartsWithSegments("/api", StringComparison.OrdinalIgnoreCase);
// Nodes without the Http role do not short-circuit /api or auth paths: the HTTP endpoints simply aren't
// mapped (see the `if (capHttp)` blocks below), so those requests fall through to static file serving
@@ -396,8 +434,10 @@ void RunInitialization()
// Per-instance daily API-egress cap. Same placement rules as the rate limiter above: after UseCraftAuth
// (so an app-only caller's AppId is resolved and CallerClassifier can tell API from UI) and after static
// serving (so a page load's assets are never charged). Only added when egress accounting is enabled —
-// see AddCraftEgressLimiter. Counts only app-only API callers; enforces (429) only once a budget is set.
-if (CraftSettings.RateLimit.Egress.ResolvedEnabled)
+// see AddCraftEgressLimiter. Flags only app-only API callers for billing (ApiEgressWireCounterMiddleware,
+// registered outside compression above, records their compressed wire bytes); enforces (429) only once a
+// budget is set.
+if (egressAccountingEnabled)
app.UseMiddleware();
// Concurrent request tracking for diagnostics. A holder object, not an int: the dispatch endpoint
diff --git a/appsettings.example.jsonc b/appsettings.example.jsonc
index cae74f4..aef6012 100644
--- a/appsettings.example.jsonc
+++ b/appsettings.example.jsonc
@@ -571,5 +571,24 @@
// // A/B measuring compressed vs raw. Env override: CRAFT_COMPRESSION=true|false (wins over this).
// "Compression": true
// }
+
+ // Dynamic /api response policy — independent of the static Frontend block above.
+ // "Api": {
+ // // Compress /api (PowerShell / native dispatch) responses on the fly: Brotli preferred, gzip
+ // // fallback, negotiated from the caller's Accept-Encoding. Default true, and INDEPENDENT of
+ // // Frontend.Compression — an origin behind a CDN that compresses the static bundle typically does
+ // // not re-compress API JSON, so the origin keeps doing it even with static compression off. A
+ // // response is only compressed when the caller advertises an accepted encoding; otherwise it is
+ // // served identity. Env override: CRAFT_API_COMPRESSION=true|false (wins over this).
+ // "Compression": true,
+ // // On-the-fly compression level for the Brotli/gzip providers: Fastest | Optimal | SmallestSize
+ // // | NoCompression. Default Optimal — measured on 2 vCPU, Brotli Optimal ~8.4x vs Fastest ~4.9x at
+ // // the SAME CPU (~14%) and latency, so it is a near-free win. Do NOT use SmallestSize for /api:
+ // // Brotli q11 pegs both cores (~180% CPU) and drives p95 into the tens of seconds. Drop to Fastest
+ // // (or NoCompression) on a very small SKU if the compressor competes with the PowerShell pool.
+ // // Applies to on-the-fly compression (dynamic /api + the static fallback for sibling-less assets);
+ // // precompressed .br/.gz siblings are unaffected. Env override: CRAFT_API_COMPRESSION_LEVEL.
+ // "CompressionLevel": "Optimal"
+ // }
}
}
diff --git a/perf-harness/README.md b/perf-harness/README.md
index 4819206..2a00b92 100644
--- a/perf-harness/README.md
+++ b/perf-harness/README.md
@@ -144,11 +144,12 @@ results/ outputs (gitignored)
── HTTP API mode (below) ──
docker-compose.api.yml http-only SUT (CRAFT_SERVE_API=true) + PerfApi module mount
-docker-compose.egress.yml overlay: egress accounting on (fast flush, known log dir) for run-egress.ps1
+docker-compose.egress.yml overlay: egress accounting on (fast flush, known log dir) for run-egress.ps1 / run-compression.ps1
api-harness/API/Modules/PerfApi/ synthetic PS HTTP endpoints (dependency-free; not for prod)
-k6/api_load.js API load test (weighted endpoint mix or single-endpoint focus)
+k6/api_load.js API load test (weighted endpoint mix or single-endpoint focus; ENC= to negotiate gzip/br)
scripts/run-api.ps1 API orchestrator (up → /healthz → warm → sample → k6 → report → down)
scripts/run-egress.ps1 e2e for the per-instance API egress cap (drive traffic → flush → read ledger → assert)
+scripts/run-compression.ps1 /api compression: correctness + accounting proof (curl) then CPU-vs-bandwidth per encoding (k6)
scripts/compare-api.ps1 before/after diff of two API result JSONs (incl. per-endpoint p95)
```
@@ -215,3 +216,38 @@ Exits non-zero on any failed check (CI-friendly). The ledger file is copied to
`results/egress-ledger-.json`. This covers what the xunit suite can't: the real file location
and serialization, real byte counting through the actual response stream, and the real
auth→classify→count path.
+
+---
+
+## /api compression (CPU vs bandwidth + accounting) (e2e)
+
+Verifies dynamic `/api` response compression end to end and measures its CPU cost. It layers
+`docker-compose.egress.yml` on the API harness (so the egress ledger is live) and runs two phases:
+
+**Phase 1 — correctness + accounting (curl, exact bytes, two isolated cases).** It reads the ledger file
+around each case to get a per-case byte delta, and measures each response's on-the-wire size
+(`curl %{size_download}`, no `--compressed`) and `Content-Encoding` (`%{header_json}`):
+
+- **without** `Accept-Encoding` → asserts the response is **not** compressed (no `Content-Encoding`), and
+ the ledger delta equals the bytes curl received;
+- **with** `Accept-Encoding: gzip` → asserts the response **is** compressed (`Content-Encoding: gzip`) and
+ materially smaller, and the ledger delta equals the bytes curl received.
+
+So the origin compresses only when asked, and the cap bills exactly the bytes that went over the wire in
+**both** cases — proving the accounting sits on the compressed side (a pre-compression counter would have
+read ≈ `Requests × identity-size` for the gzip case).
+
+**Phase 2 — CPU / ratio under load (k6 + docker stats).** Runs the same `PerfJson` payload at a fixed
+arrival rate three times — identity, gzip, br — sampling container CPU during each, and reports CPU%,
+per-response wire bytes, compression ratio, throughput and p95 side by side. This is the CPU-hit number
+to weigh against the bytes saved on a small (1–2 vCPU) container.
+
+```powershell
+docker build -f ..\build\Dockerfile -t craft:local .. # image must include the /api compression + wire counter
+
+pwsh scripts\run-compression.ps1 # 30 reqs/case proof + k6 identity/gzip/br @ rate 150
+pwsh scripts\run-compression.ps1 -JsonN 4000 -Rate 200 -Duration 30s
+pwsh scripts\run-compression.ps1 -Requests 60 -KeepUp
+```
+
+Exits non-zero if any Phase-1 assertion fails. Results go to `results/compression-.{json,md}`.
diff --git a/perf-harness/docker-compose.api.yml b/perf-harness/docker-compose.api.yml
index 7b22403..68acc09 100644
--- a/perf-harness/docker-compose.api.yml
+++ b/perf-harness/docker-compose.api.yml
@@ -25,6 +25,9 @@ services:
# Deployment role: Http (+ optionally Frontend for frontend+http cache/auth profiling — run-cache.ps1).
- CRAFT_SERVE_API=true
- CRAFT_SERVE_FRONTEND=${SERVE_FRONTEND:-}
+ # On-the-fly compression level (Fastest|Optimal|SmallestSize|NoCompression). Empty => code default
+ # (Fastest). run-compression-levels.ps1 sweeps this to measure the CPU/ratio/latency trade-off.
+ - CRAFT_API_COMPRESSION_LEVEL=${API_COMPRESSION_LEVEL:-}
# Opt-in cache/auth-middleware profiling (run-cache.ps1). Off => zero overhead.
- CRAFT_CACHE_TIMING=${CACHE_TIMING:-}
# In-memory cache body tier budget (bytes). run-cache.ps1 -DiskOnly sets 0 to A/B the disk-only "before".
diff --git a/perf-harness/k6/api_load.js b/perf-harness/k6/api_load.js
index d823376..eaea99d 100644
--- a/perf-harness/k6/api_load.js
+++ b/perf-harness/k6/api_load.js
@@ -10,6 +10,11 @@
// CPU_MS ?ms for PerfCpu (default 20)
// SLEEP_MS ?ms for PerfSleep (default 100)
// JSON_N ?n for PerfJson (default 1000)
+// ENC Accept-Encoding to negotiate: 'gzip' | 'br' | '' (default '' = identity, no header).
+// Sent as the Accept-Encoding request HEADER (NOT k6's `compression` param, which compresses
+// the request BODY and does nothing for a GET). k6 counts data_received as the on-the-wire
+// (compressed) size and still decompresses the body for the checks. Used by
+// run-compression.ps1 to A/B encodings.
//
// Output: --summary-export /out/
public int FlushSeconds { get; set; } = 60;
+ ///
+ /// Table into which per-flush accounting is mirrored as time-bucketed rows (in addition to the local
+ /// file), so the product can show usage history. Written by Craft's own table store — the same
+ /// storage account the hosted app reads with Get-CIPPTable — so CIPP can query it directly. Default
+ /// CraftEgressAccounting. Env override: CRAFT_API_EGRESS_TABLE. Blank disables the table
+ /// mirror (file-only).
+ ///
+ public string TableName { get; set; } = "CraftEgressAccounting";
+
+ ///
+ /// Width in minutes of each accounting bucket written to the table. Default 15 (→ 96 buckets/day),
+ /// which is the granularity the product surfaces. The local file keeps only running daily totals; the
+ /// bucketed time-series lives in the table. Floored at 1. Env override: CRAFT_API_EGRESS_BUCKET_MINUTES.
+ ///
+ public int BucketMinutes { get; set; } = 15;
+
+ ///
+ /// Days of bucketed rows to retain in the table before a periodic purge deletes them. Default 7.
+ /// Floored at 1. Env override: CRAFT_API_EGRESS_RETENTION_DAYS. Does not affect the local file,
+ /// which only ever holds the current UTC day.
+ ///
+ public int RetentionDays { get; set; } = 7;
+
internal const string EnabledEnv = "CRAFT_API_EGRESS_LIMIT_ENABLED";
internal const string BytesEnv = "CRAFT_API_EGRESS_LIMIT_BYTES";
internal const string FlushEnv = "CRAFT_API_EGRESS_FLUSH_SECONDS";
+ internal const string TableEnv = "CRAFT_API_EGRESS_TABLE";
+ internal const string BucketEnv = "CRAFT_API_EGRESS_BUCKET_MINUTES";
+ internal const string RetentionEnv = "CRAFT_API_EGRESS_RETENTION_DAYS";
///
/// Whether egress accounting (and therefore the middleware) should be active, resolved through
@@ -94,6 +120,32 @@ public bool ResolveEnabled(Func env)
? fromEnv
: Math.Max(1, FlushSeconds);
+ /// Resolved table name, honouring CRAFT_API_EGRESS_TABLE. Blank/whitespace = the table
+ /// mirror is off (file-only accounting).
+ public string ResolvedTableName
+ {
+ get
+ {
+ var fromEnv = Environment.GetEnvironmentVariable(TableEnv);
+ var name = string.IsNullOrWhiteSpace(fromEnv) ? TableName : fromEnv;
+ return (name ?? string.Empty).Trim();
+ }
+ }
+
+ /// Resolved bucket width in minutes, honouring CRAFT_API_EGRESS_BUCKET_MINUTES. Floored at 1.
+ public int ResolvedBucketMinutes =>
+ int.TryParse(Environment.GetEnvironmentVariable(BucketEnv), NumberStyles.Integer,
+ CultureInfo.InvariantCulture, out var fromEnv) && fromEnv > 0
+ ? fromEnv
+ : Math.Max(1, BucketMinutes);
+
+ /// Resolved retention in days, honouring CRAFT_API_EGRESS_RETENTION_DAYS. Floored at 1.
+ public int ResolvedRetentionDays =>
+ int.TryParse(Environment.GetEnvironmentVariable(RetentionEnv), NumberStyles.Integer,
+ CultureInfo.InvariantCulture, out var fromEnv) && fromEnv > 0
+ ? fromEnv
+ : Math.Max(1, RetentionDays);
+
// Tri-state flag parse (mirrors Craft.Hosting.EnvFlag, inlined to keep Configuration free of a
// dependency on Hosting): null when unset/blank, true for "true"/"1" (any casing), false otherwise.
private static bool? ParseFlag(string? value)
diff --git a/Services/Hosting/ApiEgressLimiterMiddleware.cs b/Services/Hosting/ApiEgressLimiterMiddleware.cs
index ca66a21..ba01351 100644
--- a/Services/Hosting/ApiEgressLimiterMiddleware.cs
+++ b/Services/Hosting/ApiEgressLimiterMiddleware.cs
@@ -49,21 +49,26 @@ public async Task InvokeAsync(HttpContext context)
return;
}
+ // The AppId (GUID) is the per-client accounting key — the same principal name the concurrency
+ // limiter partitions on. Carried on the charge flag below and recorded on a shed.
+ var appId = context.Request.Headers["x-ms-client-principal-name"].ToString();
+
// Already over budget for today → shed before running the (often minute-long) downstream call,
// saving both the compute and the egress. Accounting is post-hoc, so the request that tips the
// total over still completes; the NEXT one is the first to be refused. Combined with the API
// concurrency cap, the overshoot is bounded to (concurrency × largest response).
if (_ledger.ShouldReject())
{
+ _ledger.RecordShed(appId);
await RejectAsync(context);
return;
}
- // Greenlit: mark the request so ApiEgressWireCounterMiddleware (running outside the response
- // compressor) bills its on-the-wire bytes. We don't count here — a counter at this position
- // would see the pre-compression body and miss the compressor's final flush, which unwinds
- // further out. The shed body above is deliberately left unflagged, so it is never charged.
- context.Items[ApiEgressWireCounterMiddleware.ChargeItemKey] = true;
+ // Greenlit: carry the AppId so ApiEgressWireCounterMiddleware (running outside the response
+ // compressor) bills its on-the-wire bytes against this client. We don't count here — a counter at
+ // this position would see the pre-compression body and miss the compressor's final flush, which
+ // unwinds further out. The shed body above is left unflagged, so it is never charged.
+ context.Items[ApiEgressWireCounterMiddleware.ChargeItemKey] = appId;
await _next(context);
}
diff --git a/Services/Hosting/ApiEgressWireCounterMiddleware.cs b/Services/Hosting/ApiEgressWireCounterMiddleware.cs
index 41f857b..e38b735 100644
--- a/Services/Hosting/ApiEgressWireCounterMiddleware.cs
+++ b/Services/Hosting/ApiEgressWireCounterMiddleware.cs
@@ -27,9 +27,10 @@ namespace Craft.Hosting;
public sealed class ApiEgressWireCounterMiddleware
{
///
- /// Request item set by on an API request it lets through,
- /// signalling this middleware to bill the response's wire bytes. Absent for UI, anonymous, static
- /// and shed requests, which are never charged.
+ /// Request item set by on an API request it lets through:
+ /// the caller's AppId (a non-empty string), signalling this middleware to bill the response's wire
+ /// bytes against that client. Absent for UI, anonymous, static and shed requests, which are never
+ /// charged.
///
public const string ChargeItemKey = "Craft.Egress.Charge";
@@ -60,8 +61,8 @@ public async Task InvokeAsync(HttpContext context)
finally
{
context.Response.Body = original;
- if (context.Items.TryGetValue(ChargeItemKey, out var charge) && charge is true)
- _ledger.Record(counting.BytesWritten);
+ if (context.Items.TryGetValue(ChargeItemKey, out var charge) && charge is string appId && appId.Length > 0)
+ _ledger.Record(counting.BytesWritten, appId);
}
}
}
diff --git a/Services/Hosting/EgressLedger.cs b/Services/Hosting/EgressLedger.cs
index 76ce844..519627b 100644
--- a/Services/Hosting/EgressLedger.cs
+++ b/Services/Hosting/EgressLedger.cs
@@ -1,27 +1,29 @@
using System.Globalization;
using System.Text.Json;
using Craft.Configuration;
+using Craft.Storage;
namespace Craft.Hosting;
///
-/// Instance-wide daily egress counter for app-only API clients, persisted to a small local file so a
-/// restart does not reset the day's total. One number for the whole instance (deliberately NOT per
-/// client): records the outbound bytes of each API-client
-/// response into it, and asks it whether the instance is over budget before letting the next API
-/// request run.
+/// Instance-wide daily egress counter for app-only API clients, with per-client accounting. The running
+/// daily totals are the source of truth for the cap and are persisted to a small local file so a restart
+/// does not reset the day; the time-bucketed history and a durable daily audit are additionally mirrored
+/// to a table so the product can show usage over time and whether/when/how often the cap was hit.
///
-/// The request path only ever touches memory under a short lock — no per-request IO. A background timer
-/// flushes the counter to disk every FlushSeconds (and once more on shutdown), and only when it
-/// has changed, so an idle instance writes nothing. The file is loaded in the constructor (before the
-/// host serves traffic); if it is stamped with today's UTC date the total is restored, otherwise the
-/// day starts fresh. Nothing is read back from Azure — see .
+/// Request path ( / ) only touches memory under a short
+/// lock — no IO. A background timer flushes the daily totals to disk and mirrors the current 15-minute
+/// bucket plus today's daily summary to the table every FlushSeconds (and on shutdown), and only
+/// when something changed. Restart reloads the daily totals from the file (so the cap is correct
+/// immediately) and reseeds the current in-progress bucket from the table (so its row does not regress).
+/// A periodic purge drops table rows past the retention window. Table failures never affect the cap — the
+/// file is the truth; the table is a best-effort mirror.
///
///
-/// The file lives in the same directory as the log files (,
-/// e.g. {home}/logs), which is the writable area confirmed to survive restarts and crashes and is
-/// lost only on a slice/image replacement — when a fresh day is the correct behaviour anyway. Storing it
-/// there also means it follows any App__FileLogging__Directory override the deployment sets.
+/// The file lives with the log files (), the writable
+/// area confirmed to survive restarts. The table is written through Craft's own store — the same storage
+/// account the hosted app reads with Get-CIPPTable — so CIPP queries it directly. Layout: see
+/// .
///
///
public sealed class EgressLedger : BackgroundService
@@ -32,10 +34,29 @@ public sealed class EgressLedger : BackgroundService
private readonly string _filePath;
private readonly Func _utcNow;
+ // Table mirror — null store or blank table name = file-only (no mirror).
+ private readonly ICraftTableStore? _store;
+ private readonly string _tableName;
+ private readonly int _bucketMinutes;
+ private readonly int _retentionDays;
+ private bool TableEnabled => _store is not null && !string.IsNullOrWhiteSpace(_tableName);
+
private readonly object _lock = new();
private DateOnly _dateUtc;
private long _bytes;
- private bool _dirty;
+ private long _shedRequests;
+ private DateTime? _capReachedUtc;
+ private bool _dirty; // file needs rewriting
+ private bool _tableDirty; // something changed since the last table sync
+
+ private readonly Dictionary _clients = new(StringComparer.OrdinalIgnoreCase);
+ // Pending table buckets not yet finalised: bucketRowKey -> (appId -> accum).
+ private readonly Dictionary> _buckets = new(StringComparer.Ordinal);
+
+ private bool _tableReady;
+
+ private sealed class ClientTotals { public long Bytes; public long Requests; public long Shed; public DateTime LastSeenUtc; }
+ private sealed class BucketAccum { public long Bytes; public long Requests; public long Shed; }
private static readonly JsonSerializerOptions s_json = new()
{
@@ -43,26 +64,35 @@ public sealed class EgressLedger : BackgroundService
WriteIndented = false,
};
- /// Production constructor — resolves cap/flush from settings and stores the ledger file in
- /// the same directory as the log files, which is the writable area confirmed to survive restarts and
- /// crashes (and honours an App__FileLogging__Directory override).
- public EgressLedger(ILogger logger, CraftSettings settings)
+ /// Production constructor — resolves settings and injects the shared table store.
+ public EgressLedger(ILogger logger, CraftSettings settings, ICraftTableStore store)
: this(logger,
(settings ?? throw new ArgumentNullException(nameof(settings))).RateLimit.Egress.ResolvedBytesPerDay,
settings.RateLimit.Egress.ResolvedFlushSeconds,
- Path.Combine(settings.FileLogging.ResolvedDirectory, "egress-ledger.json"))
+ Path.Combine(settings.FileLogging.ResolvedDirectory, "egress-ledger.json"),
+ utcNow: null,
+ store: store,
+ tableName: settings.RateLimit.Egress.ResolvedTableName,
+ bucketMinutes: settings.RateLimit.Egress.ResolvedBucketMinutes,
+ retentionDays: settings.RateLimit.Egress.ResolvedRetentionDays)
{
}
- /// Test/explicit constructor. defaults to the real clock.
+ /// Test/explicit constructor. defaults to the real clock;
+ /// null (or a blank ) = file-only.
internal EgressLedger(ILogger logger, long capBytes, int flushSeconds, string filePath,
- Func? utcNow = null)
+ Func? utcNow = null, ICraftTableStore? store = null, string tableName = "",
+ int bucketMinutes = 15, int retentionDays = 7)
{
_logger = logger;
_capBytes = Math.Max(0, capBytes);
_flushSeconds = Math.Max(1, flushSeconds);
_filePath = filePath;
_utcNow = utcNow ?? (() => DateTime.UtcNow);
+ _store = store;
+ _tableName = tableName ?? string.Empty;
+ _bucketMinutes = Math.Max(1, bucketMinutes);
+ _retentionDays = Math.Max(1, retentionDays);
_dateUtc = DateOnly.FromDateTime(_utcNow());
Load();
}
@@ -84,26 +114,55 @@ public bool ShouldReject()
}
}
- /// Add the outbound bytes of one response to today's running total.
- public void Record(long bytes)
+ /// Add the outbound bytes of one response to today's running totals, attributed to a client.
+ public void Record(long bytes, string appId)
{
if (bytes <= 0) return;
+ appId = Normalize(appId);
lock (_lock)
{
RolloverIfNeeded();
_bytes += bytes;
+ Client(appId).Bytes += bytes;
+ Client(appId).Requests += 1;
+ Client(appId).LastSeenUtc = _utcNow();
+ Bucket(appId).Bytes += bytes;
+ Bucket(appId).Requests += 1;
+ _dirty = true;
+ _tableDirty = true;
+ }
+ }
+
+ /// Record that a request from was shed (429) once over budget. The
+ /// shed response body is deliberately not counted as bytes; this tracks the cap-hit count.
+ public void RecordShed(string appId)
+ {
+ appId = Normalize(appId);
+ lock (_lock)
+ {
+ RolloverIfNeeded();
+ _shedRequests += 1;
+ _capReachedUtc ??= _utcNow();
+ Client(appId).Shed += 1;
+ Bucket(appId).Shed += 1;
_dirty = true;
+ _tableDirty = true;
}
}
- /// Bytes served so far today (for logging / a future metrics surface).
+ /// Bytes served so far today (for logging / a metrics surface).
public long CurrentBytes
{
get { lock (_lock) { RolloverIfNeeded(); return _bytes; } }
}
- /// Seconds from now until the next UTC midnight, floored at 1 — the Retry-After a
- /// rejected caller is handed so it comes back once the daily window has reset.
+ /// Requests shed (429) so far today.
+ public long ShedRequests
+ {
+ get { lock (_lock) { RolloverIfNeeded(); return _shedRequests; } }
+ }
+
+ /// Seconds from now until the next UTC midnight, floored at 1.
public int SecondsToNextUtcMidnight()
{
var now = _utcNow();
@@ -114,6 +173,24 @@ public int SecondsToNextUtcMidnight()
public string NextUtcMidnightIso() =>
_utcNow().Date.AddDays(1).ToString("yyyy-MM-ddTHH:mm:ssZ", CultureInfo.InvariantCulture);
+ // Must be called under _lock.
+ private ClientTotals Client(string appId)
+ {
+ if (!_clients.TryGetValue(appId, out var c)) { c = new ClientTotals(); _clients[appId] = c; }
+ return c;
+ }
+
+ // Must be called under _lock. Returns the accumulator for appId in the current 15-min bucket.
+ private BucketAccum Bucket(string appId)
+ {
+ var key = EgressTableSchema.BucketRowKey(EgressTableSchema.BucketStart(_utcNow(), _bucketMinutes));
+ if (!_buckets.TryGetValue(key, out var b)) { b = new Dictionary(StringComparer.OrdinalIgnoreCase); _buckets[key] = b; }
+ if (!b.TryGetValue(appId, out var a)) { a = new BucketAccum(); b[appId] = a; }
+ return a;
+ }
+
+ private static string Normalize(string? appId) => string.IsNullOrWhiteSpace(appId) ? "unknown" : appId.Trim();
+
// Must be called under _lock.
private void RolloverIfNeeded()
{
@@ -121,32 +198,71 @@ private void RolloverIfNeeded()
if (today == _dateUtc) return;
_dateUtc = today;
_bytes = 0;
+ _shedRequests = 0;
+ _capReachedUtc = null;
+ _clients.Clear();
+ _buckets.Clear(); // yesterday's bucket rows stay in the table until retention purges them
_dirty = true;
+ _tableDirty = true;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation(
- "[Egress] API egress accounting active — cap {Cap}, flush every {Flush}s, file {File}",
+ "[Egress] API egress accounting active — cap {Cap}, flush every {Flush}s, file {File}, table {Table}",
_capBytes > 0 ? $"{_capBytes} bytes/day" : "accounting only (no cap)",
- _flushSeconds, _filePath);
+ _flushSeconds,
+ _filePath,
+ TableEnabled ? $"{_tableName} ({_bucketMinutes}-min buckets, {_retentionDays}-day retention)" : "off (file only)");
+ if (TableEnabled)
+ await InitTableAsync(stoppingToken).ConfigureAwait(false);
+
+ var lastPurgeDay = DateOnly.FromDateTime(_utcNow());
try
{
using var timer = new PeriodicTimer(TimeSpan.FromSeconds(_flushSeconds));
while (await timer.WaitForNextTickAsync(stoppingToken).ConfigureAwait(false))
+ {
Flush();
+ if (TableEnabled)
+ {
+ await SyncToTableAsync(stoppingToken).ConfigureAwait(false);
+ var today = DateOnly.FromDateTime(_utcNow());
+ if (today != lastPurgeDay) { await PurgeAsync(stoppingToken).ConfigureAwait(false); lastPurgeDay = today; }
+ }
+ }
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
- // Host shutting down (the timer await was cancelled) — expected; fall through to the final flush.
+ // Host shutting down — fall through to the final flush/sync.
}
- Flush(); // final flush so the last window of counts survives a graceful restart
+ Flush();
+ if (TableEnabled)
+ {
+ try { await SyncToTableAsync(CancellationToken.None).ConfigureAwait(false); }
+ catch (Exception ex) { _logger.LogWarning(ex, "[Egress] Final table sync failed"); }
+ }
+ }
+
+ private async Task InitTableAsync(CancellationToken ct)
+ {
+ try
+ {
+ await _store!.EnsureTableAsync(_tableName, ct).ConfigureAwait(false);
+ _tableReady = true;
+ await SeedCurrentBucketAsync(ct).ConfigureAwait(false); // backfill so a mid-bucket restart does not regress
+ await PurgeAsync(ct).ConfigureAwait(false);
+ }
+ catch (Exception ex)
+ {
+ _logger.LogWarning(ex, "[Egress] Table init failed — mirroring will retry on the next flush");
+ }
}
- /// Load the persisted counter. Restores only when the file is stamped with today's UTC
- /// date; a stale or unreadable file leaves today starting from zero.
+ /// Load the persisted daily totals. Restores only when the file is stamped with today's UTC
+ /// date; a stale or unreadable file leaves today starting from zero. Tolerates the v1 format.
internal void Load()
{
try
@@ -159,23 +275,28 @@ internal void Load()
var today = DateOnly.FromDateTime(_utcNow());
var sameDay = DateOnly.TryParseExact(state.DateUtc, "yyyy-MM-dd", CultureInfo.InvariantCulture,
DateTimeStyles.None, out var stored) && stored == today;
+ if (!sameDay || state.Bytes < 0)
+ {
+ _logger.LogInformation("[Egress] Ledger file not for today ({Stored}) — starting the day fresh", state.DateUtc);
+ return;
+ }
- if (sameDay && state.Bytes >= 0)
+ lock (_lock)
{
- lock (_lock)
+ _dateUtc = today;
+ _bytes = state.Bytes;
+ _shedRequests = Math.Max(0, state.ShedRequests);
+ _capReachedUtc = state.CapReachedUtc;
+ _clients.Clear();
+ if (state.Clients is not null)
{
- _dateUtc = today;
- _bytes = state.Bytes;
- _dirty = false;
+ foreach (var (appId, c) in state.Clients)
+ _clients[appId] = new ClientTotals { Bytes = c.Bytes, Requests = c.Requests, Shed = c.Shed, LastSeenUtc = c.LastSeenUtc };
}
- _logger.LogInformation("[Egress] Restored {Bytes} bytes already served today from {File}",
- state.Bytes, _filePath);
- }
- else
- {
- _logger.LogInformation("[Egress] Ledger file not for today ({Stored}) — starting the day fresh",
- state.DateUtc);
+ _dirty = false;
}
+ _logger.LogInformation("[Egress] Restored {Bytes} bytes across {Clients} client(s) already served today from {File}",
+ state.Bytes, state.Clients?.Count ?? 0, _filePath);
}
catch (Exception ex)
{
@@ -183,8 +304,7 @@ internal void Load()
}
}
- /// Rewrite the counter to disk via a temp file + atomic rename, but only when it has
- /// changed since the last flush. Failures are logged and retried on the next tick.
+ /// Rewrite the daily totals to disk via a temp file + atomic rename, only when changed.
internal void Flush()
{
try
@@ -196,8 +316,16 @@ internal void Flush()
RolloverIfNeeded();
snapshot = new LedgerState
{
+ Version = 2,
DateUtc = _dateUtc.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture),
Bytes = _bytes,
+ CapBytes = _capBytes,
+ CapReachedUtc = _capReachedUtc,
+ ShedRequests = _shedRequests,
+ UpdatedUtc = _utcNow(),
+ Clients = _clients.ToDictionary(
+ kv => kv.Key,
+ kv => new ClientState { Bytes = kv.Value.Bytes, Requests = kv.Value.Requests, Shed = kv.Value.Shed, LastSeenUtc = kv.Value.LastSeenUtc }),
};
_dirty = false;
}
@@ -212,13 +340,171 @@ internal void Flush()
catch (Exception ex)
{
_logger.LogWarning(ex, "[Egress] Failed to flush ledger to disk");
- lock (_lock) { _dirty = true; } // don't lose the change — try again next tick
+ lock (_lock) { _dirty = true; }
}
}
+ // ── Table mirror ─────────────────────────────────────────────────────────────────────────────────
+
+ /// Mirror the current bucket(s) and today's daily summary to the table. Only rows for clients
+ /// active since the last sync are written. Finalised buckets (older than the current one) are written
+ /// then dropped from memory; the current bucket stays and is re-written as it grows.
+ internal async Task SyncToTableAsync(CancellationToken ct)
+ {
+ if (!_tableReady) { try { await _store!.EnsureTableAsync(_tableName, ct).ConfigureAwait(false); _tableReady = true; } catch { return; } }
+
+ // Snapshot under the lock.
+ DateOnly day; long dayBytes, dayShed; DateTime? capReached; bool enforcing;
+ string currentBucketKey;
+ List<(string bucketKey, Dictionary clients)> buckets;
+ Dictionary clientDaily;
+ lock (_lock)
+ {
+ if (!_tableDirty && !_buckets.Keys.Any(k => k != CurrentBucketKey())) return;
+ RolloverIfNeeded();
+ currentBucketKey = CurrentBucketKey();
+ buckets = _buckets.Select(kv => (kv.Key, kv.Value.ToDictionary(c => c.Key, c => new BucketAccum { Bytes = c.Value.Bytes, Requests = c.Value.Requests, Shed = c.Value.Shed }, StringComparer.OrdinalIgnoreCase))).ToList();
+ var active = new HashSet(buckets.SelectMany(b => b.clients.Keys), StringComparer.OrdinalIgnoreCase);
+ clientDaily = active.Where(a => _clients.ContainsKey(a)).ToDictionary(a => a, a => new ClientTotals { Bytes = _clients[a].Bytes, Requests = _clients[a].Requests, Shed = _clients[a].Shed, LastSeenUtc = _clients[a].LastSeenUtc }, StringComparer.OrdinalIgnoreCase);
+ day = _dateUtc; dayBytes = _bytes; dayShed = _shedRequests; capReached = _capReachedUtc; enforcing = _capBytes > 0;
+ _tableDirty = false;
+ }
+
+ try
+ {
+ // Bucket rows — per client + a per-bucket instance aggregate.
+ foreach (var (bucketKey, clients) in buckets)
+ {
+ var bucketStart = ParseBucketStart(bucketKey);
+ long bBytes = 0, bReq = 0, bShed = 0;
+ foreach (var (appId, a) in clients)
+ {
+ await _store!.UpsertAsync(_tableName, EgressTableSchema.ClientBucketRow(appId, bucketStart, a.Bytes, a.Requests, a.Shed, _capBytes), ct).ConfigureAwait(false);
+ bBytes += a.Bytes; bReq += a.Requests; bShed += a.Shed;
+ }
+ await _store!.UpsertAsync(_tableName, EgressTableSchema.SystemBucketRow(bucketStart, bBytes, bReq, bShed, _capBytes), ct).ConfigureAwait(false);
+ }
+
+ // Daily audit rows — per active client + the instance summary (cap, when hit, how many shed).
+ foreach (var (appId, c) in clientDaily)
+ await _store!.UpsertAsync(_tableName, EgressTableSchema.ClientDailyRow(appId, day, c.Bytes, c.Requests, c.Shed, _capBytes), ct).ConfigureAwait(false);
+ await _store!.UpsertAsync(_tableName, EgressTableSchema.SystemDailyRow(day, dayBytes, clientDaily.Values.Sum(c => c.Requests), dayShed, _capBytes, capReached, enforcing), ct).ConfigureAwait(false);
+
+ // Drop finalised buckets we just wrote (strictly older than the current one).
+ lock (_lock)
+ {
+ foreach (var (bucketKey, _) in buckets)
+ if (string.CompareOrdinal(bucketKey, currentBucketKey) < 0)
+ _buckets.Remove(bucketKey);
+ }
+ }
+ catch (Exception ex)
+ {
+ _logger.LogWarning(ex, "[Egress] Table sync failed — will retry on the next flush");
+ lock (_lock) { _tableDirty = true; } // ensure a retry
+ }
+ }
+
+ /// Reseed the current bucket's per-client accumulators from the table so a restart mid-bucket
+ /// continues the row rather than overwriting it with a smaller value.
+ internal async Task SeedCurrentBucketAsync(CancellationToken ct)
+ {
+ try
+ {
+ var key = CurrentBucketKey();
+ var filter = $"RowKey eq '{key}'";
+ var seeded = 0;
+ await foreach (var row in _store!.QueryTableAsync(_tableName, filter, ct).ConfigureAwait(false))
+ {
+ if (row.RowKey != key || row.PartitionKey == EgressTableSchema.SystemPartition) continue; // re-apply filter; skip aggregate
+ lock (_lock)
+ {
+ if (!_buckets.TryGetValue(key, out var b)) { b = new Dictionary(StringComparer.OrdinalIgnoreCase); _buckets[key] = b; }
+ b[row.PartitionKey] = new BucketAccum { Bytes = AsLong(row[EgressTableSchema.PropBytes]), Requests = AsLong(row[EgressTableSchema.PropRequests]), Shed = AsLong(row[EgressTableSchema.PropShed]) };
+ }
+ seeded++;
+ }
+ if (seeded > 0) _logger.LogInformation("[Egress] Reseeded current bucket from table ({Count} client rows)", seeded);
+ }
+ catch (Exception ex)
+ {
+ _logger.LogWarning(ex, "[Egress] Current-bucket backfill failed — the current partial bucket may under-count until the next boundary");
+ }
+ }
+
+ /// Delete bucket and daily rows older than the retention window, partition by partition.
+ internal async Task PurgeAsync(CancellationToken ct)
+ {
+ try
+ {
+ var now = _utcNow();
+ var bucketCutoff = EgressTableSchema.BucketRetentionCutoffRowKey(now, _retentionDays, _bucketMinutes);
+ var dailyCutoff = EgressTableSchema.DailyRetentionCutoffRowKey(now, _retentionDays);
+
+ // One filtered scan for expired rows of either kind; group by partition and batch-delete.
+ var filter = $"(RowKey ge '{EgressTableSchema.BucketPrefix}' and RowKey lt '{bucketCutoff}') or " +
+ $"(RowKey ge '{EgressTableSchema.DailyPrefix}' and RowKey lt '{dailyCutoff}')";
+ var byPartition = new Dictionary>(StringComparer.Ordinal);
+ await foreach (var row in _store!.QueryTableAsync(_tableName, filter, ct).ConfigureAwait(false))
+ {
+ var rk = row.RowKey;
+ var expired = (rk.StartsWith(EgressTableSchema.BucketPrefix, StringComparison.Ordinal) && string.CompareOrdinal(rk, bucketCutoff) < 0)
+ || (rk.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal) && string.CompareOrdinal(rk, dailyCutoff) < 0);
+ if (!expired) continue; // re-apply the predicate — the store may ignore the filter
+ if (!byPartition.TryGetValue(row.PartitionKey, out var list)) { list = new List(); byPartition[row.PartitionKey] = list; }
+ list.Add(rk);
+ }
+
+ var deleted = 0;
+ foreach (var (pk, rowKeys) in byPartition)
+ {
+ await _store!.DeleteBatchAsync(_tableName, pk, rowKeys, ct).ConfigureAwait(false);
+ deleted += rowKeys.Count;
+ }
+ if (deleted > 0) _logger.LogInformation("[Egress] Purged {Count} expired row(s) beyond {Days}-day retention", deleted, _retentionDays);
+ }
+ catch (Exception ex)
+ {
+ _logger.LogWarning(ex, "[Egress] Retention purge failed — will retry on the next daily tick");
+ }
+ }
+
+ private string CurrentBucketKey() =>
+ EgressTableSchema.BucketRowKey(EgressTableSchema.BucketStart(_utcNow(), _bucketMinutes));
+
+ private static DateTime ParseBucketStart(string bucketRowKey)
+ {
+ var stamp = bucketRowKey.AsSpan(EgressTableSchema.BucketPrefix.Length); // yyyyMMddTHHmmssZ
+ return DateTime.ParseExact(stamp, "yyyyMMddTHHmmssZ", CultureInfo.InvariantCulture, DateTimeStyles.AdjustToUniversal | DateTimeStyles.AssumeUniversal);
+ }
+
+ private static long AsLong(object? value) => value switch
+ {
+ long l => l,
+ int i => i,
+ double d => (long)d,
+ string s when long.TryParse(s, NumberStyles.Integer, CultureInfo.InvariantCulture, out var r) => r,
+ JsonElement je when je.ValueKind == JsonValueKind.Number && je.TryGetInt64(out var r) => r,
+ _ => 0L,
+ };
+
private sealed class LedgerState
{
+ public int Version { get; set; }
public string DateUtc { get; set; } = "";
public long Bytes { get; set; }
+ public long CapBytes { get; set; }
+ public DateTime? CapReachedUtc { get; set; }
+ public long ShedRequests { get; set; }
+ public DateTime? UpdatedUtc { get; set; }
+ public Dictionary? Clients { get; set; }
+ }
+
+ private sealed class ClientState
+ {
+ public long Bytes { get; set; }
+ public long Requests { get; set; }
+ public long Shed { get; set; }
+ public DateTime LastSeenUtc { get; set; }
}
}
diff --git a/Services/Hosting/EgressTableSchema.cs b/Services/Hosting/EgressTableSchema.cs
new file mode 100644
index 0000000..460ef9e
--- /dev/null
+++ b/Services/Hosting/EgressTableSchema.cs
@@ -0,0 +1,128 @@
+using System.Globalization;
+using Craft.Storage;
+
+namespace Craft.Hosting;
+
+///
+/// Shapes the egress rows that mirrors into the accounting table, and the
+/// keys used to query them. Pure and side-effect-free so the layout can be unit-tested without a
+/// storage backend.
+///
+/// Two row kinds share the table, distinguished by a RowKey prefix so each can be queried on its own:
+///
+/// Buckets (bkt_) — one row per (15-min bucket, client) plus a
+/// per-bucket instance aggregate. The time-series the product charts.
+/// Daily summaries (day_) — one row per (UTC day, client) plus a
+/// per-day instance aggregate carrying the cap in force, when the cap was first hit that day, how many
+/// requests were shed, and whether enforcement was on. The durable audit of "was the cap hit, when,
+/// how many times, and what was it".
+///
+/// PartitionKey is the client's AppId (GUID) for per-client rows, or
+/// for the instance aggregate — so one client's history is a single-partition scan and the instance
+/// trend is another. RowKey embeds a lexically-sortable UTC stamp after the prefix, so "last 24h"
+/// and retention are range queries, never full scans. Azure Table key chars / \ # ? are illegal,
+/// so the prefix separator is _. The same storage account the hosted app reads with Get-CIPPTable
+/// backs this table, so CIPP queries it directly rather than through Craft.
+///
+///
+internal static class EgressTableSchema
+{
+ /// Partition holding the instance aggregate. Not a GUID, so it never collides with an AppId.
+ public const string SystemPartition = "instance-total";
+
+ public const string BucketPrefix = "bkt_";
+ public const string DailyPrefix = "day_";
+ // Upper bound of the bkt_ prefix range: 'u' is the next char after 't', so any bkt_* RowKey is < this.
+ public const string BucketPrefixEnd = "bku_";
+ public const string DailyPrefixEnd = "daz_";
+
+ // Property names on the row — kept stable, CIPP reads these.
+ public const string PropBytes = "Bytes";
+ public const string PropRequests = "Requests";
+ public const string PropShed = "Shed";
+ public const string PropCapBytes = "CapBytes";
+ public const string PropBucketStart = "BucketStartUtc";
+ public const string PropDateUtc = "DateUtc";
+ public const string PropCapReachedUtc = "CapReachedUtc";
+ public const string PropEnforcing = "Enforcing";
+ public const string PropAppId = "AppId";
+
+ /// The UTC start of the bucket falls in, floored to
+ /// .
+ public static DateTime BucketStart(DateTime utc, int bucketMinutes)
+ {
+ var minutes = Math.Max(1, bucketMinutes);
+ var day = new DateTime(utc.Year, utc.Month, utc.Day, 0, 0, 0, DateTimeKind.Utc);
+ var flooredMinutes = ((long)(utc - day).TotalMinutes / minutes) * minutes;
+ return day.AddMinutes(flooredMinutes);
+ }
+
+ private static string Stamp(DateTime utc) =>
+ utc.ToUniversalTime().ToString("yyyyMMddTHHmmssZ", CultureInfo.InvariantCulture);
+
+ private static string DayStamp(DateOnly day) =>
+ day.ToString("yyyyMMdd", CultureInfo.InvariantCulture);
+
+ public static string BucketRowKey(DateTime bucketStartUtc) => BucketPrefix + Stamp(bucketStartUtc);
+ public static string DailyRowKey(DateOnly dayUtc) => DailyPrefix + DayStamp(dayUtc);
+
+ /// Lower bound RowKey for a "last h" bucket query.
+ public static string BucketSinceRowKey(DateTime nowUtc, int hours, int bucketMinutes) =>
+ BucketPrefix + Stamp(BucketStart(nowUtc.AddHours(-Math.Abs(hours)), bucketMinutes));
+
+ /// Exclusive upper bound RowKey for purging buckets outside the retention window.
+ public static string BucketRetentionCutoffRowKey(DateTime nowUtc, int retentionDays, int bucketMinutes) =>
+ BucketPrefix + Stamp(BucketStart(nowUtc.AddDays(-Math.Max(1, retentionDays)), bucketMinutes));
+
+ /// Exclusive upper bound RowKey for purging daily rows outside the retention window.
+ public static string DailyRetentionCutoffRowKey(DateTime nowUtc, int retentionDays) =>
+ DailyPrefix + DayStamp(DateOnly.FromDateTime(nowUtc.AddDays(-Math.Max(1, retentionDays))));
+
+ public static StoreRow ClientBucketRow(string appId, DateTime bucketStartUtc,
+ long bytes, long requests, long shed, long capBytes)
+ {
+ var row = new StoreRow(appId, BucketRowKey(bucketStartUtc));
+ FillCounts(row, bytes, requests, shed, capBytes);
+ row[PropBucketStart] = bucketStartUtc;
+ row[PropAppId] = appId;
+ return row;
+ }
+
+ public static StoreRow SystemBucketRow(DateTime bucketStartUtc,
+ long bytes, long requests, long shed, long capBytes)
+ {
+ var row = new StoreRow(SystemPartition, BucketRowKey(bucketStartUtc));
+ FillCounts(row, bytes, requests, shed, capBytes);
+ row[PropBucketStart] = bucketStartUtc;
+ return row;
+ }
+
+ public static StoreRow ClientDailyRow(string appId, DateOnly dayUtc,
+ long bytes, long requests, long shed, long capBytes)
+ {
+ var row = new StoreRow(appId, DailyRowKey(dayUtc));
+ FillCounts(row, bytes, requests, shed, capBytes);
+ row[PropDateUtc] = dayUtc.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture);
+ row[PropAppId] = appId;
+ return row;
+ }
+
+ public static StoreRow SystemDailyRow(DateOnly dayUtc,
+ long bytes, long requests, long shed, long capBytes, DateTime? capReachedUtc, bool enforcing)
+ {
+ var row = new StoreRow(SystemPartition, DailyRowKey(dayUtc));
+ FillCounts(row, bytes, requests, shed, capBytes);
+ row[PropDateUtc] = dayUtc.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture);
+ row[PropEnforcing] = enforcing;
+ if (capReachedUtc.HasValue) row[PropCapReachedUtc] = capReachedUtc.Value;
+ return row;
+ }
+
+ private static void FillCounts(StoreRow row, long bytes, long requests, long shed, long capBytes)
+ {
+ row[PropBytes] = bytes;
+ row[PropRequests] = requests;
+ row[PropShed] = shed;
+ row[PropCapBytes] = capBytes;
+ }
+}
diff --git a/appsettings.example.jsonc b/appsettings.example.jsonc
index aef6012..c7c28e4 100644
--- a/appsettings.example.jsonc
+++ b/appsettings.example.jsonc
@@ -108,12 +108,21 @@
// // so you can run accounting-only first and watch real numbers before rejecting anything.
// // The counter lives in memory (no per-request IO) and is flushed to egress-ledger.json in the same
// // directory as the log files every FlushSeconds and on shutdown, so a restart doesn't reset the day.
- // // Nothing is read from Azure.
- // // Env overrides: CRAFT_API_EGRESS_LIMIT_BYTES, CRAFT_API_EGRESS_FLUSH_SECONDS.
+ // // The file holds per-client daily totals + the cap-hit state (the source of truth for the cap).
+ // // In addition, per-client + instance rows are mirrored to the TableName table — 15-minute buckets
+ // // (the time-series) and a daily audit row (cap in force, when the cap was first hit, how many
+ // // requests were shed) — so the product can chart usage and show cap hits. The table is the same
+ // // storage account the hosted app reads with Get-CIPPTable, retained RetentionDays and purged
+ // // periodically. A restart reseeds the current bucket from the table so its row doesn't regress.
+ // // Env overrides: CRAFT_API_EGRESS_LIMIT_BYTES, CRAFT_API_EGRESS_FLUSH_SECONDS, CRAFT_API_EGRESS_TABLE
+ // // (blank = file-only), CRAFT_API_EGRESS_BUCKET_MINUTES, CRAFT_API_EGRESS_RETENTION_DAYS.
// "Egress": {
// "HostedEnv": "CIPP_HOSTED", // env var whose presence enables accounting ("" = never auto-enable)
// "BytesPerDay": 0, // per-instance daily budget in bytes (0 = accounting only, no 429s; 1 GB = 1073741824)
- // "FlushSeconds": 60 // how often the counter is written to disk (also on shutdown)
+ // "FlushSeconds": 60, // how often the counter is written to disk + mirrored to the table (also on shutdown)
+ // "TableName": "CraftEgressAccounting", // table for the bucketed history + daily audit ("" = file only)
+ // "BucketMinutes": 15, // width of each history bucket (→ 96 buckets/day)
+ // "RetentionDays": 7 // days of table rows kept before the periodic purge
// }
// },
diff --git a/tests/Craft.Tests/ApiEgressLimiterMiddlewareTests.cs b/tests/Craft.Tests/ApiEgressLimiterMiddlewareTests.cs
index cc2bd0a..c6ebac5 100644
--- a/tests/Craft.Tests/ApiEgressLimiterMiddlewareTests.cs
+++ b/tests/Craft.Tests/ApiEgressLimiterMiddlewareTests.cs
@@ -55,8 +55,9 @@ private static DefaultHttpContext UiContext()
return ctx;
}
+ // The charge flag now carries the caller's AppId (a non-empty string), not a bare bool.
private static bool Charged(HttpContext ctx) =>
- ctx.Items.TryGetValue(ApiEgressWireCounterMiddleware.ChargeItemKey, out var v) && v is true;
+ ctx.Items.TryGetValue(ApiEgressWireCounterMiddleware.ChargeItemKey, out var v) && v is string s && s.Length > 0;
private static string BodyText(HttpContext ctx)
{
@@ -70,7 +71,7 @@ private static string BodyText(HttpContext ctx)
public async Task UiCaller_OverBudget_IsNeverRejectedNorFlagged()
{
var ledger = Ledger(cap: 10);
- ledger.Record(1_000_000); // instance is way over budget
+ ledger.Record(1_000_000, "seed"); // instance is way over budget
Assert.True(ledger.ShouldReject());
var called = false;
@@ -89,7 +90,7 @@ public async Task UiCaller_OverBudget_IsNeverRejectedNorFlagged()
public async Task AnonymousCaller_IsIgnored()
{
var ledger = Ledger(cap: 10);
- ledger.Record(1_000_000);
+ ledger.Record(1_000_000, "seed");
var called = false;
var mw = Middleware(ctx => { called = true; return Task.CompletedTask; }, ledger);
@@ -142,7 +143,7 @@ public async Task ApiCaller_CapZero_FlagsButNeverRejects()
public async Task ApiCaller_OverBudget_Is429_WithRetryAfter_DownstreamSkipped_AndNotFlagged()
{
var ledger = Ledger(cap: 1000);
- ledger.Record(1000); // exactly at budget → reject
+ ledger.Record(1000, "seed"); // exactly at budget → reject
Assert.True(ledger.ShouldReject());
var called = false;
diff --git a/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs b/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs
index ee2278b..2785c1a 100644
--- a/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs
+++ b/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs
@@ -42,10 +42,12 @@ private static DefaultHttpContext Context()
return ctx;
}
- // The limiter greenlights an API request by setting this on the way through; the fake handlers below
- // do the same, since in the real pipeline the flag is set inside the counter's own invocation.
+ private const string App = "11111111-2222-3333-4444-555555555555";
+
+ // The limiter greenlights an API request by setting the caller's AppId on the way through; the fake
+ // handlers below do the same, since in the real pipeline the flag is set inside the counter's own invocation.
private static void Flag(HttpContext ctx) =>
- ctx.Items[ApiEgressWireCounterMiddleware.ChargeItemKey] = true;
+ ctx.Items[ApiEgressWireCounterMiddleware.ChargeItemKey] = App;
// ── records the flagged request's wire bytes, through either write path ───────────────────────────
@@ -192,7 +194,7 @@ public async Task LimiterGreenlight_ThenCounterRecords_AdvancesLedger()
public async Task OverBudget_LimiterSheds_CounterDoesNotBillThe429()
{
var ledger = Ledger(cap: 1000);
- ledger.Record(1000); // at budget → the limiter sheds
+ ledger.Record(1000, "seed"); // at budget → the limiter sheds
var before = ledger.CurrentBytes;
var limiter = new ApiEgressLimiterMiddleware(
diff --git a/tests/Craft.Tests/EgressLedgerTableTests.cs b/tests/Craft.Tests/EgressLedgerTableTests.cs
new file mode 100644
index 0000000..b705777
--- /dev/null
+++ b/tests/Craft.Tests/EgressLedgerTableTests.cs
@@ -0,0 +1,129 @@
+using System.Globalization;
+using Craft.Hosting;
+using Craft.Storage;
+using Microsoft.Extensions.Logging.Abstractions;
+
+namespace Craft.Tests;
+
+///
+/// The table mirror is what lets the product show egress over time and whether/when/how often the cap
+/// was hit. These pin its shape and the two resilience behaviours: per-client + instance rows are written
+/// for both the 15-minute bucket and the daily audit; a restart mid-bucket reseeds from the table so the
+/// bucket row does not regress; and rows past the retention window are purged. Uses an in-memory store —
+/// the real byte counting end to end is the perf-harness's job.
+///
+public class EgressLedgerTableTests : IDisposable
+{
+ private const string App = "11111111-2222-3333-4444-555555555555";
+ private const string App2 = "99999999-8888-7777-6666-555555555555";
+ private const string Table = "CraftEgressAccounting";
+
+ private readonly string _dir =
+ Path.Combine(Path.GetTempPath(), "craft-egress-tbl-" + Guid.NewGuid().ToString("N")[..8]);
+ private string FilePath => Path.Combine(_dir, "egress-ledger.json");
+
+ public EgressLedgerTableTests() => Directory.CreateDirectory(_dir);
+ public void Dispose() { try { Directory.Delete(_dir, true); } catch { } GC.SuppressFinalize(this); }
+
+ private EgressLedger New(FakeTableStore store, long cap, Func clock, int retentionDays = 7) =>
+ new(NullLogger.Instance, cap, flushSeconds: 60, FilePath, clock, store, Table,
+ bucketMinutes: 15, retentionDays: retentionDays);
+
+ private static long Long(StoreRow r, string prop) => Convert.ToInt64(r[prop] ?? 0L, CultureInfo.InvariantCulture);
+
+ [Fact]
+ public async Task SyncToTable_WritesPerClientAndSystem_ForBucketAndDaily()
+ {
+ var store = new FakeTableStore();
+ var clock = new DateTime(2026, 9, 13, 14, 22, 0, DateTimeKind.Utc);
+ var ledger = New(store, cap: 1_000_000, () => clock);
+
+ ledger.Record(100, App);
+ ledger.Record(250, App);
+ ledger.Record(400, App2);
+ await ledger.SyncToTableAsync(CancellationToken.None);
+
+ var rows = store.All(Table);
+
+ // Bucket rows: per client + a system aggregate, all in the same 15-min bucket (bkt_...T141500Z).
+ var appBucket = rows.Single(r => r.PartitionKey == App && r.RowKey.StartsWith(EgressTableSchema.BucketPrefix, StringComparison.Ordinal));
+ Assert.Equal(350L, Long(appBucket, EgressTableSchema.PropBytes));
+ Assert.Equal(2L, Long(appBucket, EgressTableSchema.PropRequests));
+ Assert.EndsWith("T141500Z", appBucket.RowKey); // floored to the 14:15 bucket
+
+ Assert.Equal(400L, Long(rows.Single(r => r.PartitionKey == App2 && r.RowKey.StartsWith(EgressTableSchema.BucketPrefix, StringComparison.Ordinal)), EgressTableSchema.PropBytes));
+
+ var sysBucket = rows.Single(r => r.PartitionKey == EgressTableSchema.SystemPartition && r.RowKey.StartsWith(EgressTableSchema.BucketPrefix, StringComparison.Ordinal));
+ Assert.Equal(750L, Long(sysBucket, EgressTableSchema.PropBytes)); // aggregate == sum of clients
+
+ // Daily rows: per client + a system summary carrying the cap.
+ var appDaily = rows.Single(r => r.PartitionKey == App && r.RowKey.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal));
+ Assert.Equal(350L, Long(appDaily, EgressTableSchema.PropBytes));
+ var sysDaily = rows.Single(r => r.PartitionKey == EgressTableSchema.SystemPartition && r.RowKey.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal));
+ Assert.Equal(750L, Long(sysDaily, EgressTableSchema.PropBytes));
+ Assert.Equal(1_000_000L, Long(sysDaily, EgressTableSchema.PropCapBytes));
+ Assert.Equal(true, sysDaily[EgressTableSchema.PropEnforcing]);
+ }
+
+ [Fact]
+ public async Task SystemDailyRow_CarriesShedCountAndCapReached()
+ {
+ var store = new FakeTableStore();
+ var clock = new DateTime(2026, 9, 13, 9, 0, 0, DateTimeKind.Utc);
+ var ledger = New(store, cap: 1000, () => clock);
+
+ ledger.RecordShed(App);
+ ledger.RecordShed(App);
+ await ledger.SyncToTableAsync(CancellationToken.None);
+
+ var sysDaily = store.All(Table).Single(r => r.PartitionKey == EgressTableSchema.SystemPartition && r.RowKey.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal));
+ Assert.Equal(2L, Long(sysDaily, EgressTableSchema.PropShed));
+ Assert.NotNull(sysDaily[EgressTableSchema.PropCapReachedUtc]);
+ Assert.Equal(2L, Long(store.All(Table).Single(r => r.PartitionKey == App && r.RowKey.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal)), EgressTableSchema.PropShed));
+ }
+
+ [Fact]
+ public async Task Backfill_SeedsCurrentBucket_SoRestartDoesNotRegress()
+ {
+ var store = new FakeTableStore();
+ var clock = new DateTime(2026, 9, 13, 14, 22, 0, DateTimeKind.Utc);
+
+ // A previous instance already wrote 500 bytes for App into the current bucket.
+ var bucketStart = EgressTableSchema.BucketStart(clock, 15);
+ await store.UpsertAsync(Table, EgressTableSchema.ClientBucketRow(App, bucketStart, 500, 3, 0, 1000), CancellationToken.None);
+
+ // Restart: new ledger, backfill, then serve 100 more bytes.
+ var ledger = New(store, cap: 1000, () => clock);
+ await ledger.SeedCurrentBucketAsync(CancellationToken.None);
+ ledger.Record(100, App);
+ await ledger.SyncToTableAsync(CancellationToken.None);
+
+ var appBucket = store.All(Table).Single(r => r.PartitionKey == App && r.RowKey.StartsWith(EgressTableSchema.BucketPrefix, StringComparison.Ordinal));
+ Assert.Equal(600L, Long(appBucket, EgressTableSchema.PropBytes)); // 500 seeded + 100, not 100
+ }
+
+ [Fact]
+ public async Task Purge_RemovesRowsBeyondRetention_KeepsRecent()
+ {
+ var store = new FakeTableStore();
+ var clock = new DateTime(2026, 9, 13, 12, 0, 0, DateTimeKind.Utc);
+
+ // Old rows (10 days ago) — beyond a 7-day window.
+ var oldStart = EgressTableSchema.BucketStart(clock.AddDays(-10), 15);
+ await store.UpsertAsync(Table, EgressTableSchema.ClientBucketRow(App, oldStart, 111, 1, 0, 1000), CancellationToken.None);
+ await store.UpsertAsync(Table, EgressTableSchema.ClientDailyRow(App, DateOnly.FromDateTime(clock.AddDays(-10)), 111, 1, 0, 1000), CancellationToken.None);
+ // Recent rows (now).
+ var nowStart = EgressTableSchema.BucketStart(clock, 15);
+ await store.UpsertAsync(Table, EgressTableSchema.ClientBucketRow(App, nowStart, 222, 1, 0, 1000), CancellationToken.None);
+ await store.UpsertAsync(Table, EgressTableSchema.ClientDailyRow(App, DateOnly.FromDateTime(clock), 222, 1, 0, 1000), CancellationToken.None);
+
+ Assert.Equal(4, store.Count(Table));
+
+ var ledger = New(store, cap: 1000, () => clock, retentionDays: 7);
+ await ledger.PurgeAsync(CancellationToken.None);
+
+ var remaining = store.All(Table);
+ Assert.Equal(2, remaining.Count); // the two old rows are gone
+ Assert.All(remaining, r => Assert.Contains("20260913", r.RowKey)); // only today's rows survive
+ }
+}
diff --git a/tests/Craft.Tests/EgressLedgerTests.cs b/tests/Craft.Tests/EgressLedgerTests.cs
index 70f8d3a..fe15a20 100644
--- a/tests/Craft.Tests/EgressLedgerTests.cs
+++ b/tests/Craft.Tests/EgressLedgerTests.cs
@@ -1,4 +1,5 @@
using System.Globalization;
+using System.Text.Json;
using Craft.Hosting;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging.Abstractions;
@@ -8,13 +9,18 @@ namespace Craft.Tests;
///
/// The egress ledger is the instance-wide byte counter behind the daily API bandwidth cap. Its
/// correctness is what makes the cap trustworthy in both directions: under-counting lets an abuser past
-/// the budget, over-counting throttles a well-behaved caller. These pin the four behaviours that matter
-/// — it accumulates, it rejects exactly at the budget, it resets at UTC midnight, and it survives a
-/// restart within the same day without resetting (John's hard requirement) while starting fresh on a
-/// new day. The clock is injected so rollover is testable without waiting for midnight.
+/// the budget, over-counting throttles a well-behaved caller. These pin the behaviours that matter — it
+/// accumulates, it rejects exactly at the budget, it resets at UTC midnight, it survives a restart within
+/// the same day without resetting (John's hard requirement) while starting fresh on a new day, and it
+/// keeps per-client totals + a shed count + the cap-reached stamp in the persisted file. The clock is
+/// injected so rollover is testable without waiting for midnight. Table mirroring is covered separately
+/// in .
///
public class EgressLedgerTests : IDisposable
{
+ private const string App = "11111111-2222-3333-4444-555555555555";
+ private const string App2 = "99999999-8888-7777-6666-555555555555";
+
private readonly string _dir =
Path.Combine(Path.GetTempPath(), "craft-egress-test-" + Guid.NewGuid().ToString("N")[..8]);
@@ -37,8 +43,8 @@ private EgressLedger New(long cap, Func? clock = null) =>
public void Record_Accumulates()
{
var ledger = New(cap: 0);
- ledger.Record(100);
- ledger.Record(250);
+ ledger.Record(100, App);
+ ledger.Record(250, App);
Assert.Equal(350L, ledger.CurrentBytes);
}
@@ -46,8 +52,8 @@ public void Record_Accumulates()
public void Record_IgnoresZeroAndNegative()
{
var ledger = New(cap: 0);
- ledger.Record(0);
- ledger.Record(-99);
+ ledger.Record(0, App);
+ ledger.Record(-99, App);
Assert.Equal(0L, ledger.CurrentBytes);
}
@@ -55,13 +61,13 @@ public void Record_IgnoresZeroAndNegative()
public void ShouldReject_FalseBelowCap_TrueAtOrAboveCap()
{
var ledger = New(cap: 1000);
- ledger.Record(999);
+ ledger.Record(999, App);
Assert.False(ledger.ShouldReject()); // one byte under
- ledger.Record(1);
+ ledger.Record(1, App);
Assert.True(ledger.ShouldReject()); // exactly at the cap rejects
- ledger.Record(5000);
+ ledger.Record(5000, App);
Assert.True(ledger.ShouldReject()); // and stays rejecting past it
}
@@ -69,11 +75,65 @@ public void ShouldReject_FalseBelowCap_TrueAtOrAboveCap()
public void CapZero_IsAccountingOnly_NeverRejects()
{
var ledger = New(cap: 0);
- ledger.Record(long.MaxValue / 2);
+ ledger.Record(long.MaxValue / 2, App);
Assert.False(ledger.ShouldReject()); // accounting-only rollout phase: counts, never 429s
Assert.True(ledger.CurrentBytes > 0);
}
+ // ── Per-client totals + shed count + cap-reached stamp (persisted) ────────────────────────────────
+
+ [Fact]
+ public void Record_TracksPerClientTotals_InFile()
+ {
+ var ledger = New(cap: 0);
+ ledger.Record(100, App);
+ ledger.Record(250, App);
+ ledger.Record(400, App2);
+ ledger.Flush();
+
+ var state = ReadFile();
+ Assert.Equal(750L, state.RootElement.GetProperty("bytes").GetInt64());
+ var clients = state.RootElement.GetProperty("clients");
+ Assert.Equal(350L, clients.GetProperty(App).GetProperty("bytes").GetInt64());
+ Assert.Equal(2, clients.GetProperty(App).GetProperty("requests").GetInt32());
+ Assert.Equal(400L, clients.GetProperty(App2).GetProperty("bytes").GetInt64());
+ }
+
+ [Fact]
+ public void RecordShed_CountsAndStampsCapReached()
+ {
+ var clock = new DateTime(2026, 9, 11, 12, 0, 0, DateTimeKind.Utc);
+ var ledger = New(cap: 1000, () => clock);
+ ledger.RecordShed(App);
+ ledger.RecordShed(App);
+ ledger.Flush();
+
+ Assert.Equal(2L, ledger.ShedRequests);
+ var state = ReadFile();
+ Assert.Equal(2L, state.RootElement.GetProperty("shedRequests").GetInt64());
+ Assert.Equal(2, state.RootElement.GetProperty("clients").GetProperty(App).GetProperty("shed").GetInt32());
+ Assert.True(state.RootElement.TryGetProperty("capReachedUtc", out var cr) && cr.ValueKind != JsonValueKind.Null);
+ Assert.Equal(1000L, state.RootElement.GetProperty("capBytes").GetInt64());
+ }
+
+ [Fact]
+ public void PerClientTotals_RestoredAcrossRestart_SameDay()
+ {
+ var clock = new DateTime(2026, 9, 11, 6, 0, 0, DateTimeKind.Utc);
+ var first = New(cap: 0, () => clock);
+ first.Record(111, App);
+ first.Record(222, App2);
+ first.RecordShed(App);
+ first.Flush();
+
+ var restarted = New(cap: 0, () => clock);
+ restarted.Flush(); // re-persist the loaded state
+ var state = ReadFile();
+ Assert.Equal(333L, restarted.CurrentBytes);
+ Assert.Equal(1L, restarted.ShedRequests);
+ Assert.Equal(111L, state.RootElement.GetProperty("clients").GetProperty(App).GetProperty("bytes").GetInt64());
+ }
+
// ── UTC midnight rollover ───────────────────────────────────────────────────────────────────────
[Fact]
@@ -82,13 +142,13 @@ public void Rollover_ResetsCounterOnNewUtcDay()
var clock = new DateTime(2026, 9, 11, 12, 0, 0, DateTimeKind.Utc);
var ledger = New(cap: 1000, () => clock);
- ledger.Record(900);
+ ledger.Record(900, App);
Assert.Equal(900L, ledger.CurrentBytes);
clock = clock.AddDays(1); // cross UTC midnight
Assert.Equal(0L, ledger.CurrentBytes); // reads roll the day over
Assert.False(ledger.ShouldReject());
- ledger.Record(50);
+ ledger.Record(50, App);
Assert.Equal(50L, ledger.CurrentBytes); // new day counts from zero
}
@@ -103,7 +163,6 @@ public void SecondsToNextUtcMidnight_IsTimeUntilMidnight()
[Fact]
public void SecondsToNextUtcMidnight_FlooredAtOne()
{
- // A hair before midnight must never advertise Retry-After: 0 (a hot-loop invitation).
var clock = new DateTime(2026, 9, 11, 23, 59, 59, 900, DateTimeKind.Utc);
var ledger = New(cap: 0, () => clock);
Assert.Equal(1, ledger.SecondsToNextUtcMidnight());
@@ -125,10 +184,9 @@ public void Flush_ThenReload_RestoresSameDayTotal()
var clock = new DateTime(2026, 9, 11, 6, 0, 0, DateTimeKind.Utc);
var first = New(cap: 5000, () => clock);
- first.Record(1234);
+ first.Record(1234, App);
first.Flush();
- // A "restart": a brand-new ledger over the same file, same UTC day.
var restarted = New(cap: 5000, () => clock);
Assert.Equal(1234L, restarted.CurrentBytes);
}
@@ -138,7 +196,7 @@ public void Reload_OnNewUtcDay_StartsFresh()
{
var day1 = new DateTime(2026, 9, 11, 6, 0, 0, DateTimeKind.Utc);
var first = New(cap: 5000, () => day1);
- first.Record(1234);
+ first.Record(1234, App);
first.Flush();
var day2 = day1.AddDays(1);
@@ -146,6 +204,16 @@ public void Reload_OnNewUtcDay_StartsFresh()
Assert.Equal(0L, restarted.CurrentBytes); // yesterday's total is not carried into today
}
+ [Fact]
+ public void Load_V1File_RestoresBytes()
+ {
+ // A pre-per-client (v1) file has only dateUtc + bytes. It must still load.
+ var today = DateTime.UtcNow.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture);
+ File.WriteAllText(FilePath, $"{{\"dateUtc\":\"{today}\",\"bytes\":4242}}");
+ var ledger = New(cap: 0);
+ Assert.Equal(4242L, ledger.CurrentBytes);
+ }
+
[Fact]
public void Flush_SkippedWhenNothingChanged()
{
@@ -158,7 +226,7 @@ public void Flush_SkippedWhenNothingChanged()
public void Flush_LeavesNoTempFileBehind()
{
var ledger = New(cap: 1000);
- ledger.Record(10);
+ ledger.Record(10, App);
ledger.Flush();
Assert.True(File.Exists(FilePath));
Assert.False(File.Exists(FilePath + ".tmp")); // temp+rename completed cleanly
@@ -192,7 +260,7 @@ public async Task ConcurrentRecords_SumExactly()
await Task.WhenAll(Enumerable.Range(0, threads).Select(_ => Task.Run(() =>
{
- for (var i = 0; i < perThread; i++) ledger.Record(7);
+ for (var i = 0; i < perThread; i++) ledger.Record(7, App);
})));
Assert.Equal((long)threads * perThread * 7, ledger.CurrentBytes);
@@ -201,11 +269,9 @@ await Task.WhenAll(Enumerable.Range(0, threads).Select(_ => Task.Run(() =>
[Fact]
public async Task HostedService_FlushesOnShutdown()
{
- // The user's explicit requirement: a graceful stop must persist the day's total. Start the
- // background service, record, stop it — the final flush on shutdown must write the counter.
var clock = new DateTime(2026, 9, 11, 10, 0, 0, DateTimeKind.Utc);
var ledger = New(cap: 5000, () => clock);
- ledger.Record(2222);
+ ledger.Record(2222, App);
await ((IHostedService)ledger).StartAsync(CancellationToken.None);
await ((IHostedService)ledger).StopAsync(CancellationToken.None); // triggers the final flush
@@ -213,4 +279,6 @@ public async Task HostedService_FlushesOnShutdown()
var reloaded = New(cap: 5000, () => clock);
Assert.Equal(2222L, reloaded.CurrentBytes);
}
+
+ private JsonDocument ReadFile() => JsonDocument.Parse(File.ReadAllText(FilePath));
}
diff --git a/tests/Craft.Tests/FakeTableStore.cs b/tests/Craft.Tests/FakeTableStore.cs
new file mode 100644
index 0000000..f70841a
--- /dev/null
+++ b/tests/Craft.Tests/FakeTableStore.cs
@@ -0,0 +1,78 @@
+using System.Collections.Concurrent;
+using Craft.Storage;
+
+namespace Craft.Tests;
+
+///
+/// In-memory for unit tests. Deliberately IGNORES the optional
+/// filter on queries — exactly what the interface permits a backend to do — so tests also verify
+/// that callers re-apply their own predicate to the full result. Not thread-safe beyond the concurrent
+/// dictionaries; tests drive it single-threaded.
+///
+internal sealed class FakeTableStore : ICraftTableStore
+{
+ // table -> (partitionKey|rowKey) -> row
+ private readonly ConcurrentDictionary> _tables = new();
+
+ private static string Key(string pk, string rk) => pk + "" + rk;
+ private ConcurrentDictionary Table(string table) =>
+ _tables.GetOrAdd(table, _ => new ConcurrentDictionary());
+
+ /// All rows currently in a table (test helper).
+ public IReadOnlyList All(string table) => Table(table).Values.ToList();
+
+ public int Count(string table) => Table(table).Count;
+
+ public Task PingAsync(CancellationToken ct = default) => Task.CompletedTask;
+
+ public Task EnsureTableAsync(string table, CancellationToken ct = default) { Table(table); return Task.CompletedTask; }
+
+ public Task UpsertAsync(string table, StoreRow row, CancellationToken ct = default)
+ {
+ Table(table)[Key(row.PartitionKey, row.RowKey)] = row;
+ return Task.CompletedTask;
+ }
+
+ public Task UpsertBatchAsync(string table, string partitionKey, IReadOnlyList rows, CancellationToken ct = default)
+ {
+ foreach (var r in rows) Table(table)[Key(r.PartitionKey, r.RowKey)] = r;
+ return Task.CompletedTask;
+ }
+
+ public Task TryReplaceBatchAsync(string table, string partitionKey, IReadOnlyList rows, CancellationToken ct = default)
+ {
+ foreach (var r in rows) Table(table)[Key(r.PartitionKey, r.RowKey)] = r;
+ return Task.FromResult(true);
+ }
+
+ public Task GetAsync(string table, string partitionKey, string rowKey, CancellationToken ct = default) =>
+ Task.FromResult(Table(table).TryGetValue(Key(partitionKey, rowKey), out var r) ? r : null);
+
+ public async IAsyncEnumerable QueryPartitionAsync(string table, string partitionKey,
+ [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default)
+ {
+ foreach (var r in Table(table).Values.Where(r => r.PartitionKey == partitionKey)) { yield return r; }
+ await Task.CompletedTask;
+ }
+
+ public async IAsyncEnumerable QueryTableAsync(string table,
+ [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default)
+ {
+ // Snapshot so deletes during enumeration (the purge path) don't throw.
+ foreach (var r in Table(table).Values.ToList()) { yield return r; }
+ await Task.CompletedTask;
+ }
+
+ public Task DeleteAsync(string table, string partitionKey, string rowKey, CancellationToken ct = default)
+ {
+ Table(table).TryRemove(Key(partitionKey, rowKey), out _);
+ return Task.CompletedTask;
+ }
+
+ public Task DeletePartitionAsync(string table, string partitionKey, CancellationToken ct = default)
+ {
+ foreach (var k in Table(table).Where(kv => kv.Value.PartitionKey == partitionKey).Select(kv => kv.Key).ToList())
+ Table(table).TryRemove(k, out _);
+ return Task.CompletedTask;
+ }
+}