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/Configuration/EgressLimitSettings.cs b/Services/Configuration/EgressLimitSettings.cs
index 5d3a7a8..edbcd08 100644
--- a/Services/Configuration/EgressLimitSettings.cs
+++ b/Services/Configuration/EgressLimitSettings.cs
@@ -51,9 +51,35 @@ public class EgressLimitSettings
///
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 4f37b8d..ba01351 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
@@ -40,31 +49,27 @@ 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;
}
- // 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: 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);
}
private async Task RejectAsync(HttpContext context)
diff --git a/Services/Hosting/ApiEgressWireCounterMiddleware.cs b/Services/Hosting/ApiEgressWireCounterMiddleware.cs
new file mode 100644
index 0000000..e38b735
--- /dev/null
+++ b/Services/Hosting/ApiEgressWireCounterMiddleware.cs
@@ -0,0 +1,68 @@
+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:
+ /// 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";
+
+ 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 string appId && appId.Length > 0)
+ _ledger.Record(counting.BytesWritten, appId);
+ }
+ }
+}
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/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/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..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
// }
// },
@@ -571,5 +580,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/