From 11ebfb29e12d4405350320867ca5084f4b8ea157 Mon Sep 17 00:00:00 2001 From: Matthew DeVenny Date: Tue, 28 Jul 2026 08:47:17 -0700 Subject: [PATCH 1/2] #45 Write TryGetAsync payload directly into the caller buffer The IBufferWriter overload of TryGetAsync copied the payload twice: transport buffer -> CacheEntry.Data -> caller buffer. Add a per-read, destination-bound deserializer (BufferWritingCacheEntryDeserializer) that writes the payload straight into the caller's buffer on a hit, eliminating the intermediate array on the common path (a single copy; true zero-copy is impossible because the payload must land in the caller's buffer). The write only happens for a genuine hit, so the "nothing written on a miss" contract holds: undeserializable framing, absolutely-expired, and not-found all leave the buffer untouched. Sliding-expiration entries keep the materialized two-copy path because refreshing their TTL re-writes the value and therefore needs the bytes. Factor the shared header parse into CacheEntryBinarySerializer.TryReadHeader and extract UpdateEntryExpirationAsync so the array and buffer read paths agree byte-for-byte and refresh sliding TTLs through one implementation. Add integration tests: exact-payload hit, empty-payload hit, and nothing-written on miss / undeserializable / absolutely-expired / sliding-past-absolute. Co-Authored-By: Claude Opus 4.8 (1M context) Signed-off-by: Matthew DeVenny --- .../BufferWritingCacheEntryDeserializer.cs | 145 ++++++++++++++ .../CacheEntryBinarySerializer.cs | 53 ++++-- src/NatsDistributedCache/NatsCache.cs | 179 ++++++++++++++---- .../Cache/TryGetAsyncBufferTests.cs | 168 ++++++++++++++++ 4 files changed, 494 insertions(+), 51 deletions(-) create mode 100644 src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs create mode 100644 test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs diff --git a/src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs b/src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs new file mode 100644 index 0000000..6cfced8 --- /dev/null +++ b/src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs @@ -0,0 +1,145 @@ +using System.Buffers; +using NATS.Client.Core; + +namespace CodeCargo.Nats.DistributedCache; + +/// +/// Outcome of a single-copy read produced by +/// . Carries just enough for the read core to finish the +/// operation without a second copy of the payload on the common path. A result (the +/// deserializer returning null) means the stored bytes were undeserializable — exactly as +/// signals the same condition with a null entry. +/// +internal sealed class CacheEntryBufferReadResult +{ + internal CacheEntryBufferReadResult(CacheEntry entry, bool payloadWritten, bool absolutelyExpired) + { + Entry = entry; + PayloadWritten = payloadWritten; + AbsolutelyExpired = absolutelyExpired; + } + + /// + /// The decoded entry. Its is populated only when the payload was + /// not streamed to the destination (the sliding-expiration case, where the read core must + /// re-write the value to refresh its TTL); otherwise it is left null because the bytes already live in + /// the caller's buffer. + /// + internal CacheEntry Entry { get; } + + /// + /// when the payload was written straight into the caller's + /// — a confirmed hit with no sliding refresh, and the single-copy fast + /// path. The read core only needs to record the hit. + /// + internal bool PayloadWritten { get; } + + /// + /// when the entry is past its absolute expiration. Nothing was written; the read + /// core evicts the entry (using the KV revision) and reports a miss. + /// + internal bool AbsolutelyExpired { get; } +} + +/// +/// A per-read used only by NatsCache.TryGetAsync(IBufferWriter<byte>) +/// to avoid the intermediate array that the array read path allocates. It is the only +/// point in the read that holds the transport payload sequence, so it writes the payload directly into the +/// caller's destination — a single copy (transport buffer → caller buffer) instead of two (transport buffer → +/// → caller buffer). +/// +/// +/// The write only happens for a genuine hit, so the contract of "nothing +/// written on a miss" is preserved: +/// +/// undeserializable framing → returns , writes nothing; +/// absolutely expired → returns a result flagged expired, writes nothing (the core evicts + misses); +/// sliding expiration present → materializes and writes nothing, because +/// the core must re-write the whole value to refresh its TTL and therefore needs the bytes (this is the one +/// case that still costs two copies — the common absolute-only/no-expiry hit does not); +/// otherwise → writes the payload to the destination and flags the hit. +/// +/// Absolute expiry is evaluated here rather than in the core because this is the only moment the payload is +/// available, and the buffer must stay untouched when the entry turns out to be an expired miss. It uses +/// at deserialize time — the same clock and effectively the same instant +/// the core would have used immediately afterward — so the decision matches NatsCache.IsAbsolutelyExpired. +/// +internal sealed class BufferWritingCacheEntryDeserializer : INatsDeserialize +{ + private readonly IBufferWriter _destination; + private readonly TimeProvider _timeProvider; + + internal BufferWritingCacheEntryDeserializer(IBufferWriter destination, TimeProvider timeProvider) + { + _destination = destination; + _timeProvider = timeProvider; + } + + /// + public CacheEntryBufferReadResult? Deserialize(in ReadOnlySequence buffer) + { + var reader = new SequenceReader(buffer); + if (!CacheEntryBinarySerializer.TryReadHeader(ref reader, out var absoluteExpiration, out var slidingExpirationTicks)) + { + // Unknown/legacy/corrupt framing: signal an undeserializable entry (see the type doc). + return null; + } + + var remaining = reader.UnreadSequence; + + if (slidingExpirationTicks.HasValue) + { + // Sliding entries have their TTL refreshed by re-writing the value, which needs the payload + // bytes — so this path cannot stream to the caller. Materialize the full entry (one copy) and + // let the read core run its normal expiry/refresh/write, exactly like the array read path. + var data = remaining.IsEmpty ? Array.Empty() : remaining.ToArray(); + return new CacheEntryBufferReadResult( + new CacheEntry + { + AbsoluteExpiration = absoluteExpiration, + SlidingExpirationTicks = slidingExpirationTicks, + Data = data, + }, + payloadWritten: false, + absolutelyExpired: false); + } + + if (absoluteExpiration.HasValue && _timeProvider.GetUtcNow() >= absoluteExpiration.Value) + { + // Absolutely expired: write nothing. The read core evicts the entry and reports a miss. + return new CacheEntryBufferReadResult( + new CacheEntry { AbsoluteExpiration = absoluteExpiration }, + payloadWritten: false, + absolutelyExpired: true); + } + + // Genuine hit with no sliding refresh: copy the payload straight into the caller's buffer. This is + // the single copy that makes TryGet allocation-free of any intermediate array on the common path. + WritePayload(remaining); + return new CacheEntryBufferReadResult( + new CacheEntry { AbsoluteExpiration = absoluteExpiration }, + payloadWritten: true, + absolutelyExpired: false); + } + + private void WritePayload(in ReadOnlySequence payload) + { + if (payload.IsSingleSegment) + { + if (!payload.First.IsEmpty) + { + _destination.Write(payload.First.Span); + } + + return; + } + + foreach (var segment in payload) + { + if (!segment.IsEmpty) + { + _destination.Write(segment.Span); + } + } + } +} diff --git a/src/NatsDistributedCache/CacheEntryBinarySerializer.cs b/src/NatsDistributedCache/CacheEntryBinarySerializer.cs index 2a43c71..b096b66 100644 --- a/src/NatsDistributedCache/CacheEntryBinarySerializer.cs +++ b/src/NatsDistributedCache/CacheEntryBinarySerializer.cs @@ -98,21 +98,53 @@ public void Serialize(IBufferWriter bufferWriter, CacheEntry value) public CacheEntry? Deserialize(in ReadOnlySequence buffer) { var reader = new SequenceReader(buffer); + if (!TryReadHeader(ref reader, out var absoluteExpiration, out var slidingExpirationTicks)) + { + // Unknown/legacy/corrupt framing: treat as a cache miss (see TryReadHeader). + return null; + } + + var remaining = reader.UnreadSequence; + var data = remaining.IsEmpty ? Array.Empty() : remaining.ToArray(); + + return new CacheEntry + { + AbsoluteExpiration = absoluteExpiration, + SlidingExpirationTicks = slidingExpirationTicks, + Data = data, + }; + } + + /// + /// Reads and validates the fixed header (version, flags, and the optional + /// expiration fields), advancing to the first payload byte. Returns + /// for any unknown, legacy (e.g. JSON), or corrupt framing so callers can fail + /// closed to a cache miss instead of misinterpreting the remaining bytes. Shared by + /// and the single-copy read path in + /// so both agree, byte for byte, on what a valid entry + /// is and where its payload begins. + /// + internal static bool TryReadHeader( + ref SequenceReader reader, + out DateTimeOffset? absoluteExpiration, + out long? slidingExpirationTicks) + { + absoluteExpiration = null; + slidingExpirationTicks = null; if (!reader.TryRead(out var version) || version != FormatVersion) { // Unknown or legacy (e.g. JSON) entry: treat as a cache miss. - return null; + return false; } if (!reader.TryRead(out var flags) || (flags & ~KnownFlags) != 0) { // Unknown flag bits: corrupt data, or a future format that mistakenly reused this version. // Fail closed rather than misinterpreting the remaining bytes. - return null; + return false; } - DateTimeOffset? absoluteExpiration = null; if ((flags & HasAbsoluteExpiration) != 0) { if (!reader.TryReadLittleEndian(out long absoluteTicks) || @@ -120,13 +152,12 @@ public void Serialize(IBufferWriter bufferWriter, CacheEntry value) absoluteTicks > DateTimeOffset.MaxValue.Ticks) { // Missing or out-of-range ticks (corrupt entry): treat as a miss instead of throwing. - return null; + return false; } absoluteExpiration = new DateTimeOffset(absoluteTicks, TimeSpan.Zero); } - long? slidingExpirationTicks = null; if ((flags & HasSlidingExpiration) != 0) { if (!reader.TryReadLittleEndian(out long slidingTicks) || @@ -137,20 +168,12 @@ public void Serialize(IBufferWriter bufferWriter, CacheEntry value) // a miss, like the absolute-ticks bounds check above. A valid sliding window is always // a positive TimeSpan within (0, MaxTtlTicks]; the write path enforces the same ceiling, // so any legitimately written value round-trips. - return null; + return false; } slidingExpirationTicks = slidingTicks; } - var remaining = reader.UnreadSequence; - var data = remaining.IsEmpty ? Array.Empty() : remaining.ToArray(); - - return new CacheEntry - { - AbsoluteExpiration = absoluteExpiration, - SlidingExpirationTicks = slidingExpirationTicks, - Data = data, - }; + return true; } } diff --git a/src/NatsDistributedCache/NatsCache.cs b/src/NatsDistributedCache/NatsCache.cs index fec9c82..6c57d27 100644 --- a/src/NatsDistributedCache/NatsCache.cs +++ b/src/NatsDistributedCache/NatsCache.cs @@ -219,15 +219,12 @@ public async ValueTask TryGetAsync( { try { - // Reports as operation=get, not a distinct operation: the IBufferWriter overload is a zero-copy - // detail, and merging it keeps hit ratio computed over all read paths — which matters because - // HybridCache drives its L2 reads exclusively through this method. - var result = await GetAndRefreshAsync(key, CacheOperation.Get, token).ConfigureAwait(false); - if (result != null) - { - destination.Write(result); - return true; - } + // Writes the payload straight into the caller's buffer on a hit (a single copy), rather than + // fetching a byte[] and copying it a second time. Reports as operation=get, not a distinct + // operation: the IBufferWriter overload is a zero-copy detail, and merging it keeps hit ratio + // computed over all read paths — which matters because HybridCache drives its L2 reads + // exclusively through this method. + return await TryGetAndRefreshBufferAsync(key, destination, token).ConfigureAwait(false); } catch (OperationCanceledException) when (token.IsCancellationRequested) { @@ -469,7 +466,8 @@ private Lazy> CreateLazyKvStore() => } var kvEntry = natsResult.Value; - if (kvEntry.Value == null) + var entry = kvEntry.Value; + if (entry == null) { // Present entry whose bytes we cannot deserialize: a legacy JSON envelope from a // pre-binary release, or genuine corruption. Intended behavior is to treat it as a @@ -485,7 +483,7 @@ private Lazy> CreateLazyKvStore() => } // Check absolute expiration - if (IsAbsolutelyExpired(kvEntry.Value)) + if (IsAbsolutelyExpired(entry)) { // NatsKVWrongLastRevisionException is caught below var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; @@ -494,9 +492,10 @@ private Lazy> CreateLazyKvStore() => return null; } - await UpdateEntryExpirationAsync(kvEntry).ConfigureAwait(false); + await UpdateEntryExpirationAsync(kvStore, encodedKey, entry, kvEntry.Revision, token) + .ConfigureAwait(false); scope.SetHit(); - return kvEntry.Value.Data; + return entry.Data; } catch (NatsKVWrongLastRevisionException) { @@ -504,41 +503,106 @@ private Lazy> CreateLazyKvStore() => scope.SetMiss(CacheMissReason.RevisionConflict); return null; } + } + catch (Exception ex) + { + scope.SetError(ex); + throw; + } + finally + { + scope.Complete(); + } + } + + // The read core behind TryGetAsync(IBufferWriter). Mirrors GetAndRefreshAsync's telemetry, + // expiration, and revision-conflict handling exactly, but writes the payload straight into the caller's + // buffer on a hit instead of returning a byte[], eliminating the intermediate array on the common path. + // Instrumented here (not in TryGetAsync) so its swallowed failures are still recorded as errors — the + // scope closes before the exception reaches TryGetAsync's catch. + private async ValueTask TryGetAndRefreshBufferAsync( + string key, + IBufferWriter destination, + CancellationToken token) + { + var scope = NatsCacheOperationScope.Start(Telemetry, TimeProvider, CacheOperation.Get, key, token); + try + { + var encodedKey = GetEncodedKey(key); + var kvStore = await GetKvStore().ConfigureAwait(false); - // Local Functions - async Task UpdateEntryExpirationAsync(NatsKVEntry kvEntry) + // A per-read deserializer bound to this destination: on a hit with no sliding refresh it writes + // the payload directly into the buffer and reports payloadWritten; otherwise it reports enough + // for the branches below to finish without touching the buffer. See the type for the full + // "nothing written on a miss" contract. + var deserializer = new BufferWritingCacheEntryDeserializer(destination, TimeProvider); + try { - if (kvEntry.Value?.SlidingExpirationTicks == null) + var natsResult = await kvStore + .TryGetEntryAsync(encodedKey, serializer: deserializer, cancellationToken: token) + .ConfigureAwait(false); + if (!natsResult.Success) { - return; + scope.SetMiss(CacheMissReason.NotFound); + return false; } - // If we have a sliding expiration, use it as the TTL - var ttl = TimeSpan.FromTicks(kvEntry.Value.SlidingExpirationTicks.Value); + var kvEntry = natsResult.Value; + var result = kvEntry.Value; + if (result == null) + { + // Undeserializable entry: treated as a miss and left in place, identical to + // GetAndRefreshAsync (see its comment for the rolling-deploy rationale). + LogUndeserializableEntry(key); + scope.SetMiss(CacheMissReason.Undeserializable); + return false; + } - // If we also have an absolute expiration, make sure we don't exceed it - if (kvEntry.Value.AbsoluteExpiration != null) + if (result.PayloadWritten) { - var remainingTime = kvEntry.Value.AbsoluteExpiration.Value - TimeProvider.GetUtcNow(); + // Confirmed hit already streamed into the caller's buffer — the single-copy fast path. + scope.SetHit(); + return true; + } - // Use the minimum of sliding window or remaining absolute time - if (remainingTime > TimeSpan.Zero && remainingTime < ttl) - { - ttl = remainingTime; - } + if (result.AbsolutelyExpired) + { + // Nothing was written. Evict and report a miss, as in GetAndRefreshAsync. + var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; + await RemoveCoreAsync(key, natsDeleteOpts, token).ConfigureAwait(false); + scope.SetMiss(CacheMissReason.Expired); + return false; } - if (ttl > TimeSpan.Zero) + // Sliding-expiration entry: the deserializer materialized the value because refreshing the + // TTL re-writes it. Apply the same absolute-expiry, refresh, and eviction logic as + // GetAndRefreshAsync before copying the payload out, so the two read paths agree. + var entry = result.Entry; + if (IsAbsolutelyExpired(entry)) { - // Use optimistic concurrency control with the last revision - await kvStore.UpdateWithTtlAsync( - encodedKey, - kvEntry.Value, - kvEntry.Revision, - ttl, - serializer: CacheEntrySerializer, - cancellationToken: token).ConfigureAwait(false); + var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; + await RemoveCoreAsync(key, natsDeleteOpts, token).ConfigureAwait(false); + scope.SetMiss(CacheMissReason.Expired); + return false; } + + // Refresh first — a lost revision race surfaces as NatsKVWrongLastRevisionException and is + // caught below as a miss before anything is written — then copy the payload out. + await UpdateEntryExpirationAsync(kvStore, encodedKey, entry, kvEntry.Revision, token) + .ConfigureAwait(false); + if (entry.Data is { Length: > 0 } data) + { + destination.Write(data); + } + + scope.SetHit(); + return true; + } + catch (NatsKVWrongLastRevisionException) + { + // Someone else updated it during the sliding refresh; nothing was written to the buffer. + scope.SetMiss(CacheMissReason.RevisionConflict); + return false; } } catch (Exception ex) @@ -552,6 +616,49 @@ await kvStore.UpdateWithTtlAsync( } } + // Refreshes a sliding entry's TTL by re-writing it under optimistic concurrency (the last revision). + // Shared by both read cores so the array and buffer paths refresh identically. A no-op for entries + // without a sliding expiration. Extracted from GetAndRefreshAsync so there is a single implementation. + private async Task UpdateEntryExpirationAsync( + INatsKVStore kvStore, + string encodedKey, + CacheEntry entry, + ulong revision, + CancellationToken token) + { + if (entry.SlidingExpirationTicks == null) + { + return; + } + + // If we have a sliding expiration, use it as the TTL + var ttl = TimeSpan.FromTicks(entry.SlidingExpirationTicks.Value); + + // If we also have an absolute expiration, make sure we don't exceed it + if (entry.AbsoluteExpiration != null) + { + var remainingTime = entry.AbsoluteExpiration.Value - TimeProvider.GetUtcNow(); + + // Use the minimum of sliding window or remaining absolute time + if (remainingTime > TimeSpan.Zero && remainingTime < ttl) + { + ttl = remainingTime; + } + } + + if (ttl > TimeSpan.Zero) + { + // Use optimistic concurrency control with the last revision + await kvStore.UpdateWithTtlAsync( + encodedKey, + entry, + revision, + ttl, + serializer: CacheEntrySerializer, + cancellationToken: token).ConfigureAwait(false); + } + } + private async Task RemoveCoreAsync( string key, NatsKVDeleteOpts? natsKvDeleteOpts = null, diff --git a/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs b/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs new file mode 100644 index 0000000..b443722 --- /dev/null +++ b/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs @@ -0,0 +1,168 @@ +using System.Buffers; +using System.Text; +using Microsoft.Extensions.Caching.Distributed; +using NATS.Net; + +namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; + +// Behavioral coverage for the single-copy TryGetAsync(IBufferWriter) read path: a hit writes the +// exact payload into the caller's buffer, and every non-hit (miss, undeserializable, absolutely expired) +// leaves the buffer untouched — the IBufferWriter "nothing written on a miss" contract. +public class TryGetAsyncBufferTests(NatsIntegrationFixture fixture) : TestBase(fixture) +{ + // A legacy JSON envelope from a pre-binary release: the first byte is '{' (0x7B), which never matches + // the binary FormatVersion, so the entry deserializes to a miss. + private static readonly byte[] LegacyJsonEntry = + Encoding.UTF8.GetBytes("{\"absexp\":null,\"sldexp\":null,\"data\":\"AQID\"}"); + + private static CancellationToken Token => TestContext.Current.CancellationToken; + + private IBufferDistributedCache BufferCache => (IBufferDistributedCache)DistributedCache; + + [Fact] + public async Task WritesExactPayloadOnHit() + { + var key = MethodKey(); + var value = Encoding.UTF8.GetBytes($"buffer-hit-{Guid.NewGuid()}"); + await DistributedCache.SetAsync(key, value, new DistributedCacheEntryOptions(), Token); + + var destination = new ArrayBufferWriter(); + var hit = await BufferCache.TryGetAsync(key, destination, Token); + + Assert.True(hit); + Assert.Equal(value.Length, destination.WrittenCount); + Assert.Equal(value, destination.WrittenSpan.ToArray()); + } + + [Fact] + public async Task EmptyPayloadIsAHitThatWritesNothing() + { + var key = MethodKey(); + await DistributedCache.SetAsync(key, Array.Empty(), new DistributedCacheEntryOptions(), Token); + + var destination = new ArrayBufferWriter(); + var hit = await BufferCache.TryGetAsync(key, destination, Token); + + // An empty stored value is a genuine hit (true), and correctly writes zero bytes. + Assert.True(hit); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public async Task WritesNothingOnMiss() + { + var key = MethodKey(); + await DistributedCache.RemoveAsync(key, Token); // known-empty state + + var destination = new ArrayBufferWriter(); + var hit = await BufferCache.TryGetAsync(key, destination, Token); + + Assert.False(hit); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public async Task WritesNothingOnUndeserializableEntry() + { + var key = MethodKey(); + await WriteRawEntryAsync(key, LegacyJsonEntry); + + var destination = new ArrayBufferWriter(); + var hit = await BufferCache.TryGetAsync(key, destination, Token); + + // Undeserializable framing is a miss and must not leak any bytes into the caller's buffer. + Assert.False(hit); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public async Task WritesNothingOnAbsolutelyExpiredEntry() + { + var key = MethodKey(); + var value = new byte[] { 1, 2, 3, 4 }; + await DistributedCache.SetAsync( + key, + value, + new DistributedCacheEntryOptions().SetAbsoluteExpiration(TimeSpan.FromSeconds(1.1)), + Token); + + // Poll the buffer path until the entry lapses; the final read must be a miss with nothing written. + ArrayBufferWriter destination; + bool hit; + var attempts = 0; + do + { + await Task.Delay(TimeSpan.FromSeconds(0.5), Token); + destination = new ArrayBufferWriter(); + hit = await BufferCache.TryGetAsync(key, destination, Token); + } + while (hit && ++attempts < 6); + + Assert.False(hit); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public async Task SlidingExpirationHitWritesPayloadAcrossRepeatedReads() + { + var key = MethodKey(); + var value = Encoding.UTF8.GetBytes($"sliding-{Guid.NewGuid()}"); + await DistributedCache.SetAsync( + key, + value, + new DistributedCacheEntryOptions().SetSlidingExpiration(TimeSpan.FromSeconds(2)), + Token); + + // The sliding branch materializes the entry, refreshes its TTL, and copies the payload out. Read it + // twice (renewing each time) and confirm both hits return the exact bytes. + for (var i = 0; i < 2; i++) + { + var destination = new ArrayBufferWriter(); + var hit = await BufferCache.TryGetAsync(key, destination, Token); + + Assert.True(hit); + Assert.Equal(value, destination.WrittenSpan.ToArray()); + await Task.Delay(TimeSpan.FromSeconds(0.5), Token); + } + } + + [Fact] + public async Task SlidingEntryPastAbsoluteExpirationWritesNothing() + { + var key = MethodKey(); + var value = new byte[] { 9, 8, 7 }; + + // Sliding renews on access, but the absolute expiration is the hard ceiling: once it passes, the + // buffer read must evict and miss, exercising the absolute-expiry branch of the sliding read path. + await DistributedCache.SetAsync( + key, + value, + new DistributedCacheEntryOptions() + .SetSlidingExpiration(TimeSpan.FromSeconds(1.1)) + .SetAbsoluteExpiration(TimeSpan.FromSeconds(2)), + Token); + + ArrayBufferWriter destination; + bool hit; + var attempts = 0; + do + { + await Task.Delay(TimeSpan.FromSeconds(0.5), Token); + destination = new ArrayBufferWriter(); + hit = await BufferCache.TryGetAsync(key, destination, Token); + } + while (hit && ++attempts < 10); + + Assert.False(hit); + Assert.Equal(0, destination.WrittenCount); + } + + // Writes raw bytes at the key the cache reads, bypassing the binary serializer so the stored entry + // cannot be deserialized. + private async Task WriteRawEntryAsync(string key, byte[] raw) + { + var encodedKey = new NatsCacheKeyEncoder().Encode(key); + var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache"); + await kvStore.PutAsync(encodedKey, raw, cancellationToken: Token); + } +} From b902d12e797bf490e134e13f425751f2f4a85e4d Mon Sep 17 00:00:00 2001 From: Matthew DeVenny Date: Wed, 29 Jul 2026 12:26:59 -0700 Subject: [PATCH 2/2] #45 Address review: surface destination-writer failures; unify read core Requested fix: a failing destination write is no longer reported as an undeserializable miss. Moving destination.Write inside the deserializer put it where NATS swallows exceptions; now a writer failure (e.g. HybridCache's payload quota) is caught at the write site and carried out as a DestinationFailure outcome, which the read core rethrows so TryGetAsync logs a Warning and records an error -- as on main -- instead of a Debug undeserializable miss. Also, per review: - Unify the array and buffer read paths into one GetAndRefreshAsync, removing ~50 lines kept in sync by comments; the deserializer materializes for the array/sliding paths and streams for the buffer fast path. - Replace the two-bool buffer result with a CacheEntryReadOutcome enum; the hit and expired paths allocate nothing beyond the payload. - Share the absolute-expiry predicate as a static so the paths cannot drift. - Log the KV entry's error to separate corrupt framing from a legacy envelope. - Make the absolute-expiry tests deterministic with FakeTimeProvider; add unit tests for CacheEntryReadDeserializer (incl. multi-segment and the boundary) and buffer-path telemetry (destination failure, expired). - Move LegacyJsonEntry/WriteRawEntryAsync to TestBase; drop the redundant IsSingleSegment branch; use ASCII arrows; fix the "zero-copy" wording. Co-Authored-By: Claude Opus 4.8 (1M context) Signed-off-by: Matthew DeVenny --- .../BufferWritingCacheEntryDeserializer.cs | 145 ------------- .../CacheEntryReadDeserializer.cs | 172 +++++++++++++++ src/NatsDistributedCache/NatsCache.Log.cs | 3 +- src/NatsDistributedCache/NatsCache.cs | 204 ++++++++---------- test/IntegrationTests/Cache/TelemetryTests.cs | 81 +++++-- .../Cache/TryGetAsyncBufferTests.cs | 113 +++++----- .../Cache/UndeserializableEntryTests.cs | 14 -- test/IntegrationTests/TestBase.cs | 18 ++ test/TestUtils/QuotaBufferWriter.cs | 35 +++ .../CacheEntryReadDeserializerTests.cs | 167 ++++++++++++++ 10 files changed, 617 insertions(+), 335 deletions(-) delete mode 100644 src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs create mode 100644 src/NatsDistributedCache/CacheEntryReadDeserializer.cs create mode 100644 test/TestUtils/QuotaBufferWriter.cs create mode 100644 test/UnitTests/Serialization/CacheEntryReadDeserializerTests.cs diff --git a/src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs b/src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs deleted file mode 100644 index 6cfced8..0000000 --- a/src/NatsDistributedCache/BufferWritingCacheEntryDeserializer.cs +++ /dev/null @@ -1,145 +0,0 @@ -using System.Buffers; -using NATS.Client.Core; - -namespace CodeCargo.Nats.DistributedCache; - -/// -/// Outcome of a single-copy read produced by -/// . Carries just enough for the read core to finish the -/// operation without a second copy of the payload on the common path. A result (the -/// deserializer returning null) means the stored bytes were undeserializable — exactly as -/// signals the same condition with a null entry. -/// -internal sealed class CacheEntryBufferReadResult -{ - internal CacheEntryBufferReadResult(CacheEntry entry, bool payloadWritten, bool absolutelyExpired) - { - Entry = entry; - PayloadWritten = payloadWritten; - AbsolutelyExpired = absolutelyExpired; - } - - /// - /// The decoded entry. Its is populated only when the payload was - /// not streamed to the destination (the sliding-expiration case, where the read core must - /// re-write the value to refresh its TTL); otherwise it is left null because the bytes already live in - /// the caller's buffer. - /// - internal CacheEntry Entry { get; } - - /// - /// when the payload was written straight into the caller's - /// — a confirmed hit with no sliding refresh, and the single-copy fast - /// path. The read core only needs to record the hit. - /// - internal bool PayloadWritten { get; } - - /// - /// when the entry is past its absolute expiration. Nothing was written; the read - /// core evicts the entry (using the KV revision) and reports a miss. - /// - internal bool AbsolutelyExpired { get; } -} - -/// -/// A per-read used only by NatsCache.TryGetAsync(IBufferWriter<byte>) -/// to avoid the intermediate array that the array read path allocates. It is the only -/// point in the read that holds the transport payload sequence, so it writes the payload directly into the -/// caller's destination — a single copy (transport buffer → caller buffer) instead of two (transport buffer → -/// → caller buffer). -/// -/// -/// The write only happens for a genuine hit, so the contract of "nothing -/// written on a miss" is preserved: -/// -/// undeserializable framing → returns , writes nothing; -/// absolutely expired → returns a result flagged expired, writes nothing (the core evicts + misses); -/// sliding expiration present → materializes and writes nothing, because -/// the core must re-write the whole value to refresh its TTL and therefore needs the bytes (this is the one -/// case that still costs two copies — the common absolute-only/no-expiry hit does not); -/// otherwise → writes the payload to the destination and flags the hit. -/// -/// Absolute expiry is evaluated here rather than in the core because this is the only moment the payload is -/// available, and the buffer must stay untouched when the entry turns out to be an expired miss. It uses -/// at deserialize time — the same clock and effectively the same instant -/// the core would have used immediately afterward — so the decision matches NatsCache.IsAbsolutelyExpired. -/// -internal sealed class BufferWritingCacheEntryDeserializer : INatsDeserialize -{ - private readonly IBufferWriter _destination; - private readonly TimeProvider _timeProvider; - - internal BufferWritingCacheEntryDeserializer(IBufferWriter destination, TimeProvider timeProvider) - { - _destination = destination; - _timeProvider = timeProvider; - } - - /// - public CacheEntryBufferReadResult? Deserialize(in ReadOnlySequence buffer) - { - var reader = new SequenceReader(buffer); - if (!CacheEntryBinarySerializer.TryReadHeader(ref reader, out var absoluteExpiration, out var slidingExpirationTicks)) - { - // Unknown/legacy/corrupt framing: signal an undeserializable entry (see the type doc). - return null; - } - - var remaining = reader.UnreadSequence; - - if (slidingExpirationTicks.HasValue) - { - // Sliding entries have their TTL refreshed by re-writing the value, which needs the payload - // bytes — so this path cannot stream to the caller. Materialize the full entry (one copy) and - // let the read core run its normal expiry/refresh/write, exactly like the array read path. - var data = remaining.IsEmpty ? Array.Empty() : remaining.ToArray(); - return new CacheEntryBufferReadResult( - new CacheEntry - { - AbsoluteExpiration = absoluteExpiration, - SlidingExpirationTicks = slidingExpirationTicks, - Data = data, - }, - payloadWritten: false, - absolutelyExpired: false); - } - - if (absoluteExpiration.HasValue && _timeProvider.GetUtcNow() >= absoluteExpiration.Value) - { - // Absolutely expired: write nothing. The read core evicts the entry and reports a miss. - return new CacheEntryBufferReadResult( - new CacheEntry { AbsoluteExpiration = absoluteExpiration }, - payloadWritten: false, - absolutelyExpired: true); - } - - // Genuine hit with no sliding refresh: copy the payload straight into the caller's buffer. This is - // the single copy that makes TryGet allocation-free of any intermediate array on the common path. - WritePayload(remaining); - return new CacheEntryBufferReadResult( - new CacheEntry { AbsoluteExpiration = absoluteExpiration }, - payloadWritten: true, - absolutelyExpired: false); - } - - private void WritePayload(in ReadOnlySequence payload) - { - if (payload.IsSingleSegment) - { - if (!payload.First.IsEmpty) - { - _destination.Write(payload.First.Span); - } - - return; - } - - foreach (var segment in payload) - { - if (!segment.IsEmpty) - { - _destination.Write(segment.Span); - } - } - } -} diff --git a/src/NatsDistributedCache/CacheEntryReadDeserializer.cs b/src/NatsDistributedCache/CacheEntryReadDeserializer.cs new file mode 100644 index 0000000..1eaa013 --- /dev/null +++ b/src/NatsDistributedCache/CacheEntryReadDeserializer.cs @@ -0,0 +1,172 @@ +using System.Buffers; +using NATS.Client.Core; + +namespace CodeCargo.Nats.DistributedCache; + +/// +/// What a single read produced, as seen by NatsCache's read core. A +/// result (the deserializer returning null) means the stored bytes were +/// undeserializable — the same way signals it with a +/// null entry. +/// +internal enum CacheEntryReadOutcome +{ + /// + /// The payload was materialized into 's . + /// The read core checks absolute expiry, refreshes a sliding TTL, then emits the payload — returning it for + /// the array path, or writing it to the destination for a sliding buffer read. Used for every array read + /// and for sliding entries on the buffer path. + /// + Materialized, + + /// + /// Buffer fast path: the payload was written straight into the caller's . + /// A confirmed hit; the core only records it. + /// + PayloadWritten, + + /// + /// Buffer fast path: the entry is past its absolute expiration. Nothing was written; the core evicts the + /// entry and reports a miss. + /// + AbsolutelyExpired, + + /// + /// The caller's threw while the payload was being written (e.g. HybridCache's + /// payload quota). Not corrupt data: the exception is carried out so the core surfaces an error instead of + /// an undeserializable miss. + /// + DestinationFailure, +} + +/// +/// Outcome of a single read produced by . The reusable +/// and singletons let the common buffer hit and +/// the expired path allocate nothing beyond the payload itself; only the materialized and destination-failure +/// outcomes carry per-read state. +/// +internal sealed class CacheEntryReadResult +{ + internal static readonly CacheEntryReadResult PayloadWritten = + new(CacheEntryReadOutcome.PayloadWritten, entry: null, destinationFailure: null); + + internal static readonly CacheEntryReadResult AbsolutelyExpired = + new(CacheEntryReadOutcome.AbsolutelyExpired, entry: null, destinationFailure: null); + + private CacheEntryReadResult(CacheEntryReadOutcome outcome, CacheEntry? entry, Exception? destinationFailure) + { + Outcome = outcome; + Entry = entry; + DestinationFailure = destinationFailure; + } + + internal CacheEntryReadOutcome Outcome { get; } + + /// The decoded entry, non-null only for . + internal CacheEntry? Entry { get; } + + /// The writer exception, non-null only for . + internal Exception? DestinationFailure { get; } + + internal static CacheEntryReadResult Materialized(CacheEntry entry) => + new(CacheEntryReadOutcome.Materialized, entry, destinationFailure: null); + + internal static CacheEntryReadResult DestinationFailed(Exception exception) => + new(CacheEntryReadOutcome.DestinationFailure, entry: null, exception); +} + +/// +/// The per-read used by NatsCache's unified read core. +/// +/// +/// With a destination it materializes the entry for the array read path (Get/Refresh). +/// With a destination it powers the single-copy TryGetAsync(IBufferWriter<byte>) path: on a genuine +/// hit with no sliding refresh it writes the payload straight into the caller's writer, so the read costs one +/// copy (transport buffer -> caller buffer) instead of two. Nothing is written for a miss (undeserializable +/// framing or absolute expiry) or when the entry needs its sliding TTL refreshed — that re-writes the value and +/// so needs the bytes, which are materialized for the core to handle. If the caller's writer throws mid-write +/// the exception is carried out as so the core surfaces +/// an error rather than an undeserializable miss; on a multi-segment payload any segments already written stay +/// committed, so the buffer contents are unspecified when a write fails. +/// +internal sealed class CacheEntryReadDeserializer : INatsDeserialize +{ + private readonly IBufferWriter? _destination; + private readonly TimeProvider _timeProvider; + + internal CacheEntryReadDeserializer(IBufferWriter? destination, TimeProvider timeProvider) + { + _destination = destination; + _timeProvider = timeProvider; + } + + /// + public CacheEntryReadResult? Deserialize(in ReadOnlySequence buffer) + { + // Array read path: no destination to stream into, so materialize the whole entry via the canonical + // deserializer and let the read core decide expiry/refresh/eviction. + if (_destination is null) + { + var entry = CacheEntryBinarySerializer.Default.Deserialize(buffer); + return entry is null ? null : CacheEntryReadResult.Materialized(entry); + } + + var reader = new SequenceReader(buffer); + if (!CacheEntryBinarySerializer.TryReadHeader(ref reader, out var absoluteExpiration, out var slidingExpirationTicks)) + { + // Unknown/legacy/corrupt framing: signal an undeserializable entry. + return null; + } + + var remaining = reader.UnreadSequence; + + if (slidingExpirationTicks.HasValue) + { + // Sliding entries have their TTL refreshed by re-writing the value, which needs the payload + // bytes, so this path cannot stream to the caller. Materialize and let the core handle it, + // exactly like the array path. + var data = remaining.IsEmpty ? Array.Empty() : remaining.ToArray(); + return CacheEntryReadResult.Materialized(new CacheEntry + { + AbsoluteExpiration = absoluteExpiration, + SlidingExpirationTicks = slidingExpirationTicks, + Data = data, + }); + } + + if (NatsCache.IsAbsolutelyExpired(absoluteExpiration, _timeProvider.GetUtcNow())) + { + // Absolutely expired: write nothing. The core evicts and reports a miss. Expiry is evaluated + // here (the only point holding the payload) so the buffer stays untouched for an expired miss, + // using the same clock and predicate the core would. + return CacheEntryReadResult.AbsolutelyExpired; + } + + // Genuine hit with no sliding refresh: copy the payload straight into the caller's buffer — the + // single copy that makes TryGet free of any intermediate array on the common path. A writer failure + // (e.g. HybridCache's payload quota) is not corrupt data, so carry it out for the core to surface. + try + { + WritePayload(remaining); + } + catch (Exception ex) + { + return CacheEntryReadResult.DestinationFailed(ex); + } + + return CacheEntryReadResult.PayloadWritten; + } + + private void WritePayload(in ReadOnlySequence payload) + { + // The enumerator is a struct, and it yields the single segment for a single-segment sequence, so this + // covers both without a special case. + foreach (var segment in payload) + { + if (!segment.IsEmpty) + { + _destination!.Write(segment.Span); + } + } + } +} diff --git a/src/NatsDistributedCache/NatsCache.Log.cs b/src/NatsDistributedCache/NatsCache.Log.cs index 5955068..59fbfb9 100644 --- a/src/NatsDistributedCache/NatsCache.Log.cs +++ b/src/NatsDistributedCache/NatsCache.Log.cs @@ -13,9 +13,10 @@ private void LogException(Exception exception) => private void LogSwallowedException(Exception exception) => _logger.LogWarning(EventIds.Exception, exception, "NATS cache read failed in TryGetAsync; returning a cache miss"); - private void LogUndeserializableEntry(string key) => + private void LogUndeserializableEntry(string key, Exception? error) => _logger.LogDebug( EventIds.UndeserializableEntry, + error, "Cache entry for key {Key} could not be deserialized (legacy or corrupt format); returning a cache miss", key); diff --git a/src/NatsDistributedCache/NatsCache.cs b/src/NatsDistributedCache/NatsCache.cs index 6c57d27..8ff4a0a 100644 --- a/src/NatsDistributedCache/NatsCache.cs +++ b/src/NatsDistributedCache/NatsCache.cs @@ -1,5 +1,6 @@ using System.Buffers; using System.Diagnostics.Metrics; +using System.Runtime.ExceptionServices; using Microsoft.Extensions.Caching.Distributed; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; @@ -181,7 +182,7 @@ public async Task RefreshAsync(string key, CancellationToken token = default) { try { - await GetAndRefreshAsync(key, CacheOperation.Refresh, token).ConfigureAwait(false); + await GetAndRefreshAsync(key, CacheOperation.Refresh, destination: null, token).ConfigureAwait(false); } catch (Exception ex) { @@ -198,7 +199,8 @@ public async Task RefreshAsync(string key, CancellationToken token = default) { try { - return await GetAndRefreshAsync(key, CacheOperation.Get, token).ConfigureAwait(false); + return (await GetAndRefreshAsync(key, CacheOperation.Get, destination: null, token) + .ConfigureAwait(false)).Data; } catch (Exception ex) { @@ -221,10 +223,11 @@ public async ValueTask TryGetAsync( { // Writes the payload straight into the caller's buffer on a hit (a single copy), rather than // fetching a byte[] and copying it a second time. Reports as operation=get, not a distinct - // operation: the IBufferWriter overload is a zero-copy detail, and merging it keeps hit ratio + // operation: the IBufferWriter overload is a single-copy detail, and merging it keeps hit ratio // computed over all read paths — which matters because HybridCache drives its L2 reads // exclusively through this method. - return await TryGetAndRefreshBufferAsync(key, destination, token).ConfigureAwait(false); + return (await GetAndRefreshAsync(key, CacheOperation.Get, destination, token) + .ConfigureAwait(false)).Hit; } catch (OperationCanceledException) when (token.IsCancellationRequested) { @@ -243,6 +246,12 @@ public async ValueTask TryGetAsync( return false; } + // The absolute-expiry predicate over the raw instant, shared by the instance IsAbsolutelyExpired below + // and CacheEntryReadDeserializer's buffer fast path so the two read paths apply the identical rule + // (inclusive >= boundary, matching GetTtl) and cannot drift. + internal static bool IsAbsolutelyExpired(DateTimeOffset? absoluteExpiration, DateTimeOffset utcNow) => + absoluteExpiration.HasValue && utcNow >= absoluteExpiration.Value; + internal TimeSpan? GetTtl(DistributedCacheEntryOptions options) { // Maximum-value sentinels mean "never expire": normalize them to no expiration so they are not @@ -342,7 +351,7 @@ internal CacheEntry CreateCacheEntry(byte[] value, DistributedCacheEntryOptions // boundary is inclusive (>=) to match GetTtl, which treats an absolute expiration at "now" as // already elapsed. Sliding expiration is enforced separately via the NATS entry TTL. internal bool IsAbsolutelyExpired(CacheEntry entry) => - entry.AbsoluteExpiration.HasValue && TimeProvider.GetUtcNow() >= entry.AbsoluteExpiration.Value; + IsAbsolutelyExpired(entry.AbsoluteExpiration, TimeProvider.GetUtcNow()); // Builds the NatsKVConfig used when CreateBucketIfNotExists is enabled. Pure and synchronous (does not // touch the NATS connection), so it is unit-testable without a server. Cache-appropriate defaults are @@ -442,10 +451,20 @@ private Lazy> CreateLazyKvStore() => private Task GetKvStore() => _lazyKvStore.Value; // The shared read core for Get, TryGet, and Refresh (and their sync overloads). Instrumented here - // rather than in the three public methods so that TryGetAsync's swallowed failures are still recorded - // as errors — the scope closes before the exception reaches TryGetAsync's catch — and so no read path - // can be counted twice. - private async Task GetAndRefreshAsync(string key, CacheOperation operation, CancellationToken token) + // rather than in the public methods so that TryGetAsync's swallowed failures are still recorded as + // errors — the scope closes before the exception reaches TryGetAsync's catch — and so no read path can + // be counted twice. + // + // A null destination is the array path (Get/Refresh): the payload is materialized and returned as Data. + // A non-null destination is the single-copy TryGetAsync(IBufferWriter) path: on the common hit the + // payload is written straight into it and Data stays null. Both share every step below — store + // resolution, the undeserializable and absolute-expiry misses, the sliding-TTL refresh, and + // revision-conflict handling — so the two overloads cannot drift. + private async Task<(bool Hit, byte[]? Data)> GetAndRefreshAsync( + string key, + CacheOperation operation, + IBufferWriter? destination, + CancellationToken token) { var scope = NatsCacheOperationScope.Start(Telemetry, TimeProvider, operation, key, token); try @@ -454,155 +473,105 @@ private Lazy> CreateLazyKvStore() => // as errors too. This does not change which exceptions propagate. var encodedKey = GetEncodedKey(key); var kvStore = await GetKvStore().ConfigureAwait(false); + + // One deserializer for both overloads, switched on `destination`: with none it materializes the + // entry for the array path; with one it streams the payload into the caller's writer on the + // common hit and otherwise reports enough for the branches below to finish without touching the + // buffer (see the type for the full "nothing written on a miss" contract). + var deserializer = new CacheEntryReadDeserializer(destination, TimeProvider); try { var natsResult = await kvStore - .TryGetEntryAsync(encodedKey, serializer: CacheEntrySerializer, cancellationToken: token) + .TryGetEntryAsync(encodedKey, serializer: deserializer, cancellationToken: token) .ConfigureAwait(false); if (!natsResult.Success) { scope.SetMiss(CacheMissReason.NotFound); - return null; + return (false, null); } var kvEntry = natsResult.Value; - var entry = kvEntry.Value; - if (entry == null) + var result = kvEntry.Value; + if (result is null) { // Present entry whose bytes we cannot deserialize: a legacy JSON envelope from a - // pre-binary release, or genuine corruption. Intended behavior is to treat it as a - // cache miss and leave the entry in place. It self-heals when the key is next written - // (Set overwrites unconditionally), and any TTL'd entry is reaped by NATS. We - // deliberately do not evict it (an older node must not delete entries written in a - // newer format during a rolling deploy) nor throw (a cache should degrade to a miss, - // not fail the caller's operation). Logged at Debug to aid diagnosis without flooding - // logs during a JSON->binary migration, when every legacy key transiently lands here. - LogUndeserializableEntry(key); + // pre-binary release, or genuine corruption. Intended behavior is to treat it as a cache + // miss and leave the entry in place. It self-heals when the key is next written (Set + // overwrites unconditionally), and any TTL'd entry is reaped by NATS. We deliberately do + // not evict it (an older node must not delete entries written in a newer format during a + // rolling deploy) nor throw (a cache should degrade to a miss, not fail the caller's + // operation). Logged at Debug — with the KV entry's own error, which separates corrupt + // framing from a legacy envelope — to aid diagnosis without flooding logs during a + // JSON->binary migration, when every legacy key transiently lands here. + LogUndeserializableEntry(key, kvEntry.Error); scope.SetMiss(CacheMissReason.Undeserializable); - return null; + return (false, null); } - // Check absolute expiration - if (IsAbsolutelyExpired(entry)) + if (result.Outcome == CacheEntryReadOutcome.DestinationFailure) { - // NatsKVWrongLastRevisionException is caught below - var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; - await RemoveCoreAsync(key, natsDeleteOpts, token).ConfigureAwait(false); - scope.SetMiss(CacheMissReason.Expired); - return null; + // The caller's IBufferWriter threw while the payload was being written (e.g. HybridCache's + // payload quota). That is not corrupt data, so surface it: rethrowing reaches the outer + // catch (scope records error) and then TryGetAsync's catch (Warning with the exception, + // false to the caller), rather than masquerading as an undeserializable miss. + // ExceptionDispatchInfo preserves the original stack. + ExceptionDispatchInfo.Capture(result.DestinationFailure!).Throw(); } - await UpdateEntryExpirationAsync(kvStore, encodedKey, entry, kvEntry.Revision, token) - .ConfigureAwait(false); - scope.SetHit(); - return entry.Data; - } - catch (NatsKVWrongLastRevisionException) - { - // Someone else updated it; that's fine, we'll get the latest version next time - scope.SetMiss(CacheMissReason.RevisionConflict); - return null; - } - } - catch (Exception ex) - { - scope.SetError(ex); - throw; - } - finally - { - scope.Complete(); - } - } - - // The read core behind TryGetAsync(IBufferWriter). Mirrors GetAndRefreshAsync's telemetry, - // expiration, and revision-conflict handling exactly, but writes the payload straight into the caller's - // buffer on a hit instead of returning a byte[], eliminating the intermediate array on the common path. - // Instrumented here (not in TryGetAsync) so its swallowed failures are still recorded as errors — the - // scope closes before the exception reaches TryGetAsync's catch. - private async ValueTask TryGetAndRefreshBufferAsync( - string key, - IBufferWriter destination, - CancellationToken token) - { - var scope = NatsCacheOperationScope.Start(Telemetry, TimeProvider, CacheOperation.Get, key, token); - try - { - var encodedKey = GetEncodedKey(key); - var kvStore = await GetKvStore().ConfigureAwait(false); - - // A per-read deserializer bound to this destination: on a hit with no sliding refresh it writes - // the payload directly into the buffer and reports payloadWritten; otherwise it reports enough - // for the branches below to finish without touching the buffer. See the type for the full - // "nothing written on a miss" contract. - var deserializer = new BufferWritingCacheEntryDeserializer(destination, TimeProvider); - try - { - var natsResult = await kvStore - .TryGetEntryAsync(encodedKey, serializer: deserializer, cancellationToken: token) - .ConfigureAwait(false); - if (!natsResult.Success) + if (result.Outcome == CacheEntryReadOutcome.PayloadWritten) { - scope.SetMiss(CacheMissReason.NotFound); - return false; - } - - var kvEntry = natsResult.Value; - var result = kvEntry.Value; - if (result == null) - { - // Undeserializable entry: treated as a miss and left in place, identical to - // GetAndRefreshAsync (see its comment for the rolling-deploy rationale). - LogUndeserializableEntry(key); - scope.SetMiss(CacheMissReason.Undeserializable); - return false; - } - - if (result.PayloadWritten) - { - // Confirmed hit already streamed into the caller's buffer — the single-copy fast path. + // Buffer fast path: a confirmed hit already streamed into the caller's writer. scope.SetHit(); - return true; + return (true, null); } - if (result.AbsolutelyExpired) + if (result.Outcome == CacheEntryReadOutcome.AbsolutelyExpired) { - // Nothing was written. Evict and report a miss, as in GetAndRefreshAsync. - var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; - await RemoveCoreAsync(key, natsDeleteOpts, token).ConfigureAwait(false); + // Buffer fast path: nothing was written. Evict and report a miss. + await EvictExpiredAsync(key, kvEntry.Revision, token).ConfigureAwait(false); scope.SetMiss(CacheMissReason.Expired); - return false; + return (false, null); } - // Sliding-expiration entry: the deserializer materialized the value because refreshing the - // TTL re-writes it. Apply the same absolute-expiry, refresh, and eviction logic as - // GetAndRefreshAsync before copying the payload out, so the two read paths agree. - var entry = result.Entry; + // Materialized: an array read, or a sliding entry on the buffer path. The value is in hand, + // so apply absolute expiry, the sliding refresh, and eviction uniformly. + var entry = result.Entry!; if (IsAbsolutelyExpired(entry)) { - var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; - await RemoveCoreAsync(key, natsDeleteOpts, token).ConfigureAwait(false); + // NatsKVWrongLastRevisionException is caught below. + await EvictExpiredAsync(key, kvEntry.Revision, token).ConfigureAwait(false); scope.SetMiss(CacheMissReason.Expired); - return false; + return (false, null); } // Refresh first — a lost revision race surfaces as NatsKVWrongLastRevisionException and is - // caught below as a miss before anything is written — then copy the payload out. + // caught below as a miss before anything is emitted. await UpdateEntryExpirationAsync(kvStore, encodedKey, entry, kvEntry.Revision, token) .ConfigureAwait(false); + + if (destination is null) + { + scope.SetHit(); + return (true, entry.Data); + } + + // Sliding buffer hit: copy the materialized payload into the caller's writer. A writer + // failure here propagates to the outer catch and is reported as an error, matching the fast + // path's DestinationFailure handling above. if (entry.Data is { Length: > 0 } data) { destination.Write(data); } scope.SetHit(); - return true; + return (true, null); } catch (NatsKVWrongLastRevisionException) { - // Someone else updated it during the sliding refresh; nothing was written to the buffer. + // Someone else updated it (during the sliding refresh); we'll get the latest next time. + // Nothing was emitted. scope.SetMiss(CacheMissReason.RevisionConflict); - return false; + return (false, null); } } catch (Exception ex) @@ -616,6 +585,13 @@ await UpdateEntryExpirationAsync(kvStore, encodedKey, entry, kvEntry.Revision, t } } + // Evicts an absolutely-expired entry using optimistic concurrency on its last revision. A concurrent + // writer surfaces as NatsKVWrongLastRevisionException, which the read core treats as a revision-conflict + // miss. Deliberately routed through RemoveCoreAsync, which is not instrumented, so an expired read does + // not emit a phantom operation=remove. + private Task EvictExpiredAsync(string key, ulong revision, CancellationToken token) => + RemoveCoreAsync(key, new NatsKVDeleteOpts { Revision = revision }, token); + // Refreshes a sliding entry's TTL by re-writing it under optimistic concurrency (the last revision). // Shared by both read cores so the array and buffer paths refresh identically. A no-op for entries // without a sliding expiration. Extracted from GetAndRefreshAsync so there is a single implementation. diff --git a/test/IntegrationTests/Cache/TelemetryTests.cs b/test/IntegrationTests/Cache/TelemetryTests.cs index 58aeca3..c0fb3b0 100644 --- a/test/IntegrationTests/Cache/TelemetryTests.cs +++ b/test/IntegrationTests/Cache/TelemetryTests.cs @@ -1,13 +1,13 @@ using System.Buffers; using System.Diagnostics; using System.Diagnostics.Metrics; -using System.Text; using CodeCargo.Nats.DistributedCache.TestUtils; using CodeCargo.Nats.DistributedCache.TestUtils.Services.Diagnostics; using CodeCargo.Nats.DistributedCache.TestUtils.Services.Logging; using Microsoft.Extensions.Caching.Distributed; using Microsoft.Extensions.Caching.Hybrid; using Microsoft.Extensions.Diagnostics.Metrics.Testing; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Microsoft.Extensions.Time.Testing; using NATS.Net; @@ -16,11 +16,6 @@ namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; public class TelemetryTests : TestBase { - // A legacy JSON envelope from a pre-binary release: the first byte is '{' (0x7B), which never matches - // the binary FormatVersion, so the serializer returns null and the read is an undeserializable miss. - private static readonly byte[] LegacyJsonEntry = - Encoding.UTF8.GetBytes("{\"absexp\":null,\"sldexp\":null,\"data\":\"AQID\"}"); - // Held explicitly rather than captured from a primary constructor parameter: capturing a value that is // also passed to the base constructor warns under CS9107, and CI builds with warnings as errors. private readonly NatsIntegrationFixture _fixture; @@ -137,6 +132,73 @@ await bufferCache.SetAsync( Assert.Equal([("set", "ok"), ("get", "hit")], results); } + [Fact] + public async Task BufferReadWithFailingDestinationTagsErrorNotUndeserializableMiss() + { + var key = MethodKey(); + var token = TestContext.Current.CancellationToken; + + // Directly constructed so its logs are captured; it publishes on the process-wide fallback meter. + var logger = new RecordingLogger(); + var cache = new NatsCache(Options.Create(new NatsCacheOptions { BucketName = "cache" }), NatsConnection, logger); + await cache.SetAsync(key, [1, 2, 3, 4], new DistributedCacheEntryOptions(), token); + + using var duration = CreateFallbackDurationCollector(); + using var misses = CreateFallbackMissesCollector(); + + // The caller's writer rejects any payload, exactly as HybridCache's quota-limited writer does past + // its limit. The failure comes from the destination, not the stored bytes, so it must surface as an + // error rather than an undeserializable miss -- otherwise a burst of oversized entries would be + // indistinguishable from a stalled JSON->binary migration in the miss-reason signal. + var destination = new QuotaBufferWriter(maxLength: 0); + Assert.False(await ((IBufferDistributedCache)cache).TryGetAsync(key, destination, token)); + Assert.Equal(0, destination.WrittenCount); + + // Telemetry records error, not a miss. + var measurement = Assert.Single(duration.GetMeasurementSnapshot()); + Assert.Equal("get", measurement.Tags["nats.cache.operation"]); + Assert.Equal("error", measurement.Tags["nats.cache.result"]); + Assert.NotNull(measurement.Tags["error.type"]); + Assert.Empty(misses.GetMeasurementSnapshot()); + + // And it is logged once at Warning with the exception -- like any swallowed TryGet failure -- not at + // Debug as an undeserializable entry. + var record = Assert.Single(logger.Records, r => r.EventId.Name == "Exception"); + Assert.Equal(LogLevel.Warning, record.LogLevel); + Assert.NotNull(record.Exception); + } + + [Fact] + public async Task BufferReadOfExpiredEntryTagsMissReasonExpired() + { + var key = MethodKey(); + var timeProvider = new FakeTimeProvider(); + await using var provider = BuildProvider(timeProvider); + var cache = (IBufferDistributedCache)provider.GetRequiredService(); + var meterFactory = provider.GetRequiredService(); + var token = TestContext.Current.CancellationToken; + + // Five-minute absolute expiration -> five-minute real NATS TTL; the cache's own clock then jumps past + // it, so the entry is still present but absolutely expired -- the only deterministic way to reach the + // expired branch, here through the buffer overload rather than GetAsync. + await cache.SetAsync( + key, + new ReadOnlySequence(new byte[] { 1, 2, 3 }), + new DistributedCacheEntryOptions().SetAbsoluteExpiration(TimeSpan.FromMinutes(5)), + token); + + using var misses = CreateMissesCollector(meterFactory); + timeProvider.Advance(TimeSpan.FromMinutes(10)); + + var destination = new ArrayBufferWriter(); + Assert.False(await cache.TryGetAsync(key, destination, token)); + Assert.Equal(0, destination.WrittenCount); + + var miss = Assert.Single(misses.GetMeasurementSnapshot()); + Assert.Equal("get", miss.Tags["nats.cache.operation"]); + Assert.Equal("expired", miss.Tags["nats.cache.miss.reason"]); + } + [Fact] public void SyncOverloadsRecordExactlyOneMeasurementEach() { @@ -385,11 +447,4 @@ private NatsCache CreateFailingCache() => Options.Create(new NatsCacheOptions { BucketName = "does-not-exist" }), NatsConnection, new RecordingLogger()); - - private async Task WriteRawEntryAsync(string key, byte[] raw) - { - var encodedKey = new NatsCacheKeyEncoder().Encode(key); - var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache"); - await kvStore.PutAsync(encodedKey, raw, cancellationToken: TestContext.Current.CancellationToken); - } } diff --git a/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs b/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs index b443722..e53ca90 100644 --- a/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs +++ b/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs @@ -1,19 +1,24 @@ using System.Buffers; using System.Text; +using CodeCargo.Nats.DistributedCache.TestUtils; using Microsoft.Extensions.Caching.Distributed; -using NATS.Net; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Time.Testing; namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; -// Behavioral coverage for the single-copy TryGetAsync(IBufferWriter) read path: a hit writes the -// exact payload into the caller's buffer, and every non-hit (miss, undeserializable, absolutely expired) -// leaves the buffer untouched — the IBufferWriter "nothing written on a miss" contract. -public class TryGetAsyncBufferTests(NatsIntegrationFixture fixture) : TestBase(fixture) +// Behavioral coverage for the single-copy TryGetAsync(IBufferWriter) read path: a hit writes the exact +// payload into the caller's buffer, and every non-hit (miss, undeserializable, absolutely expired) leaves the +// buffer untouched -- the IBufferWriter "nothing written on a miss" contract. +public class TryGetAsyncBufferTests : TestBase { - // A legacy JSON envelope from a pre-binary release: the first byte is '{' (0x7B), which never matches - // the binary FormatVersion, so the entry deserializes to a miss. - private static readonly byte[] LegacyJsonEntry = - Encoding.UTF8.GetBytes("{\"absexp\":null,\"sldexp\":null,\"data\":\"AQID\"}"); + // Held explicitly rather than captured from the primary constructor parameter, which would also be passed + // to the base constructor and warn under CS9107 (CI builds warnings-as-errors). + private readonly NatsIntegrationFixture _fixture; + + public TryGetAsyncBufferTests(NatsIntegrationFixture fixture) + : base(fixture) => + _fixture = fixture; private static CancellationToken Token => TestContext.Current.CancellationToken; @@ -79,24 +84,22 @@ public async Task WritesNothingOnUndeserializableEntry() public async Task WritesNothingOnAbsolutelyExpiredEntry() { var key = MethodKey(); - var value = new byte[] { 1, 2, 3, 4 }; - await DistributedCache.SetAsync( + var timeProvider = new FakeTimeProvider(); + await using var provider = BuildProvider(timeProvider); + var cache = (IBufferDistributedCache)provider.GetRequiredService(); + + // Five-minute absolute expiration -> five-minute real NATS TTL, so the entry is still present; the + // cache's own clock then jumps past it. Advancing a fake clock (rather than sleeping on a ~1s real + // TTL, which risks NATS reaping the key first and landing on NotFound) pins the expired branch. + await cache.SetAsync( key, - value, - new DistributedCacheEntryOptions().SetAbsoluteExpiration(TimeSpan.FromSeconds(1.1)), + new byte[] { 1, 2, 3, 4 }, + new DistributedCacheEntryOptions().SetAbsoluteExpiration(TimeSpan.FromMinutes(5)), Token); + timeProvider.Advance(TimeSpan.FromMinutes(10)); - // Poll the buffer path until the entry lapses; the final read must be a miss with nothing written. - ArrayBufferWriter destination; - bool hit; - var attempts = 0; - do - { - await Task.Delay(TimeSpan.FromSeconds(0.5), Token); - destination = new ArrayBufferWriter(); - hit = await BufferCache.TryGetAsync(key, destination, Token); - } - while (hit && ++attempts < 6); + var destination = new ArrayBufferWriter(); + var hit = await cache.TryGetAsync(key, destination, Token); Assert.False(hit); Assert.Equal(0, destination.WrittenCount); @@ -130,39 +133,53 @@ await DistributedCache.SetAsync( public async Task SlidingEntryPastAbsoluteExpirationWritesNothing() { var key = MethodKey(); - var value = new byte[] { 9, 8, 7 }; - - // Sliding renews on access, but the absolute expiration is the hard ceiling: once it passes, the - // buffer read must evict and miss, exercising the absolute-expiry branch of the sliding read path. - await DistributedCache.SetAsync( + var timeProvider = new FakeTimeProvider(); + await using var provider = BuildProvider(timeProvider); + var cache = (IBufferDistributedCache)provider.GetRequiredService(); + + // Sliding renews on access, but the absolute expiration is the hard ceiling. The sliding flag drives + // the read down the materialized branch; advancing the fake clock past the absolute instant then + // exercises that branch's absolute-expiry eviction deterministically. + await cache.SetAsync( key, - value, + new byte[] { 9, 8, 7 }, new DistributedCacheEntryOptions() - .SetSlidingExpiration(TimeSpan.FromSeconds(1.1)) - .SetAbsoluteExpiration(TimeSpan.FromSeconds(2)), + .SetSlidingExpiration(TimeSpan.FromMinutes(1)) + .SetAbsoluteExpiration(TimeSpan.FromMinutes(5)), Token); + timeProvider.Advance(TimeSpan.FromMinutes(10)); - ArrayBufferWriter destination; - bool hit; - var attempts = 0; - do - { - await Task.Delay(TimeSpan.FromSeconds(0.5), Token); - destination = new ArrayBufferWriter(); - hit = await BufferCache.TryGetAsync(key, destination, Token); - } - while (hit && ++attempts < 10); + var destination = new ArrayBufferWriter(); + var hit = await cache.TryGetAsync(key, destination, Token); + + Assert.False(hit); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public async Task ReturnsFalseWhenDestinationWriterFails() + { + var key = MethodKey(); + await DistributedCache.SetAsync(key, new byte[] { 1, 2, 3, 4 }, new DistributedCacheEntryOptions(), Token); + + // The caller's writer throws on the payload (mirroring HybridCache's quota-limited writer). TryGet + // honors the IBufferDistributedCache contract -- it swallows and returns false rather than throwing + // -- while the read core records the failure as an error (asserted in TelemetryTests), not a miss. + var destination = new QuotaBufferWriter(maxLength: 0); + var hit = await BufferCache.TryGetAsync(key, destination, Token); Assert.False(hit); Assert.Equal(0, destination.WrittenCount); } - // Writes raw bytes at the key the cache reads, bypassing the binary serializer so the stored entry - // cannot be deserialized. - private async Task WriteRawEntryAsync(string key, byte[] raw) + // A container mirroring TestBase's own, but with a controllable clock injected so absolute-expiry branches + // can be reached by advancing time instead of sleeping on a short real TTL. + private ServiceProvider BuildProvider(TimeProvider timeProvider) { - var encodedKey = new NatsCacheKeyEncoder().Encode(key); - var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache"); - await kvStore.PutAsync(encodedKey, raw, cancellationToken: Token); + var services = new ServiceCollection(); + services.AddSingleton(timeProvider); + _fixture.ConfigureServices(services); + services.AddNatsDistributedCache(options => options.BucketName = "cache"); + return services.BuildServiceProvider(); } } diff --git a/test/IntegrationTests/Cache/UndeserializableEntryTests.cs b/test/IntegrationTests/Cache/UndeserializableEntryTests.cs index 719c904..5ae2560 100644 --- a/test/IntegrationTests/Cache/UndeserializableEntryTests.cs +++ b/test/IntegrationTests/Cache/UndeserializableEntryTests.cs @@ -9,11 +9,6 @@ namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; public class UndeserializableEntryTests(NatsIntegrationFixture fixture) : TestBase(fixture) { - // A legacy JSON envelope from a pre-binary release: the first byte is '{' (0x7B), which never - // matches the binary FormatVersion, so the serializer returns null. - private static readonly byte[] LegacyJsonEntry = - Encoding.UTF8.GetBytes("{\"absexp\":null,\"sldexp\":null,\"data\":\"AQID\"}"); - [Fact] public async Task PresentButUndeserializableEntryIsReadAsMissAndLoggedAtDebug() { @@ -57,13 +52,4 @@ public async Task UndeserializableEntryIsLeftInPlaceAndSelfHealsOnNextWrite() // ...and the entry is now readable, confirming the documented "re-populated on next write" path. Assert.Equal(value, await cache.GetAsync(key, TestContext.Current.CancellationToken)); } - - // Writes raw bytes to the "cache" bucket at the key the cache reads, bypassing the binary - // serializer so the stored entry cannot be deserialized. - private async Task WriteRawEntryAsync(string key, byte[] raw) - { - var encodedKey = new NatsCacheKeyEncoder().Encode(key); - var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache"); - await kvStore.PutAsync(encodedKey, raw, cancellationToken: TestContext.Current.CancellationToken); - } } diff --git a/test/IntegrationTests/TestBase.cs b/test/IntegrationTests/TestBase.cs index 2019daa..24505e6 100644 --- a/test/IntegrationTests/TestBase.cs +++ b/test/IntegrationTests/TestBase.cs @@ -1,5 +1,6 @@ using System.Diagnostics.Metrics; using System.Runtime.CompilerServices; +using System.Text; using CodeCargo.Nats.DistributedCache.TestUtils; using CodeCargo.Nats.DistributedCache.TestUtils.Services.Logging; using Microsoft.Extensions.Caching.Distributed; @@ -17,6 +18,12 @@ namespace CodeCargo.Nats.DistributedCache.IntegrationTests; [Collection(NatsCollection.Name)] public abstract class TestBase : IAsyncLifetime { + // A legacy JSON envelope from a pre-binary release: the first byte is '{' (0x7B), which never matches the + // binary FormatVersion, so the serializer returns null and the read is an undeserializable miss. Shared by + // the tests that seed one via WriteRawEntryAsync. + protected static readonly byte[] LegacyJsonEntry = + Encoding.UTF8.GetBytes("{\"absexp\":null,\"sldexp\":null,\"data\":\"AQID\"}"); + private int _disposed; /// @@ -106,4 +113,15 @@ public virtual async ValueTask DisposeAsync() /// Gets the key for the current test method /// protected string MethodKey([CallerMemberName] string caller = "") => caller; + + /// + /// Writes raw bytes at the key the cache reads, bypassing the binary serializer so the stored entry + /// cannot be deserialized (used to seed legacy/corrupt entries). + /// + protected async Task WriteRawEntryAsync(string key, byte[] raw) + { + var encodedKey = new NatsCacheKeyEncoder().Encode(key); + var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache"); + await kvStore.PutAsync(encodedKey, raw, cancellationToken: TestContext.Current.CancellationToken); + } } diff --git a/test/TestUtils/QuotaBufferWriter.cs b/test/TestUtils/QuotaBufferWriter.cs new file mode 100644 index 0000000..55a59bf --- /dev/null +++ b/test/TestUtils/QuotaBufferWriter.cs @@ -0,0 +1,35 @@ +using System.Buffers; + +namespace CodeCargo.Nats.DistributedCache.TestUtils; + +/// +/// An that throws once more than bytes are written, +/// mirroring the quota behaviour of the RecyclableArrayBufferWriter HybridCache hands to +/// IBufferDistributedCache.TryGetAsync (whose Advance throws +/// "Max length exceeded" past its payload limit). Used to exercise +/// the destination-writer-failure read path. +/// +public sealed class QuotaBufferWriter : IBufferWriter +{ + private readonly ArrayBufferWriter _inner = new(); + private readonly int _maxLength; + + public QuotaBufferWriter(int maxLength) => _maxLength = maxLength; + + /// Gets the number of bytes committed so far (a throwing Advance commits nothing). + public int WrittenCount => _inner.WrittenCount; + + public void Advance(int count) + { + if (_inner.WrittenCount + count > _maxLength) + { + throw new InvalidOperationException("Max length exceeded"); + } + + _inner.Advance(count); + } + + public Memory GetMemory(int sizeHint = 0) => _inner.GetMemory(sizeHint); + + public Span GetSpan(int sizeHint = 0) => _inner.GetSpan(sizeHint); +} diff --git a/test/UnitTests/Serialization/CacheEntryReadDeserializerTests.cs b/test/UnitTests/Serialization/CacheEntryReadDeserializerTests.cs new file mode 100644 index 0000000..de132e4 --- /dev/null +++ b/test/UnitTests/Serialization/CacheEntryReadDeserializerTests.cs @@ -0,0 +1,167 @@ +using System.Buffers; +using CodeCargo.Nats.DistributedCache.TestUtils; +using Microsoft.Extensions.Time.Testing; + +namespace CodeCargo.Nats.DistributedCache.UnitTests.Serialization; + +/// +/// Tests for : array-mode materialization, the buffer fast path's +/// single-copy write, and the miss/expiry/sliding/destination-failure outcomes it reports to the read core. +/// +public class CacheEntryReadDeserializerTests +{ + private static readonly DateTimeOffset Now = new(2030, 1, 1, 0, 0, 0, TimeSpan.Zero); + + // First byte is not the format version, so the header fails to parse. + private static readonly byte[] BadFraming = [0xFF, 1, 2, 3]; + + [Fact] + public void ArrayMode_MaterializesEntry() + { + var entry = new CacheEntry { AbsoluteExpiration = Now.AddHours(1), Data = [1, 2, 3] }; + + var result = Deserialize(destination: null, Serialize(entry)); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.Materialized, result.Outcome); + Assert.NotNull(result.Entry); + Assert.Equal(entry.AbsoluteExpiration, result.Entry.AbsoluteExpiration); + Assert.Equal(entry.Data, result.Entry.Data); + } + + [Fact] + public void ArrayMode_BadFramingReturnsNull() => + Assert.Null(Deserialize(destination: null, BadFraming)); + + [Fact] + public void BufferMode_HitWritesPayloadAndReportsPayloadWritten() + { + var entry = new CacheEntry { Data = [10, 20, 30] }; + var destination = new ArrayBufferWriter(); + + var result = Deserialize(destination, Serialize(entry)); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.PayloadWritten, result.Outcome); + Assert.Equal(entry.Data, destination.WrittenSpan.ToArray()); + } + + [Fact] + public void BufferMode_MultiSegmentPayloadIsWrittenWhole() + { + var entry = new CacheEntry { AbsoluteExpiration = Now.AddHours(1), Data = [1, 2, 3, 4, 5, 6] }; + var bytes = Serialize(entry); + var destination = new ArrayBufferWriter(); + + // Split within the payload so WritePayload must cross a segment boundary. + var result = Deserialize(destination, Segmented(bytes, splitAt: bytes.Length - 3)); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.PayloadWritten, result.Outcome); + Assert.Equal(entry.Data, destination.WrittenSpan.ToArray()); + } + + [Fact] + public void BufferMode_BadFramingReturnsNullAndWritesNothing() + { + var destination = new ArrayBufferWriter(); + + Assert.Null(Deserialize(destination, BadFraming)); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public void BufferMode_AbsolutelyExpiredAtBoundaryWritesNothing() + { + // Absolute instant exactly equal to "now": the inclusive >= boundary treats it as expired. + var entry = new CacheEntry { AbsoluteExpiration = Now, Data = [1, 2, 3] }; + var destination = new ArrayBufferWriter(); + + var result = Deserialize(destination, Serialize(entry), now: Now); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.AbsolutelyExpired, result.Outcome); + Assert.Equal(0, destination.WrittenCount); + } + + [Fact] + public void BufferMode_OneTickBeforeExpiryWritesPayload() + { + var entry = new CacheEntry { AbsoluteExpiration = Now.AddTicks(1), Data = [7] }; + var destination = new ArrayBufferWriter(); + + var result = Deserialize(destination, Serialize(entry), now: Now); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.PayloadWritten, result.Outcome); + Assert.Equal(entry.Data, destination.WrittenSpan.ToArray()); + } + + [Fact] + public void BufferMode_SlidingEntryIsMaterializedNotWritten() + { + var entry = new CacheEntry { SlidingExpirationTicks = TimeSpan.FromMinutes(5).Ticks, Data = [4, 5, 6] }; + var destination = new ArrayBufferWriter(); + + var result = Deserialize(destination, Serialize(entry)); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.Materialized, result.Outcome); + Assert.Equal(0, destination.WrittenCount); // core handles the refresh and the write for sliding entries + Assert.NotNull(result.Entry); + Assert.Equal(entry.SlidingExpirationTicks, result.Entry.SlidingExpirationTicks); + Assert.Equal(entry.Data, result.Entry.Data); + } + + [Fact] + public void BufferMode_DestinationFailureCarriesTheException() + { + var entry = new CacheEntry { Data = [1, 2, 3, 4] }; + var destination = new QuotaBufferWriter(maxLength: 0); + + var result = Deserialize(destination, Serialize(entry)); + + Assert.NotNull(result); + Assert.Equal(CacheEntryReadOutcome.DestinationFailure, result.Outcome); + Assert.IsType(result.DestinationFailure); + Assert.Equal(0, destination.WrittenCount); + } + + private static byte[] Serialize(CacheEntry entry) + { + var writer = new ArrayBufferWriter(); + CacheEntryBinarySerializer.Default.Serialize(writer, entry); + return writer.WrittenMemory.ToArray(); + } + + private static CacheEntryReadResult? Deserialize(IBufferWriter? destination, byte[] bytes, DateTimeOffset? now = null) => + Deserialize(destination, new ReadOnlySequence(bytes), now); + + private static CacheEntryReadResult? Deserialize( + IBufferWriter? destination, + ReadOnlySequence sequence, + DateTimeOffset? now = null) + { + var deserializer = new CacheEntryReadDeserializer(destination, new FakeTimeProvider(now ?? Now)); + return deserializer.Deserialize(sequence); + } + + private static ReadOnlySequence Segmented(byte[] data, int splitAt) + { + var first = new BufferSegment(data.AsMemory(0, splitAt)); + var second = first.Append(data.AsMemory(splitAt)); + return new ReadOnlySequence(first, 0, second, second.Memory.Length); + } + + private sealed class BufferSegment : ReadOnlySequenceSegment + { + public BufferSegment(ReadOnlyMemory memory) => Memory = memory; + + public BufferSegment Append(ReadOnlyMemory memory) + { + var segment = new BufferSegment(memory) { RunningIndex = RunningIndex + Memory.Length }; + Next = segment; + return segment; + } + } +}