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/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 fec9c82..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) { @@ -219,15 +221,13 @@ 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 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 GetAndRefreshAsync(key, CacheOperation.Get, destination, token) + .ConfigureAwait(false)).Hit; } catch (OperationCanceledException) when (token.IsCancellationRequested) { @@ -246,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 @@ -345,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 @@ -445,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 @@ -457,88 +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; - if (kvEntry.Value == 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(kvEntry.Value)) + 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(kvEntry).ConfigureAwait(false); - scope.SetHit(); - return kvEntry.Value.Data; - } - catch (NatsKVWrongLastRevisionException) - { - // Someone else updated it; that's fine, we'll get the latest version next time - scope.SetMiss(CacheMissReason.RevisionConflict); - return null; - } - - // Local Functions - async Task UpdateEntryExpirationAsync(NatsKVEntry kvEntry) - { - if (kvEntry.Value?.SlidingExpirationTicks == null) + if (result.Outcome == CacheEntryReadOutcome.PayloadWritten) { - return; + // Buffer fast path: a confirmed hit already streamed into the caller's writer. + scope.SetHit(); + return (true, null); } - // If we have a sliding expiration, use it as the TTL - var ttl = TimeSpan.FromTicks(kvEntry.Value.SlidingExpirationTicks.Value); + if (result.Outcome == CacheEntryReadOutcome.AbsolutelyExpired) + { + // 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, null); + } - // If we also have an absolute expiration, make sure we don't exceed it - if (kvEntry.Value.AbsoluteExpiration != null) + // 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 remainingTime = kvEntry.Value.AbsoluteExpiration.Value - TimeProvider.GetUtcNow(); + // NatsKVWrongLastRevisionException is caught below. + await EvictExpiredAsync(key, kvEntry.Revision, token).ConfigureAwait(false); + scope.SetMiss(CacheMissReason.Expired); + return (false, null); + } - // Use the minimum of sliding window or remaining absolute time - if (remainingTime > TimeSpan.Zero && remainingTime < ttl) - { - ttl = remainingTime; - } + // Refresh first — a lost revision race surfaces as NatsKVWrongLastRevisionException and is + // 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); } - if (ttl > TimeSpan.Zero) + // 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) { - // Use optimistic concurrency control with the last revision - await kvStore.UpdateWithTtlAsync( - encodedKey, - kvEntry.Value, - kvEntry.Revision, - ttl, - serializer: CacheEntrySerializer, - cancellationToken: token).ConfigureAwait(false); + destination.Write(data); } + + scope.SetHit(); + return (true, null); + } + catch (NatsKVWrongLastRevisionException) + { + // Someone else updated it (during the sliding refresh); we'll get the latest next time. + // Nothing was emitted. + scope.SetMiss(CacheMissReason.RevisionConflict); + return (false, null); } } catch (Exception ex) @@ -552,6 +585,56 @@ await kvStore.UpdateWithTtlAsync( } } + // 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. + 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/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 new file mode 100644 index 0000000..e53ca90 --- /dev/null +++ b/test/IntegrationTests/Cache/TryGetAsyncBufferTests.cs @@ -0,0 +1,185 @@ +using System.Buffers; +using System.Text; +using CodeCargo.Nats.DistributedCache.TestUtils; +using Microsoft.Extensions.Caching.Distributed; +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 : TestBase +{ + // 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; + + 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 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, + new byte[] { 1, 2, 3, 4 }, + new DistributedCacheEntryOptions().SetAbsoluteExpiration(TimeSpan.FromMinutes(5)), + Token); + timeProvider.Advance(TimeSpan.FromMinutes(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 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 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, + new byte[] { 9, 8, 7 }, + new DistributedCacheEntryOptions() + .SetSlidingExpiration(TimeSpan.FromMinutes(1)) + .SetAbsoluteExpiration(TimeSpan.FromMinutes(5)), + Token); + timeProvider.Advance(TimeSpan.FromMinutes(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); + } + + // 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 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; + } + } +}