Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions src/NatsDistributedCache/NatsCache.Log.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
78 changes: 56 additions & 22 deletions src/NatsDistributedCache/NatsCache.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -105,25 +106,55 @@ public async ValueTask SetAsync(
}

/// <inheritdoc />
public void Remove(string key) => RemoveAsync(key, null).GetAwaiter().GetResult();
public void Remove(string key) => RemoveAsync(key).GetAwaiter().GetResult();

/// <inheritdoc />
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;
}
}

/// <inheritdoc />
public void Refresh(string key) => RefreshAsync(key).GetAwaiter().GetResult();

/// <inheritdoc />
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;
}
}

/// <inheritdoc />
public byte[]? Get(string key) => GetAsync(key).GetAwaiter().GetResult();

/// <inheritdoc />
public Task<byte[]?> GetAsync(string key, CancellationToken token = default) =>
GetAndRefreshAsync(key, token: token);
public async Task<byte[]?> GetAsync(string key, CancellationToken token = default)
{
try
{
return await GetAndRefreshAsync(key, token).ConfigureAwait(false);
}
catch (Exception ex)
{
LogException(ex);
throw;
}
}

/// <inheritdoc />
public bool TryGet(string key, IBufferWriter<byte> destination) =>
Expand All @@ -137,16 +168,25 @@ public async ValueTask<bool> 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);
Comment thread
matthewdevenny marked this conversation as resolved.
}

return false;
Expand Down Expand Up @@ -234,12 +274,11 @@ private Lazy<Task<INatsKVStore>> 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;
}
});
Expand Down Expand Up @@ -271,7 +310,7 @@ private Lazy<Task<INatsKVStore>> 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;
}

Expand All @@ -283,11 +322,6 @@ private Lazy<Task<INatsKVStore>> 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<CacheEntry> kvEntry)
Expand Down Expand Up @@ -327,7 +361,7 @@ await kvStore.UpdateWithTtlAsync(
}
}

private async Task RemoveAsync(
private async Task RemoveCoreAsync(
string key,
NatsKVDeleteOpts? natsKvDeleteOpts = null,
CancellationToken token = default)
Expand Down
74 changes: 74 additions & 0 deletions test/IntegrationTests/Cache/TryGetAsyncTests.cs
Original file line number Diff line number Diff line change
@@ -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<byte>();

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<Exception>(
() => 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<NatsCache>();
var cache = new NatsCache(
Options.Create(new NatsCacheOptions { BucketName = "cache" }),
NatsConnection,
logger);
var destination = new ArrayBufferWriter<byte>();
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<OperationCanceledException>(
() => 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<NatsCache> logger)
{
logger = new RecordingLogger<NatsCache>();
return new NatsCache(
Options.Create(new NatsCacheOptions { BucketName = "does-not-exist" }),
NatsConnection,
logger);
}
}
55 changes: 55 additions & 0 deletions test/TestUtils/Services/Logging/RecordingLogger.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
using Microsoft.Extensions.Logging;

namespace CodeCargo.Nats.DistributedCache.TestUtils.Services.Logging;

/// <summary>
/// A captured log entry recorded by <see cref="RecordingLogger{T}" />
/// </summary>
/// <param name="LogLevel">The level the entry was logged at</param>
/// <param name="EventId">The event id associated with the entry</param>
/// <param name="Exception">The exception attached to the entry, if any</param>
/// <param name="Message">The formatted log message</param>
public sealed record LogRecord(LogLevel LogLevel, EventId EventId, Exception? Exception, string Message);

/// <summary>
/// An <see cref="ILogger{T}" /> that captures log entries in memory so tests can assert on them
/// </summary>
/// <typeparam name="T">The category type</typeparam>
public sealed class RecordingLogger<T> : ILogger<T>
{
private readonly List<LogRecord> _records = new();

/// <summary>
/// Gets a snapshot of the entries captured so far
/// </summary>
public IReadOnlyList<LogRecord> Records
{
get
{
lock (_records)
{
return _records.ToArray();
}
}
}

public void Log<TState>(
LogLevel logLevel,
EventId eventId,
TState state,
Exception? exception,
Func<TState, Exception?, string> 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>(TState state)
where TState : notnull =>
null;
}
Loading