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
14 changes: 14 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,20 @@ services.AddSingleton<TimeProvider>(new FakeTimeProvider());
services.AddNatsDistributedCache(options => options.BucketName = "cache");
```

## Cache Entry Format and Upgrades

Cache entries are stored in a compact binary envelope. When an entry cannot be deserialized — because
it was written by an incompatible release (for example a pre-binary version that used a JSON envelope)
or is otherwise corrupt — the read is treated as a **cache miss** rather than an error, and logged at
`Debug`. Because a cache's source of truth lives elsewhere, no manual migration is required:

- Entries with a TTL are reaped automatically by NATS once they expire.
- Entries without a TTL are left in place and re-populated the next time the key is written (a `Set`
overwrites the stored bytes unconditionally), which happens naturally under cache-aside usage.

Upgrading is therefore seamless in a rolling deployment: a node never deletes an entry it cannot read,
so it cannot discard entries written by a newer node still being rolled out.

## Additional Resources

* [ASP.NET Core Hybrid Cache Documentation](https://learn.microsoft.com/en-us/aspnet/core/performance/caching/hybrid?view=aspnetcore-10.0)
Expand Down
24 changes: 23 additions & 1 deletion src/NatsDistributedCache/CacheEntryBinarySerializer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,22 @@ internal sealed class CacheEntryBinarySerializer : INatsSerialize<CacheEntry>, I
/// </summary>
internal const byte FormatVersion = 1;

/// <summary>
/// The largest expiration window the cache can represent, in ticks. NATS message TTLs are encoded as
/// whole seconds in a 32-bit field (<c>(int)ttl.TotalSeconds</c> in <c>NatsExtensions.ToTtlString</c>),
/// so a window beyond <see cref="int.MaxValue"/> seconds (~68 years) would overflow the cast and emit
/// an invalid TTL header. <c>NatsCache.GetTtl</c> rejects sliding and absolute expirations that exceed
/// this on write, and the deserializer fails closed on stored sliding ticks above it, so every
/// accepted value round-trips end to end and no legitimately written entry becomes an
/// undeserializable miss.
/// </summary>
internal const long MaxTtlTicks = int.MaxValue * TimeSpan.TicksPerSecond;

/// <summary>
/// <see cref="MaxTtlTicks"/> as a <see cref="TimeSpan"/>, for range-check exception messages.
/// </summary>
internal static readonly TimeSpan MaxTtl = TimeSpan.FromTicks(MaxTtlTicks);

private const byte HasAbsoluteExpiration = 0b0000_0001;
private const byte HasSlidingExpiration = 0b0000_0010;
private const byte KnownFlags = HasAbsoluteExpiration | HasSlidingExpiration;
Expand Down Expand Up @@ -113,8 +129,14 @@ public void Serialize(IBufferWriter<byte> bufferWriter, CacheEntry value)
long? slidingExpirationTicks = null;
if ((flags & HasSlidingExpiration) != 0)
{
if (!reader.TryReadLittleEndian(out long slidingTicks))
if (!reader.TryReadLittleEndian(out long slidingTicks) ||
slidingTicks <= 0 ||
slidingTicks > MaxTtlTicks)
{
// Missing, non-positive, or out-of-range sliding ticks (corrupt entry): fail closed to
// 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;
Comment thread
matthewdevenny marked this conversation as resolved.
}

Expand Down
7 changes: 7 additions & 0 deletions src/NatsDistributedCache/NatsCache.Log.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,16 @@ 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) =>
_logger.LogDebug(
EventIds.UndeserializableEntry,
"Cache entry for key {Key} could not be deserialized (legacy or corrupt format); returning a cache miss",
key);

private static class EventIds
{
public static readonly EventId Connected = new(100, nameof(Connected));
public static readonly EventId Exception = new(101, nameof(Exception));
public static readonly EventId UndeserializableEntry = new(102, nameof(UndeserializableEntry));
}
}
105 changes: 85 additions & 20 deletions src/NatsDistributedCache/NatsCache.cs
Original file line number Diff line number Diff line change
Expand Up @@ -206,56 +206,97 @@ public async ValueTask<bool> TryGetAsync(

internal TimeSpan? GetTtl(DistributedCacheEntryOptions options)
{
if (options.AbsoluteExpiration.HasValue && options.AbsoluteExpiration.Value <= TimeProvider.GetUtcNow())
// Maximum-value sentinels mean "never expire": normalize them to no expiration so they are not
// rejected as out-of-range below. Callers commonly use DateTimeOffset.MaxValue / TimeSpan.MaxValue
// as a "cache forever" idiom, which is equivalent to omitting expiration entirely.
var absoluteExpiration = EffectiveAbsoluteExpiration(options);
var relativeToNow = EffectiveRelativeToNow(options);
var slidingExpiration = EffectiveSlidingExpiration(options);

if (absoluteExpiration.HasValue && absoluteExpiration.Value <= TimeProvider.GetUtcNow())
{
throw new ArgumentOutOfRangeException(
nameof(DistributedCacheEntryOptions.AbsoluteExpiration),
options.AbsoluteExpiration.Value,
absoluteExpiration.Value,
"The absolute expiration value must be in the future.");
}

if (options.AbsoluteExpirationRelativeToNow.HasValue &&
options.AbsoluteExpirationRelativeToNow.Value <= TimeSpan.Zero)
if (relativeToNow.HasValue && relativeToNow.Value <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(
nameof(DistributedCacheEntryOptions.AbsoluteExpirationRelativeToNow),
options.AbsoluteExpirationRelativeToNow.Value,
relativeToNow.Value,
"The relative expiration value must be positive.");
}

if (options.SlidingExpiration.HasValue && options.SlidingExpiration.Value <= TimeSpan.Zero)
if (slidingExpiration.HasValue && slidingExpiration.Value <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(
nameof(DistributedCacheEntryOptions.SlidingExpiration),
options.SlidingExpiration.Value,
slidingExpiration.Value,
"The sliding expiration value must be positive.");
}

var absoluteExpiration = ResolveAbsoluteExpiration(options);
if (!absoluteExpiration.HasValue)
// Reject sliding windows the serializer/TTL encoding could not round-trip. The read path fails
// closed above MaxTtlTicks, so accepting a larger value here would silently store an entry that
// later reads back as an undeserializable cache miss.
if (slidingExpiration.HasValue &&
slidingExpiration.Value.Ticks > CacheEntryBinarySerializer.MaxTtlTicks)
{
return options.SlidingExpiration;
throw new ArgumentOutOfRangeException(
nameof(DistributedCacheEntryOptions.SlidingExpiration),
slidingExpiration.Value,
$"The sliding expiration value is too large. The maximum is {CacheEntryBinarySerializer.MaxTtl}.");
}

var resolvedAbsolute = ResolveAbsoluteExpiration(options);
if (!resolvedAbsolute.HasValue)
{
return slidingExpiration;
}

var ttl = absoluteExpiration.Value - TimeProvider.GetUtcNow();
var ttl = resolvedAbsolute.Value - TimeProvider.GetUtcNow();
if (ttl.TotalMilliseconds <= 0)
{
// Value is in the past, remove it
return TimeSpan.Zero;
}

// If there's also a sliding expiration, use the minimum of the two
return options.SlidingExpiration.HasValue
? TimeSpan.FromTicks(Math.Min(ttl.Ticks, options.SlidingExpiration.Value.Ticks))
: ttl;
// If there's also a sliding expiration, use the minimum of the two. Sliding is bounded to
// MaxTtlTicks above, so the minimum is always within the encodable range.
if (slidingExpiration.HasValue)
{
return TimeSpan.FromTicks(Math.Min(ttl.Ticks, slidingExpiration.Value.Ticks));
}

// Absolute-only: the TTL spans the full window to the absolute instant. Reject windows the NATS
// TTL encoding cannot represent (see CacheEntryBinarySerializer.MaxTtlTicks) rather than emitting
// an overflowed header.
if (ttl.Ticks > CacheEntryBinarySerializer.MaxTtlTicks)
{
if (relativeToNow.HasValue)
{
throw new ArgumentOutOfRangeException(
nameof(DistributedCacheEntryOptions.AbsoluteExpirationRelativeToNow),
relativeToNow.Value,
$"The relative expiration value is too large. The maximum is {CacheEntryBinarySerializer.MaxTtl}.");
}

throw new ArgumentOutOfRangeException(
nameof(DistributedCacheEntryOptions.AbsoluteExpiration),
absoluteExpiration!.Value,
$"The absolute expiration is too far in the future. The maximum window is {CacheEntryBinarySerializer.MaxTtl}.");
}

return ttl;
}

internal CacheEntry CreateCacheEntry(byte[] value, DistributedCacheEntryOptions options) =>
new CacheEntry
{
Data = value,
AbsoluteExpiration = ResolveAbsoluteExpiration(options),
SlidingExpirationTicks = options.SlidingExpiration?.Ticks
SlidingExpirationTicks = EffectiveSlidingExpiration(options)?.Ticks
};

// An entry is absolutely expired once the clock reaches its absolute expiration instant. The
Expand Down Expand Up @@ -291,12 +332,27 @@ internal NatsKVConfig BuildBucketConfig()
return config;
}

// "Never expire" sentinels: a DateTimeOffset.MaxValue absolute instant or a TimeSpan.MaxValue window
// is normalized to no expiration, so it is not rejected as out-of-range and the entry lives forever.
private static DateTimeOffset? EffectiveAbsoluteExpiration(DistributedCacheEntryOptions options) =>
options.AbsoluteExpiration == DateTimeOffset.MaxValue ? null : options.AbsoluteExpiration;

private static TimeSpan? EffectiveRelativeToNow(DistributedCacheEntryOptions options) =>
options.AbsoluteExpirationRelativeToNow == TimeSpan.MaxValue ? null : options.AbsoluteExpirationRelativeToNow;

private static TimeSpan? EffectiveSlidingExpiration(DistributedCacheEntryOptions options) =>
options.SlidingExpiration == TimeSpan.MaxValue ? null : options.SlidingExpiration;

// Resolves the effective absolute expiration instant: a relative expiration (offset from the
// current clock) takes precedence over an explicit absolute expiration when both are set.
private DateTimeOffset? ResolveAbsoluteExpiration(DistributedCacheEntryOptions options) =>
options.AbsoluteExpirationRelativeToNow.HasValue
? TimeProvider.GetUtcNow().Add(options.AbsoluteExpirationRelativeToNow.Value)
: options.AbsoluteExpiration;
// Maximum-value sentinels are treated as "no expiration" (see the Effective* helpers).
private DateTimeOffset? ResolveAbsoluteExpiration(DistributedCacheEntryOptions options)
{
var relativeToNow = EffectiveRelativeToNow(options);
return relativeToNow.HasValue
? TimeProvider.GetUtcNow().Add(relativeToNow.Value)
: EffectiveAbsoluteExpiration(options);
}

private string GetEncodedKey(string key) =>
string.IsNullOrEmpty(_keyPrefix)
Expand Down Expand Up @@ -363,6 +419,15 @@ private Lazy<Task<INatsKVStore>> CreateLazyKvStore() =>
var kvEntry = natsResult.Value;
if (kvEntry.Value == 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);
return null;
}

Expand Down
69 changes: 69 additions & 0 deletions test/IntegrationTests/Cache/UndeserializableEntryTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
using System.Text;
using CodeCargo.Nats.DistributedCache.TestUtils.Services.Logging;
using Microsoft.Extensions.Caching.Distributed;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using NATS.Net;

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()
{
var key = MethodKey();
var logger = new RecordingLogger<NatsCache>();
var cache = new NatsCache(Options.Create(new NatsCacheOptions { BucketName = "cache" }), NatsConnection, logger);
await WriteRawEntryAsync(key, LegacyJsonEntry);

// The undeserializable entry reads as a cache miss rather than throwing...
var result = await cache.GetAsync(key, TestContext.Current.CancellationToken);
Assert.Null(result);

// ...and is surfaced once at Debug to aid diagnosis without failing the caller's operation.
var record = Assert.Single(logger.Records, r => r.EventId.Name == "UndeserializableEntry");
Assert.Equal(LogLevel.Debug, record.LogLevel);
}

[Fact]
public async Task UndeserializableEntryIsLeftInPlaceAndSelfHealsOnNextWrite()
{
var key = MethodKey();
var cache = new NatsCache(
Options.Create(new NatsCacheOptions { BucketName = "cache" }),
NatsConnection,
new RecordingLogger<NatsCache>());
await WriteRawEntryAsync(key, LegacyJsonEntry);

// The read leaves the entry untouched (no eviction): the original legacy bytes are still in
// the bucket afterwards, and it is still a miss on a second read.
Assert.Null(await cache.GetAsync(key, TestContext.Current.CancellationToken));
var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache");
var storedEntry = await kvStore.GetEntryAsync<byte[]>(
new NatsCacheKeyEncoder().Encode(key),
cancellationToken: TestContext.Current.CancellationToken);
Assert.Equal(LegacyJsonEntry, storedEntry.Value);

// Writing the key (a no-TTL entry, as in the migration case) overwrites the legacy bytes...
var value = Encoding.UTF8.GetBytes($"healed-{Guid.NewGuid()}");
await cache.SetAsync(key, value, new DistributedCacheEntryOptions(), TestContext.Current.CancellationToken);

// ...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);
}
}
57 changes: 57 additions & 0 deletions test/UnitTests/Cache/TimeExpirationUnitTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -91,4 +91,61 @@ public void ZeroSlidingExpirationThrows()
"The sliding expiration value must be positive.",
TimeSpan.Zero);
}

[Fact]
public void TooLargeSlidingExpirationThrows()
{
var key = MethodKey();
var value = new byte[1];

// A window beyond the NATS TTL encoding ceiling (~68 years) but short of the TimeSpan.MaxValue
// "never expire" sentinel must be rejected rather than stored as an entry that later reads back
// as an undeserializable miss or emits an overflowed TTL header.
var sliding = TimeSpan.FromDays(365 * 100);
ExceptionAssert.ThrowsArgumentOutOfRange(
() =>
{
Cache.Set(key, value, new DistributedCacheEntryOptions().SetSlidingExpiration(sliding));
},
nameof(DistributedCacheEntryOptions.SlidingExpiration),
"The sliding expiration value is too large.",
sliding);
}

[Fact]
public void TooFarRelativeExpirationThrows()
{
var key = MethodKey();
var value = new byte[1];

// A relative expiration beyond the NATS TTL encoding limit (int.MaxValue seconds, ~68 years)
// would overflow the (int) cast in ToTtlString, so it must be rejected on write.
var relative = TimeSpan.FromDays(365 * 100);
ExceptionAssert.ThrowsArgumentOutOfRange(
() =>
{
Cache.Set(key, value, new DistributedCacheEntryOptions().SetAbsoluteExpiration(relative));
},
nameof(DistributedCacheEntryOptions.AbsoluteExpirationRelativeToNow),
"The relative expiration value is too large.",
relative);
}

[Fact]
public void TooFarAbsoluteExpirationThrows()
{
var key = MethodKey();
var value = new byte[1];

// Same limit for an absolute instant: more than ~68 years out overflows the TTL header.
var absolute = TimeProvider.GetUtcNow().AddYears(100);
ExceptionAssert.ThrowsArgumentOutOfRange(
() =>
{
Cache.Set(key, value, new DistributedCacheEntryOptions().SetAbsoluteExpiration(absolute));
},
nameof(DistributedCacheEntryOptions.AbsoluteExpiration),
"The absolute expiration is too far in the future.",
absolute);
}
}
Loading
Loading