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
76 changes: 57 additions & 19 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,17 @@ A .NET 8+ library (tested on .NET 8 and .NET 10) for using NATS with `HybridCach
## Requirements

- NATS 2.11 or later
- A NATS KV bucket with `LimitMarkerTTL` set for per-key TTL support. Example:
- A NATS KV bucket with `LimitMarkerTTL` set for per-key TTL support. Either enable
[automatic bucket creation](#automatic-bucket-creation) (`options.CreateBucketIfNotExists = true`), or
pre-create the bucket yourself:
```csharp
using NATS.Client.KeyValueStore;
using NATS.Net;

// assuming an INatsConnection natsConnection
var kvContext = natsConnection.CreateKeyValueStoreContext();
await kvContext.CreateOrUpdateStoreAsync(new NatsKVConfig("cache") { LimitMarkerTTL = TimeSpan.FromSeconds(1) });
await kvContext.CreateOrUpdateStoreAsync(
new NatsKVConfig("cache") { LimitMarkerTTL = TimeSpan.FromSeconds(1), History = 1 });
Comment thread
matthewdevenny marked this conversation as resolved.
```

## Use with `HybridCache`
Expand All @@ -37,8 +43,6 @@ dotnet add package NATS.Extensions.Microsoft.DependencyInjection
using CodeCargo.Nats.HybridCacheExtensions;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NATS.Client.Core;
using NATS.Client.KeyValueStore;
using NATS.Extensions.Microsoft.DependencyInjection;
using NATS.Net;

Expand All @@ -56,19 +60,15 @@ builder.ConfigureServices(services =>
services.AddNatsHybridCache(options =>
{
options.BucketName = "cache";

// Create the KV bucket on first use if it doesn't already exist.
// Omit this if you pre-create the bucket yourself (see Requirements).
options.CreateBucketIfNotExists = true;
});
});

var host = builder.Build();

// Ensure that the KV Store is created
var natsConnection = host.Services.GetRequiredService<INatsConnection>();
var kvContext = natsConnection.CreateKeyValueStoreContext();
await kvContext.CreateOrUpdateStoreAsync(new NatsKVConfig("cache")
{
LimitMarkerTTL = TimeSpan.FromSeconds(1)
});

// Start the host
await host.RunAsync();
```
Expand All @@ -88,8 +88,6 @@ dotnet add package NATS.Extensions.Microsoft.DependencyInjection
using CodeCargo.Nats.DistributedCache;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NATS.Client.Core;
using NATS.Client.KeyValueStore;
using NATS.Extensions.Microsoft.DependencyInjection;
using NATS.Net;

Expand All @@ -107,20 +105,60 @@ builder.ConfigureServices(services =>
services.AddNatsDistributedCache(options =>
{
options.BucketName = "cache";

// Create the KV bucket on first use if it doesn't already exist.
// Omit this if you pre-create the bucket yourself (see Requirements).
options.CreateBucketIfNotExists = true;
});
});

var host = builder.Build();

// Ensure that the KV Store is created
var natsConnection = host.Services.GetRequiredService<INatsConnection>();
var kvContext = natsConnection.CreateKeyValueStoreContext();
await kvContext.CreateOrUpdateStoreAsync(new NatsKVConfig("cache") { LimitMarkerTTL = TimeSpan.FromSeconds(1) });

// Start the host
await host.RunAsync();
```

## Automatic bucket creation

By default the KV bucket must already exist. Set `CreateBucketIfNotExists = true` to have the cache
create it on first use if it is missing, with the settings per-key TTL requires —
`History = 1` and a non-zero `LimitMarkerTTL`:

```csharp
services.AddNatsDistributedCache(options =>
{
options.BucketName = "cache";
options.CreateBucketIfNotExists = true;
});
```

To customize storage, replication, or size limits, use `ConfigureBucketOnCreate`. `NatsKVConfig` is an
immutable record, so return a modified copy with a `with` expression:

```csharp
using NATS.Client.KeyValueStore;

services.AddNatsDistributedCache(options =>
{
options.BucketName = "cache";
options.CreateBucketIfNotExists = true;
options.ConfigureBucketOnCreate = config => config with
{
Storage = NatsKVStorageType.File,
NumberOfReplicas = 3,
};
});
```

Notes:

- Only a missing bucket is created; an existing bucket is used as-is and never modified, so
operator-managed settings are preserved. `ConfigureBucketOnCreate` therefore only applies when the bucket is
first created.
- Creating a bucket requires JetStream stream-management permissions.
- Overriding `History` (away from `1`) or clearing `LimitMarkerTTL` in `ConfigureBucketOnCreate` disables
reliable per-key TTL.

## Controlling Expiration Timing

Expiration is computed from a [`TimeProvider`](https://learn.microsoft.com/dotnet/api/system.timeprovider),
Expand Down
63 changes: 62 additions & 1 deletion src/NatsDistributedCache/NatsCache.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using NATS.Client.Core;
using NATS.Client.JetStream;
using NATS.Client.KeyValueStore;
using NATS.Net;

Expand All @@ -26,10 +27,19 @@ public class CacheEntry
/// </summary>
public partial class NatsCache : IBufferDistributedCache
{
// JetStream "stream not found" error code (JSStreamNotFoundErr), returned by GetStoreAsync when the
// KV bucket does not exist.
private const int StreamNotFoundErrCode = 10059;

// Compact binary serializer for the CacheEntry envelope (replaces the previous JSON+base64 format).
private static readonly CacheEntryBinarySerializer CacheEntrySerializer = CacheEntryBinarySerializer.Default;

// Non-zero LimitMarkerTTL enables per-key TTL (NATS 2.11+); the actual per-key expiry is set per Put.
private static readonly TimeSpan DefaultLimitMarkerTtl = TimeSpan.FromSeconds(1);

private readonly string _bucketName;
private readonly bool _createBucketIfNotExists;
private readonly Func<NatsKVConfig, NatsKVConfig>? _configureBucketOnCreate;
private readonly INatsCacheKeyEncoder _keyEncoder;
private readonly string _keyPrefix;
private readonly ILogger _logger;
Expand All @@ -49,6 +59,8 @@ public NatsCache(
_keyPrefix = string.IsNullOrEmpty(options.CacheKeyPrefix)
? string.Empty
: options.CacheKeyPrefix.TrimEnd('.');
_createBucketIfNotExists = options.CreateBucketIfNotExists;
_configureBucketOnCreate = options.ConfigureBucketOnCreate;
_lazyKvStore = CreateLazyKvStore();
_natsConnection = natsConnection;
_logger = logger ?? NullLogger<NatsCache>.Instance;
Expand Down Expand Up @@ -252,6 +264,33 @@ internal CacheEntry CreateCacheEntry(byte[] value, DistributedCacheEntryOptions
internal bool IsAbsolutelyExpired(CacheEntry entry) =>
entry.AbsoluteExpiration.HasValue && TimeProvider.GetUtcNow() >= entry.AbsoluteExpiration.Value;

// 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
// applied first, then the user hook (a record `with` transform) may override them. The bucket name is
// re-asserted afterward so the hook cannot retarget creation to a bucket other than the one GetKvStore
// reads from.
internal NatsKVConfig BuildBucketConfig()
{
var config = new NatsKVConfig(_bucketName)
{
History = 1, // required for well-defined per-key TTL behavior
LimitMarkerTTL = DefaultLimitMarkerTtl, // non-zero => enables per-key TTL (NATS 2.11+)
};

if (_configureBucketOnCreate != null)
{
config = _configureBucketOnCreate(config)
?? throw new InvalidOperationException(
$"{nameof(NatsCacheOptions)}.{nameof(NatsCacheOptions.ConfigureBucketOnCreate)} must not return null.");
if (config.Bucket != _bucketName)
{
config = config with { Bucket = _bucketName };
}
}

return config;
}

// 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) =>
Expand All @@ -264,13 +303,35 @@ private string GetEncodedKey(string key) =>
? _keyEncoder.Encode(key)
: _keyEncoder.Encode($"{_keyPrefix}.{key}");

// Returns the KV store for the configured bucket, creating it only if it does not already exist. An
// existing (operator-managed) bucket is used as-is and never modified — matching the
// CreateBucketIfNotExists contract — so CreateStoreAsync is used rather than CreateOrUpdateStoreAsync.
// GetStoreAsync is attempted first (rather than listing every bucket) so this stays O(1) and needs no
// stream-list permission; only a genuine "stream not found" triggers creation, and any other error
// (connectivity, auth) propagates unchanged. A create that loses a race with a concurrent creator
// surfaces as a failure and is retried via the reset-on-failure path in CreateLazyKvStore.
Comment thread
matthewdevenny marked this conversation as resolved.
private async Task<INatsKVStore> GetOrCreateStoreAsync(INatsKVContext kv)
{
try
{
return await kv.GetStoreAsync(_bucketName).ConfigureAwait(false);
}
catch (NatsJSApiException ex) when (ex.Error is { ErrCode: StreamNotFoundErrCode })
{
// Null-safe pattern: a null Error simply doesn't match, so the filter is false rather than throwing.
return await kv.CreateStoreAsync(BuildBucketConfig()).ConfigureAwait(false);
}
}

private Lazy<Task<INatsKVStore>> CreateLazyKvStore() =>
new(async () =>
{
try
{
var kv = _natsConnection.CreateKeyValueStoreContext();
var store = await kv.GetStoreAsync(_bucketName).ConfigureAwait(false);
var store = _createBucketIfNotExists
? await GetOrCreateStoreAsync(kv).ConfigureAwait(false)
: await kv.GetStoreAsync(_bucketName).ConfigureAwait(false);
LogConnected(_bucketName);
return store;
}
Expand Down
35 changes: 35 additions & 0 deletions src/NatsDistributedCache/NatsCacheOptions.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using Microsoft.Extensions.Options;
using NATS.Client.KeyValueStore;

namespace CodeCargo.Nats.DistributedCache
{
Expand All @@ -22,6 +23,40 @@ public class NatsCacheOptions : IOptions<NatsCacheOptions>
/// </summary>
public string? CacheKeyPrefix { get; set; }

/// <summary>
/// When <see langword="true"/>, the <see cref="BucketName"/> KV bucket is created the first time the
/// cache is used, if it does not already exist. Defaults to <see langword="false"/>, in which case the
/// bucket must be pre-created by the operator.
/// </summary>
/// <remarks>
/// Only a missing bucket is created; an existing bucket is used as-is and never modified, so
/// operator-managed settings are preserved. Creating a bucket requires JetStream stream-management
/// permissions.
/// </remarks>
public bool CreateBucketIfNotExists { get; set; }
Comment thread
matthewdevenny marked this conversation as resolved.

/// <summary>
/// Optional hook to customize the <see cref="NatsKVConfig"/> used when
/// <see cref="CreateBucketIfNotExists"/> is enabled and a missing bucket is created (for example
/// <c>Storage</c>, <c>NumberOfReplicas</c>, <c>MaxBytes</c>, or <c>MaxAge</c>). Ignored when
/// <see cref="CreateBucketIfNotExists"/> is <see langword="false"/> or the bucket already exists.
/// </summary>
/// <remarks>
/// <para>
/// <see cref="NatsKVConfig"/> is an immutable record, so the hook receives the pre-populated config
/// and returns a modified copy using a <c>with</c> expression, for example
/// <c>options.ConfigureBucketOnCreate = cfg =&gt; cfg with { Storage = NatsKVStorageType.Memory };</c>.
/// </para>
/// <para>
/// The library pre-populates the config with cache-appropriate defaults — <c>History = 1</c> and a
/// non-zero <c>LimitMarkerTTL</c>, both required for per-key TTL on NATS 2.11+ — before this hook runs,
/// so the hook can override any property. The <c>Bucket</c> name is always forced back to
/// <see cref="BucketName"/> afterward. Overriding <c>History</c> to a value other than 1, or clearing
/// <c>LimitMarkerTTL</c>, disables reliable per-key TTL.
/// </para>
/// </remarks>
public Func<NatsKVConfig, NatsKVConfig>? ConfigureBucketOnCreate { get; set; }

NatsCacheOptions IOptions<NatsCacheOptions>.Value => this;
}
}
Loading
Loading