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;
+}