From e34c0a3cad09faacdf400f20695a08debcadea82 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Thu, 17 Sep 2026 17:14:01 +0800 Subject: [PATCH 1/6] fix(http): handle aborted PowerShell requests Link HTTP request cancellation to PowerShell execution so client disconnects stop the pipeline promptly. Return a 499-style result for aborted requests, skip cache/response writes for partial output, and still invalidate write caches because a non-GET operation may have already applied. --- .../Endpoints/PowerShellDispatchEndpoint.cs | 11 +++++++ .../PowerShellHost/PowerShellRunnerService.cs | 29 +++++++++++++++---- 2 files changed, 35 insertions(+), 5 deletions(-) diff --git a/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs b/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs index 72394d8..9121ce4 100644 --- a/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs +++ b/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs @@ -100,8 +100,19 @@ await ServeCachedAsync(context, cache, psRunner, orchestrator, logger, var result = await psRunner.ExecuteHttpScript(endpoint, context); // Writes can invalidate anything, so they clear the cache regardless of endpoint name. + // Do this even on a client abort: a write may have partially applied before the + // connection dropped, so dropping the possibly-stale cache is the safe choice. if (context.Request.Method != "GET") InvalidateForWrite(context, cache); + // Client hung up mid-request (navigated away / hit Cancel). The pipeline was already + // stopped and the worker reclaimed; there is no live connection to write to, and a + // partial result must not be cached. Bail before the cache write and response write. + if (context.RequestAborted.IsCancellationRequested) + { + logger.LogInformation("[HTTP] /API/{Endpoint} cancelled by client", endpoint); + return; + } + if (useCache && cacheKey is not null && result.StatusCode is >= 200 and < 400) await cache.Set(cacheKey, result); diff --git a/Services/PowerShellHost/PowerShellRunnerService.cs b/Services/PowerShellHost/PowerShellRunnerService.cs index f162254..4e7c68d 100644 --- a/Services/PowerShellHost/PowerShellRunnerService.cs +++ b/Services/PowerShellHost/PowerShellRunnerService.cs @@ -231,7 +231,7 @@ public async Task ExecuteHttpScript(string route, HttpContext http if (!DispatchProfiler.Enabled) { var req = await BuildRequestObject(httpContext); - return await ExecuteHttpScriptInternal(route, req, isHttp: true); + return await ExecuteHttpScriptInternal(route, req, isHttp: true, clientAborted: httpContext.RequestAborted); } // Profiling path: time request marshaling + the runner-side segments (checkout/invoke/extract). @@ -241,7 +241,8 @@ public async Task ExecuteHttpScript(string route, HttpContext http var request = await BuildRequestObject(httpContext); var marshalTicks = Stopwatch.GetTimestamp() - mStart; var timing = new DispatchTiming(); - var result = await ExecuteHttpScriptInternal(route, request, isHttp: true, timing); + var result = await ExecuteHttpScriptInternal(route, request, isHttp: true, timing, + clientAborted: httpContext.RequestAborted); DispatchProfiler.Record(marshalTicks, timing.CheckoutTicks, timing.InvokeTicks, timing.ExtractTicks, Stopwatch.GetTimestamp() - totalStart); return result; @@ -258,7 +259,7 @@ public async Task ExecuteHttpScript(string route, Hashtable reques } private async Task ExecuteHttpScriptInternal(string route, Hashtable request, bool isHttp, - DispatchTiming? timing = null) + DispatchTiming? timing = null, CancellationToken clientAborted = default) { var sw = Stopwatch.StartNew(); var entry = _repo.GetByRoute(route); @@ -365,9 +366,16 @@ private async Task ExecuteHttpScriptInternal(string route, Hashtab worker.Streams.Verbose.DataAdded += onVerbose; var timeoutSeconds = isHttp ? _workerSettings.HttpTimeoutSeconds : _workerSettings.BgTimeoutSeconds; - using var cts = timeoutSeconds > 0 + using var timeoutCts = timeoutSeconds > 0 ? new CancellationTokenSource(TimeSpan.FromSeconds(timeoutSeconds)) : null; + // Link the client-disconnect token (present on live HTTP requests) with the timeout so + // either one stops the pipeline. When the client hangs up — navigated away or hit Cancel — + // RequestAborted fires, InvokeAsync calls PowerShell.Stop(), and the worker is freed instead + // of paging on for a response nobody is waiting for. Background cache refresh passes default + // (no live client), so only the timeout applies there. + using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource( + timeoutCts?.Token ?? CancellationToken.None, clientAborted); // When Scripts.HttpHandler is set, ALL HTTP routes dispatch through that single // function instead of invoking the route's function directly. The endpoint name @@ -378,7 +386,7 @@ private async Task ExecuteHttpScriptInternal(string route, Hashtab ? _scriptsSettings.HttpHandler : entry.FunctionName; var invokeStart = timing != null ? Stopwatch.GetTimestamp() : 0; - var results = await worker.InvokeAsync(targetFunction, parameters, cts?.Token ?? default); + var results = await worker.InvokeAsync(targetFunction, parameters, linkedCts.Token); if (timing != null) timing.InvokeTicks = Stopwatch.GetTimestamp() - invokeStart; var extractStart = timing != null ? Stopwatch.GetTimestamp() : 0; @@ -397,6 +405,17 @@ private async Task ExecuteHttpScriptInternal(string route, Hashtab catch (OperationCanceledException) when (sw.ElapsedMilliseconds > 0) { sw.Stop(); + + // Client hung up (navigated away / hit Cancel) vs the request exceeding its time budget. + // A cancel is normal and expected — log it quietly and return 499 (never actually written; + // the connection is gone) so the dispatcher skips caching this partial result. + if (clientAborted.IsCancellationRequested) + { + _logger.LogInformation("[{Pool}] {Function} cancelled by client after {Ms}ms", + poolLabel, entry?.FunctionName ?? route, sw.ElapsedMilliseconds); + return new ScriptResult { StatusCode = 499, Body = string.Empty }; + } + var timeoutSeconds = isHttp ? _workerSettings.HttpTimeoutSeconds : _workerSettings.BgTimeoutSeconds; _logger.LogWarning("[{Pool}] {Function} timed out after {Ms}ms (limit: {Limit}s)", poolLabel, entry?.FunctionName ?? route, sw.ElapsedMilliseconds, timeoutSeconds); From 627cd05d67b29694efa0d631e5f6bfe9b8ad5b81 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Thu, 17 Sep 2026 17:23:31 +0800 Subject: [PATCH 2/6] feat: gate http worker cancelation behind only get requests --- Services/PowerShellHost/PowerShellRunnerService.cs | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/Services/PowerShellHost/PowerShellRunnerService.cs b/Services/PowerShellHost/PowerShellRunnerService.cs index 4e7c68d..66d3354 100644 --- a/Services/PowerShellHost/PowerShellRunnerService.cs +++ b/Services/PowerShellHost/PowerShellRunnerService.cs @@ -228,10 +228,17 @@ public Dictionary DiscoverHttpEndpoints() /// public async Task ExecuteHttpScript(string route, HttpContext httpContext) { + // Only reads honour client-disconnect cancellation. A write (POST/PUT/DELETE/PATCH) runs to + // completion even if the client hangs up — a half-applied mutation with nobody listening is + // worse than one that finished. Reads mutate nothing, so aborting them mid-flight is safe. + var clientAborted = HttpMethods.IsGet(httpContext.Request.Method) + ? httpContext.RequestAborted + : CancellationToken.None; + if (!DispatchProfiler.Enabled) { var req = await BuildRequestObject(httpContext); - return await ExecuteHttpScriptInternal(route, req, isHttp: true, clientAborted: httpContext.RequestAborted); + return await ExecuteHttpScriptInternal(route, req, isHttp: true, clientAborted: clientAborted); } // Profiling path: time request marshaling + the runner-side segments (checkout/invoke/extract). @@ -242,7 +249,7 @@ public async Task ExecuteHttpScript(string route, HttpContext http var marshalTicks = Stopwatch.GetTimestamp() - mStart; var timing = new DispatchTiming(); var result = await ExecuteHttpScriptInternal(route, request, isHttp: true, timing, - clientAborted: httpContext.RequestAborted); + clientAborted: clientAborted); DispatchProfiler.Record(marshalTicks, timing.CheckoutTicks, timing.InvokeTicks, timing.ExtractTicks, Stopwatch.GetTimestamp() - totalStart); return result; From e16ebc4e64a534d4dc8cac945d6f8595a09cccdd Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Fri, 18 Sep 2026 17:10:46 +0800 Subject: [PATCH 3/6] Add Azure Table large-entity splitting Store oversized table entities across multiple properties and rows, then reassemble them on read so task parameters and other payloads no longer fail Azure Table size limits. This includes stale part cleanup, conditional-write safety, and regression tests covering large payloads, batch writes, deletes, and orchestrator task rehydration. --- Services/Storage/AzureTableStore.cs | 350 +++++++- Services/Storage/EntitySplitter.cs | 789 ++++++++++++++++++ Services/Storage/IncompleteEntityException.cs | 29 + .../AzureTableStoreLargeEntityTests.cs | 240 ++++++ tests/Craft.Tests/EntitySplitterTests.cs | 140 ++++ ...estratorTaskLargeParametersAzuriteTests.cs | 109 +++ 6 files changed, 1612 insertions(+), 45 deletions(-) create mode 100644 Services/Storage/EntitySplitter.cs create mode 100644 Services/Storage/IncompleteEntityException.cs create mode 100644 tests/Craft.Tests/AzureTableStoreLargeEntityTests.cs create mode 100644 tests/Craft.Tests/EntitySplitterTests.cs create mode 100644 tests/Craft.Tests/OrchestratorTaskLargeParametersAzuriteTests.cs diff --git a/Services/Storage/AzureTableStore.cs b/Services/Storage/AzureTableStore.cs index e994093..2bd6a25 100644 --- a/Services/Storage/AzureTableStore.cs +++ b/Services/Storage/AzureTableStore.cs @@ -25,8 +25,8 @@ public sealed class AzureTableStore : ICraftTableStore private readonly TableClientOptions _clientOptions; // Azure Table transaction limits: at most 100 entities, all sharing a partition key, ~4 MB total. + // The payload budget lives in EntitySplitter.MaxTransactionPayload, used by SubmitSizedAsync. private const int MaxBatch = 100; - private const int MaxBatchChars = 1_600_000; // ≈3.2 MB UTF-16, safely under the 4 MB cap private static readonly HashSet s_systemKeys = new(StringComparer.Ordinal) { @@ -150,39 +150,57 @@ private async Task RecreateTableAsync(string table, CancellationToken ct) public async Task UpsertAsync(string table, StoreRow row, CancellationToken ct = default) { var client = Client(table); - var entity = ToEntity(row); - try - { - await client.UpsertEntityAsync(entity, TableUpdateMode.Replace, ct); - } - catch (RequestFailedException ex) when (IsTableNotFound(ex)) + var split = EntitySplitter.Split(ToEntity(row)); + + // Fast path: the entity fits one row unchanged — a single unconditional upsert, exactly as + // before large-entity splitting existed. A key that WAS split earlier and is now small leaves + // its extra "{RowKey}-part{n}" rows behind, but Reassemble's plain-row precedence means the + // read still returns this value, and the next engaged write of the key removes them. + if (!split.Engaged) { - await RecreateTableAsync(table, ct); - await client.UpsertEntityAsync(entity, TableUpdateMode.Replace, ct); + var entity = split.Rows[0]; + try + { + await client.UpsertEntityAsync(entity, TableUpdateMode.Replace, ct); + } + catch (RequestFailedException ex) when (IsTableNotFound(ex)) + { + await RecreateTableAsync(table, ct); + await client.UpsertEntityAsync(entity, TableUpdateMode.Replace, ct); + } + return; } + + var actions = split.Rows + .Select(r => new TableTransactionAction(TableTransactionActionType.UpsertReplace, r)) + .ToList(); + await SubmitSizedAsync(table, client, actions, ct); + await RemoveStalePartRowsAsync(table, row.PartitionKey, row.RowKey, + new HashSet(split.Rows.Select(r => r.RowKey), StringComparer.Ordinal), ct); } public async Task UpsertBatchAsync(string table, string partitionKey, IReadOnlyList rows, CancellationToken ct = default) { var client = Client(table); - var batch = new List(MaxBatch); - var chars = 0; + var actions = new List(rows.Count); + var engaged = new List<(string PartitionKey, string RowKey, HashSet Live)>(); foreach (var row in rows) { - var rowChars = EstimateChars(row); - if (batch.Count > 0 && (batch.Count >= MaxBatch || chars + rowChars > MaxBatchChars)) - { - await SubmitAsync(table, client, batch, ct); - batch.Clear(); - chars = 0; - } - batch.Add(new TableTransactionAction(TableTransactionActionType.UpsertReplace, ToEntity(row))); - chars += rowChars; + var split = EntitySplitter.Split(ToEntity(row)); + foreach (var r in split.Rows) + actions.Add(new TableTransactionAction(TableTransactionActionType.UpsertReplace, r)); + if (split.Engaged) + engaged.Add((row.PartitionKey, row.RowKey, + new HashSet(split.Rows.Select(r => r.RowKey), StringComparer.Ordinal))); } - if (batch.Count > 0) - await SubmitAsync(table, client, batch, ct); + await SubmitSizedAsync(table, client, actions, ct); + + // Only an engaged (split) write can leave stale part rows behind, so the common all-small batch + // does no extra reads at all. + foreach (var (pk, rk, live) in engaged) + await RemoveStalePartRowsAsync(table, pk, rk, live, ct); } public async Task TryReplaceBatchAsync(string table, string partitionKey, IReadOnlyList rows, @@ -190,13 +208,9 @@ public async Task TryReplaceBatchAsync(string table, string partitionKey, { if (rows.Count == 0) return true; - // One transaction, so one round-trip and one atomic outcome. A caller claiming more than a - // transaction can hold would silently get partial application, which for a claim means rows - // marked as owned by a worker that never receives them. - if (rows.Count > MaxBatch) - throw new ArgumentException($"Conditional batch is limited to {MaxBatch} rows, got {rows.Count}.", nameof(rows)); - var actions = new List(rows.Count); + var engaged = new List<(string PartitionKey, string RowKey, HashSet Live)>(); + foreach (var row in rows) { // A row with no ETag was never read from storage, so there is nothing to guard against and @@ -204,14 +218,32 @@ public async Task TryReplaceBatchAsync(string table, string partitionKey, if (string.IsNullOrEmpty(row.ETag)) throw new ArgumentException($"Row {row.PartitionKey}/{row.RowKey} has no ETag to guard the write.", nameof(rows)); - actions.Add(new TableTransactionAction( - TableTransactionActionType.UpdateReplace, ToEntity(row), new ETag(row.ETag))); + var split = EntitySplitter.Split(ToEntity(row)); + foreach (var r in split.Rows) + { + // The concurrency guard belongs on the entity's own row — splitting never changes that + // RowKey. The extra "{RowKey}-part{n}" rows carry no independent token and ride along + // unconditionally in the SAME atomic transaction, so the claim stays all-or-nothing. + if (string.Equals(r.RowKey, row.RowKey, StringComparison.Ordinal)) + actions.Add(new TableTransactionAction(TableTransactionActionType.UpdateReplace, r, new ETag(row.ETag))); + else + actions.Add(new TableTransactionAction(TableTransactionActionType.UpsertReplace, r)); + } + if (split.Engaged) + engaged.Add((row.PartitionKey, row.RowKey, + new HashSet(split.Rows.Select(r => r.RowKey), StringComparer.Ordinal))); } + // One transaction, so one round-trip and one atomic outcome. A caller claiming more than a + // transaction can hold would silently get partial application, which for a claim means rows + // marked as owned by a worker that never receives them. Splitting a large guarded row can turn + // one logical claim into several physical rows, so the cap is checked after splitting. + if (actions.Count > MaxBatch) + throw new ArgumentException($"Conditional batch is limited to {MaxBatch} rows, got {actions.Count} after large-entity splitting.", nameof(rows)); + try { await Client(table).SubmitTransactionAsync(actions, ct); - return true; } catch (RequestFailedException ex) when (IsTableNotFound(ex)) { @@ -227,6 +259,14 @@ public async Task TryReplaceBatchAsync(string table, string partitionKey, // Deliberately NO fallback to unconditional upserts: that is what would steal the row. return false; } + + // The claim landed. Remove any part rows left by an earlier, larger version of a claimed entity + // so a later read cannot merge stale fragments. Best-effort and post-commit: the claim already + // succeeded and stale parts only ever affect a subsequent read. + foreach (var (pk, rk, live) in engaged) + await RemoveStalePartRowsAsync(table, pk, rk, live, ct); + + return true; } private async Task SubmitAsync(string table, TableClient client, List batch, CancellationToken ct) @@ -311,10 +351,10 @@ public async Task DeleteBatchAsync(string table, string partitionKey, IReadOnlyL public async Task GetAsync(string table, string partitionKey, string rowKey, CancellationToken ct = default) { + TableEntity entity; try { - var response = await Client(table).GetEntityAsync(partitionKey, rowKey, cancellationToken: ct); - return ToRow(response.Value); + entity = (await Client(table).GetEntityAsync(partitionKey, rowKey, cancellationToken: ct)).Value; } catch (RequestFailedException ex) when (IsTableNotFound(ex)) { @@ -327,6 +367,21 @@ public async Task DeleteBatchAsync(string table, string partitionKey, IReadOnlyL { return null; } + + // No split markers → this row is the whole entity. The hot path, unchanged. + if (!HasSplitMarkers(entity)) + return ToRow(entity); + + // A column-only chunked row reassembles from itself; a cross-row root needs its sibling + // "{RowKey}-part{n}" rows fetched first. Try in place, and only re-query the partition when this + // one row is not the whole entity. + IncompleteEntityException? incomplete = null; + var self = EntitySplitter.Reassemble(new[] { entity }, OnReassemblyWarning, ex => incomplete = ex).FirstOrDefault(); + if (incomplete is null && self is not null) + return ToRow(self); + + var full = await ReadLogicalEntityAsync(table, partitionKey, rowKey, ct); + return full is null ? null : ToRow(full); } /// @@ -367,14 +422,14 @@ public async IAsyncEnumerable QueryPartitionAsync(string table, string [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default) { var filter = $"PartitionKey eq '{Escape(partitionKey)}'"; - await foreach (var entity in EnumerateAsync(table, () => Client(table).QueryAsync(filter: filter, cancellationToken: ct), ct)) + await foreach (var entity in StreamReassembledAsync(table, () => Client(table).QueryAsync(filter: filter, cancellationToken: ct), ct)) yield return ToRow(entity); } public async IAsyncEnumerable QueryTableAsync(string table, [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default) { - await foreach (var entity in EnumerateAsync(table, () => Client(table).QueryAsync(cancellationToken: ct), ct)) + await foreach (var entity in StreamReassembledAsync(table, () => Client(table).QueryAsync(cancellationToken: ct), ct)) yield return ToRow(entity); } @@ -386,7 +441,7 @@ public async IAsyncEnumerable QueryTableAsync(string table, public async IAsyncEnumerable QueryTableAsync(string table, string? filter, [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default) { - await foreach (var entity in EnumerateAsync(table, () => Client(table).QueryAsync(filter: filter, cancellationToken: ct), ct)) + await foreach (var entity in StreamReassembledAsync(table, () => Client(table).QueryAsync(filter: filter, cancellationToken: ct), ct)) yield return ToRow(entity); } @@ -394,6 +449,11 @@ public async IAsyncEnumerable QueryTableAsync(string table, string? fi /// The filtered scan with a $select, for callers that want keys and a stamp rather than the /// row. The retention sweep reads every Results row's partition this way, and a Results row is a /// 64 KiB chunk of payload it has no use for. + /// + /// Unlike the full-row scans, this projected path does NOT reassemble large entities: the projection + /// strips the split markers reassembly needs, and its only caller wants the physical keys anyway (it + /// groups by partition to find and delete whole orphaned partitions, part rows included). Reassembly + /// here would drop every split entity as "incomplete" and hide the very keys the sweep must see. /// public async IAsyncEnumerable QueryTableAsync(string table, string? filter, IReadOnlyList? properties, [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default) @@ -412,6 +472,14 @@ public async Task DeleteAsync(string table, string partitionKey, string rowKey, { // Already gone — not an error. } + + // If this entity had been split across rows, its "{rowKey}-part{n}" rows are addressed by + // nothing else, so a plain delete would orphan them — and a later partition scan would then + // resurrect the deleted entity by reassembling those leftover parts. Remove them too (an empty + // "live" set marks every part row stale). This is a keys-only scan that finds nothing for the + // common unsplit row. + await RemoveStalePartRowsAsync(table, partitionKey, rowKey, + new HashSet(StringComparer.Ordinal), ct); } public async Task DeletePartitionAsync(string table, string partitionKey, CancellationToken ct = default) @@ -440,6 +508,207 @@ public async Task DeletePartitionAsync(string table, string partitionKey, Cancel } } + // ── Large-entity splitting ───────────────────────────────────────────────── + // + // Azure Table Storage caps a string property at 64 KiB (32K UTF-16 units) and an entity at 1 MiB. + // A host payload that exceeds either — a scheduled task whose Parameters embed a whole policy + // template, say — cannot be stored as one entity: the write 400s with PropertyValueTooLarge and the + // row is lost, so the task can never be dispatched. EntitySplitter (vendored from AzBobbyTables via + // CIPP.TableClient) breaks such an entity across extra properties and, if still too big, extra rows, + // and reassembles it on read. This is contained here for the same reason the rest of the Azure + // specifics are: nothing above ICraftTableStore knows a row was ever split. + // + // The orchestrator Results table does its own hand-rolled chunking with the marker names + // "OriginalEntityId"/"PartIndex"; EntitySplitter deliberately uses "_Craft"-prefixed marker names so + // the two schemes never touch. A Results part row is therefore invisible to the reassembly below + // (its own reader still sees the physical rows it expects), and an already-chunked Results property + // is under the size limit so it is never re-split. + + private void OnReassemblyWarning(string message) => + _logger?.LogWarning("[TableStore] {Message}", message); + + private static bool HasSplitMarkers(TableEntity entity) => + entity.ContainsKey(EntitySplitter.SplitOverPropsKey) || + entity.ContainsKey(EntitySplitter.OriginalEntityIdKey) || + entity.ContainsKey(EntitySplitter.PartCountKey); + + /// + /// Submit upsert actions grouped by partition and packed by BOTH count and estimated payload size. + /// Split rows each approach the per-row budget, so 100 of them in one transaction would blow the + /// batch's ~4 MB payload cap; the size bound is what a plain count cannot see. Delegates each packed + /// batch to so the table-recreate and per-entity fallback still apply. + /// + private async Task SubmitSizedAsync(string table, TableClient client, List actions, + CancellationToken ct) + { + foreach (var group in actions.GroupBy(a => a.Entity.PartitionKey)) + { + var batch = new List(); + long size = 0; + foreach (var action in group) + { + var actionSize = action.Entity is TableEntity te ? EntitySplitter.EstimateEntitySize(te) : 0; + if (batch.Count > 0 && (batch.Count >= MaxBatch || size + actionSize > EntitySplitter.MaxTransactionPayload)) + { + await SubmitAsync(table, client, batch, ct); + batch.Clear(); + size = 0; + } + batch.Add(action); + size += actionSize; + } + if (batch.Count > 0) + await SubmitAsync(table, client, batch, ct); + } + } + + /// + /// Stream a query's rows, reassembling large entities, without materialising the whole result set: + /// plain rows (the overwhelming majority) are yielded as they page in, and only split-part rows — + /// those carrying a cross-row or column-split marker — are buffered, their missing siblings fetched, + /// and the group reassembled at the end. A plain row supersedes leftover parts of the same identity, + /// so a reassembled group whose identity a plain row already emitted is dropped, matching + /// EntitySplitter's plain-row precedence. This is what keeps the Results streaming reader from having + /// to hold a 50–150 MB run in memory just because reassembly exists. + /// + private async IAsyncEnumerable StreamReassembledAsync(string table, + Func> query, + [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken ct = default) + { + var plainKeys = new HashSet<(string, string)>(); + List? parts = null; + + await foreach (var entity in EnumerateAsync(table, query, ct)) + { + if (entity.ContainsKey(EntitySplitter.OriginalEntityIdKey) || entity.ContainsKey(EntitySplitter.SplitOverPropsKey)) + { + (parts ??= new List()).Add(entity); + } + else + { + plainKeys.Add((entity.PartitionKey, entity.RowKey)); + yield return entity; + } + } + + if (parts is null) yield break; + + var recovered = await RecoverMissingPartRowsAsync(table, parts, ct); + foreach (var reassembled in EntitySplitter.Reassemble(recovered, OnReassemblyWarning, inc => OnReassemblyWarning(inc.Message))) + { + if (plainKeys.Contains((reassembled.PartitionKey, reassembled.RowKey))) continue; + yield return reassembled; + } + } + + /// + /// Fetch the rows of any entity in that is only partly present (its query + /// filter had no reason to match every "{RowKey}-part{n}" row it was split over) and return the set + /// with them added. Part rows are matched by a RowKey range within the entity's partition — an index + /// seek that needs no marker property to be selected. + /// + private async Task> RecoverMissingPartRowsAsync(string table, List rows, + CancellationToken ct) + { + var incomplete = EntitySplitter.FindIncompleteGroups(rows); + if (incomplete.Count == 0) return rows; + + var seen = new HashSet<(string, string)>(rows.Select(r => (r.PartitionKey, r.RowKey))); + + foreach (var partitionGroup in incomplete.GroupBy(g => g.PartitionKey, StringComparer.Ordinal)) + { + var partitionClause = $"PartitionKey eq '{Escape(partitionGroup.Key)}'"; + foreach (var entityId in partitionGroup.Select(g => g.EntityId).Distinct(StringComparer.Ordinal)) + { + var filter = $"{partitionClause} and {BuildRowKeyPrefixClause(entityId)}"; + await foreach (var row in Client(table).QueryAsync(filter: filter, cancellationToken: ct)) + { + if (seen.Add((row.PartitionKey, row.RowKey))) + rows.Add(row); + } + } + } + + return rows; + } + + /// + /// Read one logical entity by (partition, row) key, reassembling it from its physical rows. Used by + /// when a point read lands on a split entity's root row. Returns null when the + /// entity does not exist or cannot be reassembled from the rows present. + /// + private async Task ReadLogicalEntityAsync(string table, string partitionKey, string rowKey, + CancellationToken ct) + { + // The prefix range is an index seek catching the root row and every "{rowKey}-part{n}" row, but + // also unrelated rows sharing the prefix; filter to rows that belong to this entity before + // reassembling. + var filter = $"PartitionKey eq '{Escape(partitionKey)}' and {BuildRowKeyPrefixClause(rowKey)}"; + var rows = new List(); + await foreach (var row in Client(table).QueryAsync(filter: filter, cancellationToken: ct)) + { + if (row.RowKey == rowKey || + (row.TryGetValue(EntitySplitter.OriginalEntityIdKey, out var id) && id?.ToString() == rowKey)) + rows.Add(row); + } + + IncompleteEntityException? incomplete = null; + var entity = EntitySplitter.Reassemble(rows, OnReassemblyWarning, ex => incomplete = ex) + .FirstOrDefault(e => e.RowKey == rowKey); + + if (incomplete is not null) + { + OnReassemblyWarning(incomplete.Message); + return null; + } + + return entity; + } + + /// + /// Delete the part rows of a split entity that this write did NOT rewrite — leftovers from an + /// earlier, larger version, which would otherwise merge stale fragments into a later read. Part rows + /// are found by the cross-row marker; the root row (RowKey == originalRowKey) is in + /// for a current split write, and an empty (from a + /// delete) marks every part row stale. + /// + private async Task RemoveStalePartRowsAsync(string table, string partitionKey, string originalRowKey, + HashSet live, CancellationToken ct) + { + var filter = $"PartitionKey eq '{Escape(partitionKey)}' and {EntitySplitter.OriginalEntityIdKey} eq '{Escape(originalRowKey)}'"; + var stale = new List(); + await foreach (var row in Client(table).QueryAsync(filter: filter, select: select, cancellationToken: ct)) + { + if (!live.Contains(row.RowKey)) + stale.Add(row.RowKey); + } + + if (stale.Count > 0) + await DeleteBatchAsync(table, partitionKey, stale, ct); + } + + /// + /// An OData clause matching every RowKey beginning with : the smallest + /// string sorting above the prefix is the prefix with its last character incremented. + /// + private static string BuildRowKeyPrefixClause(string prefix) + { + var lower = $"RowKey ge '{Escape(prefix)}'"; + var bound = prefix.ToCharArray(); + for (var i = bound.Length - 1; i >= 0; i--) + { + if (bound[i] < char.MaxValue) + { + bound[i]++; + var upper = new string(bound, 0, i + 1); + return $"({lower} and RowKey lt '{Escape(upper)}')"; + } + } + + // Every character is already the maximum, so nothing sorts above the prefix. + return $"({lower})"; + } + // ── Conversion ──────────────────────────────────────────────────────────── private static TableEntity ToEntity(StoreRow row) @@ -467,15 +736,6 @@ private static StoreRow ToRow(TableEntity entity) return row; } - private static int EstimateChars(StoreRow row) - { - var total = row.PartitionKey.Length + row.RowKey.Length + 128; - foreach (var value in row.Properties.Values) - if (value is string s) total += s.Length + 20; - else total += 24; - return total; - } - /// Escape single quotes in OData filter values to prevent injection. private static string Escape(string value) => value.Replace("'", "''"); } diff --git a/Services/Storage/EntitySplitter.cs b/Services/Storage/EntitySplitter.cs new file mode 100644 index 0000000..f9362dd --- /dev/null +++ b/Services/Storage/EntitySplitter.cs @@ -0,0 +1,789 @@ +using Azure.Data.Tables; +using System.Text; +using System.Text.Json; + +namespace Craft.Storage; + +/// +/// Splits entities that exceed the Azure Table Storage size limits (64 KiB per string +/// property, 1 MiB per entity) into multiple properties and rows, and reassembles +/// them on read. +/// +/// Cross-column splitting: an oversized string property "Data" is cut into chunks of +/// at most characters, stored as "Data_Part0", +/// "Data_Part1", ... A JSON manifest in the property +/// records which chunk properties belong to which original property, in order: +/// [{"OriginalHeader":"Data","SplitHeaders":["Data_Part0","Data_Part1"]}]. The +/// manifest is read as either an array or a single object. +/// +/// Cross-row splitting: when the entity is still over after +/// column splitting, its properties are distributed over multiple rows. The first row +/// keeps the original RowKey, subsequent rows use "{RowKey}-part{i}". Every row +/// carries (the original RowKey) and +/// (the zero-based row index). Chunks of one property may +/// land on different rows; reassembly merges rows before joining chunks, so the +/// combination is well-defined. +/// +/// Reassembly groups rows by (PartitionKey, OriginalEntityId), merges each group in +/// PartIndex order, then rejoins chunked properties per the manifest. A plain row +/// whose RowKey equals the OriginalEntityId of leftover part rows takes precedence, +/// so an entity rewritten small after having been split still reads correctly. +/// +/// Vendored from CIPP's CIPP.TableClient, itself derived from AzBobbyTables +/// (MIT-licensed, © Björn Sundling). The only change from the upstream copy is the +/// marker property names, which carry a "_Craft" prefix here so they cannot collide +/// with the orchestrator Results table's own hand-rolled chunking scheme (which reuses +/// the plain "OriginalEntityId"/"PartIndex" names). Craft's orchestrator tables are +/// private to Craft, so the on-disk format does not need to match CIPP's. +/// +public static class EntitySplitter +{ + /// Property holding the JSON manifest describing column-split properties. + public const string SplitOverPropsKey = "_CraftSplitOverProps"; + + /// Property on part rows holding the RowKey of the original entity. + public const string OriginalEntityIdKey = "_CraftPartOf"; + + /// Property on part rows holding the zero-based row index. + public const string PartIndexKey = "_CraftPartIndex"; + + /// + /// Property on part rows holding how many rows the entity was split into, so a + /// reader can tell a complete group from a truncated one. + /// + public const string PartCountKey = "_CraftPartCount"; + + /// Default for . + public const int DefaultMaxPropertyChars = 32_256; + + /// Default for . + public const int DefaultMaxRowSize = 900_000; + + /// Default for . + public const int DefaultMaxTransactionPayload = 2_500_000; + + /// + /// Maximum number of characters per chunk when splitting an oversized string + /// property. The service limit is 64 KiB per string property, i.e. 32768 UTF-16 + /// units regardless of encoding; the default keeps a 512-unit margin below it. + /// Settable for tuning and testing; the reader accepts any chunk size. + /// + public static int MaxPropertyChars { get; set; } = DefaultMaxPropertyChars; + + /// + /// Estimated entity size in bytes at which an entity is distributed over multiple + /// rows. The service allows 1 MiB per entity, but measures it differently per + /// implementation: the real service uses the documented UTF-16 size formula while + /// the storage emulator measures the escaped JSON payload. The estimator counts + /// the worst case of both, so the default of 900k keeps ~15% headroom against + /// either. Settable for tuning and testing. + /// + public static int MaxRowSize { get; set; } = DefaultMaxRowSize; + + /// + /// Estimated payload budget in bytes for a single transaction. Batch requests cap + /// at 4 MiB (rejected at ~3.8 MiB empirically); split rows approach + /// each, so batches must be packed by size as well as + /// count. The default keeps a comfortable margin for the worst-case encoding. + /// Settable for tuning and testing. + /// + public static int MaxTransactionPayload { get; set; } = DefaultMaxTransactionPayload; + + /// + /// Properties that describe a row rather than the logical entity, excluded when + /// merging part rows back together. + /// + private static readonly HashSet MergeExcludedProperties = new(StringComparer.Ordinal) + { + OriginalEntityIdKey, + PartIndexKey, + PartCountKey, + "PartitionKey", + "RowKey", + "Timestamp", + "odata.etag", + "ETag", + }; + + /// + /// The result of splitting one logical entity into the rows to write. + /// + public sealed class SplitResult + { + /// The physical rows to write. A single row when no splitting was needed. + public List Rows { get; } + + /// + /// Whether the splitter rewrote the entity (column chunks and/or multiple rows). + /// When false, contains the original entity untouched. + /// + public bool Engaged { get; } + + public SplitResult(List rows, bool engaged) + { + Rows = rows; + Engaged = engaged; + } + } + + /// + /// Split an entity into one or more rows that each fit within the storage limits. + /// Entities within the limits are passed through untouched. + /// + public static SplitResult Split(TableEntity entity) + { + var hasOversizedProperty = entity.Any(p => p.Value is string s && s.Length > MaxPropertyChars); + + if (!hasOversizedProperty && EstimateEntitySize(entity) <= MaxRowSize) + { + return new SplitResult(new List { entity }, false); + } + + // Rebuild the entity without row-scoped metadata. Split rows are new rows, + // the source row's Timestamp and ETag do not apply to them + var working = new TableEntity(entity.PartitionKey, entity.RowKey); + foreach (var property in entity) + { + if (property.Key is "PartitionKey" or "RowKey" or "Timestamp" or "odata.etag" or "ETag") + { + continue; + } + working[property.Key] = property.Value; + } + + SplitOversizedProperties(working); + + if (EstimateEntitySize(working) <= MaxRowSize) + { + return new SplitResult(new List { working }, true); + } + + return new SplitResult(SplitAcrossRows(working), true); + } + + /// + /// Cross-column splitting: replace every oversized string property with chunk + /// properties and record the manifest in . + /// + private static void SplitOversizedProperties(TableEntity working) + { + var manifest = new List<(string OriginalHeader, List SplitHeaders)>(); + + foreach (var propertyName in working.Keys.ToList()) + { + if (working[propertyName] is not string value || value.Length <= MaxPropertyChars) + { + continue; + } + + var chunks = ChunkString(value); + var chunkNames = new List(chunks.Count); + + working.Remove(propertyName); + for (var i = 0; i < chunks.Count; i++) + { + var chunkName = $"{propertyName}_Part{i}"; + chunkNames.Add(chunkName); + working[chunkName] = chunks[i]; + } + + manifest.Add((propertyName, chunkNames)); + } + + if (manifest.Count > 0) + { + working[SplitOverPropsKey] = SerializeManifest(manifest); + } + } + + /// + /// Cross-row splitting: distribute the properties of an oversized entity over + /// multiple rows, filling each row up to . + /// + private static List SplitAcrossRows(TableEntity working) + { + var originalRowKey = working.RowKey; + var rows = new List(); + + TableEntity CreateRow(int index) + { + var rowKey = index == 0 ? originalRowKey : $"{originalRowKey}-part{index}"; + var row = new TableEntity(working.PartitionKey, rowKey) + { + [OriginalEntityIdKey] = originalRowKey, + [PartIndexKey] = index, + }; + return row; + } + + var rowIndex = 0; + var currentRow = CreateRow(rowIndex); + var currentSize = EstimateEntitySize(currentRow); + + // The manifest goes on the root row after the loop; charge it here so it cannot + // push that row over the limit. + if (working.TryGetValue(SplitOverPropsKey, out var manifestForBudget)) + { + currentSize += EstimatePropertySize(SplitOverPropsKey, manifestForBudget); + } + + foreach (var property in working) + { + if (property.Key is "PartitionKey" or "RowKey") + { + continue; + } + + // Placed on the root row below instead. Left to the loop it lands on the + // last row, and a reader holding only the first row would then have every + // chunk and nothing describing them. + if (property.Key == SplitOverPropsKey) + { + continue; + } + + var propertySize = EstimatePropertySize(property.Key, property.Value); + + // Start a new row when this property would push the current one over the + // budget. A row always accepts at least one property so a single property + // larger than the budget cannot loop forever. + if (currentSize + propertySize > MaxRowSize && HasPayload(currentRow)) + { + rows.Add(currentRow); + rowIndex++; + currentRow = CreateRow(rowIndex); + currentSize = EstimateEntitySize(currentRow); + } + + currentRow[property.Key] = property.Value; + currentSize += propertySize; + } + + rows.Add(currentRow); + + if (working.TryGetValue(SplitOverPropsKey, out var manifest)) + { + rows[0][SplitOverPropsKey] = manifest; + } + + // On every row, so any subset of them knows how many there should be. + foreach (var row in rows) + { + row[PartCountKey] = rows.Count; + } + + return rows; + } + + /// + /// Whether a row under construction holds any property beyond keys and part markers. + /// + private static bool HasPayload(TableEntity row) => + row.Keys.Any(k => k is not ("PartitionKey" or "RowKey" or OriginalEntityIdKey or PartIndexKey)); + + /// + /// Split a string into chunks of at most characters, + /// never cutting between the halves of a surrogate pair. + /// + private static List ChunkString(string value) + { + var chunks = new List((value.Length + MaxPropertyChars - 1) / MaxPropertyChars); + var position = 0; + + while (position < value.Length) + { + var length = Math.Min(MaxPropertyChars, value.Length - position); + + // Move the boundary back one character rather than splitting a surrogate pair. + if (position + length < value.Length && char.IsHighSurrogate(value[position + length - 1]) && length > 1) + { + length--; + } + + chunks.Add(value.Substring(position, length)); + position += length; + } + + return chunks; + } + + /// + /// Serialize the manifest of column-split properties as a JSON array. + /// + private static string SerializeManifest(List<(string OriginalHeader, List SplitHeaders)> manifest) + { + using var stream = new MemoryStream(); + using (var writer = new Utf8JsonWriter(stream)) + { + writer.WriteStartArray(); + foreach (var (originalHeader, splitHeaders) in manifest) + { + writer.WriteStartObject(); + writer.WriteString("OriginalHeader", originalHeader); + writer.WritePropertyName("SplitHeaders"); + writer.WriteStartArray(); + foreach (var header in splitHeaders) + { + writer.WriteStringValue(header); + } + writer.WriteEndArray(); + writer.WriteEndObject(); + } + writer.WriteEndArray(); + } + + return Encoding.UTF8.GetString(stream.ToArray()); + } + + /// + /// Reassemble physical rows into logical entities: merge part rows in PartIndex + /// order, prefer plain rows over leftover part rows with the same identity, and + /// rejoin column-split properties according to their manifest. + /// + /// The physical rows, e.g. the result of a table query. + /// Called with a message when a malformed manifest is skipped. + /// + /// Called for each entity whose rows or split-property chunks are missing. That entity + /// is left out of the results and the rest are still returned, so one entity with rows + /// missing does not hide the entities returned alongside it. + /// + public static IEnumerable Reassemble(IEnumerable entities, Action? onWarning = null, Action? onIncomplete = null) + { + var (order, groups) = GroupRows(entities); + + foreach (var key in order) + { + var group = groups[key]; + + TableEntity? result; + if (group.Root is not null) + { + result = group.Root; + } + else if (group.Parts.Count > 0) + { + if (DescribeIncompleteness(group) is { } reason) + { + onIncomplete?.Invoke(new IncompleteEntityException( + key.PartitionKey, + key.EntityId, + $"Skipped entity with PartitionKey='{key.PartitionKey}' and RowKey='{key.EntityId}': {reason}")); + continue; + } + + result = MergeParts(group.Parts, key.PartitionKey, key.EntityId); + } + else + { + continue; + } + + try + { + JoinSplitProperties(result, onWarning); + } + catch (IncompleteEntityException ex) + { + onIncomplete?.Invoke(ex); + continue; + } + + yield return result; + } + } + + /// + /// Group physical rows by the logical entity they belong to, preserving first-seen + /// order. + /// + private static (List<(string PartitionKey, string EntityId)> Order, Dictionary<(string PartitionKey, string EntityId), EntityGroup> Groups) GroupRows(IEnumerable entities) + { + var order = new List<(string PartitionKey, string EntityId)>(); + var groups = new Dictionary<(string PartitionKey, string EntityId), EntityGroup>(); + + foreach (var entity in entities) + { + var isPart = TryGetOriginalEntityId(entity, out var originalEntityId); + var key = (entity.PartitionKey, isPart ? originalEntityId : entity.RowKey); + + if (!groups.TryGetValue(key, out var group)) + { + group = new EntityGroup(); + groups[key] = group; + order.Add(key); + } + + if (!isPart) + { + // A plain row is the authoritative version of this identity: any part + // rows with the same identity are stale leftovers from an earlier, + // larger version of the entity and are dropped. + group.Root = entity; + group.Parts.Clear(); + } + else if (group.Root is null) + { + group.Parts.Add(entity); + } + } + + return (order, groups); + } + + /// + /// The identities of entities that are present only in part, i.e. rows were split + /// over more rows than the given set contains. + /// + /// + /// Lets a reader fetch the rows it is missing before reassembling, instead of + /// producing a truncated entity from what it has. + /// + /// The physical rows, e.g. the result of a table query. + public static IReadOnlyList<(string PartitionKey, string EntityId)> FindIncompleteGroups(IEnumerable entities) + { + var (order, groups) = GroupRows(entities); + var incomplete = new List<(string PartitionKey, string EntityId)>(); + + foreach (var key in order) + { + var group = groups[key]; + if (group.Root is null && group.Parts.Count > 0 && DescribeIncompleteness(group) is not null) + { + incomplete.Add(key); + } + } + + return incomplete; + } + + /// + /// Why a group of part rows is not the whole entity, or null when it is. + /// + /// + /// answers this outright. Rows written before it existed + /// fall back to two weaker signals: indexes must run from zero without gaps, and a + /// group holding chunk properties must also hold the manifest describing them, + /// which older writers put on the last row. + /// + private static string? DescribeIncompleteness(EntityGroup group) + { + var indexes = group.Parts.Select(GetPartIndex).OrderBy(i => i).ToList(); + + var declaredCount = group.Parts + .Select(GetPartCount) + .Where(c => c > 0) + .DefaultIfEmpty(0) + .Max(); + + if (declaredCount > 0 && group.Parts.Count != declaredCount) + { + return $"the entity was split over {declaredCount} rows but only {group.Parts.Count} were returned."; + } + + if (indexes[0] != 0) + { + return $"the first row of the entity (PartIndex 0) was not returned."; + } + + for (var i = 1; i < indexes.Count; i++) + { + if (indexes[i] != indexes[i - 1] + 1) + { + return $"the rows returned skip from PartIndex {indexes[i - 1]} to {indexes[i]}."; + } + } + + if (declaredCount == 0 && + !group.Parts.Any(p => p.ContainsKey(SplitOverPropsKey)) && + group.Parts.Any(p => p.Keys.Any(IsChunkPropertyName))) + { + return "the entity holds split property chunks but no manifest describing them, so the row carrying the manifest was not returned."; + } + + return null; + } + + /// + /// Whether a property name looks like one of the chunks a column-split property was + /// broken into, i.e. ends in _Part followed by digits. + /// + private static bool IsChunkPropertyName(string name) + { + var separator = name.LastIndexOf("_Part", StringComparison.Ordinal); + if (separator < 1 || separator + 5 >= name.Length) + { + return false; + } + + for (var i = separator + 5; i < name.Length; i++) + { + if (name[i] < '0' || name[i] > '9') + { + return false; + } + } + + return true; + } + + /// + /// The PartCount of a part row, tolerating the integer widening that occurs when + /// rows round-trip through JSON. Rows written before the property existed report 0. + /// + private static long GetPartCount(TableEntity entity) + { + if (!entity.TryGetValue(PartCountKey, out var value)) + { + return 0; + } + + return value switch + { + int i => i, + long l => l, + string s when long.TryParse(s, out var parsed) => parsed, + _ => 0, + }; + } + + private sealed class EntityGroup + { + public TableEntity? Root { get; set; } + public List Parts { get; } = new(); + } + + /// + /// Whether the row is a part of a split entity, i.e. carries a non-empty + /// property. + /// + private static bool TryGetOriginalEntityId(TableEntity entity, out string originalEntityId) + { + if (entity.TryGetValue(OriginalEntityIdKey, out var value) && value?.ToString() is { Length: > 0 } id) + { + originalEntityId = id; + return true; + } + + originalEntityId = string.Empty; + return false; + } + + /// + /// The PartIndex of a part row, tolerating the integer widening that occurs when + /// rows round-trip through JSON. Rows without a usable index sort first. + /// + private static long GetPartIndex(TableEntity entity) + { + if (!entity.TryGetValue(PartIndexKey, out var value)) + { + return long.MinValue; + } + + return value switch + { + int i => i, + long l => l, + string s when long.TryParse(s, out var parsed) => parsed, + _ => long.MinValue, + }; + } + + /// + /// Merge part rows into one logical entity in PartIndex order. Row-scoped metadata + /// is dropped; the merged entity takes its RowKey from the original entity id and + /// its Timestamp and ETag from the first part row. + /// + private static TableEntity MergeParts(List parts, string partitionKey, string originalEntityId) + { + var ordered = parts.OrderBy(GetPartIndex).ToList(); + var merged = new TableEntity(partitionKey, originalEntityId); + + foreach (var part in ordered) + { + foreach (var property in part) + { + if (MergeExcludedProperties.Contains(property.Key)) + { + continue; + } + + // The writer places each property on exactly one row; string values + // colliding across rows are treated as fragments and concatenated. + if (merged.TryGetValue(property.Key, out var existing) && existing is string left && property.Value is string right) + { + merged[property.Key] = left + right; + } + else + { + merged[property.Key] = property.Value; + } + } + } + + var first = ordered[0]; + merged.Timestamp = first.Timestamp; + if (first.TryGetValue("odata.etag", out var etag)) + { + merged["odata.etag"] = etag; + } + + return merged; + } + + /// + /// Rejoin column-split properties according to the + /// manifest, removing the chunk properties and the manifest itself. Malformed + /// manifests are reported through and skipped; the + /// manifest property is removed either way. + /// + private static void JoinSplitProperties(TableEntity entity, Action? onWarning) + { + if (!entity.TryGetValue(SplitOverPropsKey, out var manifestValue) || manifestValue is not string manifestJson || manifestJson.Length == 0) + { + return; + } + + try + { + using var document = JsonDocument.Parse(manifestJson); + + // The manifest may be a bare object instead of a single-element array. + var entries = document.RootElement.ValueKind switch + { + JsonValueKind.Array => document.RootElement.EnumerateArray().ToList(), + JsonValueKind.Object => new List { document.RootElement }, + _ => throw new JsonException($"Unexpected manifest root of kind {document.RootElement.ValueKind}."), + }; + + foreach (var entry in entries) + { + if (!entry.TryGetProperty("OriginalHeader", out var originalHeaderElement) || + originalHeaderElement.GetString() is not { Length: > 0 } originalHeader || + !entry.TryGetProperty("SplitHeaders", out var splitHeadersElement)) + { + throw new JsonException("Manifest entry is missing OriginalHeader or SplitHeaders."); + } + + var splitHeaders = splitHeadersElement.ValueKind switch + { + JsonValueKind.Array => splitHeadersElement.EnumerateArray() + .Select(h => h.GetString()) + .Where(h => h is { Length: > 0 }) + .Select(h => h!) + .ToList(), + // Single-element unwrapping tolerance, should not occur in practice. + JsonValueKind.String when splitHeadersElement.GetString() is { Length: > 0 } single => new List { single }, + _ => throw new JsonException("Manifest SplitHeaders is neither an array nor a string."), + }; + + var joined = new StringBuilder(); + foreach (var header in splitHeaders) + { + // Skipping a missing chunk would store the rest as though it were the + // whole value - a truncation nothing downstream can detect. + if (!entity.TryGetValue(header, out var chunk) || chunk is null) + { + throw new IncompleteEntityException( + entity.PartitionKey, + entity.RowKey, + $"Cannot reassemble property '{originalHeader}' of entity with PartitionKey='{entity.PartitionKey}' and RowKey='{entity.RowKey}': chunk property '{header}' is missing. The query did not return every row the entity was split over."); + } + + joined.Append(chunk); + } + + entity[originalHeader] = joined.ToString(); + foreach (var header in splitHeaders) + { + entity.Remove(header); + } + } + } + catch (Exception ex) when (ex is JsonException or InvalidOperationException) + { + onWarning?.Invoke($"Failed to process {SplitOverPropsKey} for entity with PartitionKey='{entity.PartitionKey}' and RowKey='{entity.RowKey}': {ex.Message}"); + } + finally + { + entity.Remove(SplitOverPropsKey); + } + } + + /// + /// Estimate the stored size of an entity in bytes, based on the documented Azure + /// Table Storage entity size formula (property names and string values count as + /// UTF-16, i.e. two bytes per character). + /// + public static int EstimateEntitySize(TableEntity entity) + { + var size = 4; + size += (entity.PartitionKey?.Length ?? 0) * 2; + size += (entity.RowKey?.Length ?? 0) * 2; + + foreach (var property in entity) + { + if (property.Key is "PartitionKey" or "RowKey" or "Timestamp" or "odata.etag") + { + continue; + } + size += EstimatePropertySize(property.Key, property.Value); + } + + return size; + } + + /// + /// Estimate the stored size of a single property in bytes: fixed per-property + /// overhead, the name, and the value by type. Sizes are the maximum of the + /// documented storage size formula and the JSON wire representation, since the + /// service enforces the entity limit against whichever is relevant to it (the + /// emulator measures the request payload, where non-ASCII characters are escaped + /// as \uXXXX sequences). Non-string types include headroom for the OData type + /// annotation property that accompanies them on the wire. + /// + private static int EstimatePropertySize(string name, object? value) + { + var valueSize = value switch + { + null => 0, + string s => EstimateStringSize(s), + byte[] b => EstimateBinarySize(b.Length), + BinaryData b => EstimateBinarySize((int)Math.Min(b.ToMemory().Length, int.MaxValue)), + bool => 8, + int => 16, + long => 48, + double => 48, + Guid => 64, + DateTime => 72, + DateTimeOffset => 72, + _ => EstimateStringSize(value.ToString() ?? string.Empty), + }; + + return 8 + name.Length * 2 + valueSize; + } + + /// + /// Estimate the stored size of a string value in bytes. Per UTF-16 unit, the cost + /// is the maximum of the storage formula (2 bytes per unit) and the escaped JSON + /// wire form (6 bytes for control and non-ASCII units). + /// + private static int EstimateStringSize(string value) + { + var size = 4; + foreach (var c in value) + { + size += c < 0x20 || c >= 0x80 ? 6 : 2; + } + return size; + } + + /// + /// Estimate the stored size of a binary value in bytes: the maximum of the storage + /// formula (raw bytes + 4 overhead) and the base64 JSON wire form. Base64 encodes + /// every 3 bytes as 4 characters; ceiling division ((length + 2) / 3 * 4) + /// is used so a length not divisible by 3 is not underestimated. The +44 covers + /// the OData type annotation property that accompanies binary values on the wire. + /// + private static int EstimateBinarySize(int length) => + Math.Max(length + 4, (length + 2) / 3 * 4 + 44); +} diff --git a/Services/Storage/IncompleteEntityException.cs b/Services/Storage/IncompleteEntityException.cs new file mode 100644 index 0000000..f514427 --- /dev/null +++ b/Services/Storage/IncompleteEntityException.cs @@ -0,0 +1,29 @@ +namespace Craft.Storage; + +/// +/// Thrown when a split entity cannot be reassembled because some of the rows or chunk +/// properties it was split into are not present. +/// +/// +/// An incomplete read fails rather than returning what it managed to reassemble, which +/// would be indistinguishable from the real value. +/// +/// Vendored from CIPP's CIPP.TableClient (itself derived from AzBobbyTables). See +/// for the storage format and attribution. +/// +public class IncompleteEntityException : Exception +{ + /// The PartitionKey of the entity that could not be reassembled. + public string EntityPartitionKey { get; } + + /// The RowKey of the entity that could not be reassembled. + public string EntityRowKey { get; } + + /// Create an exception for an entity that could not be reassembled. + public IncompleteEntityException(string partitionKey, string rowKey, string message) + : base(message) + { + EntityPartitionKey = partitionKey; + EntityRowKey = rowKey; + } +} diff --git a/tests/Craft.Tests/AzureTableStoreLargeEntityTests.cs b/tests/Craft.Tests/AzureTableStoreLargeEntityTests.cs new file mode 100644 index 0000000..20a57c4 --- /dev/null +++ b/tests/Craft.Tests/AzureTableStoreLargeEntityTests.cs @@ -0,0 +1,240 @@ +using Azure; +using Azure.Data.Tables; +using Craft.Configuration; +using Craft.Storage; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// Tests here allocate multi-MB strings to force the storage size limits, and heap-delta measurement +/// tests (see ) cannot share a process with concurrent allocation. +/// Marking this collection non-parallel puts it in the same sequential phase as those, so the two never +/// run at once. See the note on for the failure mode this avoids. +/// +[CollectionDefinition(LargeAllocationSerialTests.Name, DisableParallelization = true)] +public class LargeAllocationSerialTests +{ + public const string Name = "large-allocation-serial"; +} + +/// +/// Proves the large-entity split/reassemble path against a real table backend: a value too big for one +/// property (or one entity) is stored transparently and read back byte-identical, claims still work on a +/// split row, shrinking a split entity does not resurrect it from stale parts, and deleting one takes its +/// parts with it. Azurite by default, a real account via CRAFT_TEST_TABLE_CONNECTION; skipped, not failed, +/// when neither is reachable (a skip reads as a pass). +/// +/// This is the regression guard for the failure it was written for: a scheduled task whose Parameters +/// embed a whole policy template exceeded Azure Table's 64 KiB-per-property limit, so the orchestrator +/// task row 400'd with PropertyValueTooLarge, was dropped, and the task could never be dispatched. +/// +[Collection(LargeAllocationSerialTests.Name)] +public class AzureTableStoreLargeEntityTests +{ + private sealed class Fixture : IAsyncDisposable + { + public required AzureTableStore Store { get; init; } + public required string Table { get; init; } + public required string Connection { get; init; } + + public static async Task TryConnectAsync() + { + var settings = new CraftSettings(); + var connection = Environment.GetEnvironmentVariable("CRAFT_TEST_TABLE_CONNECTION"); + if (!string.IsNullOrWhiteSpace(connection)) + settings.Auth.UserStorageConnection = connection; + else + { + settings.Storage.AllowDevelopmentStorage = true; + connection = "UseDevelopmentStorage=true"; + } + + var store = new AzureTableStore(settings, NullLogger.Instance); + try + { + using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(3)); + await store.PingAsync(cts.Token); + } + catch + { + return null; + } + + var table = "azle" + Guid.NewGuid().ToString("N")[..8]; + await store.EnsureTableAsync(table); + return new Fixture { Store = store, Table = table, Connection = connection }; + } + + /// Count the physical rows in a partition, bypassing reassembly (raw client). + public async Task PhysicalRowCountAsync(string partitionKey) + { + var count = 0; + var client = new TableClient(Connection, Table); + await foreach (var _ in client.QueryAsync(filter: $"PartitionKey eq '{partitionKey}'")) + count++; + return count; + } + + public async ValueTask DisposeAsync() + { + try { await new TableServiceClient(Connection).DeleteTableAsync(Table); } + catch (RequestFailedException) { /* never created, or already gone */ } + } + } + + private static string Text(int chars, char c = 'x') => new(c, chars); + + private static StoreRow Row(string pk, string rk, string bigValue) => + new(pk, rk) { Properties = { ["Status"] = "Pending", ["ParametersJson"] = bigValue } }; + + [Fact] + public async Task ColumnSplit_OversizedProperty_RoundTripsThroughUpsertAndGetAndQuery() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + var value = Text(80_000); // > 64 KiB per-property limit; this is the exact failure mode + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", value)); + + var got = await fx.Store.GetAsync(fx.Table, "p", "task1"); + Assert.Equal(value, got?.GetString("ParametersJson")); + Assert.Equal("Pending", got?.GetString("Status")); + + var queried = new List(); + await foreach (var r in fx.Store.QueryPartitionAsync(fx.Table, "p")) queried.Add(r); + Assert.Single(queried); + Assert.Equal(value, queried[0].GetString("ParametersJson")); + + // A column split stays one physical row. + Assert.Equal(1, await fx.PhysicalRowCountAsync("p")); + } + + [Fact] + public async Task CrossRowSplit_HugeProperty_RoundTripsAndUsesMultipleRows() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + var value = Text(1_200_000); // over the 1 MiB entity cap → spills onto extra rows + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", value)); + + Assert.True(await fx.PhysicalRowCountAsync("p") > 1, "expected the entity to occupy more than one physical row"); + + var got = await fx.Store.GetAsync(fx.Table, "p", "task1"); + Assert.Equal(value, got?.GetString("ParametersJson")); + + var queried = new List(); + await foreach (var r in fx.Store.QueryPartitionAsync(fx.Table, "p")) queried.Add(r); + Assert.Single(queried); + Assert.Equal(value, queried[0].GetString("ParametersJson")); + } + + [Fact] + public async Task UpsertBatch_MixesLargeAndSmallRows_AllRoundTrip() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + var big = Text(1_200_000); + var mid = Text(80_000); + await fx.Store.UpsertBatchAsync(fx.Table, "p", new[] + { + new StoreRow("p", "small") { Properties = { ["ParametersJson"] = "tiny" } }, + Row("p", "mid", mid), + Row("p", "big", big), + }); + + var byKey = new Dictionary(); + await foreach (var r in fx.Store.QueryPartitionAsync(fx.Table, "p")) + byKey[r.RowKey] = r.GetString("ParametersJson"); + + Assert.Equal(3, byKey.Count); + Assert.Equal("tiny", byKey["small"]); + Assert.Equal(mid, byKey["mid"]); + Assert.Equal(big, byKey["big"]); + } + + [Fact] + public async Task ShrinkingASplitEntityToASmallValue_ReadsTheNewValue_WithoutResurrection() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", Text(1_200_000))); // large: many rows + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", "small-now")); // small: one plain row + + // The correctness guarantee: both read paths return the new small value, never the old one + // reassembled from leftover part rows (plain-row precedence). Physical cleanup of those leftovers + // happens on the next ENGAGED write or on delete, not on a shrink-to-small — matching the + // AzBobbyTables write path, which would otherwise add a scan to every small upsert. + var got = await fx.Store.GetAsync(fx.Table, "p", "task1"); + Assert.Equal("small-now", got?.GetString("ParametersJson")); + + var queried = new List(); + await foreach (var r in fx.Store.QueryPartitionAsync(fx.Table, "p")) queried.Add(r); + Assert.Single(queried); + Assert.Equal("small-now", queried[0].GetString("ParametersJson")); + } + + [Fact] + public async Task EngagedOverwrite_ReclaimsStalePartRows_FromTheLargerVersion() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", Text(3_000_000))); // many rows + var smaller = Row("p", "task1", Text(1_100_000)); // still split, fewer rows + await fx.Store.UpsertAsync(fx.Table, smaller); + + var got = await fx.Store.GetAsync(fx.Table, "p", "task1"); + Assert.Equal(Text(1_100_000), got?.GetString("ParametersJson")); + + // The engaged write removed the higher-index part rows the larger version had left behind: + // exactly the new version's rows remain. + var expected = EntitySplitter.Split(new TableEntity("p", "task1") + { + ["Status"] = "Pending", + ["ParametersJson"] = Text(1_100_000) + }).Rows.Count; + Assert.Equal(expected, await fx.PhysicalRowCountAsync("p")); + } + + [Fact] + public async Task ConditionalReplace_WorksOnALargeChunkedRow() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", Text(80_000))); + + var current = await fx.Store.GetAsync(fx.Table, "p", "task1"); + Assert.NotNull(current); + current!["Status"] = "Running"; + + Assert.True(await fx.Store.TryReplaceBatchAsync(fx.Table, "p", new[] { current })); + Assert.Equal("Running", (await fx.Store.GetAsync(fx.Table, "p", "task1"))?.GetString("Status")); + + // The now-stale ETag must be rejected — no silent unconditional overwrite. + current["Status"] = "Cancelled"; + Assert.False(await fx.Store.TryReplaceBatchAsync(fx.Table, "p", new[] { current })); + } + + [Fact] + public async Task Delete_RemovesEveryPartRow_OfASplitEntity() + { + await using var fx = await Fixture.TryConnectAsync(); + if (fx == null) return; + + await fx.Store.UpsertAsync(fx.Table, Row("p", "task1", Text(1_200_000))); + Assert.True(await fx.PhysicalRowCountAsync("p") > 1); + + await fx.Store.DeleteAsync(fx.Table, "p", "task1"); + + Assert.Equal(0, await fx.PhysicalRowCountAsync("p")); + Assert.Null(await fx.Store.GetAsync(fx.Table, "p", "task1")); + var any = false; + await foreach (var _ in fx.Store.QueryPartitionAsync(fx.Table, "p")) any = true; + Assert.False(any); + } +} diff --git a/tests/Craft.Tests/EntitySplitterTests.cs b/tests/Craft.Tests/EntitySplitterTests.cs new file mode 100644 index 0000000..b40a743 --- /dev/null +++ b/tests/Craft.Tests/EntitySplitterTests.cs @@ -0,0 +1,140 @@ +using Azure.Data.Tables; +using Craft.Storage; + +namespace Craft.Tests; + +/// +/// Pure round-trip tests for — no backend. These prove the split/reassemble +/// invariant that relies on: whatever goes in comes back byte-identical, +/// however the storage limits forced it to be broken up, and a partial set of rows is refused rather than +/// silently truncated. Uses real large strings so it never mutates the splitter's global tuning knobs. +/// +[Collection(LargeAllocationSerialTests.Name)] +public class EntitySplitterTests +{ + private static string Text(int chars, char c = 'x') => new(c, chars); + + private static TableEntity Entity(string pk, string rk, params (string Key, object Value)[] props) + { + var e = new TableEntity(pk, rk); + foreach (var (k, v) in props) e[k] = v; + return e; + } + + private static TableEntity RoundTrip(TableEntity entity) + { + var split = EntitySplitter.Split(entity); + // Reassemble sees exactly the physical rows a query would return. + var back = EntitySplitter.Reassemble(split.Rows).ToList(); + Assert.Single(back); + return back[0]; + } + + [Fact] + public void SmallEntity_IsNotEngaged_AndPassesThroughUnchanged() + { + var entity = Entity("p", "r", ("Value", "hello"), ("N", 7)); + var split = EntitySplitter.Split(entity); + + Assert.False(split.Engaged); + Assert.Single(split.Rows); + Assert.Same(entity, split.Rows[0]); + } + + [Fact] + public void OversizedProperty_SplitsAcrossColumns_InOneRow_AndRejoins() + { + var value = Text(80_000); // > 32K limit, well under the 1 MiB entity cap + var entity = Entity("p", "r", ("Status", "Pending"), ("ParametersJson", value)); + + var split = EntitySplitter.Split(entity); + + Assert.True(split.Engaged); + Assert.Single(split.Rows); // column split stays one physical row + Assert.False(split.Rows[0].ContainsKey("ParametersJson")); // original replaced by chunks + Assert.True(split.Rows[0].ContainsKey("ParametersJson_Part0")); + Assert.True(split.Rows[0].ContainsKey(EntitySplitter.SplitOverPropsKey)); + + var back = EntitySplitter.Reassemble(split.Rows).Single(); + Assert.Equal(value, back.GetString("ParametersJson")); + Assert.Equal("Pending", back.GetString("Status")); + Assert.False(back.ContainsKey(EntitySplitter.SplitOverPropsKey)); + } + + [Fact] + public void HugeProperty_SplitsAcrossRows_AndRejoins() + { + var value = Text(1_200_000); // forces the entity over the 1 MiB cap → cross-row + var entity = Entity("p", "r", ("Status", "Pending"), ("ParametersJson", value)); + + var split = EntitySplitter.Split(entity); + + Assert.True(split.Engaged); + Assert.True(split.Rows.Count > 1); // spilled onto extra rows + Assert.Equal("r", split.Rows[0].RowKey); // root keeps the original RowKey + Assert.All(split.Rows, r => Assert.Equal("r", r.GetString(EntitySplitter.OriginalEntityIdKey))); + + var back = EntitySplitter.Reassemble(split.Rows).Single(); + Assert.Equal("r", back.RowKey); + Assert.Equal(value, back.GetString("ParametersJson")); + Assert.Equal("Pending", back.GetString("Status")); + } + + [Fact] + public void SurrogatePairs_AreNotSplitDownTheMiddle() + { + // A string of astral-plane code points (each a surrogate pair) sized past the chunk boundary: + // a naive cut would produce an invalid half-surrogate and corrupt the round-trip. + var value = string.Concat(Enumerable.Repeat("\U0001F600", 40_000)); // 80k UTF-16 units + var entity = Entity("p", "r", ("Data", value)); + + var back = RoundTrip(entity); + + Assert.Equal(value, back.GetString("Data")); + } + + [Fact] + public void PlainRow_SupersedesLeftoverPartRows_OfTheSameIdentity() + { + // The entity was large (parts written), then rewritten small (a plain root at the same RowKey). + // Reassembly must return the small value and ignore the stale parts. + var big = EntitySplitter.Split(Entity("p", "r", ("Data", Text(1_200_000)))).Rows; + var plain = Entity("p", "r", ("Data", "small-now")); + + var rows = new List { plain }; + rows.AddRange(big.Where(r => r.RowKey != "r")); // leftover "-part{n}" rows only + + var back = EntitySplitter.Reassemble(rows).Single(); + Assert.Equal("small-now", back.GetString("Data")); + } + + [Fact] + public void MissingPartRow_IsReportedIncomplete_AndDropped_WithoutAffectingOthers() + { + var incomplete = EntitySplitter.Split(Entity("p", "gone", ("Data", Text(1_200_000)))).Rows; + var whole = Entity("p", "here", ("Data", "fine")); + + // Drop the last part row of the "gone" entity to simulate a query that missed it. + var rows = new List { whole }; + rows.AddRange(incomplete.Take(incomplete.Count - 1)); + + IncompleteEntityException? reported = null; + var back = EntitySplitter.Reassemble(rows, onIncomplete: ex => reported = ex).ToList(); + + Assert.NotNull(reported); + Assert.Equal("gone", reported!.EntityRowKey); + Assert.DoesNotContain(back, e => e.RowKey == "gone"); + Assert.Contains(back, e => e.RowKey == "here" && e.GetString("Data") == "fine"); + } + + [Fact] + public void FindIncompleteGroups_IdentifiesEntitiesMissingRows() + { + var parts = EntitySplitter.Split(Entity("p", "r", ("Data", Text(1_200_000)))).Rows; + var missingLast = parts.Take(parts.Count - 1).ToList(); + + var incomplete = EntitySplitter.FindIncompleteGroups(missingLast); + + Assert.Contains(("p", "r"), incomplete); + } +} diff --git a/tests/Craft.Tests/OrchestratorTaskLargeParametersAzuriteTests.cs b/tests/Craft.Tests/OrchestratorTaskLargeParametersAzuriteTests.cs new file mode 100644 index 0000000..b44a86c --- /dev/null +++ b/tests/Craft.Tests/OrchestratorTaskLargeParametersAzuriteTests.cs @@ -0,0 +1,109 @@ +using System.Text.Json; +using Craft.Configuration; +using Craft.Orchestration; +using Craft.Storage; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// End-to-end guard for the failure this whole change exists for: a scheduled task whose Parameters +/// embed a whole policy template serialize to more than Azure Table's 64 KiB-per-property limit, so the +/// orchestrator's Tasks row used to 400 with PropertyValueTooLarge, get dropped, and the task could never +/// be rehydrated at dispatch ("Parameters could not be rehydrated at dispatch — the Tasks-table row is +/// missing"). With large-entity splitting in the backing store, the task row persists and both read paths +/// the orchestrator uses — GetRunAsync (partition scan) and GetTaskParametersAsync (point read, the +/// dispatch rehydrate) — return the Parameters byte-for-byte. +/// +/// Azurite by default, a real account via CRAFT_TEST_TABLE_CONNECTION; skipped, not failed, when neither +/// is reachable. +/// +[Collection(LargeAllocationSerialTests.Name)] +public class OrchestratorTaskLargeParametersAzuriteTests +{ + private static async Task TryConnectAsync() + { + var settings = new CraftSettings(); + var connection = Environment.GetEnvironmentVariable("CRAFT_TEST_TABLE_CONNECTION"); + if (!string.IsNullOrWhiteSpace(connection)) + settings.Auth.UserStorageConnection = connection; + else + settings.Storage.AllowDevelopmentStorage = true; + + settings.Orchestrator.TablePrefix = "aztp" + Guid.NewGuid().ToString("N")[..8]; + + var backing = new AzureTableStore(settings); + try + { + using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(3)); + await backing.PingAsync(cts.Token); + } + catch + { + return null; + } + + var store = new OrchestratorTableStore(NullLogger.Instance, settings, backing); + await store.InitializeAsync(); + return store; + } + + private static string AsString(object? v) => v switch + { + string s => s, + JsonElement { ValueKind: JsonValueKind.String } je => je.GetString() ?? "", + _ => v?.ToString() ?? "" + }; + + [Theory] + [InlineData(80_000)] // > 64 KiB per-property: the exact reported failure (column split) + [InlineData(1_200_000)] // > 1 MiB entity: forces a cross-row split too + public async Task ATaskWhoseParametersExceedTheLimit_PersistsAndRehydrates(int settingsChars) + { + var store = await TryConnectAsync(); + if (store == null) return; + + const string run = "UserTaskOrchestrator_contoso.com"; + var big = new string('T', settingsChars); + var parameters = new Dictionary + { + ["Tenant"] = "contoso.com", + ["Settings"] = big, + }; + + try + { + await store.UpsertRunAsync(new OrchestratorRun + { + Name = run, + Status = "Running", + Priority = 2, + StartedUtc = DateTime.UtcNow, + TaskScriptName = "ExecScheduledCommand", + Tasks = [new OrchestratorTaskItem { Id = "task-0", Status = "Pending" }] + }); + + // The exact write path Start-UserTasksOrchestrator uses to enqueue a task's payload. + await store.UpsertTaskBatchAsync(run, new List + { + new() { Id = "task-0", Status = "Pending", Parameters = parameters } + }); + + // Dispatch rehydrate: a point read of the one task's Parameters. + var rehydrated = await store.GetTaskParametersAsync(run, "task-0"); + Assert.NotNull(rehydrated); + Assert.Equal("contoso.com", AsString(rehydrated!["Tenant"])); + Assert.Equal(big, AsString(rehydrated["Settings"])); + + // Whole-run read: the partition scan reassembles the same task. + var loaded = await store.GetRunAsync(run); + Assert.NotNull(loaded); + var task = Assert.Single(loaded!.Tasks); + Assert.Equal(big, AsString(task.Parameters["Settings"])); + } + finally + { + await store.CleanupRunAsync(run); + } + } +} From 0a41401d0d9a3c482968b803f4771774f6a549c3 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Fri, 18 Sep 2026 17:19:49 +0800 Subject: [PATCH 4/6] add env for craft version to image for downstream app and internal usage --- Services/Storage/EntitySplitter.cs | 2 +- build/Dockerfile | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/Services/Storage/EntitySplitter.cs b/Services/Storage/EntitySplitter.cs index f9362dd..0f44750 100644 --- a/Services/Storage/EntitySplitter.cs +++ b/Services/Storage/EntitySplitter.cs @@ -1,6 +1,6 @@ -using Azure.Data.Tables; using System.Text; using System.Text.Json; +using Azure.Data.Tables; namespace Craft.Storage; diff --git a/build/Dockerfile b/build/Dockerfile index bfe1837..3997ae8 100644 --- a/build/Dockerfile +++ b/build/Dockerfile @@ -135,6 +135,7 @@ ARG IMAGE_TAG=latest ENV APP_VERSION=$APP_VERSION ENV COMMIT_SHA=$COMMIT_SHA ENV IMAGE_TAG=$IMAGE_TAG +ENV CRAFT_VERSION=$APP_VERSION ENV CRAFT_VERBOSE="false" From 5320c48d59d46a989df9fa7bb7e67d3e70dd0c63 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Thu, 24 Sep 2026 01:39:59 +0800 Subject: [PATCH 5/6] perf(hosting): coalesce response-compression input into 64 KiB blocks String bodies (Response.WriteAsync) and CraftResult.Stream writers reach the encoder through the response pipe, one ~4 KiB segment per write. Brotli's fast qualities compress each write as an isolated fragment, and every encoder pays per-call overhead. Both providers now sit behind a 64 KiB BufferedStream: output still streams every 64 KiB and flushes pass straight through. Real 300 KB ListLogs body: br Fastest 85 KB -> 36 KB and 2.5x faster; br Optimal (the default) ~10-20% less CPU with identical bytes; gzip unchanged. Providers are now constructed with their level instead of via IOptions. Tests: per-encoding negotiation and round-trip, br winning a browser-style Accept-Encoding, and a real-Kestrel check that a varied string body is not fed to the encoder as fragments (fails without the coalescing). Harness: PerfFile serves a captured real response verbatim from PAYLOAD_DIR; run-compression-levels.ps1 gains -Url/-PayloadDir, p50, and a discarded warm-up pass so the first-measured encoding no longer absorbs JIT and pool growth. --- .../Hosting/CoalescingCompressionProvider.cs | 28 ++++++++ .../Hosting/CraftHostBuilderExtensions.cs | 14 ++-- .../API/Modules/PerfApi/PerfApi.psd1 | 2 +- .../API/Modules/PerfApi/PerfApi.psm1 | 9 +++ perf-harness/api-harness/payloads/.gitkeep | 0 perf-harness/docker-compose.api.yml | 3 + perf-harness/k6/api_load.js | 9 ++- .../scripts/run-compression-levels.ps1 | 30 ++++++--- .../ApiCompressionPipelineTests.cs | 66 +++++++++++++++++-- 9 files changed, 139 insertions(+), 22 deletions(-) create mode 100644 Services/Hosting/CoalescingCompressionProvider.cs create mode 100644 perf-harness/api-harness/payloads/.gitkeep diff --git a/Services/Hosting/CoalescingCompressionProvider.cs b/Services/Hosting/CoalescingCompressionProvider.cs new file mode 100644 index 0000000..861fd2e --- /dev/null +++ b/Services/Hosting/CoalescingCompressionProvider.cs @@ -0,0 +1,28 @@ +using Microsoft.AspNetCore.ResponseCompression; + +namespace Craft.Hosting; + +/// +/// Feeds an encoder in 64 KiB blocks however the response is written, while still streaming. +/// +/// +/// A string body (Response.WriteAsync) or a CraftResult.Stream writer reaches the encoder +/// through the response pipe, which hands it one ~4 KiB segment per write. Brotli's fast qualities +/// compress each write as an isolated fragment, so a 300 KB ListLogs body at Fastest came out 85 KB +/// instead of 29 KB, and every encoder paid per-call overhead (Brotli Optimal ~10-15% CPU). Coalescing to +/// 64 KiB recovers the ratio (36 KB; exact for bodies up to 64 KiB) and the CPU, and output still leaves +/// every 64 KiB — the body is never materialised. Flushes pass straight through, so an explicit flush +/// (SSE, progressive output) behaves exactly as before; writes of 64 KiB or more bypass the buffer. +/// +public sealed class CoalescingCompressionProvider(ICompressionProvider inner) : ICompressionProvider +{ + // Under the 85 KB large-object threshold, so the per-response buffer is a cheap gen0 allocation. + private const int BlockSize = 64 * 1024; + + public string EncodingName => inner.EncodingName; + public bool SupportsFlush => inner.SupportsFlush; + + // ponytail: BufferedStream allocates its 64 KiB buffer per compressed response (gen0). Swap for an + // ArrayPool-backed buffer if allocation rate ever shows up in the GC diagnostics. + public Stream CreateStream(Stream outputStream) => new BufferedStream(inner.CreateStream(outputStream), BlockSize); +} diff --git a/Services/Hosting/CraftHostBuilderExtensions.cs b/Services/Hosting/CraftHostBuilderExtensions.cs index 48531b6..2189af0 100644 --- a/Services/Hosting/CraftHostBuilderExtensions.cs +++ b/Services/Hosting/CraftHostBuilderExtensions.cs @@ -223,15 +223,19 @@ public static IServiceCollection AddCraftResponseCompression( services.AddResponseCompression(options => { options.EnableForHttps = true; - options.Providers.Add(); - options.Providers.Add(); + // Order is preference: equal-q ties go to the earliest provider, so br wins over gzip. + // Every provider is coalesced — see CoalescingCompressionProvider. + ICompressionProvider[] providers = + [ + new BrotliCompressionProvider(Options.Create(new BrotliCompressionProviderOptions { Level = level })), + new GzipCompressionProvider(Options.Create(new GzipCompressionProviderOptions { Level = level })), + ]; + foreach (var provider in providers) + options.Providers.Add(new CoalescingCompressionProvider(provider)); options.MimeTypes = ResponseCompressionDefaults.MimeTypes.Concat( second); }); - services.Configure(o => o.Level = level); - services.Configure(o => o.Level = level); - return services; } diff --git a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 index 95a728f..28b8e31 100644 --- a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 +++ b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 @@ -5,7 +5,7 @@ Author = 'CRAFT perf-harness' Description = 'Synthetic HTTP endpoints for load-testing CRAFT in http-only mode. Not for production.' PowerShellVersion = '7.2' - FunctionsToExport = @('Invoke-PerfPing', 'Invoke-PerfEcho', 'Invoke-PerfCpu', 'Invoke-PerfSleep', 'Invoke-PerfJson', 'Invoke-PerfBgEnqueue', 'Push-PerfBg', 'Push-PerfBgLeaf', 'Invoke-PerfManyRuns', 'Push-PerfHold', 'Invoke-PerfThreads', 'Invoke-PerfTableOp', 'Push-PerfCheck', 'Invoke-PerfCheckCounts', 'Push-PerfSeq', 'Invoke-PerfSeqResult', 'Push-PerfSeqWorker', 'Invoke-PerfSeqWorkerEnqueue', 'Invoke-PerfSeqWorkerResult', 'Invoke-PerfThreadBreakdown', 'Invoke-PerfSeedRuns', 'Invoke-ListPerf', 'Invoke-PerfWhoami', 'Invoke-PerfTimerTick', 'Invoke-PerfTimerCount', 'Invoke-PerfPublish', 'Invoke-PerfAllocation', 'Invoke-PerfRuns') + FunctionsToExport = @('Invoke-PerfPing', 'Invoke-PerfEcho', 'Invoke-PerfCpu', 'Invoke-PerfSleep', 'Invoke-PerfJson', 'Invoke-PerfFile', 'Invoke-PerfBgEnqueue', 'Push-PerfBg', 'Push-PerfBgLeaf', 'Invoke-PerfManyRuns', 'Push-PerfHold', 'Invoke-PerfThreads', 'Invoke-PerfTableOp', 'Push-PerfCheck', 'Invoke-PerfCheckCounts', 'Push-PerfSeq', 'Invoke-PerfSeqResult', 'Push-PerfSeqWorker', 'Invoke-PerfSeqWorkerEnqueue', 'Invoke-PerfSeqWorkerResult', 'Invoke-PerfThreadBreakdown', 'Invoke-PerfSeedRuns', 'Invoke-ListPerf', 'Invoke-PerfWhoami', 'Invoke-PerfTimerTick', 'Invoke-PerfTimerCount', 'Invoke-PerfPublish', 'Invoke-PerfAllocation', 'Invoke-PerfRuns') CmdletsToExport = @() VariablesToExport = @() AliasesToExport = @() diff --git a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 index cdf4f8f..8b3661f 100644 --- a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 +++ b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 @@ -539,6 +539,15 @@ function Invoke-PerfJson { return @{ StatusCode = 200; Body = @{ ok = $true; endpoint = 'PerfJson'; count = $n; items = @($items) } } } +# Serves a captured real API response verbatim (?name= under /payloads, mounted by the compose file +# from PAYLOAD_DIR). A string body that parses as JSON goes out byte-for-byte, so compression sweeps run +# against real production-shaped JSON rather than the synthetic PerfJson rows. +function Invoke-PerfFile { + param($Request, $TriggerMetadata) + $name = [IO.Path]::GetFileName([string]$Request.Query.name) + return @{ StatusCode = 200; Body = [IO.File]::ReadAllText("/payloads/$name") } +} + # Realtime bridge driver: publishes a job event via the C# RealtimeBridge so an SSE consumer of # /.craft/events can observe it. userId is taken from the caller's identity so it is delivered back to # the same principal. Query: ?jobId=&mode=start|update|end&size=. diff --git a/perf-harness/api-harness/payloads/.gitkeep b/perf-harness/api-harness/payloads/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/perf-harness/docker-compose.api.yml b/perf-harness/docker-compose.api.yml index 68acc09..b2d169f 100644 --- a/perf-harness/docker-compose.api.yml +++ b/perf-harness/docker-compose.api.yml @@ -19,6 +19,9 @@ services: volumes: # Overlay the synthetic API module onto the base image's empty /app/API. Read-only; live-editable. - ./api-harness/API:/app/API:ro + # Captured real API responses for PerfFile (?name=). Default is an empty in-repo dir; point + # PAYLOAD_DIR at a folder of captured JSON (kept out of git - it is real tenant data). + - ${PAYLOAD_DIR:-./api-harness/payloads}:/payloads:ro environment: - ASPNETCORE_ENVIRONMENT=Production - ASPNETCORE_URLS=http://+:8080 diff --git a/perf-harness/k6/api_load.js b/perf-harness/k6/api_load.js index eaea99d..6ddd80a 100644 --- a/perf-harness/k6/api_load.js +++ b/perf-harness/k6/api_load.js @@ -10,7 +10,8 @@ // 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). +// URL drive this one path instead of the catalogue (e.g. /API/PerfFile?name=listlogs.json) +// ENC Accept-Encoding to negotiate: 'gzip' | 'br' | '' (default '' = identity). // 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 @@ -41,13 +42,15 @@ const endpoints = [ { name: 'PerfJson', url: `/API/PerfJson?n=${JSON_N}`, weight: 2 }, ]; -const active = ONLY ? endpoints.filter((e) => e.name === ONLY) : endpoints; +if (__ENV.URL) endpoints.push({ name: 'Custom', url: __ENV.URL, weight: 1 }); +const ONLY_EFF = __ENV.URL ? 'Custom' : ONLY; +const active = ONLY_EFF ? endpoints.filter((e) => e.name === ONLY_EFF) : endpoints; if (active.length === 0) throw new Error(`ONLY=${ONLY} matched no endpoint`); // Weighted pick list (each endpoint repeated `weight` times); single-endpoint mode picks it every time. const pick = []; for (const e of active) { - const n = ONLY ? 1 : e.weight; + const n = ONLY_EFF ? 1 : e.weight; for (let i = 0; i < n; i++) pick.push(e); } diff --git a/perf-harness/scripts/run-compression-levels.ps1 b/perf-harness/scripts/run-compression-levels.ps1 index b6dff31..c13afb8 100644 --- a/perf-harness/scripts/run-compression-levels.ps1 +++ b/perf-harness/scripts/run-compression-levels.ps1 @@ -21,6 +21,9 @@ pwsh scripts\run-compression-levels.ps1 -Build pwsh scripts\run-compression-levels.ps1 -Levels Fastest,Optimal,SmallestSize -Encodings gzip,br pwsh scripts\run-compression-levels.ps1 -JsonN 8000 -Rate 20 -Pool 4 + # A captured real response (served verbatim by PerfFile from -PayloadDir), all four encodings: + pwsh scripts\run-compression-levels.ps1 -PayloadDir C:\captures -Url '/API/PerfFile?name=listlogs.json' ` + -Encodings br,gzip,identity #> [CmdletBinding()] param( @@ -33,6 +36,8 @@ param( [int]$Rate = 10, [string]$Duration = '20s', [int]$JsonN = 2000, + [string]$Url = '', # path to drive instead of PerfJson, e.g. /API/PerfFile?name=listlogs.json + [string]$PayloadDir = '', # host folder mounted at /payloads for PerfFile [int]$ReadyTimeoutSec = 120, [switch]$Build, [switch]$KeepUp @@ -60,6 +65,8 @@ if($Build){ if($LASTEXITCODE -ne 0){ throw "docker build failed" } } +if(-not $Url){ $Url = "/API/PerfJson?n=$JsonN" } +if($PayloadDir){ $env:PAYLOAD_DIR = (Resolve-Path $PayloadDir).Path } $env:SUT_IMAGE = $SutImage; $env:SUT_PORT = "$Port"; $env:SUT_CPUS = "$Cpus"; $env:POOL = "$Pool" # One curl to /API/PerfJson at the given encoding ('' = identity). Returns the on-the-wire body size @@ -67,7 +74,7 @@ $env:SUT_IMAGE = $SutImage; $env:SUT_PORT = "$Port"; $env:SUT_CPUS = "$Cpus"; $e function WireBytes([string]$Enc){ $a = @('-s','-o','NUL','--max-time','30','-w','%{size_download}') if($Enc){ $a += @('-H', "Accept-Encoding: $Enc") } - $a += "$base/API/PerfJson?n=$JsonN" + $a += "$base$Url" [long](& curl.exe @a 2>$null) } @@ -96,7 +103,7 @@ function LoadOne([string]$Enc, [string]$tag){ $summary = Join-Path $resultsDir "levels-$tag-$stamp.k6.json" & docker run --rm --network $network ` -e BASE="http://sut:8080" -e RATE="$Rate" -e DURATION="$Duration" ` - -e ONLY='PerfJson' -e JSON_N="$JsonN" -e ENC="$Enc" ` + -e URL="$Url" -e ENC="$Enc" ` -v "${k6DirD}:/scripts:ro" -v "${resultsDirD}:/out" ` grafana/k6 run /scripts/api_load.js --summary-export "/out/$(Split-Path $summary -Leaf)" 2>&1 | Out-Null @@ -112,6 +119,7 @@ function LoadOne([string]$Enc, [string]$tag){ cpuAvgPct = if($cpu){ [math]::Round(($cpu|Measure-Object -Average).Average,1) } else { $null } cpuMaxPct = if($cpu){ [math]::Round(($cpu|Measure-Object -Maximum).Maximum,1) } else { $null } reqPerSec = [math]::Round(([double](MetricVal $k6 'http_reqs' 'rate')),1) + p50Ms = [math]::Round(([double](MetricVal $k6 'http_req_duration' 'med')),2) p95Ms = [math]::Round(([double](MetricVal $k6 'http_req_duration' 'p(95)')),2) wirePerResp = if($reqs -gt 0){ [long]($recv/$reqs) } else { 0 } } @@ -133,6 +141,10 @@ try { # Warm once, capture the identity baseline once (level-independent). [void](WireBytes '') + # Discarded load pass over every encoding: a fresh container JITs its compressors and grows its + # runspace pool under the first load, which otherwise lands on whichever encoding is measured first. + Info " warm-up pass (discarded) ..." + foreach($enc in $Encodings){ [void](LoadOne $enc "$level-warm-$enc") } if($rawBytes -eq 0){ $rawBytes = WireBytes '' } foreach($enc in $Encodings){ @@ -142,7 +154,7 @@ try { $load = LoadOne $enc "$level-$enc" $rows.Add([ordered]@{ level=$level; encoding=$enc; wireBytes=$wire; ratio=$ratio - cpuAvgPct=$load.cpuAvgPct; cpuMaxPct=$load.cpuMaxPct; p95Ms=$load.p95Ms; reqPerSec=$load.reqPerSec + cpuAvgPct=$load.cpuAvgPct; cpuMaxPct=$load.cpuMaxPct; p50Ms=$load.p50Ms; p95Ms=$load.p95Ms; reqPerSec=$load.reqPerSec }) } } @@ -150,7 +162,7 @@ try { # ── report ──────────────────────────────────────────────────────────────────── $result = [ordered]@{ label='compression-levels'; timestamp=$stamp; sutImage=$SutImage - config=@{ pool=$Pool; cpus=$Cpus; jsonN=$JsonN; rate=$Rate; duration=$Duration; levels=$Levels; encodings=$Encodings } + config=@{ pool=$Pool; cpus=$Cpus; url=$Url; rate=$Rate; duration=$Duration; levels=$Levels; encodings=$Encodings } identityBytes=$rawBytes rows=$rows } @@ -161,13 +173,13 @@ try { $sb = [System.Text.StringBuilder]::new() [void]$sb.AppendLine("# /api compression level sweep ($stamp)") [void]$sb.AppendLine("") - [void]$sb.AppendLine("- **SUT:** ``$SutImage`` **CPUs:** $Cpus **Pool:** $Pool **Payload:** PerfJson n=$JsonN ($rawBytes B raw)") + [void]$sb.AppendLine("- **SUT:** ``$SutImage`` **CPUs:** $Cpus **Pool:** $Pool **Payload:** ``$Url`` ($rawBytes B raw)") [void]$sb.AppendLine("- **Load:** rate=$Rate for $Duration per (level x encoding)") [void]$sb.AppendLine("") - [void]$sb.AppendLine("| level | encoding | wire B/resp | ratio | CPU% avg | CPU% max | p95 ms | req/s |") - [void]$sb.AppendLine("|---|---|---:|---:|---:|---:|---:|---:|") + [void]$sb.AppendLine("| level | encoding | wire B/resp | ratio | CPU% avg | CPU% max | p50 ms | p95 ms | req/s |") + [void]$sb.AppendLine("|---|---|---:|---:|---:|---:|---:|---:|---:|") foreach($r in $rows){ - [void]$sb.AppendLine("| $($r.level) | $($r.encoding) | $($r.wireBytes) | $($r.ratio)x | $($r.cpuAvgPct) | $($r.cpuMaxPct) | $($r.p95Ms) | $($r.reqPerSec) |") + [void]$sb.AppendLine("| $($r.level) | $($r.encoding) | $($r.wireBytes) | $($r.ratio)x | $($r.cpuAvgPct) | $($r.cpuMaxPct) | $($r.p50Ms) | $($r.p95Ms) | $($r.reqPerSec) |") } $sb.ToString() | Set-Content $mdOut -Encoding utf8 @@ -177,7 +189,7 @@ try { Get-Content $mdOut | Write-Host } finally { - Remove-Item Env:\API_COMPRESSION_LEVEL -ErrorAction SilentlyContinue + Remove-Item Env:\API_COMPRESSION_LEVEL, Env:\PAYLOAD_DIR -ErrorAction SilentlyContinue if($KeepUp){ Warn "leaving containers up (-KeepUp). Tear down: docker compose -f `"$composeApi`" down -v" } else { Info "tearing down ..."; docker compose -f $composeApi down -v 2>&1 | Out-Null } } diff --git a/tests/Craft.Tests/ApiCompressionPipelineTests.cs b/tests/Craft.Tests/ApiCompressionPipelineTests.cs index ea0b4d2..60938e5 100644 --- a/tests/Craft.Tests/ApiCompressionPipelineTests.cs +++ b/tests/Craft.Tests/ApiCompressionPipelineTests.cs @@ -1,3 +1,4 @@ +using System.IO.Compression; using System.Net; using System.Text; using Craft.Hosting; @@ -55,6 +56,15 @@ private EgressLedger Ledger() => private async Task<(string encoding, long bytes)> Request( Action configure, string contentType = "application/json", bool routed = false) { + var (enc, body) = await RequestBody(configure, contentType, routed); + return (enc, body.LongLength); + } + + private async Task<(string encoding, byte[] body)> RequestBody( + Action configure, string contentType = "application/json", bool routed = false, + string acceptEncoding = "gzip", string? payload = null) + { + payload ??= Payload; var builder = WebApplication.CreateBuilder(); builder.WebHost.ConfigureKestrel(o => o.Listen(IPAddress.Loopback, 0)); // dynamic port builder.Logging.ClearProviders(); @@ -70,7 +80,7 @@ async Task Terminal(HttpContext ctx) ctx.Response.StatusCode = 200; ctx.Response.ContentType = contentType; ctx.Response.Headers["X-Cache"] = "MISS"; - await ctx.Response.WriteAsync(Payload); + await ctx.Response.WriteAsync(payload); } // routed = the real shape: a mapped endpoint (MapMethods) executed by the endpoint middleware, @@ -87,16 +97,64 @@ async Task Terminal(HttpContext ctx) using var handler = new HttpClientHandler { AutomaticDecompression = DecompressionMethods.None }; using var client = new HttpClient(handler); var req = new HttpRequestMessage(HttpMethod.Get, $"{addr}/API/thing"); - req.Headers.TryAddWithoutValidation("Accept-Encoding", "gzip"); + req.Headers.TryAddWithoutValidation("Accept-Encoding", acceptEncoding); var resp = await client.SendAsync(req); var enc = resp.Content.Headers.ContentEncoding.FirstOrDefault() ?? ""; - var bytes = (await resp.Content.ReadAsByteArrayAsync()).LongLength; - return (enc, bytes); + return (enc, await resp.Content.ReadAsByteArrayAsync()); } finally { await app.StopAsync(); } } + [Theory] + [InlineData("br")] + [InlineData("gzip")] + public async Task EachEncoding_NegotiatesAlone_AndRoundTrips(string encoding) + { + var (enc, body) = await RequestBody(app => app.UseResponseCompression(), acceptEncoding: encoding); + Assert.Equal(encoding, enc); + Assert.True(body.LongLength < RawLength, $"expected compressed < {RawLength}, got {body.LongLength}"); + + using var input = new MemoryStream(body); + using Stream decoder = encoding switch + { + "br" => new BrotliStream(input, CompressionMode.Decompress), + _ => new GZipStream(input, CompressionMode.Decompress), + }; + using var reader = new StreamReader(decoder, Encoding.UTF8); + Assert.Equal(Payload, await reader.ReadToEndAsync()); + } + + [Fact] + public async Task StringBody_ReachesEncoderCoalesced_NotAsPipeFragments() + { + // Response.WriteAsync(string) arrives at the encoder as ~4 KiB pipe segments. Brotli Fastest + // compresses each write as an isolated fragment, which made a real 300 KB body ~2.9x the size of + // compressing it in one go. Coalescing keeps it close to one-shot. Varied JSON, not a repeated + // string, so the fragment penalty would actually show. + var varied = "[" + string.Join(",", Enumerable.Range(0, 6000).Select(i => + $"{{\"id\":{i},\"tenant\":\"t{i % 97}.onmicrosoft.com\",\"msg\":\"event {i * 7919 % 10007} on {i % 13}\"}}")) + "]"; + var raw = Encoding.UTF8.GetBytes(varied); + var oneShot = new byte[BrotliEncoder.GetMaxCompressedLength(raw.Length)]; + Assert.True(BrotliEncoder.TryCompress(raw, oneShot, out var oneShotLength, quality: 1, window: 22)); + + var (enc, body) = await RequestBody(app => app.UseResponseCompression(), acceptEncoding: "br", payload: varied); + + Assert.Equal("br", enc); + Assert.True(body.Length < oneShotLength * 1.35, + $"br Fastest wire {body.Length} B vs one-shot {oneShotLength} B — encoder is being fed fragments"); + using var decoded = new StreamReader(new BrotliStream(new MemoryStream(body), CompressionMode.Decompress)); + Assert.Equal(varied, await decoded.ReadToEndAsync()); + } + + [Fact] + public async Task BrowserAcceptEncoding_StillPrefersBrotli() + { + // Chrome/Firefox send all four at equal q; br must win the tie. + var (enc, _) = await RequestBody(app => app.UseResponseCompression(), acceptEncoding: "gzip, deflate, br, zstd"); + Assert.Equal("br", enc); + } + // ── pipeline shapes, outer→inner, that isolate where compression is lost ────────────────────────── [Fact] From 2fce627f009a2063bd89bc53b1ea1ba8c72dc5a4 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Thu, 24 Sep 2026 01:25:42 +0800 Subject: [PATCH 6/6] perf(build): run gzip on zlib-ng instead of the system zlib .NET 8 on Linux does gzip through the system libz.so.1. Stage zlib-ng 2.3.3 (zlib-compat build, pinned + sha256-verified, runtime CPU dispatch) and copy it onto the runtime's libz.so.1, the same way mimalloc is staged. Azure Linux runtime, 2 vCPU, default level (Optimal = zlib 6): - 10 MB ListLogs body: 165 -> 84 ms per encode, 3.8% smaller on the wire - under load (3 rps of 10 MB): gzip CPU over identity ~60 -> ~37 points - 32 KB: 296 -> 187 us; <= 2 KB unchanged Every case round-trips byte-identical. Brotli is unaffected. --- build/Dockerfile | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/build/Dockerfile b/build/Dockerfile index 3997ae8..e801793 100644 --- a/build/Dockerfile +++ b/build/Dockerfile @@ -102,6 +102,26 @@ RUN set -eux; \ esac; \ install -D /usr/lib/${TRIPLET}/libmimalloc.so.2 /mimalloc/libmimalloc.so.2 +# zlib-ng built in zlib-compat mode, dropped in over the runtime's libz.so.1. .NET 8 on Linux does gzip +# through the system zlib (libSystem.IO.Compression.Native links libz.so.1), so this is what /api gzip +# runs on. Measured on this runtime, 2 vCPU, at the default level (Optimal = zlib 6): a 10 MB ListLogs +# body 165 -> 84 ms and 4% smaller; 32 KB 296 -> 187 us; <=2 KB unchanged. Fastest (level 1) gets ~3x +# faster but ~24% bigger (zlib-ng's quick strategy). Same drop-in Fedora ships as its system zlib. +# Built from a pinned, checksummed release (not packaged for Azure Linux or bookworm); runtime CPU +# dispatch picks SIMD paths per host. Bump ZLIBNG_VERSION + ZLIBNG_SHA256 together. +FROM ${DOTNET_REGISTRY}/dotnet/aspnet:8.0-bookworm-slim AS zlibng-src +ARG ZLIBNG_VERSION=2.3.3 +ARG ZLIBNG_SHA256=f9c65aa9c852eb8255b636fd9f07ce1c406f061ec19a2e7d508b318ca0c907d1 +RUN set -eux; \ + apt-get update && apt-get install -y --no-install-recommends ca-certificates curl cmake gcc make libc6-dev; \ + curl -fsSL -o /tmp/zlib-ng.tar.gz "https://github.com/zlib-ng/zlib-ng/archive/refs/tags/${ZLIBNG_VERSION}.tar.gz"; \ + echo "${ZLIBNG_SHA256} /tmp/zlib-ng.tar.gz" | sha256sum -c -; \ + tar -xzf /tmp/zlib-ng.tar.gz -C /tmp; \ + cmake -S /tmp/zlib-ng-${ZLIBNG_VERSION} -B /tmp/zlib-ng-build \ + -DZLIB_COMPAT=ON -DBUILD_TESTING=OFF -DZLIB_ENABLE_TESTS=OFF -DCMAKE_BUILD_TYPE=Release; \ + cmake --build /tmp/zlib-ng-build -j"$(nproc)"; \ + install -D "$(readlink -f /tmp/zlib-ng-build/libz.so.1)" /zlibng/libz.so.1 + # distroless-extra = distroless + icu/tzdata. PowerShell needs globalization, so # the plain -distroless tag (no ICU) would break it — the -extra tag is required. FROM ${DOTNET_REGISTRY}/dotnet/aspnet:8.0-azurelinux3.0-distroless-extra AS runtime @@ -119,6 +139,9 @@ FROM ${DOTNET_REGISTRY}/dotnet/aspnet:8.0-azurelinux3.0-distroless-extra AS runt # ENV PATH="/usr/share/powershell:${PATH}" COPY --from=mimalloc-src /mimalloc/libmimalloc.so.2 /usr/lib/libmimalloc.so.2 +# Onto the soname path the loader resolves, not the versioned file behind it: the base tag floats, so a +# zlib bump renaming libz.so.1.3.x must not silently strand this copy. +COPY --from=zlibng-src /zlibng/libz.so.1 /usr/lib/libz.so.1 # Create /app owned by the non-root user. COPYing with NO preceding WORKDIR makes Docker CREATE the # destination /app with this ownership; a `WORKDIR /app` first would pre-create it root-owned, leaving