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/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..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); + return await ExecuteHttpScriptInternal(route, req, isHttp: true, clientAborted: clientAborted); } // Profiling path: time request marshaling + the runner-side segments (checkout/invoke/extract). @@ -241,7 +248,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: clientAborted); DispatchProfiler.Record(marshalTicks, timing.CheckoutTicks, timing.InvokeTicks, timing.ExtractTicks, Stopwatch.GetTimestamp() - totalStart); return result; @@ -258,7 +266,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 +373,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 +393,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 +412,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); 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..0f44750 --- /dev/null +++ b/Services/Storage/EntitySplitter.cs @@ -0,0 +1,789 @@ +using System.Text; +using System.Text.Json; +using Azure.Data.Tables; + +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/build/Dockerfile b/build/Dockerfile index bfe1837..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 @@ -135,6 +158,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" 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] 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); + } + } +}