diff --git a/README.md b/README.md index 81e0795..0c8761e 100644 --- a/README.md +++ b/README.md @@ -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 }); ``` ## Use with `HybridCache` @@ -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; @@ -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(); -var kvContext = natsConnection.CreateKeyValueStoreContext(); -await kvContext.CreateOrUpdateStoreAsync(new NatsKVConfig("cache") -{ - LimitMarkerTTL = TimeSpan.FromSeconds(1) -}); - // Start the host await host.RunAsync(); ``` @@ -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; @@ -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(); -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), diff --git a/src/NatsDistributedCache/NatsCache.cs b/src/NatsDistributedCache/NatsCache.cs index f35765a..a37d187 100644 --- a/src/NatsDistributedCache/NatsCache.cs +++ b/src/NatsDistributedCache/NatsCache.cs @@ -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; @@ -26,10 +27,19 @@ public class CacheEntry /// 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? _configureBucketOnCreate; private readonly INatsCacheKeyEncoder _keyEncoder; private readonly string _keyPrefix; private readonly ILogger _logger; @@ -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.Instance; @@ -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) => @@ -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. + private async Task 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> 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; } diff --git a/src/NatsDistributedCache/NatsCacheOptions.cs b/src/NatsDistributedCache/NatsCacheOptions.cs index 1fd79f6..67e95b8 100644 --- a/src/NatsDistributedCache/NatsCacheOptions.cs +++ b/src/NatsDistributedCache/NatsCacheOptions.cs @@ -1,4 +1,5 @@ using Microsoft.Extensions.Options; +using NATS.Client.KeyValueStore; namespace CodeCargo.Nats.DistributedCache { @@ -22,6 +23,40 @@ public class NatsCacheOptions : IOptions /// public string? CacheKeyPrefix { get; set; } + /// + /// When , the KV bucket is created the first time the + /// cache is used, if it does not already exist. Defaults to , in which case the + /// bucket must be pre-created by the operator. + /// + /// + /// 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. + /// + public bool CreateBucketIfNotExists { get; set; } + + /// + /// Optional hook to customize the used when + /// is enabled and a missing bucket is created (for example + /// Storage, NumberOfReplicas, MaxBytes, or MaxAge). Ignored when + /// is or the bucket already exists. + /// + /// + /// + /// is an immutable record, so the hook receives the pre-populated config + /// and returns a modified copy using a with expression, for example + /// options.ConfigureBucketOnCreate = cfg => cfg with { Storage = NatsKVStorageType.Memory };. + /// + /// + /// The library pre-populates the config with cache-appropriate defaults — History = 1 and a + /// non-zero LimitMarkerTTL, both required for per-key TTL on NATS 2.11+ — before this hook runs, + /// so the hook can override any property. The Bucket name is always forced back to + /// afterward. Overriding History to a value other than 1, or clearing + /// LimitMarkerTTL, disables reliable per-key TTL. + /// + /// + public Func? ConfigureBucketOnCreate { get; set; } + NatsCacheOptions IOptions.Value => this; } } diff --git a/test/IntegrationTests/Cache/BucketAutoCreationTests.cs b/test/IntegrationTests/Cache/BucketAutoCreationTests.cs new file mode 100644 index 0000000..4353ad4 --- /dev/null +++ b/test/IntegrationTests/Cache/BucketAutoCreationTests.cs @@ -0,0 +1,138 @@ +using CodeCargo.Nats.DistributedCache.TestUtils; +using Microsoft.Extensions.Caching.Distributed; +using Microsoft.Extensions.Logging; +using NATS.Client.KeyValueStore; +using NATS.Net; + +namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; + +/// +/// Verifies opt-in bucket auto-creation (issue #38). Uses buckets distinct from the shared "cache" +/// bucket so they never collide with 's KV_cache purge, and deletes any +/// bucket it touches on teardown so nothing leaks across the collection lifetime. +/// +[Collection(NatsCollection.Name)] +public class BucketAutoCreationTests(NatsIntegrationFixture fixture) : IAsyncLifetime +{ + private const string Key = "auto-create-key"; + + private readonly List _serviceProviders = new(); + private readonly List _bucketsToDelete = new(); + + [Fact] + public async Task FirstOperation_CreatesMissingBucket_WithHistoryOneAndLimitMarkerTtl() + { + const string bucketName = "auto-created-cache"; + var ct = TestContext.Current.CancellationToken; + + // Memory storage keeps the test light and exercises the ConfigureBucketOnCreate hook end-to-end. + var cache = BuildCache(bucketName, cfg => cfg with { Storage = NatsKVStorageType.Memory }); + var value = new byte[] { 1, 2, 3 }; + + // The first cache operation triggers creation of the missing bucket. + await cache.SetAsync( + Key, + value, + new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromMinutes(1) }, + ct); + + Assert.Equal(value, await cache.GetAsync(Key, ct)); + + // The bucket now exists with the cache-required config: History = 1 and a non-zero LimitMarkerTTL. + var status = await GetStatusAsync(bucketName, ct); + Assert.Equal(bucketName, status.Bucket); + Assert.NotEqual(TimeSpan.Zero, status.LimitMarkerTTL); + Assert.Equal(1, status.Info.Config.MaxMsgsPerSubject); // KV History maps to MaxMsgsPerSubject + } + + [Fact] + public async Task ExistingBucket_IsUsedAsIs_AndNeverModified() + { + const string bucketName = "operator-managed-cache"; + var ct = TestContext.Current.CancellationToken; + + // Operator pre-creates the bucket with a distinctive LimitMarkerTTL (5s). + var kv = fixture.NatsConnection.CreateKeyValueStoreContext(); + _bucketsToDelete.Add(bucketName); + await kv.CreateStoreAsync( + new NatsKVConfig(bucketName) + { + History = 1, + LimitMarkerTTL = TimeSpan.FromSeconds(5), + Storage = NatsKVStorageType.Memory, + }, + ct); + + // Enable auto-create with a hook that WOULD set a different LimitMarkerTTL (1s) if it created the bucket. + var cache = BuildCache( + bucketName, + cfg => cfg with { LimitMarkerTTL = TimeSpan.FromSeconds(1), Storage = NatsKVStorageType.Memory }, + trackForDeletion: false); + + await cache.SetAsync( + Key, + new byte[] { 9 }, + new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromMinutes(1) }, + ct); + + // The pre-existing bucket must be left untouched: LimitMarkerTTL is still the operator's 5s, not 1s. + var status = await GetStatusAsync(bucketName, ct); + Assert.Equal(TimeSpan.FromSeconds(5), status.LimitMarkerTTL); + } + + public ValueTask InitializeAsync() => ValueTask.CompletedTask; + + public async ValueTask DisposeAsync() + { + var kv = fixture.NatsConnection.CreateKeyValueStoreContext(); + foreach (var bucket in _bucketsToDelete) + { + try + { + await kv.DeleteStoreAsync(bucket, TestContext.Current.CancellationToken); + } + catch + { + // Best-effort cleanup; memory-backed buckets are discarded when the server stops regardless. + } + } + + foreach (var sp in _serviceProviders) + { + await sp.DisposeAsync(); + } + + GC.SuppressFinalize(this); + } + + private IDistributedCache BuildCache( + string bucketName, + Func configureBucket, + bool trackForDeletion = true) + { + var services = new ServiceCollection(); + services.AddLogging(); + fixture.ConfigureServices(services); + services.AddNatsDistributedCache(options => + { + options.BucketName = bucketName; + options.CreateBucketIfNotExists = true; + options.ConfigureBucketOnCreate = configureBucket; + }); + + var serviceProvider = services.BuildServiceProvider(); + _serviceProviders.Add(serviceProvider); + if (trackForDeletion) + { + _bucketsToDelete.Add(bucketName); + } + + return serviceProvider.GetRequiredService(); + } + + private async Task GetStatusAsync(string bucketName, CancellationToken ct) + { + var kv = fixture.NatsConnection.CreateKeyValueStoreContext(); + return await (await kv.GetStoreAsync(bucketName, ct)).GetStatusAsync(ct); + } +} diff --git a/test/UnitTests/Cache/BucketConfigUnitTests.cs b/test/UnitTests/Cache/BucketConfigUnitTests.cs new file mode 100644 index 0000000..7041716 --- /dev/null +++ b/test/UnitTests/Cache/BucketConfigUnitTests.cs @@ -0,0 +1,90 @@ +using Microsoft.Extensions.Options; +using Moq; +using NATS.Client.Core; +using NATS.Client.KeyValueStore; + +namespace CodeCargo.Nats.DistributedCache.UnitTests.Cache; + +public class BucketConfigUnitTests +{ + private const string BucketName = "cache"; + + [Fact] + public void BuildBucketConfig_SetsBucketName() + { + var config = CreateCache().BuildBucketConfig(); + + Assert.Equal(BucketName, config.Bucket); + } + + [Fact] + public void BuildBucketConfig_DefaultsHistoryToOne() + { + var config = CreateCache().BuildBucketConfig(); + + Assert.Equal(1, config.History); + } + + [Fact] + public void BuildBucketConfig_SetsNonZeroLimitMarkerTtl() + { + var config = CreateCache().BuildBucketConfig(); + + Assert.NotEqual(TimeSpan.Zero, config.LimitMarkerTTL); + Assert.Equal(TimeSpan.FromSeconds(1), config.LimitMarkerTTL); + } + + [Fact] + public void BuildBucketConfig_InvokesConfigureBucketOnCreateHook() + { + var config = CreateCache(o => o.ConfigureBucketOnCreate = cfg => cfg with + { + Storage = NatsKVStorageType.Memory, + NumberOfReplicas = 3, + }).BuildBucketConfig(); + + Assert.Equal(NatsKVStorageType.Memory, config.Storage); + Assert.Equal(3, config.NumberOfReplicas); + } + + [Fact] + public void BuildBucketConfig_HookCanOverrideDefaults() + { + var config = CreateCache(o => o.ConfigureBucketOnCreate = cfg => cfg with + { + History = 5, + LimitMarkerTTL = TimeSpan.FromSeconds(30), + }).BuildBucketConfig(); + + // The hook runs after the defaults, so it wins (documented as caller responsibility for TTL correctness). + Assert.Equal(5, config.History); + Assert.Equal(TimeSpan.FromSeconds(30), config.LimitMarkerTTL); + } + + [Fact] + public void BuildBucketConfig_ReassertsBucketName_WhenHookChangesIt() + { + // The hook must not be able to retarget creation to a different bucket than the cache reads from. + var config = CreateCache(o => o.ConfigureBucketOnCreate = cfg => cfg with { Bucket = "some-other-bucket" }) + .BuildBucketConfig(); + + Assert.Equal(BucketName, config.Bucket); + } + + [Fact] + public void BuildBucketConfig_ThrowsClearException_WhenHookReturnsNull() + { + var cache = CreateCache(o => o.ConfigureBucketOnCreate = _ => null!); + + var ex = Assert.Throws(() => cache.BuildBucketConfig()); + Assert.Contains(nameof(NatsCacheOptions.ConfigureBucketOnCreate), ex.Message); + } + + // BuildBucketConfig never touches the connection, so a bare mock is sufficient and no server is needed. + private static NatsCache CreateCache(Action? configure = null) + { + var options = new NatsCacheOptions { BucketName = BucketName, CreateBucketIfNotExists = true }; + configure?.Invoke(options); + return new NatsCache(Options.Create(options), new Mock().Object); + } +} diff --git a/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs b/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs index 2af5542..352ba69 100644 --- a/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs +++ b/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs @@ -5,6 +5,7 @@ using Microsoft.Extensions.Time.Testing; using Moq; using NATS.Client.Core; +using NATS.Client.KeyValueStore; namespace CodeCargo.Nats.DistributedCache.UnitTests.Extensions; @@ -79,6 +80,41 @@ public void AddNatsCache_SetsCacheOptions() Assert.Equal("cache", options.BucketName); } + [Fact] + public void AddNatsCache_SetsBucketCreationOptions() + { + // Arrange + var services = new ServiceCollection(); + services.AddSingleton(_mockNatsConnection.Object); + Func configureBucket = cfg => cfg; + + // Act + services.AddNatsDistributedCache(options => + { + options.BucketName = "cache"; + options.CreateBucketIfNotExists = true; + options.ConfigureBucketOnCreate = configureBucket; + }); + + // Build the provider to verify options + var provider = services.BuildServiceProvider(); + var options = provider.GetRequiredService>().Value; + + // Assert + Assert.True(options.CreateBucketIfNotExists); + Assert.Same(configureBucket, options.ConfigureBucketOnCreate); + } + + [Fact] + public void NatsCacheOptions_BucketCreation_DefaultsToDisabled() + { + // Guards the "disabled by default / existing behavior unchanged" contract for issue #38. + var options = new NatsCacheOptions(); + + Assert.False(options.CreateBucketIfNotExists); + Assert.Null(options.ConfigureBucketOnCreate); + } + [Fact] public void AddNatsCache_UsesCacheOptionsAction() { diff --git a/util/ReadmeExample/DistributedCache.Example.cs b/util/ReadmeExample/DistributedCache.Example.cs index ddedcdc..28f7e74 100644 --- a/util/ReadmeExample/DistributedCache.Example.cs +++ b/util/ReadmeExample/DistributedCache.Example.cs @@ -1,8 +1,6 @@ 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; @@ -25,16 +23,15 @@ public static async Task Run(string[] args) 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(); - var kvContext = natsConnection.CreateKeyValueStoreContext(); - await kvContext.CreateOrUpdateStoreAsync(new NatsKVConfig("cache") { LimitMarkerTTL = TimeSpan.FromSeconds(1) }); - // Start the host await host.RunAsync(); } diff --git a/util/ReadmeExample/DistributedCache.cs b/util/ReadmeExample/DistributedCache.cs index 616e743..48f43de 100644 --- a/util/ReadmeExample/DistributedCache.cs +++ b/util/ReadmeExample/DistributedCache.cs @@ -5,10 +5,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; -using NATS.Client.Core; -using NATS.Client.KeyValueStore; using NATS.Extensions.Microsoft.DependencyInjection; -using NATS.Net; namespace CodeCargo.ReadmeExample; @@ -47,6 +44,9 @@ public static async Task RunAsync(string[] args) services.AddNatsDistributedCache(options => { options.BucketName = "cache"; + + // Create the KV bucket on first use if it doesn't already exist. + options.CreateBucketIfNotExists = true; }); services.AddScoped(); @@ -55,14 +55,6 @@ public static async Task RunAsync(string[] args) var host = builder.Build(); var lifetime = host.Services.GetRequiredService(); - // Ensure that the KV Store is created - Console.WriteLine("Creating KV store..."); - var natsConnection = host.Services.GetRequiredService(); - var kvContext = natsConnection.CreateKeyValueStoreContext(); - await kvContext.CreateOrUpdateStoreAsync( - new NatsKVConfig("cache") { LimitMarkerTTL = TimeSpan.FromSeconds(1) }, startupCts.Token); - Console.WriteLine("KV store created"); - // Start the host Console.WriteLine("Starting app..."); using var appCts = new CancellationTokenSource(); diff --git a/util/ReadmeExample/HybridCache.Example.cs b/util/ReadmeExample/HybridCache.Example.cs index 53a8cff..cb514a6 100644 --- a/util/ReadmeExample/HybridCache.Example.cs +++ b/util/ReadmeExample/HybridCache.Example.cs @@ -1,8 +1,6 @@ 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; @@ -25,19 +23,15 @@ public static async Task Run(string[] args) 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(); - var kvContext = natsConnection.CreateKeyValueStoreContext(); - await kvContext.CreateOrUpdateStoreAsync(new NatsKVConfig("cache") - { - LimitMarkerTTL = TimeSpan.FromSeconds(1) - }); - // Start the host await host.RunAsync(); } diff --git a/util/ReadmeExample/HybridCache.cs b/util/ReadmeExample/HybridCache.cs index ffbe8af..cd9fe19 100644 --- a/util/ReadmeExample/HybridCache.cs +++ b/util/ReadmeExample/HybridCache.cs @@ -5,10 +5,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; -using NATS.Client.Core; -using NATS.Client.KeyValueStore; using NATS.Extensions.Microsoft.DependencyInjection; -using NATS.Net; namespace CodeCargo.ReadmeExample; @@ -48,6 +45,9 @@ public static async Task RunAsync(string[] args) services.AddNatsHybridCache(options => { options.BucketName = "cache"; + + // Create the KV bucket on first use if it doesn't already exist. + options.CreateBucketIfNotExists = true; }); services.AddScoped(); @@ -56,14 +56,6 @@ public static async Task RunAsync(string[] args) var host = builder.Build(); var lifetime = host.Services.GetRequiredService(); - // Ensure that the KV Store is created - Console.WriteLine("Creating KV store..."); - var natsConnection = host.Services.GetRequiredService(); - var kvContext = natsConnection.CreateKeyValueStoreContext(); - await kvContext.CreateOrUpdateStoreAsync( - new NatsKVConfig("cache") { LimitMarkerTTL = TimeSpan.FromSeconds(1) }, startupCts.Token); - Console.WriteLine("KV store created"); - // Start the host Console.WriteLine("Starting app..."); using var appCts = new CancellationTokenSource();