diff --git a/docs/design/core/shared/streaming.md b/docs/design/core/shared/streaming.md index 3f8494b6e..d7c3ec146 100644 --- a/docs/design/core/shared/streaming.md +++ b/docs/design/core/shared/streaming.md @@ -21,14 +21,15 @@ flowchart LR tee -. "tee for inline verify" .-> ver[RoundTripVerifier] ``` -- **`ProgressStream(inner, IProgress)`** wraps the *source* at the top of the chain so progress tracks logical bytes read off disk, not compressed bytes on the wire. `Length` is delegated to the inner `FileStream`, so the consumer knows the total up front and can compute a percentage. It reports after each read; a zero-length read reports nothing. The download path reuses it the same way, wrapping the blob's read stream. +- **`ProgressStream(inner, IProgress)`** wraps the *source* at the top of the chain so progress tracks logical bytes read off disk, not compressed bytes on the wire. `Length` is delegated to the inner `FileStream`, so the consumer knows the total up front and can compute a percentage. Reports are **coalesced to at most one per 500 ms**: the first read reports immediately, later reads wait for the interval, and the true total is emitted at EOF or disposal. Empty sources report nothing. The download path reuses it the same way, wrapping the blob's read stream. - **`CountingStream(inner)`** sits at the *bottom* of the chain, directly above `OpenWriteAsync`, and increments `BytesWritten` on every write. It is read *after* the chain is disposed to capture the final compressed-and-encrypted blob size, which `UploadChunkAsync` then writes into blob metadata (`chunk-size`). The note in `UploadChunkAsync` is load-bearing here: the encryption stream is disposed *explicitly* before reading `BytesWritten`, because GCM flushes its final auth tag on dispose — reading the count earlier would undercount by the tag bytes. ## Key invariants -- **No buffering proportional to file size.** Both wrappers hold only a `long` counter; the chain streams a multi-GB file without an O(file-size) allocation. See [memory-boundedness](../../cross-cutting/memory-boundedness.md). +- **No buffering proportional to file size.** Both wrappers hold only a few `long` counters; the chain streams a multi-GB file without an O(file-size) allocation. See [memory-boundedness](../../cross-cutting/memory-boundedness.md). +- **A coalesced `ProgressStream` still ends on the true total.** Throttling may drop intermediate reports, but the EOF and dispose flushes must stay, or a consumer is left short of 100% by up to one interval's worth of bytes. - **`CountingStream.BytesWritten` is valid only after the whole chain is finalized.** Every layer above it (compression frame, GCM tag) must be flushed/disposed first, or the recorded `chunk-size` is short. - **`ProgressStream` is read-only, `CountingStream` is write-only.** They guard their unsupported direction with `NotSupportedException` rather than silently no-op'ing, so a misuse fails loudly. - **Counting reflects what was actually persisted.** `CountingStream` wraps the Azure write stream, so its total is the blob's real stored size, not a pre-computed estimate. @@ -39,5 +40,5 @@ Splitting "report read progress" and "count written bytes" into separate one-lin ## Open seams / future -- `ProgressStream` reports raw read counts; smoothing/throttling of the `IProgress` callback is left to the consumer (`UploadChunkAsync` already de-dupes non-increasing reports via a `CallbackProgress`). +- `ProgressStream` throttles reports at the source to one per 500 ms. `UploadChunkAsync`'s `CallbackProgress` still de-dupes non-increasing reports on top of it. The interval is a fixed constant, not a per-consumer setting. - Any future upload layer (e.g. a second integrity tee) slots into the same push chain in `UploadChunkAsync` between source and `OpenWriteAsync`; these two wrappers stay unchanged as the progress/size endpoints. diff --git a/docs/design/cross-cutting/memory-boundedness.md b/docs/design/cross-cutting/memory-boundedness.md index ae2965d08..5e2011983 100644 --- a/docs/design/cross-cutting/memory-boundedness.md +++ b/docs/design/cross-cutting/memory-boundedness.md @@ -46,7 +46,7 @@ Earlier the [chunk index](../../glossary.md#chunk-index) kept an in-memory LRU o ### 4. Streaming up/download — bound the single-file byte axis -A single multi-GB file must not be buffered. `ChunkStorageService.UploadChunkAsync` assembles a push chain (source → `ProgressStream` → zstd → encryption → `CountingStream` → `OpenWriteAsync`) with no `MemoryStream` and no intermediate temp file (the tar bundle being the one streamed-in exception). The `Streaming` decorators hold only a `long` counter, so the chain moves a file of any size with a working set of roughly one buffer. This is the byte-scale complement to the file-count techniques above — see [streaming](../core/shared/streaming.md). +A single multi-GB file must not be buffered. `ChunkStorageService.UploadChunkAsync` assembles a push chain (source → `ProgressStream` → zstd → encryption → `CountingStream` → `OpenWriteAsync`) with no `MemoryStream` and no intermediate temp file (the tar bundle being the one streamed-in exception). The `Streaming` decorators hold only a few counters, so the chain moves a file of any size with a working set of roughly one buffer. This is the byte-scale complement to the file-count techniques above — see [streaming](../core/shared/streaming.md). ## Key invariants diff --git a/src/Arius.Api.FakeTestHost/CanonicalScenarios.cs b/src/Arius.Api.FakeTestHost/CanonicalScenarios.cs index a35e8cb12..98c1a4ee4 100644 --- a/src/Arius.Api.FakeTestHost/CanonicalScenarios.cs +++ b/src/Arius.Api.FakeTestHost/CanonicalScenarios.cs @@ -18,7 +18,7 @@ public static class CanonicalScenarios new FileHashingEvent(RelativePath.Parse("big.bin"), 100_000_000), new FileHashedEvent(RelativePath.Parse("big.bin"), ContentHash.Parse(new string('a', 64)), FastHashReused: false, FastHashRehashed: true, FileSize: 100_000_000), new ChunkUploadedEvent(ChunkHash.Parse(new string('a', 64)), StoredSize: 60_000_000, OriginalSize: 100_000_000), - new FileDedupedEvent(ContentHash.Parse(new string('b', 64)), OriginalSize: 48_000_000), + new FileDedupedEvent(RelativePath.Parse("deduped.bin"), ContentHash.Parse(new string('b', 64)), OriginalSize: 48_000_000), new RoutingCompleteEvent(NewByteTotal: 100_000_000), // dedup/route drained → exact new-byte total new SnapshotCreatedEvent(default, DateTimeOffset.UnixEpoch, 3122), ], diff --git a/src/Arius.Api.Integration.Tests/RepresentationScenarioTests.cs b/src/Arius.Api.Integration.Tests/RepresentationScenarioTests.cs index 5fd60ebb9..67906be9b 100644 --- a/src/Arius.Api.Integration.Tests/RepresentationScenarioTests.cs +++ b/src/Arius.Api.Integration.Tests/RepresentationScenarioTests.cs @@ -25,7 +25,7 @@ public async Task Pointer_heavy_archive_reports_additive_new_bytes_not_underflow Events: [ new ScanCompleteEvent(TotalFiles: 1001, TotalBytes: 100_000_000), // pointer-only files scanned as 0 - new FileDedupedEvent(ContentHash.Parse(new string('b', 64)), OriginalSize: 1_000_000_000), // pointer-only dedup, full size + new FileDedupedEvent(RelativePath.Parse("deduped.bin"), ContentHash.Parse(new string('b', 64)), OriginalSize: 1_000_000_000), // pointer-only dedup, full size new ChunkUploadingEvent(ChunkHash.Parse(new string('c', 64)), 100_000_000), // one new chunk queued new ChunkUploadedEvent(ChunkHash.Parse(new string('c', 64)), StoredSize: 60_000_000, OriginalSize: 100_000_000), ], diff --git a/src/Arius.Api.Tests/Jobs/JobSinkAggregateTests.cs b/src/Arius.Api.Tests/Jobs/JobSinkAggregateTests.cs index 6022fdb2f..87ae960bd 100644 --- a/src/Arius.Api.Tests/Jobs/JobSinkAggregateTests.cs +++ b/src/Arius.Api.Tests/Jobs/JobSinkAggregateTests.cs @@ -112,7 +112,7 @@ public async Task Archive_forwarders_populate_byte_layers() await new ScanCompleteForwarder(s).Handle(new ScanCompleteEvent(2, 3000), default); await new FileScannedForwarder(s).Handle(new FileScannedEvent(RelativePath.Parse("a"), 2000), default); await new FileHashingForwarder(s).Handle(new FileHashingEvent(RelativePath.Parse("a"), 2000), default); - await new FileDedupedForwarder(s).Handle(new FileDedupedEvent(ContentHash.Parse(new string('b', 64)), 1000), default); + await new FileDedupedForwarder(s).Handle(new FileDedupedEvent(RelativePath.Parse("deduped.bin"), ContentHash.Parse(new string('b', 64)), 1000), default); await new ChunkUploadingForwarder(s).Handle(new ChunkUploadingEvent(ChunkHash.Parse(new string('d', 64)), 2000), default); await new ChunkUploadedForwarder(s).Handle(new ChunkUploadedEvent(ChunkHash.Parse(new string('c', 64)), 300, 2000), default); diff --git a/src/Arius.Benchmarks/AllocationBenchmarks.cs b/src/Arius.Benchmarks/AllocationBenchmarks.cs new file mode 100644 index 000000000..03dfaa78d --- /dev/null +++ b/src/Arius.Benchmarks/AllocationBenchmarks.cs @@ -0,0 +1,269 @@ +using Arius.Core.Features.ArchiveCommand; +using Arius.Core.Shared.ChunkIndex; +using Arius.Core.Shared.Encryption; +using Arius.Core.Shared.FileSystem; +using Arius.Core.Shared.FileTree; +using Arius.Core.Shared.HashCache; +using Arius.Core.Shared.Hashes; +using Arius.Core.Shared.Storage; +using Arius.Tests.Shared; +using BenchmarkDotNet.Attributes; + +namespace Arius.Benchmarks; + +/// +/// In-process allocation benchmarks for selected Arius.Core components. +/// Run with micro (e.g. micro --filter '*TarBuilder*'); no Azurite or Docker is required. +/// +[MemoryDiagnoser] +public class AllocationBenchmarks +{ + private const int EntryCount = 1_000; + + /// + /// Digest containing hexadecimal letters, ensuring the lowercase conversion path is exercised. + /// + private readonly byte[] _digest = CreateHighNibbleDigest(); + + private const string CanonicalHex = "00112233445566778899aabbccddeeff00112233445566778899aabbccddeeff"; + private const string UppercaseHex = "00112233445566778899AABBCCDDEEFF00112233445566778899AABBCCDDEEFF"; + + [Benchmark(Description = "HashCodec.ToLowerHex")] + public string HashCodec_ToLowerHex() => HashCodec.ToLowerHex(_digest); + + /// Parses an already canonical lowercase value. + [Benchmark(Description = "HashCodec.NormalizeHex (already canonical)")] + public string HashCodec_NormalizeHex_Canonical() => HashCodec.NormalizeHex(CanonicalHex); + + [Benchmark(Description = "HashCodec.NormalizeHex (uppercase input)")] + public string HashCodec_NormalizeHex_Uppercase() => HashCodec.NormalizeHex(UppercaseHex); + + // ── Sparse fingerprint ─────────────────────────────────────────────────────── + + private const long SmallFileSize = 200L * 1024; // one whole-file region + private const long LargeFileSize = 64L * 1024 * 1024 * 1024; // k = MaxBlocks = 64 regions + + private byte[] _readBuffer = null!; + + /// Exercises the single-region sampler path for a small file. + [Benchmark(Description = "SparseFingerprint.Sampler small file (200 KB)")] + public byte[] SparseFingerprint_Sampler_SmallFile() + { + using var sampler = new SparseFingerprint.Sampler(SmallFileSize); + var position = 0L; + + while (position < SmallFileSize) + { + var length = (int)Math.Min(_readBuffer.Length, SmallFileSize - position); + sampler.Capture(position, _readBuffer.AsSpan(0, length)); + position += length; + } + + return sampler.Finish(); + } + + /// Exercises the maximum sampled-region buffer size. + [Benchmark(Description = "SparseFingerprint.Sampler large file (64 GB logical)")] + public byte[] SparseFingerprint_Sampler_LargeFile() + { + using var sampler = new SparseFingerprint.Sampler(LargeFileSize); + + foreach (var (offset, length) in SparseFingerprint.Regions(LargeFileSize)) + sampler.Capture(offset, _readBuffer.AsSpan(0, length)); + + return sampler.Finish(); + } + + // ── Filetree serialization ─────────────────────────────────────────────────── + + private IReadOnlyList _fileTreeEntries = null!; + private byte[] _fileTreeBytes = null!; + + [Benchmark(Description = "FileTreeSerializer.Serialize (1000 entries)")] + public byte[] FileTreeSerializer_Serialize() => FileTreeSerializer.Serialize(_fileTreeEntries); + + [Benchmark(Description = "FileTreeSerializer.Deserialize (1000 entries)")] + public IReadOnlyList FileTreeSerializer_Deserialize() => FileTreeSerializer.Deserialize(_fileTreeBytes); + + // ── Tar builder ────────────────────────────────────────────────────────────── + + private const long TarTargetSize = 64L * 1024 * 1024; + private const int TarEntrySize = 64 * 1024; + + private byte[] _tarEntryPayload = null!; + private byte[] _smallEntryPayload = null!; + private IEncryptionService _encryption = null!; + + /// Builds and seals one 64 MB TAR bundle. + [Benchmark(Description = "TarBuilder seal one 64 MB bundle")] + public async Task TarBuilder_Seal_64MB() + { + await using var builder = new TarBuilder(TarTargetSize, _encryption); + + var sealedCount = 0; + for (var i = 0; i < TarTargetSize / TarEntrySize; i++) + { + var source = new MemoryStream(_tarEntryPayload, writable: false); + if (await builder.AddAsync(CreateUpload(i, TarEntrySize), source, CancellationToken.None) is not null) + sealedCount++; + } + + return sealedCount; + } + + /// A handful of small files: the shape of a tiny repository or a small tail bundle. + [Benchmark(Description = "TarBuilder seal one small bundle (5 x 1 KB)")] + public async Task TarBuilder_Seal_SmallBundle() + { + await using var builder = new TarBuilder(TarTargetSize, _encryption); + + for (var i = 0; i < 5; i++) + { + var source = new MemoryStream(_smallEntryPayload, writable: false); + await builder.AddAsync(CreateUpload(i, _smallEntryPayload.Length), source, CancellationToken.None); + } + + return (await builder.SealAsync(CancellationToken.None))!.Entries.Count; + } + + // ── Chunk-index local store ────────────────────────────────────────────────── + + private LocalDirectory _storeRoot = default; + private ChunkIndexLocalStore _store = null!; + private ShardEntry[] _shardEntries = null!; + private ContentHash[] _lookupHashes = null!; + private PathSegment _rangePrefix = default; + + [Benchmark(Description = "ChunkIndexLocalStore.UpsertRemoteBacked (1000 rows)")] + public void ChunkIndexLocalStore_UpsertRemoteBacked() => _store.UpsertRemoteBacked(_shardEntries); + + [Benchmark(Description = "ChunkIndexLocalStore.ReadRangeEntries (1000 rows)")] + public int ChunkIndexLocalStore_ReadRangeEntries() + { + var count = 0; + _store.ReadRangeEntries(_rangePrefix, _ => count++); + return count; + } + + /// Measures the per-hash lookup shape for a 256-hash deduplication batch. + [Benchmark(Description = "ChunkIndexLocalStore.FindEntry x256 (one dedup batch)")] + public int ChunkIndexLocalStore_FindEntry_256() + { + var found = 0; + foreach (var hash in _lookupHashes) + if (_store.FindEntry(hash) is not null) + found++; + + return found; + } + + /// Measures the batched lookup for a 256-hash deduplication batch. + [Benchmark(Description = "ChunkIndexLocalStore.FindEntries x1 (one dedup batch)")] + public int ChunkIndexLocalStore_FindEntries_Batch() => _store.FindEntries(_lookupHashes).Count; + + // ── Setup ──────────────────────────────────────────────────────────────────── + + [GlobalSetup] + public void Setup() + { + _readBuffer = new byte[256 * 1024]; + _tarEntryPayload = new byte[TarEntrySize]; + _smallEntryPayload = new byte[1024]; + Random.Shared.NextBytes(_readBuffer); + Random.Shared.NextBytes(_tarEntryPayload); + Random.Shared.NextBytes(_smallEntryPayload); + + _encryption = IEncryptionService.EncryptedInstance; + + _fileTreeEntries = BuildFileTreeEntries(EntryCount); + _fileTreeBytes = FileTreeSerializer.Serialize(_fileTreeEntries); + + _shardEntries = BuildShardEntries(EntryCount); + _lookupHashes = _shardEntries.Take(256).Select(e => e.ContentHash).ToArray(); + _rangePrefix = PathSegment.Parse("00"); + + _storeRoot = TestTempRoots.CreateDirectory("benchmark-chunkindex"); + _store = new ChunkIndexLocalStore(_storeRoot); + _store.UpsertRemoteBacked(_shardEntries); + } + + [GlobalCleanup] + public void Cleanup() + { + try + { + RelativeFileSystem.DeleteDirectory(_storeRoot, RelativePath.Root, recursive: true); + } + catch (IOException) + { + // Best effort: the SQLite connection pool may still hold the file. TestTempRoots sweeps stale dirs. + } + } + + // ── Deterministic fixtures ─────────────────────────────────────────────────── + + /// + /// Creates distinct digests under the "00" range prefix used by the range benchmark. + /// + private static byte[] CreateDigest(int seed) + { + var digest = new byte[32]; + BitConverter.TryWriteBytes(digest.AsSpan(1), seed); + return digest; + } + + private static ContentHash CreateContentHash(int seed) => ContentHash.FromDigest(CreateDigest(seed)); + + /// Creates a digest containing hexadecimal letters. + private static byte[] CreateHighNibbleDigest() + { + var digest = new byte[32]; + for (var i = 0; i < digest.Length; i++) + digest[i] = (byte)(0xA0 | (i & 0x0F)); + + return digest; + } + + private static IReadOnlyList BuildFileTreeEntries(int count) + { + var entries = new List(count); + var created = new DateTimeOffset(2026, 1, 1, 0, 0, 0, TimeSpan.Zero); + + for (var i = 0; i < count; i++) + { + entries.Add(new FileEntry + { + Name = PathSegment.Parse($"file-{i:D6}.bin"), + ContentHash = CreateContentHash(i), + Created = created.AddSeconds(i), + Modified = created.AddSeconds(i * 2), + }); + } + + return entries; + } + + private static ShardEntry[] BuildShardEntries(int count) + { + var entries = new ShardEntry[count]; + for (var i = 0; i < count; i++) + { + var contentHash = CreateContentHash(i); + entries[i] = new ShardEntry( + ContentHash: contentHash, + ChunkHash: ChunkHash.Parse(contentHash), // large chunk: chunk hash == content hash + OriginalSize: 4096 + i, + ChunkSize: 2048 + i, + StorageTierHint: BlobTier.Cool); + } + + return entries; + } + + private static FileToUpload CreateUpload(int seed, long size) + { + var filePair = new FilePair { RelativePath = RelativePath.Parse($"file-{seed:D6}.bin") }; + var hashed = new HashedFilePair(filePair, CreateContentHash(seed), DateTimeOffset.UnixEpoch, DateTimeOffset.UnixEpoch); + return new FileToUpload(hashed, size); + } +} diff --git a/src/Arius.Benchmarks/BenchmarkRunOptions.cs b/src/Arius.Benchmarks/BenchmarkRunOptions.cs index 9d22312de..c7cd98c47 100644 --- a/src/Arius.Benchmarks/BenchmarkRunOptions.cs +++ b/src/Arius.Benchmarks/BenchmarkRunOptions.cs @@ -71,6 +71,7 @@ static string FindRepositoryRoot() static void PrintHelp() { Console.WriteLine("Runs the canonical representative workflow benchmark on Azurite."); + Console.WriteLine("Pass 'micro [BenchmarkDotNet options]' instead for the in-process allocation benchmarks."); Console.WriteLine(); Console.WriteLine("Options:"); Console.WriteLine(" --raw-output Folder where per-run raw BenchmarkDotNet output is saved."); diff --git a/src/Arius.Benchmarks/Program.cs b/src/Arius.Benchmarks/Program.cs index dbbfc1a20..0f9b70be5 100644 --- a/src/Arius.Benchmarks/Program.cs +++ b/src/Arius.Benchmarks/Program.cs @@ -5,6 +5,14 @@ using BenchmarkDotNet.Loggers; using BenchmarkDotNet.Running; +// In-process allocation micro-benchmarks run on BenchmarkDotNet's default job and take its own options, +// e.g. `micro --filter '*TarBuilder*'`. They need no Docker and do not append to the tail log. +if (args is ["micro", .. var microArgs]) +{ + BenchmarkSwitcher.FromTypes([typeof(AllocationBenchmarks)]).Run(microArgs); + return; +} + var options = BenchmarkRunOptions.Parse(args); var runStartedAt = DateTimeOffset.UtcNow; var runId = runStartedAt.ToString("yyyyMMddTHHmmss.fffZ"); @@ -12,6 +20,8 @@ Directory.CreateDirectory(rawOutputDirectory); Directory.CreateDirectory(Path.GetDirectoryName(options.TailLogPath)!); +// The archive step runs for tens of seconds and mutates real fixtures per iteration, so it is measured +// with one invocation per iteration and no warmup. var config = ManualConfig .Create(DefaultConfig.Instance) .AddJob(Job.Default diff --git a/src/Arius.Benchmarks/benchmark-tail.md b/src/Arius.Benchmarks/benchmark-tail.md index f9f2f62c4..f7631579a 100644 --- a/src/Arius.Benchmarks/benchmark-tail.md +++ b/src/Arius.Benchmarks/benchmark-tail.md @@ -1,5 +1,5 @@ -| ComputerName | DateTimeUtc | Git Head | RepresentativeScaleDivisor | Iterations | Mean | Error | StdDev | Gen 0 | Gen 1 | Gen 2 | Allocated | Completed Work Items | Lock Contentions | RawOutputPath | -| ------------- | --------------------------------- | ---------------------------------------- | -------------------------- | ---------- | -------- | --------- | -------- | ----------- | ----------- | ---------- | --------- | -------------------- | ---------------- | --------------------------------------------- | +| ComputerName | DateTimeUtc | Git Head | RepresentativeScaleDivisor | Iterations | Mean | Error | StdDev | Gen 0 | Gen 1 | Gen 2 | Allocated | Completed Work Items | Lock Contentions | RawOutputPath | Notes | +| ------------- | --------------------------------- | ---------------------------------------- | -------------------------- | ---------- | -------- | --------- | -------- | ----------- | ----------- | ---------- | --------- | -------------------- | ---------------- | --------------------------------------------- | ----- | | woutbook6 | 2026-04-28T04:05:05.9286260+00:00 | 391911db0ae4c2bd3f7871be4e1bd38295770521 | 8 | 3 | 25.75 s | 3.873 s | 0.212 s | 120000.0000 | 27000.0000 | 15000.0000 | 1.08 GB | 74874.0000 | - | src/Arius.Benchmarks/raw/20260428T040505.928Z | # Baseline woutbook run | | runnervmeorf1 | 2026-04-28T04:26:53.0615861+00:00 | 4eb6d183d1984113f94c69d2e1e5832aecefdc0f | 8 | 3 | 39.12 s | 5.755 s | 0.315 s | 66000.0000 | 17000.0000 | 13000.0000 | 1.07 GB | 72894.0000 | 2.0000 | src/Arius.Benchmarks/raw/20260428T042653.061Z | # Baseline GH runner | | woutbook6 | 2026-04-28T04:45:04.1518940+00:00 | 0eccc408f7eb421a76563d69b13d4de9b8a86c0a | 8 | 3 | 25.68 s | 2.462 s | 0.135 s | 87000.0000 | 23000.0000 | 16000.0000 | 840.47 MB | 72540.0000 | 2.0000 | src/Arius.Benchmarks/raw/20260428T044504.151Z | # Run on Woutbook after refactor to hashes | @@ -12,3 +12,5 @@ | woutbook6 | 2026-05-01T15:23:58.2515390+00:00 | 1eedb04c4fe77ae90f7faa384c2c7b7dedbe0600 | 1 | 3 | 1.453 m | 0.1986 m | 0.0109 m | 434000.0000 | 54000.0000 | 4000.0000 | 3.52 GB | 506427.0000 | 15.0000 | src/Arius.Benchmarks/raw/20260501T152358.251Z | | woutbook6 | 2026-05-01T16:44:07.0942170+00:00 | 1eedb04c4fe77ae90f7faa384c2c7b7dedbe0600 | 1 | 3 | 33.16 s | 1.135 s | 0.062 s | 56000.0000 | 5000.0000 | | 444.88 MB | 39835.0000 | 2.0000 | src/Arius.Benchmarks/raw/20260501T164407.094Z | | woutbook6 | 2026-05-01T19:11:23.3076330+00:00 | 468d2479d64cef776b02c14f98ed9b667ecd2b46 | 8 | 3 | 7.352 s | 2.350 s | 0.1288 s | 14000.0000 | 3000.0000 | | 115.94 MB | 7893.0000 | 1.0000 | src/Arius.Benchmarks/raw/20260501T191123.307Z | +| woutbook6 | 2026-09-08T08:52:24.1656440+00:00 | 10947e7bbb38e1161669cdb9aebe90af17384759 | 1 | 3 | 3.808 s | 3.631 s | 0.1990 s | 33000.0000 | 11000.0000 | 2000.0000 | 407.74 MB | 28623.0000 | 429.0000 | src/Arius.Benchmarks/raw/20260908T085224.165Z | # Mid-branch (10947e7b): includes 0ae0798c and a9bdefa7, both reverted later, so it is not the merged code; superseded by the next row. Controlled base run on the same machine/session was 478.22 MB (ee4e9f3e). Time/GC counts on this host are too noisy to compare (3 iterations, no warmup, Error +/-31 s on the base). | +| woutbook6 | 2026-09-29T11:49:40.9965940+00:00 | b4e5a01171380b121bb12619d193e67dba48446f | 1 | 3 | 3.445 s | 7.194 s | 0.3943 s | 29000.0000 | 9000.0000 | 1000.0000 | 392.1 MB | 28159.0000 | 491.0000 | src/Arius.Benchmarks/raw/20260929T114940.996Z | # Branch head after the review fixes: 392.1 MB vs 478.22 MB at the base (ee4e9f3e), -18%. Allocated only; time/GC counts are noise at 3 iterations. | diff --git a/src/Arius.Benchmarks/raw/20260908T085224.165Z/Arius.Benchmarks.ArchiveStepBenchmarks-20260908-105224.log b/src/Arius.Benchmarks/raw/20260908T085224.165Z/Arius.Benchmarks.ArchiveStepBenchmarks-20260908-105224.log new file mode 100644 index 000000000..e69de29bb diff --git a/src/Arius.Benchmarks/raw/20260908T085224.165Z/benchmark-output.log b/src/Arius.Benchmarks/raw/20260908T085224.165Z/benchmark-output.log new file mode 100644 index 000000000..93d0807cc --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260908T085224.165Z/benchmark-output.log @@ -0,0 +1,159 @@ +// Validating benchmarks: +// ***** BenchmarkRunner: Start ***** +// ***** Found 1 benchmark(s) in total ***** +// ***** Building 1 exe(s) in Parallel: Start ***** +// start dotnet restore --nodeReuse:false /p:UseSharedCompilation=false /p:Deterministic=true /p:Optimize=true /p:ArtifactsPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/" /p:OutDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:OutputPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:PublishDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/publish/" in /Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1 +// command took 1.09 sec and exited with 0 +// start dotnet build -c Release --no-restore --nodeReuse:false /p:UseSharedCompilation=false /p:Deterministic=true /p:Optimize=true /p:ArtifactsPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/" /p:OutDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:OutputPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:PublishDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/publish/" --output "/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" in /Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1 +// command took 11.54 sec and exited with 0 +// ***** Done, took 00:00:12 (12.68 sec) ***** +// Found 1 benchmarks: +// ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) + +// ************************** +// Benchmark: ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) +// *** Execute *** +// Launch: 1 / 1 +// Execute: dotnet Arius.Benchmarks-Job-HILDPN-1.dll --anonymousPipes 126 127 --benchmarkName Arius.Benchmarks.ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite --job "InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0" --benchmarkId 0 in /Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0 +// Failed to set up high priority (Permission denied). In order to run benchmarks with high priority, make sure you have the right permissions. +// BeforeAnythingElse + +// Benchmark Process Environment Information: +// BenchmarkDotNet v0.15.8 +// Runtime=.NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a +// GC=Concurrent Workstation +// HardwareIntrinsics=ArmBase+AdvSimd,AES,CRC32,DP,RDM,SHA1,SHA256 VectorSize=128 +// Job: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) + +[testcontainers.org 00:00:00.08] Connected to Docker: + Host: unix:///var/run/docker.sock + Server Version: 29.3.1 + Kernel Version: 6.12.76-linuxkit + API Version: 1.54 + Operating System: Docker Desktop + Total Memory: 7.65 GB + Labels: + com.docker.desktop.address=unix:///Users/wouter/Library/Containers/com.docker.docker/Data/docker-cli.sock +[testcontainers.org 00:00:00.19] Docker container db7721aecdbe created +[testcontainers.org 00:00:00.21] Start Docker container db7721aecdbe +[testcontainers.org 00:00:00.30] Wait for Docker container db7721aecdbe to complete readiness checks +[testcontainers.org 00:00:01.33] Docker container db7721aecdbe ready +[testcontainers.org 00:00:01.41] Docker container a79377d769d1 created +[testcontainers.org 00:00:01.42] Start Docker container a79377d769d1 +[testcontainers.org 00:00:01.51] Wait for Docker container a79377d769d1 to complete readiness checks +[testcontainers.org 00:00:02.58] Docker container a79377d769d1 ready +OverheadJitting 1: 1 op, 61167.00 ns, 61.1670 us/op +WorkloadJitting 1: 1 op, 3657928917.00 ns, 3.6579 s/op + +OverheadWarmup 1: 1 op, 417.00 ns, 417.0000 ns/op +OverheadWarmup 2: 1 op, 84.00 ns, 84.0000 ns/op +OverheadWarmup 3: 1 op, 0.00 ns, 0.0000 ns/op +OverheadWarmup 4: 1 op, 0.00 ns, 0.0000 ns/op +OverheadWarmup 5: 1 op, 41.00 ns, 41.0000 ns/op +OverheadWarmup 6: 1 op, 0.00 ns, 0.0000 ns/op +OverheadWarmup 7: 1 op, 0.00 ns, 0.0000 ns/op + +OverheadActual 1: 1 op, 167.00 ns, 167.0000 ns/op +OverheadActual 2: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 3: 1 op, 209.00 ns, 209.0000 ns/op +OverheadActual 4: 1 op, 42.00 ns, 42.0000 ns/op +OverheadActual 5: 1 op, 250.00 ns, 250.0000 ns/op +OverheadActual 6: 1 op, 250.00 ns, 250.0000 ns/op +OverheadActual 7: 1 op, 42.00 ns, 42.0000 ns/op +OverheadActual 8: 1 op, 41.00 ns, 41.0000 ns/op +OverheadActual 9: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 10: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 11: 1 op, 42.00 ns, 42.0000 ns/op +OverheadActual 12: 1 op, 42.00 ns, 42.0000 ns/op +OverheadActual 13: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 14: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 15: 1 op, 208.00 ns, 208.0000 ns/op +OverheadActual 16: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 17: 1 op, 84.00 ns, 84.0000 ns/op +OverheadActual 18: 1 op, 41.00 ns, 41.0000 ns/op +OverheadActual 19: 1 op, 416.00 ns, 416.0000 ns/op +OverheadActual 20: 1 op, 83.00 ns, 83.0000 ns/op + + +// BeforeActualRun +WorkloadActual 1: 1 op, 3805428292.00 ns, 3.8054 s/op +WorkloadActual 2: 1 op, 4008702500.00 ns, 4.0087 s/op +WorkloadActual 3: 1 op, 3610687125.00 ns, 3.6107 s/op + +// AfterActualRun +WorkloadResult 1: 1 op, 3805428250.00 ns, 3.8054 s/op +WorkloadResult 2: 1 op, 4008702458.00 ns, 4.0087 s/op +WorkloadResult 3: 1 op, 3610687083.00 ns, 3.6107 s/op +// GC: 33 11 2 427543360 1 +// Threading: 28623 429 1 +// Exceptions: 4 + +[testcontainers.org 00:00:25.04] Delete Docker container a79377d769d1 +// AfterAll +// Benchmark Process 75165 has exited with code 0. + +Mean = 3.808 s, StdErr = 0.115 s (3.02%), N = 3, StdDev = 0.199 s +Min = 3.611 s, Q1 = 3.708 s, Median = 3.805 s, Q3 = 3.907 s, Max = 4.009 s +IQR = 0.199 s, LowerFence = 3.410 s, UpperFence = 4.206 s +ConfidenceInterval = [0.177 s; 7.439 s] (CI 99.9%), Margin = 3.631 s (95.34% of Mean) +Skewness = 0.01, Kurtosis = 0.67, MValue = 2 + +// ** Remained 0 (0,0%) benchmark(s) to run. Estimated finish 2026-09-08 10:53 (0h 0m from now) ** +// ***** BenchmarkRunner: Finish ***** + +// * Export * + src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv + src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md + src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html + +// * Detailed results * +ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) +Runtime = .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a; GC = Concurrent Workstation +Mean = 3.808 s, StdErr = 0.115 s (3.02%), N = 3, StdDev = 0.199 s +Min = 3.611 s, Q1 = 3.708 s, Median = 3.805 s, Q3 = 3.907 s, Max = 4.009 s +IQR = 0.199 s, LowerFence = 3.410 s, UpperFence = 4.206 s +ConfidenceInterval = [0.177 s; 7.439 s] (CI 99.9%), Margin = 3.631 s (95.34% of Mean) +Skewness = 0.01, Kurtosis = 0.67, MValue = 2 +-------------------- Histogram -------------------- +[3.527 s ; 3.889 s) | @@ +[3.889 s ; 4.190 s) | @ +--------------------------------------------------- + +// * Summary * + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M4, 1 CPU, 10 logical and 10 physical cores +.NET SDK 10.0.201 + [Host] : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + Job-HILDPN : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + +InvocationCount=1 IterationCount=3 LaunchCount=1 +UnrollFactor=1 WarmupCount=0 + +| Method | Mean | Error | StdDev | Gen0 | Completed Work Items | Lock Contentions | Gen1 | Gen2 | Allocated | +|--------------------------------------- |--------:|--------:|---------:|-----------:|---------------------:|-----------------:|-----------:|----------:|----------:| +| Archive_Step_V1_Representative_Azurite | 3.808 s | 3.631 s | 0.1990 s | 33000.0000 | 28623.0000 | 429.0000 | 11000.0000 | 2000.0000 | 407.74 MB | + +// * Legends * + Mean : Arithmetic mean of all measurements + Error : Half of 99.9% confidence interval + StdDev : Standard deviation of all measurements + Gen0 : GC Generation 0 collects per 1000 operations + Completed Work Items : The number of work items that have been processed in ThreadPool (per single operation) + Lock Contentions : The number of times there was contention upon trying to take a Monitor's lock (per single operation) + Gen1 : GC Generation 1 collects per 1000 operations + Gen2 : GC Generation 2 collects per 1000 operations + Allocated : Allocated memory per single operation (managed only, inclusive, 1KB = 1024B) + 1 s : 1 Second (1 sec) + +// * Diagnostic Output - MemoryDiagnoser * + +// * Diagnostic Output - ThreadingDiagnoser * + + +// ***** BenchmarkRunner: End ***** +Run time: 00:00:25 (25.83 sec), executed benchmarks: 1 + +Global total time: 00:00:38 (38.78 sec), executed benchmarks: 1 +// * Artifacts cleanup * +Artifacts cleanup is finished diff --git a/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md b/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md new file mode 100644 index 000000000..26e495fc8 --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md @@ -0,0 +1,15 @@ +``` + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M4, 1 CPU, 10 logical and 10 physical cores +.NET SDK 10.0.201 + [Host] : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + Job-HILDPN : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + +InvocationCount=1 IterationCount=3 LaunchCount=1 +UnrollFactor=1 WarmupCount=0 + +``` +| Method | Mean | Error | StdDev | Gen0 | Completed Work Items | Lock Contentions | Gen1 | Gen2 | Allocated | +|--------------------------------------- |--------:|--------:|---------:|-----------:|---------------------:|-----------------:|-----------:|----------:|----------:| +| Archive_Step_V1_Representative_Azurite | 3.808 s | 3.631 s | 0.1990 s | 33000.0000 | 28623.0000 | 429.0000 | 11000.0000 | 2000.0000 | 407.74 MB | diff --git a/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv b/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv new file mode 100644 index 000000000..7127f63e1 --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv @@ -0,0 +1,2 @@ +Method,Job,AnalyzeLaunchVariance,EvaluateOverhead,MaxAbsoluteError,MaxRelativeError,MinInvokeCount,MinIterationTime,OutlierMode,Affinity,EnvironmentVariables,Jit,LargeAddressAware,Platform,PowerPlanMode,Runtime,AllowVeryLargeObjects,Concurrent,CpuGroups,Force,HeapAffinitizeMask,HeapCount,NoAffinitize,RetainVm,Server,Arguments,BuildConfiguration,Clock,EngineFactory,NuGetReferences,Toolchain,IsMutator,InvocationCount,IterationCount,IterationTime,LaunchCount,MaxIterationCount,MaxWarmupIterationCount,MemoryRandomization,MinIterationCount,MinWarmupIterationCount,RunStrategy,UnrollFactor,WarmupCount,Mean,Error,StdDev,Gen0,Completed Work Items,Lock Contentions,Gen1,Gen2,Allocated +Archive_Step_V1_Representative_Azurite,Job-HILDPN,False,Default,Default,Default,Default,Default,Default,0000000000,Empty,RyuJit,Default,Arm64,8c5e7fda-e8bf-4a96-9a85-a6e23a8c635c,.NET 10.0,False,True,False,True,Default,Default,False,False,False,Default,Default,Default,Default,Default,Default,Default,1,3,Default,1,Default,Default,Default,Default,Default,Default,1,0,3.808 s,3.631 s,0.1990 s,33000.0000,28623.0000,429.0000,11000.0000,2000.0000,407.74 MB diff --git a/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html b/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html new file mode 100644 index 000000000..9221acfaf --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260908T085224.165Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html @@ -0,0 +1,32 @@ + + + + +Arius.Benchmarks.ArchiveStepBenchmarks-20260908-105237 + + + + +

+BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0]
+Apple M4, 1 CPU, 10 logical and 10 physical cores
+.NET SDK 10.0.201
+  [Host]     : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a
+  Job-HILDPN : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a
+
+
InvocationCount=1  IterationCount=3  LaunchCount=1  
+UnrollFactor=1  WarmupCount=0  
+
+ + + + + +
Method MeanErrorStdDevGen0Completed Work ItemsLock ContentionsGen1Gen2Allocated
Archive_Step_V1_Representative_Azurite3.808 s3.631 s0.1990 s33000.000028623.0000429.000011000.00002000.0000407.74 MB
+ + diff --git a/src/Arius.Benchmarks/raw/20260929T114940.996Z/Arius.Benchmarks.ArchiveStepBenchmarks-20260929-134941.log b/src/Arius.Benchmarks/raw/20260929T114940.996Z/Arius.Benchmarks.ArchiveStepBenchmarks-20260929-134941.log new file mode 100644 index 000000000..e69de29bb diff --git a/src/Arius.Benchmarks/raw/20260929T114940.996Z/benchmark-output.log b/src/Arius.Benchmarks/raw/20260929T114940.996Z/benchmark-output.log new file mode 100644 index 000000000..e9ed7c079 --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260929T114940.996Z/benchmark-output.log @@ -0,0 +1,159 @@ +// Validating benchmarks: +// ***** BenchmarkRunner: Start ***** +// ***** Found 1 benchmark(s) in total ***** +// ***** Building 1 exe(s) in Parallel: Start ***** +// start dotnet restore --nodeReuse:false /p:UseSharedCompilation=false /p:Deterministic=true /p:Optimize=true /p:ArtifactsPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/" /p:OutDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:OutputPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:PublishDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/publish/" in /Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1 +// command took 0.88 sec and exited with 0 +// start dotnet build -c Release --no-restore --nodeReuse:false /p:UseSharedCompilation=false /p:Deterministic=true /p:Optimize=true /p:ArtifactsPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/" /p:OutDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:OutputPath="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" /p:PublishDir="/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/publish/" --output "/Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0/" in /Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1 +// command took 10.44 sec and exited with 0 +// ***** Done, took 00:00:11 (11.45 sec) ***** +// Found 1 benchmarks: +// ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) + +// ************************** +// Benchmark: ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) +// *** Execute *** +// Launch: 1 / 1 +// Execute: dotnet Arius.Benchmarks-Job-HILDPN-1.dll --anonymousPipes 126 127 --benchmarkName Arius.Benchmarks.ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite --job "InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0" --benchmarkId 0 in /Users/wouter/.superset/worktrees/288a93b0-804e-460f-999d-8c3f427fff33/memory-opt/src/Arius.Benchmarks/bin/Release/net10.0/Arius.Benchmarks-Job-HILDPN-1/bin/Release/net10.0 +// Failed to set up high priority (Permission denied). In order to run benchmarks with high priority, make sure you have the right permissions. +// BeforeAnythingElse + +// Benchmark Process Environment Information: +// BenchmarkDotNet v0.15.8 +// Runtime=.NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a +// GC=Concurrent Workstation +// HardwareIntrinsics=ArmBase+AdvSimd,AES,CRC32,DP,RDM,SHA1,SHA256 VectorSize=128 +// Job: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) + +[testcontainers.org 00:00:00.07] Connected to Docker: + Host: unix:///var/run/docker.sock + Server Version: 29.3.1 + Kernel Version: 6.12.76-linuxkit + API Version: 1.54 + Operating System: Docker Desktop + Total Memory: 7.65 GB + Labels: + com.docker.desktop.address=unix:///Users/wouter/Library/Containers/com.docker.docker/Data/docker-cli.sock +[testcontainers.org 00:00:00.19] Docker container 02e615d7dc78 created +[testcontainers.org 00:00:00.21] Start Docker container 02e615d7dc78 +[testcontainers.org 00:00:00.30] Wait for Docker container 02e615d7dc78 to complete readiness checks +[testcontainers.org 00:00:01.33] Docker container 02e615d7dc78 ready +[testcontainers.org 00:00:01.40] Docker container 4574fa6d9a39 created +[testcontainers.org 00:00:01.41] Start Docker container 4574fa6d9a39 +[testcontainers.org 00:00:01.49] Wait for Docker container 4574fa6d9a39 to complete readiness checks +[testcontainers.org 00:00:02.54] Docker container 4574fa6d9a39 ready +OverheadJitting 1: 1 op, 54708.00 ns, 54.7080 us/op +WorkloadJitting 1: 1 op, 3172194458.00 ns, 3.1722 s/op + +OverheadWarmup 1: 1 op, 333.00 ns, 333.0000 ns/op +OverheadWarmup 2: 1 op, 42.00 ns, 42.0000 ns/op +OverheadWarmup 3: 1 op, 125.00 ns, 125.0000 ns/op +OverheadWarmup 4: 1 op, 83.00 ns, 83.0000 ns/op +OverheadWarmup 5: 1 op, 41.00 ns, 41.0000 ns/op +OverheadWarmup 6: 1 op, 42.00 ns, 42.0000 ns/op +OverheadWarmup 7: 1 op, 42.00 ns, 42.0000 ns/op + +OverheadActual 1: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 2: 1 op, 167.00 ns, 167.0000 ns/op +OverheadActual 3: 1 op, 41.00 ns, 41.0000 ns/op +OverheadActual 4: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 5: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 6: 1 op, 41.00 ns, 41.0000 ns/op +OverheadActual 7: 1 op, 42.00 ns, 42.0000 ns/op +OverheadActual 8: 1 op, 167.00 ns, 167.0000 ns/op +OverheadActual 9: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 10: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 11: 1 op, 84.00 ns, 84.0000 ns/op +OverheadActual 12: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 13: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 14: 1 op, 125.00 ns, 125.0000 ns/op +OverheadActual 15: 1 op, 84.00 ns, 84.0000 ns/op +OverheadActual 16: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 17: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 18: 1 op, 0.00 ns, 0.0000 ns/op +OverheadActual 19: 1 op, 83.00 ns, 83.0000 ns/op +OverheadActual 20: 1 op, 84.00 ns, 84.0000 ns/op + + +// BeforeActualRun +WorkloadActual 1: 1 op, 3713232958.00 ns, 3.7132 s/op +WorkloadActual 2: 1 op, 2992224084.00 ns, 2.9922 s/op +WorkloadActual 3: 1 op, 3629489209.00 ns, 3.6295 s/op + +// AfterActualRun +WorkloadResult 1: 1 op, 3713232875.00 ns, 3.7132 s/op +WorkloadResult 2: 1 op, 2992224001.00 ns, 2.9922 s/op +WorkloadResult 3: 1 op, 3629489126.00 ns, 3.6295 s/op +// GC: 29 9 1 411149920 1 +// Threading: 28159 491 1 +// Exceptions: 4 + +[testcontainers.org 00:00:22.89] Delete Docker container 4574fa6d9a39 +// AfterAll +// Benchmark Process 49638 has exited with code 0. + +Mean = 3.445 s, StdErr = 0.228 s (6.61%), N = 3, StdDev = 0.394 s +Min = 2.992 s, Q1 = 3.311 s, Median = 3.629 s, Q3 = 3.671 s, Max = 3.713 s +IQR = 0.361 s, LowerFence = 2.770 s, UpperFence = 4.212 s +ConfidenceInterval = [-3.749 s; 10.639 s] (CI 99.9%), Margin = 7.194 s (208.83% of Mean) +Skewness = -0.37, Kurtosis = 0.67, MValue = 2 + +// ** Remained 0 (0,0%) benchmark(s) to run. Estimated finish 2026-09-29 13:50 (0h 0m from now) ** +// ***** BenchmarkRunner: Finish ***** + +// * Export * + raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv + raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md + raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html + +// * Detailed results * +ArchiveStepBenchmarks.Archive_Step_V1_Representative_Azurite: Job-HILDPN(InvocationCount=1, IterationCount=3, LaunchCount=1, UnrollFactor=1, WarmupCount=0) +Runtime = .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a; GC = Concurrent Workstation +Mean = 3.445 s, StdErr = 0.228 s (6.61%), N = 3, StdDev = 0.394 s +Min = 2.992 s, Q1 = 3.311 s, Median = 3.629 s, Q3 = 3.671 s, Max = 3.713 s +IQR = 0.361 s, LowerFence = 2.770 s, UpperFence = 4.212 s +ConfidenceInterval = [-3.749 s; 10.639 s] (CI 99.9%), Margin = 7.194 s (208.83% of Mean) +Skewness = -0.37, Kurtosis = 0.67, MValue = 2 +-------------------- Histogram -------------------- +[2.633 s ; 3.313 s) | @ +[3.313 s ; 4.072 s) | @@ +--------------------------------------------------- + +// * Summary * + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M4, 1 CPU, 10 logical and 10 physical cores +.NET SDK 10.0.201 + [Host] : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + Job-HILDPN : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + +InvocationCount=1 IterationCount=3 LaunchCount=1 +UnrollFactor=1 WarmupCount=0 + +| Method | Mean | Error | StdDev | Gen0 | Completed Work Items | Lock Contentions | Gen1 | Gen2 | Allocated | +|--------------------------------------- |--------:|--------:|---------:|-----------:|---------------------:|-----------------:|----------:|----------:|----------:| +| Archive_Step_V1_Representative_Azurite | 3.445 s | 7.194 s | 0.3943 s | 29000.0000 | 28159.0000 | 491.0000 | 9000.0000 | 1000.0000 | 392.1 MB | + +// * Legends * + Mean : Arithmetic mean of all measurements + Error : Half of 99.9% confidence interval + StdDev : Standard deviation of all measurements + Gen0 : GC Generation 0 collects per 1000 operations + Completed Work Items : The number of work items that have been processed in ThreadPool (per single operation) + Lock Contentions : The number of times there was contention upon trying to take a Monitor's lock (per single operation) + Gen1 : GC Generation 1 collects per 1000 operations + Gen2 : GC Generation 2 collects per 1000 operations + Allocated : Allocated memory per single operation (managed only, inclusive, 1KB = 1024B) + 1 s : 1 Second (1 sec) + +// * Diagnostic Output - MemoryDiagnoser * + +// * Diagnostic Output - ThreadingDiagnoser * + + +// ***** BenchmarkRunner: End ***** +Run time: 00:00:23 (23.57 sec), executed benchmarks: 1 + +Global total time: 00:00:35 (35.21 sec), executed benchmarks: 1 +// * Artifacts cleanup * +Artifacts cleanup is finished diff --git a/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md b/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md new file mode 100644 index 000000000..d35a75319 --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report-github.md @@ -0,0 +1,15 @@ +``` + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M4, 1 CPU, 10 logical and 10 physical cores +.NET SDK 10.0.201 + [Host] : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + Job-HILDPN : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a + +InvocationCount=1 IterationCount=3 LaunchCount=1 +UnrollFactor=1 WarmupCount=0 + +``` +| Method | Mean | Error | StdDev | Gen0 | Completed Work Items | Lock Contentions | Gen1 | Gen2 | Allocated | +|--------------------------------------- |--------:|--------:|---------:|-----------:|---------------------:|-----------------:|----------:|----------:|----------:| +| Archive_Step_V1_Representative_Azurite | 3.445 s | 7.194 s | 0.3943 s | 29000.0000 | 28159.0000 | 491.0000 | 9000.0000 | 1000.0000 | 392.1 MB | diff --git a/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv b/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv new file mode 100644 index 000000000..6f408da36 --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.csv @@ -0,0 +1,2 @@ +Method,Job,AnalyzeLaunchVariance,EvaluateOverhead,MaxAbsoluteError,MaxRelativeError,MinInvokeCount,MinIterationTime,OutlierMode,Affinity,EnvironmentVariables,Jit,LargeAddressAware,Platform,PowerPlanMode,Runtime,AllowVeryLargeObjects,Concurrent,CpuGroups,Force,HeapAffinitizeMask,HeapCount,NoAffinitize,RetainVm,Server,Arguments,BuildConfiguration,Clock,EngineFactory,NuGetReferences,Toolchain,IsMutator,InvocationCount,IterationCount,IterationTime,LaunchCount,MaxIterationCount,MaxWarmupIterationCount,MemoryRandomization,MinIterationCount,MinWarmupIterationCount,RunStrategy,UnrollFactor,WarmupCount,Mean,Error,StdDev,Gen0,Completed Work Items,Lock Contentions,Gen1,Gen2,Allocated +Archive_Step_V1_Representative_Azurite,Job-HILDPN,False,Default,Default,Default,Default,Default,Default,0000000000,Empty,RyuJit,Default,Arm64,8c5e7fda-e8bf-4a96-9a85-a6e23a8c635c,.NET 10.0,False,True,False,True,Default,Default,False,False,False,Default,Default,Default,Default,Default,Default,Default,1,3,Default,1,Default,Default,Default,Default,Default,Default,1,0,3.445 s,7.194 s,0.3943 s,29000.0000,28159.0000,491.0000,9000.0000,1000.0000,392.1 MB diff --git a/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html b/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html new file mode 100644 index 000000000..a925c1e2d --- /dev/null +++ b/src/Arius.Benchmarks/raw/20260929T114940.996Z/results/Arius.Benchmarks.ArchiveStepBenchmarks-report.html @@ -0,0 +1,32 @@ + + + + +Arius.Benchmarks.ArchiveStepBenchmarks-20260929-134952 + + + + +

+BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0]
+Apple M4, 1 CPU, 10 logical and 10 physical cores
+.NET SDK 10.0.201
+  [Host]     : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a
+  Job-HILDPN : .NET 10.0.5 (10.0.5, 10.0.526.15411), Arm64 RyuJIT armv8.0-a
+
+
InvocationCount=1  IterationCount=3  LaunchCount=1  
+UnrollFactor=1  WarmupCount=0  
+
+ + + + + +
Method MeanErrorStdDevGen0Completed Work ItemsLock ContentionsGen1Gen2Allocated
Archive_Step_V1_Representative_Azurite3.445 s7.194 s0.3943 s29000.000028159.0000491.00009000.00001000.0000392.1 MB
+ + diff --git a/src/Arius.Cli.Tests/Commands/Archive/NotificationHandlerTests.cs b/src/Arius.Cli.Tests/Commands/Archive/NotificationHandlerTests.cs index 47a2caf5d..59f149ddd 100644 --- a/src/Arius.Cli.Tests/Commands/Archive/NotificationHandlerTests.cs +++ b/src/Arius.Cli.Tests/Commands/Archive/NotificationHandlerTests.cs @@ -145,6 +145,45 @@ public async Task FileSkippedHandler_DuringUpload_RemovesTrackedFileWithoutIncre state.TrackedFiles.ContainsKey(path).ShouldBeFalse(); } + [Test] + public async Task FileDedupedHandler_DuplicateOfFileBeingUploaded_KeepsTheUploadingRow() + { + var state = new ProgressState(); + var hashingH = new FileHashingHandler(state); + var hashedH = new FileHashedHandler(state); + var dedupedH = new FileDedupedHandler(state); + var uploadingH = new ChunkUploadingHandler(state); + var original = RelativePath.Parse("original.bin"); + var duplicate = RelativePath.Parse("copy/original.bin"); + + await hashingH.Handle(new FileHashingEvent(original, 5000), CancellationToken.None); + await hashingH.Handle(new FileHashingEvent(duplicate, 5000), CancellationToken.None); + await hashedH.Handle(new FileHashedEvent(original, FakeContentHash('7'), false, false), CancellationToken.None); + await hashedH.Handle(new FileHashedEvent(duplicate, FakeContentHash('7'), false, false), CancellationToken.None); + await dedupedH.Handle(new FileDedupedEvent(duplicate, FakeContentHash('7'), 5000), CancellationToken.None); + await uploadingH.Handle(new ChunkUploadingEvent(FakeChunkHash('7'), 5000), CancellationToken.None); + + state.TrackedFiles.ContainsKey(duplicate).ShouldBeFalse(); + state.TrackedFiles[original].State.ShouldBe(FileState.Uploading); + state.FilesUnique.ShouldBe(1L); + } + + [Test] + public async Task FileDedupedHandler_LastTrackedCopy_ForgetsTheContentHash() + { + var state = new ProgressState(); + var hashingH = new FileHashingHandler(state); + var hashedH = new FileHashedHandler(state); + var dedupedH = new FileDedupedHandler(state); + var path = RelativePath.Parse("known.bin"); + + await hashingH.Handle(new FileHashingEvent(path, 5000), CancellationToken.None); + await hashedH.Handle(new FileHashedEvent(path, FakeContentHash('6'), false, false), CancellationToken.None); + await dedupedH.Handle(new FileDedupedEvent(path, FakeContentHash('6'), 5000), CancellationToken.None); + + state.ContentHashToPath.ContainsKey(FakeContentHash('6')).ShouldBeFalse(); + } + [Test] public async Task TarBundleStartedHandler_CreatesTrackedTar() { @@ -190,6 +229,7 @@ public async Task TarEntryAddedHandler_RemovesTrackedFileAndUpdatesTrackedTar() await tarEntryH.Handle(new TarEntryAddedEvent(FakeContentHash('b'), 1, 500), CancellationToken.None); state.TrackedFiles.ContainsKey(path).ShouldBeFalse(); + state.ContentHashToPath.ContainsKey(FakeContentHash('b')).ShouldBeFalse(); state.TrackedTars[1].FileCount.ShouldBe(1); state.TrackedTars[1].AccumulatedBytes.ShouldBe(500L); } @@ -281,6 +321,7 @@ public async Task ChunkUploadedHandler_RemovesFileAndIncrementsChunksUploaded() await uploadedH.Handle(new ChunkUploadedEvent(FakeChunkHash('9'), 4000, 5000), CancellationToken.None); state.TrackedFiles.ContainsKey(path).ShouldBeFalse(); + state.ContentHashToPath.ContainsKey(FakeContentHash('9')).ShouldBeFalse(); state.ChunksUploaded.ShouldBe(1L); state.BytesUploaded.ShouldBe(4000L); } diff --git a/src/Arius.Cli/Commands/Archive/ArchiveProgressHandlers.cs b/src/Arius.Cli/Commands/Archive/ArchiveProgressHandlers.cs index 1b5b16325..3a315cf64 100644 --- a/src/Arius.Cli/Commands/Archive/ArchiveProgressHandlers.cs +++ b/src/Arius.Cli/Commands/Archive/ArchiveProgressHandlers.cs @@ -91,6 +91,7 @@ public ValueTask Handle(TarBundleStartedEvent notification, CancellationToken ca var bundleNumber = state.NextBundleNumber(); var tar = new TrackedTar(bundleNumber, state.TarTargetSize); state.TrackedTars.TryAdd(bundleNumber, tar); + state.AccumulatingTar = tar; return ValueTask.CompletedTask; } } @@ -106,16 +107,9 @@ public sealed class TarEntryAddedHandler(ProgressState state) : INotificationHan { public ValueTask Handle(TarEntryAddedEvent notification, CancellationToken cancellationToken) { - if (state.ContentHashToPath.TryGetValue(notification.ContentHash, out var paths)) - foreach (var path in paths) - state.RemoveFile(path); + state.RemoveContentHash(notification.ContentHash); - var tar = state.TrackedTars.Values - .Where(t => t.State == TarState.Accumulating) - .OrderByDescending(t => t.BundleNumber) - .FirstOrDefault(); - - if (tar != null) + if (state.AccumulatingTar is { } tar) { var addedBytes = notification.CurrentTarSize - tar.AccumulatedBytes; tar.AddEntry(addedBytes > 0 ? addedBytes : 0); @@ -134,12 +128,7 @@ public sealed class TarBundleSealingHandler(ProgressState state) : INotification { public ValueTask Handle(TarBundleSealingEvent notification, CancellationToken cancellationToken) { - var tar = state.TrackedTars.Values - .Where(t => t.State == TarState.Accumulating || t.State == TarState.Sealing) - .OrderByDescending(t => t.BundleNumber) - .FirstOrDefault(); - - if (tar != null) + if (state.AccumulatingTar is { } tar) { tar.TarHash = notification.TarHash; // Use the sealed tar's archive byte size (headers + padding included), not the sum of file @@ -149,6 +138,22 @@ public ValueTask Handle(TarBundleSealingEvent notification, CancellationToken ca tar.State = TarState.Sealing; } + // The bundle is no longer accumulating; the next TarBundleStartedEvent supplies the next one. + state.AccumulatingTar = null; + + return ValueTask.CompletedTask; + } +} + +/// +/// Removes tracked files when their content is deduplicated. +/// +public sealed class FileDedupedHandler(ProgressState state) : INotificationHandler +{ + public ValueTask Handle(FileDedupedEvent notification, CancellationToken cancellationToken) + { + state.RemoveDeduplicatedFile(notification.RelativePath, notification.ContentHash); + return ValueTask.CompletedTask; } } @@ -192,9 +197,7 @@ public sealed class ChunkUploadedHandler(ProgressState state) : INotificationHan { public ValueTask Handle(ChunkUploadedEvent notification, CancellationToken cancellationToken) { - if (state.ContentHashToPath.TryGetValue(ContentHash.Parse(notification.ChunkHash), out var paths)) - foreach (var path in paths) - state.RemoveFile(path); + state.RemoveContentHash(ContentHash.Parse(notification.ChunkHash)); state.IncrementChunksUploaded(notification.StoredSize); return ValueTask.CompletedTask; } diff --git a/src/Arius.Cli/ProgressState.cs b/src/Arius.Cli/ProgressState.cs index c4973275d..f2e053b84 100644 --- a/src/Arius.Cli/ProgressState.cs +++ b/src/Arius.Cli/ProgressState.cs @@ -253,7 +253,8 @@ public sealed class ProgressState /// /// Reverse lookup: ContentHash → one or more RelativePaths. /// One-to-many because two files can legitimately hash to the same content in one run. - /// Populated when FileHashedEvent fires; used by downstream content-hash-keyed events. + /// Populated when FileHashedEvent fires; used by downstream content-hash-keyed events, and + /// removed once the content's rows are done. /// public ConcurrentDictionary> ContentHashToPath { get; } = new(); @@ -307,6 +308,25 @@ public bool SetFileUploading(ContentHash contentHash) public void RemoveFile(RelativePath relativePath) => TrackedFiles.TryRemove(relativePath, out _); + /// Removes the reverse lookup for and the rows of all its paths. + public void RemoveContentHash(ContentHash contentHash) + { + if (ContentHashToPath.TryRemove(contentHash, out var paths)) + foreach (var path in paths) + RemoveFile(path); + } + + /// + /// Removes the row of one deduplicated path, and the reverse lookup once none of its paths is tracked. + /// Other paths keep theirs: the first copy of the content may still be uploading. + /// + public void RemoveDeduplicatedFile(RelativePath relativePath, ContentHash contentHash) + { + RemoveFile(relativePath); + if (ContentHashToPath.TryGetValue(contentHash, out var paths) && !paths.Any(TrackedFiles.ContainsKey)) + ContentHashToPath.TryRemove(contentHash, out _); + } + /// /// Removes the tracked row for a skipped file. Only files still in /// count toward hashing completion; later-stage skips still clear the row but must not double-count. @@ -322,6 +342,16 @@ public void SkipFileDuringHashing(RelativePath relativePath) /// TAR bundles currently tracked, keyed by bundle number. public ConcurrentDictionary TrackedTars { get; } = new(); + /// + /// The bundle currently receiving entries, or null between bundles. + /// + public TrackedTar? AccumulatingTar + { + get => Volatile.Read(ref _accumulatingTar); + set => Volatile.Write(ref _accumulatingTar, value); + } + private TrackedTar? _accumulatingTar; + /// Monotonically increasing bundle counter; call to allocate a new ID. private long _bundleCounter; diff --git a/src/Arius.Core.Tests/Features/ArchiveCommand/TarBuilderTests.cs b/src/Arius.Core.Tests/Features/ArchiveCommand/TarBuilderTests.cs index 80ed525b3..6df8aaccd 100644 --- a/src/Arius.Core.Tests/Features/ArchiveCommand/TarBuilderTests.cs +++ b/src/Arius.Core.Tests/Features/ArchiveCommand/TarBuilderTests.cs @@ -103,6 +103,34 @@ public async Task SealAsync_TarHash_MatchesHashOfBundleContent() bundle.TarHash.ShouldBe(ChunkHashOf(bundle.Content, IEncryptionService.PlaintextInstance)); } + // ── Buffer sizing ───────────────────────────────────────────────────────── + + [Test] + public async Task AddAsync_FullBundleOfTinyFiles_BufferFitsTheirTarFraming() + { + // 4 KB files carry ~37% tar framing (headers + padding), far above a fixed headroom over the target. + await using var builder = new TarBuilder(targetSize: 8 * 1024 * 1024, IEncryptionService.PlaintextInstance); + + SealedTar? bundle = null; + for (var i = 0; bundle is null; i++) + bundle = await builder.AddAsync(Upload($"f{i}", UniqueHash(i), 4 * 1024), new MemoryStream(new byte[4 * 1024]), CancellationToken.None); + + bundle.Content.Array!.Length.ShouldBeLessThan((int)(bundle.Content.Count * 1.25)); + } + + [Test] + public async Task SealAsync_PartialBundle_DoesNotHoldTheFullTargetBuffer() + { + await using var builder = new TarBuilder(targetSize: 64 * 1024 * 1024, IEncryptionService.PlaintextInstance); + + for (var i = 0; i < 768; i++) // 3 MB of 4 KB files: past any presize threshold, far short of the target + await builder.AddAsync(Upload($"f{i}", UniqueHash(i), 4 * 1024), new MemoryStream(new byte[4 * 1024]), CancellationToken.None); + + var bundle = await builder.SealAsync(CancellationToken.None); + + bundle!.Content.Array!.Length.ShouldBeLessThanOrEqualTo(bundle.Content.Count * 2); + } + // ── Lifecycle callbacks ─────────────────────────────────────────────────── [Test] @@ -176,6 +204,13 @@ private static (byte[] Content, ContentHash Hash) Content(byte fill, int count) return (bytes, IEncryptionService.PlaintextInstance.ComputeHash(bytes)); } + private static ContentHash UniqueHash(int seed) + { + var digest = new byte[32]; + BitConverter.TryWriteBytes(digest, seed); + return ContentHash.FromDigest(digest); + } + /// Reads a sealed tar bundle back into a map of entry-name → entry-bytes. private static async Task> ReadBundleAsync(SealedTar bundle) { diff --git a/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexLocalStoreTests.cs b/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexLocalStoreTests.cs index 428d3899f..4249943bb 100644 --- a/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexLocalStoreTests.cs +++ b/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexLocalStoreTests.cs @@ -23,6 +23,33 @@ public void Initialize_CreatesSchemaVersionAndWalMode() journalMode.ExecuteScalar().ShouldBe("wal"); } + [Test] + public void FindEntries_WithMoreHashesThanSqliteVariableLimit_ReturnsMatches() + { + // Wider than SQLITE_MAX_VARIABLE_NUMBER (32766), as a large directory listing can be. + var repositoryKey = $"acct-local-store-wide-{Guid.NewGuid():N}"; + var root = RepositoryLocalStatePaths.GetChunkIndexCacheRoot(repositoryKey, repositoryKey); + var store = new ChunkIndexLocalStore(root); + + var stored = new ShardEntry(FakeContentHash('a'), FakeChunkHash('b'), 10, 5, BlobTier.Cool); + store.UpsertPendingFlush(stored); + + var hashes = new List(40_001) { stored.ContentHash }; + hashes.AddRange(Enumerable.Range(0, 40_000).Select(WideLookupHash)); + + var found = store.FindEntries(hashes); + + found.Count.ShouldBe(1); + found[stored.ContentHash].ShouldBe(stored); + } + + private static ContentHash WideLookupHash(int seed) + { + var digest = new byte[32]; + BitConverter.TryWriteBytes(digest.AsSpan(1), seed); + return ContentHash.FromDigest(digest); + } + [Test] public void UpsertPendingFlush_AndLookup_RoundTripsEntry() { diff --git a/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexServiceLookupTests.cs b/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexServiceLookupTests.cs index cea970055..57a95f2d4 100644 --- a/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexServiceLookupTests.cs +++ b/src/Arius.Core.Tests/Shared/ChunkIndex/ChunkIndexServiceLookupTests.cs @@ -299,6 +299,30 @@ public async Task LookupAsync_BatchedHashes_LoadEachTouchedPrefixOnce() blobs.RequestedBlobNames.Count(name => name == BlobPaths.ChunkIndexShardPath(ChunkIndexRouter.GetRootPrefix(otherHash))).ShouldBe(1); } + [Test] + public async Task LookupAsync_SameRootBatchLargerThan256_ResolvesAllEntries() + { + var blobs = new FakeInMemoryBlobContainerService(); + var hashes = Enumerable.Range(0, 300) + .Select(index => ContentHash.Parse($"aa{index:X62}")) + .ToArray(); + var entries = hashes + .Select((hash, index) => new ShardEntry(hash, FakeChunkHash((char)('0' + index % 10)), index + 1, index + 1, BlobTier.Cool)) + .ToArray(); + + blobs.SeedBlob( + BlobPaths.ChunkIndexShardPath(PathSegment.Parse("aa")), + await ShardSerializer.SerializeAsync(CreateShard(entries), IEncryptionService.PlaintextInstance, ICompressionService.ZtdInstance), + BlobTier.Cool); + using var index = CreateIndex(blobs, "large-same-root-batch"); + + var result = await index.LookupAsync(hashes); + + result.Count.ShouldBe(300); + foreach (var entry in entries) + result[entry.ContentHash].ShouldBe(entry); + } + [Test] public async Task LookupAsync_CorruptCleanSqlite_FailsWithLocalStoreRecoveryGuidance() { diff --git a/src/Arius.Core.Tests/Shared/ChunkStorage/ChunkStorageServiceUploadTests.cs b/src/Arius.Core.Tests/Shared/ChunkStorage/ChunkStorageServiceUploadTests.cs index d4d4f29ff..6a7e48dbd 100644 --- a/src/Arius.Core.Tests/Shared/ChunkStorage/ChunkStorageServiceUploadTests.cs +++ b/src/Arius.Core.Tests/Shared/ChunkStorage/ChunkStorageServiceUploadTests.cs @@ -223,7 +223,12 @@ public async Task UploadLargeAsync_RetryAfterMetadataConflict_ReportsSingleProgr result.AlreadyExisted.ShouldBeFalse(); blobs.DeletedBlobNames.ShouldContain(blobName); - reports.ShouldBe([512L, 1024L, 1536L, 2048L]); + + // Verify that retry progress is monotonic and ends at the true total; intermediate reports are throttled. + reports.ShouldNotBeEmpty(); + reports.ShouldBeInOrder(); + reports.Distinct().Count().ShouldBe(reports.Count); + reports[^1].ShouldBe((long)content.Length); } [Test] diff --git a/src/Arius.Core.Tests/Shared/FileSystem/RelativePathTests.cs b/src/Arius.Core.Tests/Shared/FileSystem/RelativePathTests.cs index ab9a75621..898af7818 100644 --- a/src/Arius.Core.Tests/Shared/FileSystem/RelativePathTests.cs +++ b/src/Arius.Core.Tests/Shared/FileSystem/RelativePathTests.cs @@ -28,6 +28,14 @@ public void Parse_DotSegment_Throws() Should.Throw(() => RelativePath.Parse("photos/../pic.jpg")); } + [Test] + [Arguments("photos/pi\u0001c.jpg")] + [Arguments("pho\ntos/pic.jpg")] + public void Parse_ControlCharacter_Throws(string value) + { + Should.Throw(() => RelativePath.Parse(value)); + } + [Test] public void FromPlatformRelativePath_NormalizesDirectorySeparators() { diff --git a/src/Arius.Core.Tests/Shared/HashCache/SparseFingerprintTests.cs b/src/Arius.Core.Tests/Shared/HashCache/SparseFingerprintTests.cs index 059a04d4f..d6c5841c9 100644 --- a/src/Arius.Core.Tests/Shared/HashCache/SparseFingerprintTests.cs +++ b/src/Arius.Core.Tests/Shared/HashCache/SparseFingerprintTests.cs @@ -61,7 +61,7 @@ public void ComputeBySeeking_WholeFileRegion_LargerThanBlockSize_DoesNotThrow() fp.Length.ShouldBe(32); // And it must still agree with the streaming Sampler over the same content. - var sampler = new SparseFingerprint.Sampler(size); + using var sampler = new SparseFingerprint.Sampler(size); var pos = 0; const int chunk = 64 * 1024; while (pos < data.Length) @@ -83,7 +83,7 @@ public void Sampler_MatchesSeekingFingerprint_ForSameContent() var seekFp = SparseFingerprint.ComputeBySeeking(fs, path, size); // Drive the sampler the way a sequential read would. - var sampler = new SparseFingerprint.Sampler(size); + using var sampler = new SparseFingerprint.Sampler(size); var pos = 0; const int chunk = 64 * 1024; while (pos < data.Length) @@ -103,4 +103,15 @@ private static (RelativeFileSystem fs, RelativePath path) WriteTempFile(byte[] d fs.WriteAllBytes(path, data); return (fs, path); } + + [Test] + public void Sampler_FinishAfterDispose_Throws() + { + // Dispose returns the capture buffer to ArrayPool; fingerprinting it afterwards would hash bytes + // another consumer now owns. + var sampler = new SparseFingerprint.Sampler(1024); + sampler.Dispose(); + + Should.Throw(() => sampler.Finish()); + } } diff --git a/src/Arius.Core.Tests/Shared/Streaming/FrozenTimeProvider.cs b/src/Arius.Core.Tests/Shared/Streaming/FrozenTimeProvider.cs new file mode 100644 index 000000000..bdf8a775e --- /dev/null +++ b/src/Arius.Core.Tests/Shared/Streaming/FrozenTimeProvider.cs @@ -0,0 +1,7 @@ +namespace Arius.Core.Tests.Shared.Streaming; + +/// A clock that never advances, so every throttle window stays open. +internal sealed class FrozenTimeProvider : TimeProvider +{ + public override long GetTimestamp() => 0; +} diff --git a/src/Arius.Core.Tests/Shared/Streaming/ProgressStreamTests.cs b/src/Arius.Core.Tests/Shared/Streaming/ProgressStreamTests.cs index 9e2094e14..6b6e127b4 100644 --- a/src/Arius.Core.Tests/Shared/Streaming/ProgressStreamTests.cs +++ b/src/Arius.Core.Tests/Shared/Streaming/ProgressStreamTests.cs @@ -5,21 +5,19 @@ namespace Arius.Core.Tests.Shared.Streaming; public class ProgressStreamTests { [Test] - public void Read_ReportsProgressAfterEachChunk() + public void Read_CoalescesReports_ButStillReportsTheTotal() { var data = new byte[1024]; Random.Shared.NextBytes(data); using var src = new MemoryStream(data); var reports = new List(); var progress = new SyncProgress(v => reports.Add(v)); - using var ps = new ProgressStream(src, progress); + using var ps = new ProgressStream(src, progress, new FrozenTimeProvider()); var buf = new byte[256]; while (ps.Read(buf, 0, buf.Length) > 0) { } - reports.Count.ShouldBe(4); - reports[^1].ShouldBe(1024); - reports.ShouldBeInOrder(); + reports.ShouldBe([256, 1024]); } [Test] @@ -93,6 +91,22 @@ public void Read_ZeroLengthSource_NoProgressReported() reportCount.ShouldBe(0); } + [Test] + public void Read_ZeroLengthBuffer_DoesNotReportFinalProgress() + { + using var src = new MemoryStream(new byte[100]); + var reports = new List(); + var progress = new SyncProgress(value => reports.Add(value)); + using var ps = new ProgressStream(src, progress); + + var buffer = new byte[100]; + ps.Read(buffer, 0, buffer.Length).ShouldBe(100); + reports.ShouldBe([100]); + + ps.Read(buffer, 0, 0).ShouldBe(0); + reports.ShouldBe([100]); + } + [Test] public void ReadSpan_ReportsProgress() { diff --git a/src/Arius.Core/Features/ArchiveCommand/ArchiveCommandHandler.cs b/src/Arius.Core/Features/ArchiveCommand/ArchiveCommandHandler.cs index f9c36df9c..075ec5e91 100644 --- a/src/Arius.Core/Features/ArchiveCommand/ArchiveCommandHandler.cs +++ b/src/Arius.Core/Features/ArchiveCommand/ArchiveCommandHandler.cs @@ -481,7 +481,7 @@ async ValueTask FullHashAndRecordAsync(RelativePath relativePath, l Interlocked.Increment(ref filesDeduped); var size = isKnown ? known[hashed.ContentHash].OriginalSize : inFlightHashes[hashed.ContentHash]; Interlocked.Add(ref originalSize, size); - await _mediator.Publish(new FileDedupedEvent(hashed.ContentHash, size), cancellationToken); + await _mediator.Publish(new FileDedupedEvent(hashed.FilePair.RelativePath, hashed.ContentHash, size), cancellationToken); continue; } @@ -493,7 +493,7 @@ async ValueTask FullHashAndRecordAsync(RelativePath relativePath, l Interlocked.Increment(ref filesDeduped); var size = fs.GetFileSize(hashed.FilePair.RelativePath); Interlocked.Add(ref originalSize, size); - await _mediator.Publish(new FileDedupedEvent(hashed.ContentHash, size), cancellationToken); + await _mediator.Publish(new FileDedupedEvent(hashed.FilePair.RelativePath, hashed.ContentHash, size), cancellationToken); } else { diff --git a/src/Arius.Core/Features/ArchiveCommand/Events.cs b/src/Arius.Core/Features/ArchiveCommand/Events.cs index 505589023..8603f973a 100644 --- a/src/Arius.Core/Features/ArchiveCommand/Events.cs +++ b/src/Arius.Core/Features/ArchiveCommand/Events.cs @@ -88,9 +88,11 @@ public sealed record ChunkUploadedEvent(ChunkHash ChunkHash, long StoredSize, lo /// A file's contents are already present in the remote. /// Contrast , which fires for content that is uploaded. /// +/// The deduplicated file. Consumers tracking per-file state need this: a content +/// hash can map to several paths in one run, and the other copies may still be uploading. /// Content hash of the deduplicated file (already present in the repository). /// Uncompressed size in bytes of the file whose content was not re-uploaded. -public sealed record FileDedupedEvent(ContentHash ContentHash, long OriginalSize) : INotification; +public sealed record FileDedupedEvent(RelativePath RelativePath, ContentHash ContentHash, long OriginalSize) : INotification; /// A tar bundle is being sealed. /// Number of entries in the sealed tar. diff --git a/src/Arius.Core/Features/ArchiveCommand/TarBuilder.cs b/src/Arius.Core/Features/ArchiveCommand/TarBuilder.cs index dea243497..30a7ddfe6 100644 --- a/src/Arius.Core/Features/ArchiveCommand/TarBuilder.cs +++ b/src/Arius.Core/Features/ArchiveCommand/TarBuilder.cs @@ -31,6 +31,15 @@ internal sealed class TarBuilder : IAsyncDisposable private MemoryStream? _tarStream; private long _currentSize; + /// + /// Tar length past which the buffer grows in one step to the bundle's projected full size. Below it, + /// MemoryStream keeps doubling, so a bundle of a few small files never reserves the full target. + /// + private const int PresizeThreshold = 1024 * 1024; + + /// Headroom over the projected size, for the entry that overshoots the target. + private const int ProjectionHeadroomDivisor = 8; + /// A bundle is sealed once its accumulated size reaches this threshold. /// Used to hash the sealed tar body. /// Invoked when a new bundle is opened (its first entry). @@ -71,6 +80,7 @@ public TarBuilder( } // Write the entry named by content-hash (not original path). + var lengthBefore = _tarStream!.Length; var tarEntry = new PaxTarEntry(TarEntryType.RegularFile, upload.HashedPair.ContentHash.ToString()); await using (source) { @@ -82,6 +92,15 @@ public TarBuilder( _entries.Add(new TarEntry(upload.HashedPair.ContentHash, upload.FileSize, upload.HashedPair)); _currentSize += upload.FileSize; + // Rather than doubling all the way to the target, grow once, sized by the tar bytes written per file + // byte so far: per-entry framing makes a bundle of tiny files far longer than its summed file sizes. + if (lengthBefore < PresizeThreshold && _tarStream.Length >= PresizeThreshold) + { + var projected = _tarStream.Length * (_targetSize + _targetSize / ProjectionHeadroomDivisor) / _currentSize; + if (projected > _tarStream.Capacity) + _tarStream.Capacity = (int)Math.Min(projected, int.MaxValue / 2); + } + await (_onEntryAdded?.Invoke(upload.HashedPair.ContentHash, _entries.Count, _currentSize) ?? ValueTask.CompletedTask); return _currentSize >= _targetSize ? await SealAsync(cancellationToken) : null; @@ -101,6 +120,11 @@ public TarBuilder( _tarWriter = null; var body = _tarStream!; // set together with _tarWriter in AddAsync, so non-null whenever the writer was + + // A bundle sealed well short of its presized buffer (the run's tail) must not hold it through the upload. + if (body.Capacity > 2 * body.Length) + body.Capacity = (int)body.Length; + body.Position = 0; var tarHash = ChunkHash.Parse(await _encryption.ComputeHashAsync(body, cancellationToken)); diff --git a/src/Arius.Core/Shared/ChunkIndex/ChunkIndexLocalStore.cs b/src/Arius.Core/Shared/ChunkIndex/ChunkIndexLocalStore.cs index 2ea87f55d..cbed3509f 100644 --- a/src/Arius.Core/Shared/ChunkIndex/ChunkIndexLocalStore.cs +++ b/src/Arius.Core/Shared/ChunkIndex/ChunkIndexLocalStore.cs @@ -1,3 +1,5 @@ +using System.Collections.ObjectModel; +using System.Text; using System.Text.Json; using Arius.Core.Shared.Storage; using Microsoft.Data.Sqlite; @@ -80,6 +82,89 @@ private void Initialize() /// public ShardEntry? FindPendingFlushEntry(ContentHash contentHash) => FindEntryCore(contentHash, pendingFlushOnly: true); + /// + /// Batched twin of : one query for a whole set of hashes. + /// Hashes with no stored entry are simply absent from the result. + /// + public IReadOnlyDictionary FindEntries(IReadOnlyCollection contentHashes) + => FindEntriesCore(contentHashes, pendingFlushOnly: false); + + /// + /// Batched twin of : one query for a whole set of hashes. + /// Returns only entries still pending local flush. + /// + public IReadOnlyDictionary FindPendingFlushEntries(IReadOnlyCollection contentHashes) + => FindEntriesCore(contentHashes, pendingFlushOnly: true); + + /// + /// Looks up a set of hashes with a single IN (...) query per page instead of one query per hash. + /// + /// + /// Paged because callers are unbounded (ls looks up one hash per file of a directory), and one + /// parameter per hash would otherwise exceed SQLite's variable limit. + /// + private IReadOnlyDictionary FindEntriesCore(IReadOnlyCollection contentHashes, bool pendingFlushOnly) + { + if (contentHashes.Count == 0) + return ReadOnlyDictionary.Empty; + + try + { + var results = new Dictionary(contentHashes.Count); + + using var connection = OpenConnection(); + foreach (var page in contentHashes.Chunk(MaxLookupPageSize)) + { + using var command = connection.CreateCommand(); + command.CommandText = BuildSql(page.Length, pendingFlushOnly); + + for (var i = 0; i < page.Length; i++) + { + var digest = CreateDigestBuffer(); + WriteDigest(page[i].ToString(), digest); + command.Parameters.Add($"$h{i}", SqliteType.Blob).Value = digest; + } + + using var reader = command.ExecuteReader(); + while (reader.Read()) + { + var entry = ReadEntry(reader); + results[entry.ContentHash] = entry; + } + } + + _logger.LogDebug("[chunk-index-local] FindEntries: requested={Requested} pendingFlushOnly={PendingFlushOnly} found={Found}", contentHashes.Count, pendingFlushOnly, results.Count); + return results; + } + catch (SqliteException ex) + { + throw CreateLocalStoreException(ex); + } + + static string BuildSql(int slots, bool pendingFlushOnly) + { + var sql = new StringBuilder("SELECT content_hash, chunk_hash, original_size, chunk_size, storage_tier_hint FROM chunk_index_entries WHERE content_hash IN ("); + for (var i = 0; i < slots; i++) + { + if (i > 0) + sql.Append(", "); + sql.Append("$h").Append(i); + } + + sql.Append(')'); + if (pendingFlushOnly) + sql.Append(" AND pending_flush = 1"); + + return sql.Append(';').ToString(); + } + } + + /// + /// Largest number of hashes bound into one IN (...) query, keeping the parameter count well + /// inside SQLite's limit regardless of how wide a caller's lookup is. + /// + private const int MaxLookupPageSize = 256; + private ShardEntry? FindEntryCore(ContentHash contentHash, bool pendingFlushOnly) { try @@ -89,7 +174,9 @@ private void Initialize() command.CommandText = pendingFlushOnly ? "SELECT content_hash, chunk_hash, original_size, chunk_size, storage_tier_hint FROM chunk_index_entries WHERE content_hash = $contentHash AND pending_flush = 1;" : "SELECT content_hash, chunk_hash, original_size, chunk_size, storage_tier_hint FROM chunk_index_entries WHERE content_hash = $contentHash;"; - command.Parameters.Add("$contentHash", SqliteType.Blob).Value = ParseHashBytes(contentHash.ToString()); + var digest = CreateDigestBuffer(); + WriteDigest(contentHash.ToString(), digest); + command.Parameters.Add("$contentHash", SqliteType.Blob).Value = digest; using var reader = command.ExecuteReader(); var entry = reader.Read() ? ReadEntry(reader) : null; _logger.LogDebug("[chunk-index-local] FindEntry: contentHash={ContentHash} pendingFlushOnly={PendingFlushOnly} found={Found}", contentHash.Short8, pendingFlushOnly, entry is not null); @@ -385,7 +472,7 @@ public void UpsertPendingFlush(ShardEntry entry) using var connection = OpenConnection(); using var transaction = connection.BeginTransaction(); using var command = CreateUpsertCommand(connection, transaction, pendingFlush: true, preservePendingFlushRows: false); - BindEntry(command, entry); + BindEntry(command, entry, CreateDigestBuffer(), CreateDigestBuffer()); var rowsAffected = command.ExecuteNonQuery(); transaction.Commit(); @@ -415,10 +502,12 @@ public void UpsertPendingFlush(IEnumerable entries) using var connection = OpenConnection(); using var transaction = connection.BeginTransaction(); using var command = CreateUpsertCommand(connection, transaction, pendingFlush: true, preservePendingFlushRows: false); + var contentDigest = CreateDigestBuffer(); + var chunkDigest = CreateDigestBuffer(); var rowsAffected = 0; foreach (var entry in materialized) { - BindEntry(command, entry); + BindEntry(command, entry, contentDigest, chunkDigest); rowsAffected += command.ExecuteNonQuery(); } @@ -452,10 +541,12 @@ public void UpsertRemoteBacked(IEnumerable entries) using var connection = OpenConnection(); using var transaction = connection.BeginTransaction(); using var command = CreateUpsertCommand(connection, transaction, pendingFlush: false, preservePendingFlushRows: false); + var contentDigest = CreateDigestBuffer(); + var chunkDigest = CreateDigestBuffer(); var rowsAffected = 0; foreach (var entry in batch) { - BindEntry(command, entry); + BindEntry(command, entry, contentDigest, chunkDigest); rowsAffected += command.ExecuteNonQuery(); } @@ -497,9 +588,11 @@ public void EnrichThinChunks(IReadOnlyDictionary + /// Binds one row using reusable digest buffers; each execution completes before the buffers are overwritten. + /// + private static void BindEntry(SqliteCommand command, ShardEntry entry, byte[] contentDigest, byte[] chunkDigest) { - command.Parameters["$contentHash"].Value = ParseHashBytes(entry.ContentHash.ToString()); - command.Parameters["$chunkHash"].Value = ParseHashBytes(entry.ChunkHash.ToString()); + WriteDigest(entry.ContentHash.ToString(), contentDigest); + WriteDigest(entry.ChunkHash.ToString(), chunkDigest); + command.Parameters["$contentHash"].Value = contentDigest; + command.Parameters["$chunkHash"].Value = chunkDigest; command.Parameters["$originalSize"].Value = entry.OriginalSize; command.Parameters["$chunkSize"].Value = entry.ChunkSize; command.Parameters["$storageTierHint"].Value = ShardEntry.SerializeTier(entry.StorageTierHint); @@ -907,6 +1007,14 @@ private static ShardEntry ReadEntry(SqliteDataReader reader) reader.GetInt64(3), ShardEntry.DeserializeTier(reader.GetInt32(4))); - private static byte[] ParseHashBytes(string value) - => Convert.FromHexString(value); + /// + /// Writes a canonical hexadecimal hash into a 32-byte digest buffer. The hash types guarantee exactly + /// 64 canonical hex characters, so the conversion cannot fail. + /// + private static void WriteDigest(string hex, byte[] destination) + => Convert.FromHexString(hex, destination, out _, out _); + + /// Creates a reusable 32-byte digest buffer. + private static byte[] CreateDigestBuffer() => new byte[HashCodec.Sha256ByteLength]; + } diff --git a/src/Arius.Core/Shared/ChunkIndex/ChunkIndexService.cs b/src/Arius.Core/Shared/ChunkIndex/ChunkIndexService.cs index d6a2bc0e3..e2762ec32 100644 --- a/src/Arius.Core/Shared/ChunkIndex/ChunkIndexService.cs +++ b/src/Arius.Core/Shared/ChunkIndex/ChunkIndexService.cs @@ -104,14 +104,16 @@ public async Task> LookupAsync(IEnu if (hashes.Length == 0) return result; + // Probe pending-flush entries in one batch before remote validation. + var pendingFlush = _localStore.FindPendingFlushEntries(hashes); + var validationWork = new List<(PathSegment Root, List Hashes)>(); foreach (var rootGroup in hashes.GroupBy(ChunkIndexRouter.GetRootPrefix)) { var hashesNeedingValidation = new List(); foreach (var contentHash in rootGroup) { - var pendingFlushEntry = _localStore.FindPendingFlushEntry(contentHash); - if (pendingFlushEntry is not null) + if (pendingFlush.TryGetValue(contentHash, out var pendingFlushEntry)) { // Entry is local-only / dirty result[contentHash] = pendingFlushEntry; @@ -143,15 +145,11 @@ await Parallel.ForEachAsync( await EnsureCoverageForHashesAsync(item.Root, item.Hashes, latestSnapshotName, ct); }); - // Construct the result from the validated shards + // Populate the result from validated shards using batched lookups. foreach (var item in validationWork) { - foreach (var contentHash in item.Hashes) - { - var entry = _localStore.FindEntry(contentHash); - if (entry is not null) - result[contentHash] = entry; - } + foreach (var (contentHash, entry) in _localStore.FindEntries(item.Hashes)) + result[contentHash] = entry; } return result; diff --git a/src/Arius.Core/Shared/ChunkStorage/ChunkStorageService.cs b/src/Arius.Core/Shared/ChunkStorage/ChunkStorageService.cs index 987967b62..d6a07bd94 100644 --- a/src/Arius.Core/Shared/ChunkStorage/ChunkStorageService.cs +++ b/src/Arius.Core/Shared/ChunkStorage/ChunkStorageService.cs @@ -306,6 +306,13 @@ private sealed class ChunkDownloadStream(Stream inner) : Stream public override long Position { get => inner.Position; set => inner.Position = value; } public override void Flush() => inner.Flush(); public override int Read(byte[] buffer, int offset, int count) => inner.Read(buffer, offset, count); + + // Override the async members so restore uses the inner stream's asynchronous read and copy paths. + public override int Read(Span buffer) => inner.Read(buffer); + public override Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken) => inner.ReadAsync(buffer, offset, count, cancellationToken); + public override ValueTask ReadAsync(Memory buffer, CancellationToken cancellationToken = default) => inner.ReadAsync(buffer, cancellationToken); + public override Task CopyToAsync(Stream destination, int bufferSize, CancellationToken cancellationToken) => inner.CopyToAsync(destination, bufferSize, cancellationToken); + public override long Seek(long offset, SeekOrigin origin) => inner.Seek(offset, origin); public override void SetLength(long value) => inner.SetLength(value); public override void Write(byte[] buffer, int offset, int count) => inner.Write(buffer, offset, count); @@ -337,4 +344,4 @@ private void DisposeResources() inner.Dispose(); } } -} \ No newline at end of file +} diff --git a/src/Arius.Core/Shared/Compression/ZstdCompressionService.cs b/src/Arius.Core/Shared/Compression/ZstdCompressionService.cs index e63ebc6c8..0b69d6dd4 100644 --- a/src/Arius.Core/Shared/Compression/ZstdCompressionService.cs +++ b/src/Arius.Core/Shared/Compression/ZstdCompressionService.cs @@ -13,7 +13,7 @@ internal sealed class ZstdCompressionService(int compressionLevel = ZstdCompress : ICompressionService { // Magic number at the start of a zstd frame (little-endian on disk). - private static readonly byte[] ZstdMagic = [0x28, 0xB5, 0x2F, 0xFD]; // 0xFD2FB528 + private static ReadOnlySpan ZstdMagic => [0x28, 0xB5, 0x2F, 0xFD]; // 0xFD2FB528 // Legacy "+gzip" blobs are decoded by the gzip codec; reads delegate to it when a gzip frame is detected. private static readonly GZipCompressionService LegacyGzip = new(); diff --git a/src/Arius.Core/Shared/Encryption/PassphraseEncryptionService.cs b/src/Arius.Core/Shared/Encryption/PassphraseEncryptionService.cs index cbe638ae3..b35048197 100644 --- a/src/Arius.Core/Shared/Encryption/PassphraseEncryptionService.cs +++ b/src/Arius.Core/Shared/Encryption/PassphraseEncryptionService.cs @@ -35,7 +35,7 @@ internal sealed class PassphraseEncryptionService : IEncryptionService // key for anything written today (production writes AES-GCM, see GcmPbkdf2Iter = 100_000). It // must match the blobs on disk, so it cannot be raised. private const int CbcPbkdf2Iter = 10_000; - private static readonly byte[] SaltedMagic = "Salted__"u8.ToArray(); + private static ReadOnlySpan SaltedMagic => "Salted__"u8; // ── GCM constants ──────────────────────────────────────────────────────────── private const int GcmSaltSize = 16; @@ -45,7 +45,7 @@ internal sealed class PassphraseEncryptionService : IEncryptionService private const int GcmBlockSize = 64 * 1024; // 64 KiB private const int GcmPbkdf2Iter = 100_000; private const uint GcmMaxPbkdf2Iter = 10_000_000; // sanity cap: reject crafted blobs - private static readonly byte[] GcmMagic = "ArGCM1"u8.ToArray(); // 6 bytes + private static ReadOnlySpan GcmMagic => "ArGCM1"u8; // 6 bytes private readonly byte[] _passphraseBytes; diff --git a/src/Arius.Core/Shared/FileSystem/PathSegment.cs b/src/Arius.Core/Shared/FileSystem/PathSegment.cs index 1e925ed99..4a34833c7 100644 --- a/src/Arius.Core/Shared/FileSystem/PathSegment.cs +++ b/src/Arius.Core/Shared/FileSystem/PathSegment.cs @@ -38,7 +38,7 @@ public static bool TryParse(string? value, out PathSegment segment) return false; } - if (value.Contains('/') || value.Contains('\\') || value.Any(char.IsControl)) + if (value.Contains('/') || value.Contains('\\') || ContainsControlCharacter(value)) { segment = default; return false; @@ -46,6 +46,16 @@ public static bool TryParse(string? value, out PathSegment segment) segment = new PathSegment(value); return true; + + // Allocation-free replacement for value.Any(char.IsControl). + static bool ContainsControlCharacter(string value) + { + foreach (var c in value) + if (char.IsControl(c)) + return true; + + return false; + } } public bool Contains(string value, StringComparison comparisonType) => diff --git a/src/Arius.Core/Shared/FileSystem/RelativeFileSystem.cs b/src/Arius.Core/Shared/FileSystem/RelativeFileSystem.cs index 57d0a425b..ac13532ea 100644 --- a/src/Arius.Core/Shared/FileSystem/RelativeFileSystem.cs +++ b/src/Arius.Core/Shared/FileSystem/RelativeFileSystem.cs @@ -273,6 +273,17 @@ public async Task WriteAllBytesAsync(RelativePath path, byte[] content, Cancella await File.WriteAllBytesAsync(fullPath, content, cancellationToken); } + /// + /// Writes to without copying it to an array first. + /// Mirrors . + /// + public async Task WriteAllBytesAsync(RelativePath path, ReadOnlyMemory content, CancellationToken cancellationToken) + { + var fullPath = root.Resolve(path); + CreateDirectory(path.Parent ?? RelativePath.Root); + await File.WriteAllBytesAsync(fullPath, content, cancellationToken); + } + public void ReplaceFileAtomically(RelativePath source, RelativePath destination) { var sourcePath = root.Resolve(source); diff --git a/src/Arius.Core/Shared/FileSystem/RelativePath.cs b/src/Arius.Core/Shared/FileSystem/RelativePath.cs index 59fc4d875..ca2a8d083 100644 --- a/src/Arius.Core/Shared/FileSystem/RelativePath.cs +++ b/src/Arius.Core/Shared/FileSystem/RelativePath.cs @@ -95,7 +95,7 @@ public static bool TryParse(string? value, out RelativePath path) return false; } - if (value.Contains('\\') || value.Contains("//", StringComparison.Ordinal) || value.Any(char.IsControl)) + if (value.Contains('\\') || value.Contains("//", StringComparison.Ordinal)) { path = default; return false; diff --git a/src/Arius.Core/Shared/FileTree/FileTreeSerializer.cs b/src/Arius.Core/Shared/FileTree/FileTreeSerializer.cs index 1869f8e45..adf84b8c8 100644 --- a/src/Arius.Core/Shared/FileTree/FileTreeSerializer.cs +++ b/src/Arius.Core/Shared/FileTree/FileTreeSerializer.cs @@ -11,10 +11,8 @@ internal static class FileTreeSerializer { private static readonly Encoding s_utf8 = new UTF8Encoding(encoderShouldEmitUTF8Identifier: false); - public static IReadOnlyList Deserialize(byte[] bytes) + public static IReadOnlyList Deserialize(ReadOnlySpan bytes) { - ArgumentNullException.ThrowIfNull(bytes); - var text = s_utf8.GetString(bytes); return ParsePersistedLines(text.Split('\n')); } diff --git a/src/Arius.Core/Shared/FileTree/FileTreeService.cs b/src/Arius.Core/Shared/FileTree/FileTreeService.cs index 694915392..817ed76ca 100644 --- a/src/Arius.Core/Shared/FileTree/FileTreeService.cs +++ b/src/Arius.Core/Shared/FileTree/FileTreeService.cs @@ -1,6 +1,7 @@ using System.Collections.Concurrent; using Arius.Core.Shared.Compression; using Arius.Core.Shared.Encryption; +using Arius.Core.Shared.Extensions; using Arius.Core.Shared.Snapshot; using Arius.Core.Shared.Storage; using Microsoft.Extensions.Logging; @@ -243,6 +244,7 @@ private async Task SerializeStorageAsync(ReadOnlyMemory plaintext, await compressionStream.WriteAsync(plaintext, cancellationToken); } + // The codec chain closes ms; ToArray remains valid after close, unlike buffer-based accessors. return ms.ToArray(); } @@ -252,7 +254,8 @@ private async Task> DeserializeStorageAsync(Stream await using var decompressStream = _compression.WrapForDecompression(decStream); using var ms = new MemoryStream(); await decompressStream.CopyToAsync(ms, cancellationToken); - return FileTreeSerializer.Deserialize(ms.ToArray()); + var buffer = ms.ToArraySegment(); + return FileTreeSerializer.Deserialize(buffer.AsSpan()); } private async Task WriteCacheAtomicallyAsync(RelativePath diskPath, ReadOnlyMemory plaintext, CancellationToken cancellationToken) @@ -261,7 +264,7 @@ private async Task WriteCacheAtomicallyAsync(RelativePath diskPath, ReadOnlyMemo try { - await _diskCacheFileSystem.WriteAllBytesAsync(tempPath, plaintext.ToArray(), cancellationToken); + await _diskCacheFileSystem.WriteAllBytesAsync(tempPath, plaintext, cancellationToken); _diskCacheFileSystem.ReplaceFileAtomically(tempPath, diskPath); } finally diff --git a/src/Arius.Core/Shared/FileTree/FileTreeStagingWriter.cs b/src/Arius.Core/Shared/FileTree/FileTreeStagingWriter.cs index a06a09409..4f7a3dba5 100644 --- a/src/Arius.Core/Shared/FileTree/FileTreeStagingWriter.cs +++ b/src/Arius.Core/Shared/FileTree/FileTreeStagingWriter.cs @@ -96,7 +96,7 @@ private async Task AppendDirectoryEntriesAsync(RelativePath filePath, Cancellati private async Task AppendLineAsync(RelativePath path, string line, CancellationToken cancellationToken) { - var nodeLock = _lockStripes[(uint)StringComparer.Ordinal.GetHashCode(path) % (uint)_lockStripes.Length]; + var nodeLock = _lockStripes[(uint)path.GetHashCode() % (uint)_lockStripes.Length]; await nodeLock.WaitAsync(cancellationToken); try diff --git a/src/Arius.Core/Shared/HashCache/SparseFingerprint.cs b/src/Arius.Core/Shared/HashCache/SparseFingerprint.cs index 95d139f95..5bd4d55e0 100644 --- a/src/Arius.Core/Shared/HashCache/SparseFingerprint.cs +++ b/src/Arius.Core/Shared/HashCache/SparseFingerprint.cs @@ -1,3 +1,5 @@ +using System.Buffers; +using System.Buffers.Binary; using System.Security.Cryptography; namespace Arius.Core.Shared.HashCache; @@ -48,21 +50,36 @@ internal static class SparseFingerprint public static byte[] ComputeBySeeking(RelativeFileSystem fs, RelativePath path, long size) { using var sha = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); - sha.AppendData(BitConverter.GetBytes(size)); + AppendSize(sha, size); using var stream = fs.OpenRead(path); var regions = Regions(size); + if (regions.Count == 0) + return sha.GetHashAndReset(); + // Buffer the largest region: the single whole-file region for a small file can reach k×BlockSize // (up to 1 MiB at k=MinBlocks), which is larger than BlockSize — a fixed BlockSize buffer would // overflow ReadExactly for files in (BlockSize, k×BlockSize]. - var buffer = new byte[regions.Count == 0 ? 0 : regions.Max(r => r.Length)]; - foreach (var (offset, length) in regions) + var longest = 0; + for (var i = 0; i < regions.Count; i++) + longest = Math.Max(longest, regions[i].Length); + + var buffer = ArrayPool.Shared.Rent(longest); + try + { + foreach (var (offset, length) in regions) + { + stream.Seek(offset, SeekOrigin.Begin); + stream.ReadExactly(buffer, 0, length); + sha.AppendData(buffer, 0, length); + } + + return sha.GetHashAndReset(); + } + finally { - stream.Seek(offset, SeekOrigin.Begin); - stream.ReadExactly(buffer, 0, length); - sha.AppendData(buffer, 0, length); + ArrayPool.Shared.Return(buffer); } - return sha.GetHashAndReset(); } /// @@ -72,17 +89,38 @@ public static byte[] ComputeBySeeking(RelativeFileSystem fs, RelativePath path, /// for the same content (same , same size ‖ region-bytes framing) — keep /// the two in sync. /// - public sealed class Sampler + public sealed class Sampler : IDisposable { private readonly long _size; private readonly IReadOnlyList<(long Off, int Len)> _regions; - private readonly byte[][] _captured; + + /// + /// Stores sampled regions contiguously so can hash them in one pass. + /// + private readonly byte[] _buffer; + private readonly int[] _bufferOffsets; + private readonly int _capturedLength; + + private bool _disposed; public Sampler(long size) { - _size = size; - _regions = Regions(size); - _captured = _regions.Select(r => new byte[r.Len]).ToArray(); + _size = size; + _regions = Regions(size); + + _bufferOffsets = new int[_regions.Count]; + var total = 0; + for (var i = 0; i < _regions.Count; i++) + { + _bufferOffsets[i] = total; + total += _regions[i].Len; + } + + _capturedLength = total; + _buffer = ArrayPool.Shared.Rent(total); + + // Clear uncaptured regions because pooled buffers are not initialized. + _buffer.AsSpan(0, _capturedLength).Clear(); } /// Offer the bytes read at ; overlapping region bytes are copied out. @@ -97,17 +135,39 @@ public void Capture(long position, ReadOnlySpan buffer) continue; var srcStart = (int)(from - position); var dstStart = (int)(from - off); - buffer.Slice(srcStart, (int)(to - from)).CopyTo(_captured[i].AsSpan(dstStart)); + buffer.Slice(srcStart, (int)(to - from)).CopyTo(_buffer.AsSpan(_bufferOffsets[i] + dstStart)); } } public byte[] Finish() { + // Hashing a returned buffer would fingerprint another consumer's bytes. + ObjectDisposedException.ThrowIf(_disposed, this); + using var sha = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); - sha.AppendData(BitConverter.GetBytes(_size)); - foreach (var region in _captured) - sha.AppendData(region); + AppendSize(sha, _size); + sha.AppendData(_buffer.AsSpan(0, _capturedLength)); return sha.GetHashAndReset(); } + + public void Dispose() + { + if (_disposed) + return; + + _disposed = true; + ArrayPool.Shared.Return(_buffer); + } + } + + /// + /// Frames the file size into the digest as 8 little-endian bytes. Both compute paths must agree on this + /// exactly — it is the size ‖ region-bytes prefix. + /// + private static void AppendSize(IncrementalHash sha, long size) + { + Span sizeBytes = stackalloc byte[sizeof(long)]; + BinaryPrimitives.WriteInt64LittleEndian(sizeBytes, size); + sha.AppendData(sizeBytes); } } diff --git a/src/Arius.Core/Shared/HashCache/SparseSamplingStream.cs b/src/Arius.Core/Shared/HashCache/SparseSamplingStream.cs index 5845a4140..897841988 100644 --- a/src/Arius.Core/Shared/HashCache/SparseSamplingStream.cs +++ b/src/Arius.Core/Shared/HashCache/SparseSamplingStream.cs @@ -63,9 +63,13 @@ public override long Position protected override void Dispose(bool disposing) { - if (disposing) - _inner.Dispose(); - + if (disposing) + { + // Fingerprint must be called before disposal releases the sampler buffer. + _sampler.Dispose(); + _inner.Dispose(); + } + base.Dispose(disposing); } } diff --git a/src/Arius.Core/Shared/Hashes/HashCodec.cs b/src/Arius.Core/Shared/Hashes/HashCodec.cs index 7e7e286b5..e1434b244 100644 --- a/src/Arius.Core/Shared/Hashes/HashCodec.cs +++ b/src/Arius.Core/Shared/Hashes/HashCodec.cs @@ -17,19 +17,27 @@ public static string NormalizeHex(string value) throw new FormatException($"Expected {Sha256HexLength} hex characters but got {value.Length}."); Span chars = stackalloc char[Sha256HexLength]; + var alreadyCanonical = true; for (var i = 0; i < value.Length; i++) { var c = value[i]; - chars[i] = c switch + switch (c) { - >= '0' and <= '9' => c, - >= 'a' and <= 'f' => c, - >= 'A' and <= 'F' => char.ToLowerInvariant(c), - _ => throw new FormatException($"Invalid hex character '{c}'.") - }; + case >= '0' and <= '9': + case >= 'a' and <= 'f': + chars[i] = c; + break; + case >= 'A' and <= 'F': + chars[i] = char.ToLowerInvariant(c); + alreadyCanonical = false; + break; + default: + throw new FormatException($"Invalid hex character '{c}'."); + } } - return new string(chars); + // Preserve canonical input to avoid allocating a duplicate string. + return alreadyCanonical ? value : new string(chars); } public static string ToLowerHex(ReadOnlySpan digest) @@ -37,6 +45,6 @@ public static string ToLowerHex(ReadOnlySpan digest) if (digest.Length != Sha256ByteLength) throw new ArgumentException($"Expected {Sha256ByteLength}-byte SHA-256 digest.", nameof(digest)); - return Convert.ToHexString(digest).ToLowerInvariant(); + return Convert.ToHexStringLower(digest); } } diff --git a/src/Arius.Core/Shared/Storage/BlobConstants.cs b/src/Arius.Core/Shared/Storage/BlobConstants.cs index 4a5f27941..e3b7c7fe0 100644 --- a/src/Arius.Core/Shared/Storage/BlobConstants.cs +++ b/src/Arius.Core/Shared/Storage/BlobConstants.cs @@ -77,25 +77,25 @@ public static class BlobPaths // NOTE: These methods require a strong domain type (unless there is none). For string convenience overloads used in Test suites, see Arius.Tests.Shared.BlobPathsExtensions. /// Content-addressable chunks: large files and thin pointers. - public static RelativePath ChunksPrefix => RelativePath.Root / PathSegment.Parse("chunks"); + public static readonly RelativePath ChunksPrefix = RelativePath.Root / PathSegment.Parse("chunks"); /// Temporary Hot-tier copies for in-progress rehydration. - public static RelativePath ChunksRehydratedPrefix => RelativePath.Root / PathSegment.Parse("chunks-rehydrated"); + public static readonly RelativePath ChunksRehydratedPrefix = RelativePath.Root / PathSegment.Parse("chunks-rehydrated"); /// /// Chunk metadata sidecars for chunks whose own metadata cannot be written — Archive-tier blobs migrated from v5, where Azure forbids Set Blob Metadata. /// NOTE: if chunk pruning/GC is ever added, deleting a chunk must also delete its metadata sidecar. /// - public static RelativePath V5LegacySideCarPrefix => RelativePath.Root / PathSegment.Parse("chunks-v5legacy-metadata"); + public static readonly RelativePath V5LegacySideCarPrefix = RelativePath.Root / PathSegment.Parse("chunks-v5legacy-metadata"); /// Merkle tree blobs (one per directory). - public static RelativePath FileTreesPrefix => RelativePath.Root / PathSegment.Parse("filetrees"); + public static readonly RelativePath FileTreesPrefix = RelativePath.Root / PathSegment.Parse("filetrees"); /// Snapshot manifests. - public static RelativePath SnapshotsPrefix => RelativePath.Root / PathSegment.Parse("snapshots"); + public static readonly RelativePath SnapshotsPrefix = RelativePath.Root / PathSegment.Parse("snapshots"); /// Chunk index shards. - public static RelativePath ChunkIndexPrefix => RelativePath.Root / PathSegment.Parse("chunk-index"); + public static readonly RelativePath ChunkIndexPrefix = RelativePath.Root / PathSegment.Parse("chunk-index"); public static RelativePath ChunkPath(ChunkHash hash) => ChunksPrefix / PathSegment.Parse(hash.ToString()); public static RelativePath ThinChunkPath(ContentHash hash) => ChunksPrefix / PathSegment.Parse(hash.ToString()); diff --git a/src/Arius.Core/Shared/Streaming/ProgressStream.cs b/src/Arius.Core/Shared/Streaming/ProgressStream.cs index 701079d8e..91fff7d8c 100644 --- a/src/Arius.Core/Shared/Streaming/ProgressStream.cs +++ b/src/Arius.Core/Shared/Streaming/ProgressStream.cs @@ -1,82 +1,109 @@ namespace Arius.Core.Shared.Streaming; /// -/// Read-mode stream wrapper that reports cumulative bytes read via . -/// Delegates all reads to the inner stream and reports progress after each read operation. -/// Does not buffer any data. +/// Read-mode stream wrapper that reports cumulative bytes read without buffering. +/// The first read reports immediately; later reports are limited to one per +/// . The final total is emitted at EOF or disposal. /// public sealed class ProgressStream : Stream { + /// Minimum wall-clock gap between two progress callbacks. + private static readonly TimeSpan ReportInterval = TimeSpan.FromMilliseconds(500); + private readonly Stream _inner; private readonly IProgress _progress; + private readonly TimeProvider _timeProvider; private long _bytesRead; + private long _reportedBytes; + private long _lastReportTimestamp; /// The readable source stream. - /// Receives cumulative bytes read after each read call. + /// Receives the cumulative bytes read, throttled as described on the type. public ProgressStream(Stream inner, IProgress progress) + : this(inner, progress, TimeProvider.System) + { + } + + internal ProgressStream(Stream inner, IProgress progress, TimeProvider timeProvider) { ArgumentNullException.ThrowIfNull(inner); ArgumentNullException.ThrowIfNull(progress); if (!inner.CanRead) throw new ArgumentException("Inner stream must be readable.", nameof(inner)); - _inner = inner; - _progress = progress; + _inner = inner; + _progress = progress; + _timeProvider = timeProvider; } - public override bool CanRead => true; - public override bool CanWrite => false; - public override bool CanSeek => false; - - public override int Read(byte[] buffer, int offset, int count) + /// Accounts for one inner read of bytes out of . + private int OnRead(int read, int requested) { - var n = _inner.Read(buffer, offset, count); - if (n > 0) + if (read > 0) { - _bytesRead += n; - _progress.Report(_bytesRead); + _bytesRead += read; + ReportThrottled(); } - return n; - } - - public override int Read(Span buffer) - { - var n = _inner.Read(buffer); - if (n > 0) + else if (requested != 0) // a zero-length read is not EOF { - _bytesRead += n; - _progress.Report(_bytesRead); + ReportFinal(); } - return n; + + return read; } - public override async Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken ct) + /// + /// Reports the running total at most once per using monotonic time. + /// + private void ReportThrottled() { - var n = await _inner.ReadAsync(buffer, offset, count, ct); - if (n > 0) - { - _bytesRead += n; - _progress.Report(_bytesRead); - } - return n; + var now = _timeProvider.GetTimestamp(); + + if (_reportedBytes > 0 && _timeProvider.GetElapsedTime(_lastReportTimestamp, now) < ReportInterval) + return; + + _lastReportTimestamp = now; + _reportedBytes = _bytesRead; + _progress.Report(_bytesRead); } - public override async ValueTask ReadAsync(Memory buffer, CancellationToken ct = default) + /// + /// Reports any total withheld by throttling at EOF or disposal. + /// + private void ReportFinal() { - var n = await _inner.ReadAsync(buffer, ct); - if (n > 0) - { - _bytesRead += n; - _progress.Report(_bytesRead); - } - return n; + // Also keeps an empty source from emitting a spurious zero. + if (_reportedBytes == _bytesRead) + return; + + _reportedBytes = _bytesRead; + _progress.Report(_bytesRead); } + public override bool CanRead => true; + public override bool CanWrite => false; + public override bool CanSeek => false; + + public override int Read(byte[] buffer, int offset, int count) => OnRead(_inner.Read(buffer, offset, count), count); + + public override int Read(Span buffer) => OnRead(_inner.Read(buffer), buffer.Length); + + public override async Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken ct) + => OnRead(await _inner.ReadAsync(buffer, offset, count, ct), count); + + public override async ValueTask ReadAsync(Memory buffer, CancellationToken ct = default) + => OnRead(await _inner.ReadAsync(buffer, ct), buffer.Length); + public override void Flush() => _inner.Flush(); protected override void Dispose(bool disposing) { - if (disposing) _inner.Dispose(); + if (disposing) + { + ReportFinal(); + _inner.Dispose(); + } + base.Dispose(disposing); }