diff --git a/src/NatsDistributedCache/NatsCache.Log.cs b/src/NatsDistributedCache/NatsCache.Log.cs index 0ffe146..19c51a0 100644 --- a/src/NatsDistributedCache/NatsCache.Log.cs +++ b/src/NatsDistributedCache/NatsCache.Log.cs @@ -10,6 +10,9 @@ private void LogConnected(string bucketName) => private void LogException(Exception exception) => _logger.LogError(EventIds.Exception, exception, "Exception in NatsDistributedCache"); + private void LogSwallowedException(Exception exception) => + _logger.LogWarning(EventIds.Exception, exception, "NATS cache read failed in TryGetAsync; returning a cache miss"); + private static class EventIds { public static readonly EventId Connected = new(100, nameof(Connected)); diff --git a/src/NatsDistributedCache/NatsCache.cs b/src/NatsDistributedCache/NatsCache.cs index 1ac9a60..f35765a 100644 --- a/src/NatsDistributedCache/NatsCache.cs +++ b/src/NatsDistributedCache/NatsCache.cs @@ -73,10 +73,11 @@ public async Task SetAsync( { var ttl = GetTtl(options); var entry = CreateCacheEntry(value, options); - var kvStore = await GetKvStore().ConfigureAwait(false); try { + var kvStore = await GetKvStore().ConfigureAwait(false); + // todo: remove cast after https://github.com/nats-io/nats.net/pull/852 is released await ((NatsKVStore)kvStore) .PutWithTtlAsync(GetEncodedKey(key), entry, ttl ?? TimeSpan.Zero, CacheEntrySerializer, token) @@ -105,25 +106,55 @@ public async ValueTask SetAsync( } /// - public void Remove(string key) => RemoveAsync(key, null).GetAwaiter().GetResult(); + public void Remove(string key) => RemoveAsync(key).GetAwaiter().GetResult(); /// - public async Task RemoveAsync(string key, CancellationToken token = default) => - await RemoveAsync(key, null, token).ConfigureAwait(false); + public async Task RemoveAsync(string key, CancellationToken token = default) + { + try + { + await RemoveCoreAsync(key, null, token).ConfigureAwait(false); + } + catch (Exception ex) + { + LogException(ex); + throw; + } + } /// public void Refresh(string key) => RefreshAsync(key).GetAwaiter().GetResult(); /// - public Task RefreshAsync(string key, CancellationToken token = default) => - GetAndRefreshAsync(key, token: token); + public async Task RefreshAsync(string key, CancellationToken token = default) + { + try + { + await GetAndRefreshAsync(key, token).ConfigureAwait(false); + } + catch (Exception ex) + { + LogException(ex); + throw; + } + } /// public byte[]? Get(string key) => GetAsync(key).GetAwaiter().GetResult(); /// - public Task GetAsync(string key, CancellationToken token = default) => - GetAndRefreshAsync(key, token: token); + public async Task GetAsync(string key, CancellationToken token = default) + { + try + { + return await GetAndRefreshAsync(key, token).ConfigureAwait(false); + } + catch (Exception ex) + { + LogException(ex); + throw; + } + } /// public bool TryGet(string key, IBufferWriter destination) => @@ -137,16 +168,25 @@ public async ValueTask TryGetAsync( { try { - var result = await GetAsync(key, token).ConfigureAwait(false); + var result = await GetAndRefreshAsync(key, token).ConfigureAwait(false); if (result != null) { destination.Write(result); return true; } } - catch + catch (OperationCanceledException) when (token.IsCancellationRequested) { - // Ignore failures here; they will surface later + // Cooperative cancellation is not a cache failure; let it propagate rather than + // masquerading as a cache miss. + throw; + } + catch (Exception ex) + { + // A read failure (e.g. NATS connectivity or a corrupt entry) is swallowed to honor the + // IBufferDistributedCache contract (return false), but logged at warning so it stays + // visible in production and is distinguishable from a normal cache miss. + LogSwallowedException(ex); } return false; @@ -234,12 +274,11 @@ private Lazy> CreateLazyKvStore() => LogConnected(_bucketName); return store; } - catch (Exception ex) + catch (Exception) { - // Reset the lazy initializer on failure for next attempt + // Reset the lazy initializer on failure so the next attempt retries. The exception + // propagates to the calling operation, which logs it once at the appropriate level. _lazyKvStore = CreateLazyKvStore(); - - LogException(ex); throw; } }); @@ -271,7 +310,7 @@ private Lazy> CreateLazyKvStore() => { // NatsKVWrongLastRevisionException is caught below var natsDeleteOpts = new NatsKVDeleteOpts { Revision = kvEntry.Revision }; - await RemoveAsync(key, natsDeleteOpts, token).ConfigureAwait(false); + await RemoveCoreAsync(key, natsDeleteOpts, token).ConfigureAwait(false); return null; } @@ -283,11 +322,6 @@ private Lazy> CreateLazyKvStore() => // Someone else updated it; that's fine, we'll get the latest version next time return null; } - catch (Exception ex) - { - LogException(ex); - throw; - } // Local Functions async Task UpdateEntryExpirationAsync(NatsKVEntry kvEntry) @@ -327,7 +361,7 @@ await kvStore.UpdateWithTtlAsync( } } - private async Task RemoveAsync( + private async Task RemoveCoreAsync( string key, NatsKVDeleteOpts? natsKvDeleteOpts = null, CancellationToken token = default) diff --git a/test/IntegrationTests/Cache/TryGetAsyncTests.cs b/test/IntegrationTests/Cache/TryGetAsyncTests.cs new file mode 100644 index 0000000..7bec7ba --- /dev/null +++ b/test/IntegrationTests/Cache/TryGetAsyncTests.cs @@ -0,0 +1,74 @@ +using System.Buffers; +using CodeCargo.Nats.DistributedCache.TestUtils.Services.Logging; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; + +namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; + +public class TryGetAsyncTests(NatsIntegrationFixture fixture) : TestBase(fixture) +{ + [Fact] + public async Task TryGetAsyncSwallowsFailureAndLogsOnceAtWarning() + { + var cache = CreateFailingCache(out var logger); + var destination = new ArrayBufferWriter(); + + var result = await cache.TryGetAsync(MethodKey(), destination, TestContext.Current.CancellationToken); + + Assert.False(result); + Assert.Equal(0, destination.WrittenCount); + + // The failure is logged exactly once, at warning, with no redundant error-level entry. + var record = Assert.Single(logger.Records); + Assert.Equal(LogLevel.Warning, record.LogLevel); + Assert.Equal("Exception", record.EventId.Name); + Assert.NotNull(record.Exception); + } + + [Fact] + public async Task GetAsyncPropagatesFailureAndLogsOnceAtError() + { + var cache = CreateFailingCache(out var logger); + + await Assert.ThrowsAnyAsync( + () => cache.GetAsync(MethodKey(), TestContext.Current.CancellationToken)); + + // A propagating read logs exactly once, at error, and is not double-logged by shared helpers. + var record = Assert.Single(logger.Records); + Assert.Equal(LogLevel.Error, record.LogLevel); + Assert.Equal("Exception", record.EventId.Name); + Assert.NotNull(record.Exception); + } + + [Fact] + public async Task TryGetAsyncPropagatesCancellationWithoutLogging() + { + // Use the real bucket so the store resolves; the cancelled token then fails the read itself. + var logger = new RecordingLogger(); + var cache = new NatsCache( + Options.Create(new NatsCacheOptions { BucketName = "cache" }), + NatsConnection, + logger); + var destination = new ArrayBufferWriter(); + using var cts = new CancellationTokenSource(); + await cts.CancelAsync(); + + // Cancellation is not a cache failure: it propagates instead of becoming a false miss... + await Assert.ThrowsAnyAsync( + () => cache.TryGetAsync(MethodKey(), destination, cts.Token).AsTask()); + + // ...and it is not logged as an exception (a benign "Connected" info entry may be present). + Assert.DoesNotContain(logger.Records, r => r.EventId.Name == "Exception"); + } + + // Points a cache at a bucket that does not exist so the read path throws when it resolves the KV + // store, without depending on the shared fixture connection being torn down. + private NatsCache CreateFailingCache(out RecordingLogger logger) + { + logger = new RecordingLogger(); + return new NatsCache( + Options.Create(new NatsCacheOptions { BucketName = "does-not-exist" }), + NatsConnection, + logger); + } +} diff --git a/test/TestUtils/Services/Logging/RecordingLogger.cs b/test/TestUtils/Services/Logging/RecordingLogger.cs new file mode 100644 index 0000000..b2725d0 --- /dev/null +++ b/test/TestUtils/Services/Logging/RecordingLogger.cs @@ -0,0 +1,55 @@ +using Microsoft.Extensions.Logging; + +namespace CodeCargo.Nats.DistributedCache.TestUtils.Services.Logging; + +/// +/// A captured log entry recorded by +/// +/// The level the entry was logged at +/// The event id associated with the entry +/// The exception attached to the entry, if any +/// The formatted log message +public sealed record LogRecord(LogLevel LogLevel, EventId EventId, Exception? Exception, string Message); + +/// +/// An that captures log entries in memory so tests can assert on them +/// +/// The category type +public sealed class RecordingLogger : ILogger +{ + private readonly List _records = new(); + + /// + /// Gets a snapshot of the entries captured so far + /// + public IReadOnlyList Records + { + get + { + lock (_records) + { + return _records.ToArray(); + } + } + } + + public void Log( + LogLevel logLevel, + EventId eventId, + TState state, + Exception? exception, + Func formatter) + { + var record = new LogRecord(logLevel, eventId, exception, formatter(state, exception)); + lock (_records) + { + _records.Add(record); + } + } + + public bool IsEnabled(LogLevel logLevel) => true; + + public IDisposable? BeginScope(TState state) + where TState : notnull => + null; +}