diff --git a/README.md b/README.md index 84efa0a..f953225 100644 --- a/README.md +++ b/README.md @@ -159,6 +159,52 @@ Notes: - Overriding `History` (away from `1`) or clearing `LimitMarkerTTL` in `ConfigureBucketOnCreate` disables reliable per-key TTL. +## Multi-tenant Key Prefixes and Bulk Purge + +Set `CacheKeyPrefix` to partition a single bucket across apps, services, or tenants. Every key is then +stored as `.`, so several consumers can safely share one KV bucket: + +```csharp +services.AddNatsDistributedCache(options => +{ + options.BucketName = "cache"; + options.CacheKeyPrefix = "tenant-42"; +}); +``` + +To evict every entry beneath a sub-prefix — for example, all keys for one tenant — resolve +`INatsCacheMaintenance` and call `PurgeByPrefixAsync`. It is implemented by the same singleton that backs +`IDistributedCache`, so it shares the bucket, key prefix, and key encoding: + +```csharp +using CodeCargo.Nats.DistributedCache; + +var maintenance = serviceProvider.GetRequiredService(); + +// Purges every key stored as "orders.<...>" beneath the configured CacheKeyPrefix. +// Returns the number of stream messages purged (see the count note below). +long purged = await maintenance.PurgeByPrefixAsync("orders"); +``` + +The supplied prefix is relative to `CacheKeyPrefix` and matches keys hierarchically: `"orders"` purges +`orders.a` and `orders.a.b`, but not `orders-archive.a`. + +Notes: + +- This is a bulk, **irreversible** maintenance operation, not a per-request cache call. It removes every + matching entry in a single JetStream stream purge (a subject-filtered purge of the bucket's backing + `KV_` stream), so the messages are deleted outright rather than left as purge-marker tombstones. +- This requires JetStream stream-purge permission on the `KV_` stream, in addition to the ordinary + KV access the cache uses. +- The prefix must be non-empty and not consist solely of whitespace or `.` characters; otherwise + `PurgeByPrefixAsync` throws `ArgumentException`. This guards against accidentally purging the entire + prefix space (or, with no `CacheKeyPrefix`, the whole bucket). +- Purging is scoped to children of the prefix (`.<...>`); the cache never stores a bare key equal + to the prefix itself. +- The returned count is the number of stream messages purged. For the cache's single-revision + (`History = 1`) buckets this equals the number of live entries removed, but it can also include + not-yet-compacted delete markers left by earlier evictions, so treat it as an approximate count. + ## Controlling Expiration Timing Expiration is computed from a [`TimeProvider`](https://learn.microsoft.com/dotnet/api/system.timeprovider), diff --git a/src/NatsDistributedCache/INatsCacheMaintenance.cs b/src/NatsDistributedCache/INatsCacheMaintenance.cs new file mode 100644 index 0000000..5caf2f2 --- /dev/null +++ b/src/NatsDistributedCache/INatsCacheMaintenance.cs @@ -0,0 +1,42 @@ +namespace CodeCargo.Nats.DistributedCache; + +/// +/// Maintenance operations for a NATS-backed distributed cache that fall outside the +/// surface. +/// +/// +/// Resolve this from dependency injection to reach the same cache instance registered by +/// . +/// +public interface INatsCacheMaintenance +{ + /// + /// Purges every cache entry whose key begins with , treated as a sub-prefix + /// beneath the configured (entries stored as + /// {prefix}.{key} are matched). This is a bulk, irreversible maintenance operation intended for + /// scenarios such as evicting every entry belonging to a single tenant. + /// + /// + /// The matching entries are removed in a single JetStream stream purge (a subject-filtered purge of the + /// bucket's backing KV_<bucket> stream), so the messages are deleted outright rather than + /// left as purge-marker tombstones. This requires JetStream stream-purge permission on the + /// KV_<bucket> stream, in addition to the ordinary KV access the cache uses. + /// + /// + /// The key sub-prefix to purge, relative to the configured cache key prefix. Must be non-empty and not + /// consist solely of whitespace or '.' characters, so a scoped purge cannot collapse into a purge + /// of the entire prefix space (or, when no is configured, + /// the whole bucket). + /// + /// A token used to cancel the operation. + /// + /// The number of stream messages purged beneath the prefix. For the cache's single-revision + /// (History = 1) buckets this is the number of live entries removed, though it can also include + /// not-yet-compacted delete markers left by earlier evictions, so treat it as an approximate count rather + /// than an exact live-entry total. + /// + /// + /// is null, empty, whitespace, or consists solely of '.' characters. + /// + Task PurgeByPrefixAsync(string prefix, CancellationToken cancellationToken = default); +} diff --git a/src/NatsDistributedCache/NatsCache.Maintenance.cs b/src/NatsDistributedCache/NatsCache.Maintenance.cs new file mode 100644 index 0000000..8e9cd7e --- /dev/null +++ b/src/NatsDistributedCache/NatsCache.Maintenance.cs @@ -0,0 +1,60 @@ +using NATS.Client.JetStream.Models; + +namespace CodeCargo.Nats.DistributedCache; + +public partial class NatsCache : INatsCacheMaintenance +{ + /// + public async Task PurgeByPrefixAsync(string prefix, CancellationToken cancellationToken = default) + { + // Reject an empty/whitespace prefix outright. Without it the subject filter below would collapse to + // the entire configured-prefix space (or, with no CacheKeyPrefix, the whole bucket), turning a scoped + // maintenance call into an accidental full purge. + if (string.IsNullOrWhiteSpace(prefix)) + { + throw new ArgumentException("Prefix must not be null, empty, or whitespace.", nameof(prefix)); + } + + // Trailing '.' is trimmed to mirror the constructor's normalization of _keyPrefix, and so the raw + // prefix never ends in '.'. That matters because the key encoder only escapes a trailing '.' at the + // very end of a string; leaving one on the prefix would encode the separator differently here than it + // is encoded mid-key inside a stored full key, and the filter would stop matching. + var trimmedPrefix = prefix.TrimEnd('.'); + if (trimmedPrefix.Length == 0) + { + throw new ArgumentException("Prefix must not consist solely of '.' characters.", nameof(prefix)); + } + + // Compose the caller's sub-prefix with the configured CacheKeyPrefix exactly as GetEncodedKey composes + // a full key, so the raw prefix is a genuine leading segment of every stored key beneath it. + var rawPrefix = string.IsNullOrEmpty(_keyPrefix) ? trimmedPrefix : $"{_keyPrefix}.{trimmedPrefix}"; + + // The encoder URL-encodes per character and leaves '.' unescaped (it is RFC 3986 unreserved) while + // escaping the NATS wildcards '*' and '>', so the encoded prefix is a byte-for-byte leading segment of + // every encoded full key beneath it and cannot itself contain a wildcard. A KV key is stored on the + // subject '$KV..', so appending the NATS multi-token wildcard '>' after the '.' + // separator yields a subject filter that matches every child of the prefix namespace. The cache never + // stores a bare-prefix key (keys are always '{prefix}.{userKey}'), so scoping the filter to children + // with '$KV..{encodedPrefix}.>' -- rather than also matching a key exactly equal to the prefix + // -- is correct. + var encodedPrefix = _keyEncoder.Encode(rawPrefix); + var subjectFilter = $"$KV.{_bucketName}.{encodedPrefix}.>"; + + var store = await GetKvStore().ConfigureAwait(false); + + // Purge every matching message from the KV bucket's backing JetStream stream ('KV_') in a + // single server round-trip, scoped to the prefix's subject space. This deletes the messages outright + // rather than enumerating the keys and issuing a per-key KV purge -- which would take N round-trips and + // leave purge-marker tombstones that linger until a PurgeDeletes() compaction. With the bucket's + // single-revision history (History = 1) there is exactly one message per live key, so the returned + // Purged count is the number of entries removed. + var response = await store.JetStreamContext + .PurgeStreamAsync( + $"KV_{_bucketName}", + new StreamPurgeRequest { Filter = subjectFilter }, + cancellationToken) + .ConfigureAwait(false); + + return response.Purged; + } +} diff --git a/src/NatsDistributedCache/NatsDistributedCacheExtensions.cs b/src/NatsDistributedCache/NatsDistributedCacheExtensions.cs index fff7e0a..0ffb5f9 100644 --- a/src/NatsDistributedCache/NatsDistributedCacheExtensions.cs +++ b/src/NatsDistributedCache/NatsDistributedCacheExtensions.cs @@ -30,7 +30,11 @@ public static IServiceCollection AddNatsDistributedCache( .Configure(configureOptions) .Validate(o => !string.IsNullOrWhiteSpace(o.BucketName), NatsCacheOptions.BucketNameRequiredMessage) .ValidateOnStart(); - services.AddSingleton(sp => + + // Register the concrete NatsCache as the single instance, then forward both public service surfaces + // to it. IDistributedCache and INatsCacheMaintenance must resolve to the SAME instance so a purge and + // a cache hit share one KV store, key prefix, and key encoder. + services.AddSingleton(sp => { var optionsAccessor = sp.GetRequiredService>(); var natsConnection = connectionServiceKey == null @@ -52,6 +56,8 @@ public static IServiceCollection AddNatsDistributedCache( MeterFactory = meterFactory, }; }); + services.AddSingleton(sp => sp.GetRequiredService()); + services.AddSingleton(sp => sp.GetRequiredService()); return services; } diff --git a/test/IntegrationTests/Cache/PurgeByPrefixTests.cs b/test/IntegrationTests/Cache/PurgeByPrefixTests.cs new file mode 100644 index 0000000..ef0373a --- /dev/null +++ b/test/IntegrationTests/Cache/PurgeByPrefixTests.cs @@ -0,0 +1,73 @@ +using Microsoft.Extensions.Caching.Distributed; + +namespace CodeCargo.Nats.DistributedCache.IntegrationTests.Cache; + +public class PurgeByPrefixTests(NatsIntegrationFixture fixture) : TestBase(fixture) +{ + // The maintenance surface and IDistributedCache resolve to the same NatsCache singleton, so seeding via + // DistributedCache and purging via Maintenance operate on one bucket, key prefix, and key encoder. + private INatsCacheMaintenance Maintenance => ServiceProvider.GetRequiredService(); + + [Fact] + public async Task PurgeByPrefixRemovesMatchingKeysAndLeavesOthers() + { + var token = TestContext.Current.CancellationToken; + var value = new byte[] { 1 }; + + string[] aKeys = ["A.one", "A.two", "A.three"]; + string[] bKeys = ["B.one", "B.two"]; + + foreach (var key in aKeys.Concat(bKeys)) + { + await DistributedCache.SetAsync(key, value, token); + } + + var purged = await Maintenance.PurgeByPrefixAsync("A", token); + + Assert.Equal(aKeys.Length, (int)purged); + foreach (var key in aKeys) + { + Assert.Null(await DistributedCache.GetAsync(key, token)); + } + + foreach (var key in bKeys) + { + Assert.NotNull(await DistributedCache.GetAsync(key, token)); + } + + // A second purge finds nothing left under the prefix. + var purgedAgain = await Maintenance.PurgeByPrefixAsync("A", token); + Assert.Equal(0, (int)purgedAgain); + } + + [Fact] + public async Task PurgeByPrefixMatchesKeysThatRequireEncoding() + { + // '#' is not a valid unencoded key character, so the encoder escapes it (# -> %23 -> =23). This + // exercises the encoding nuance the subject filter relies on: the encoder leaves the '.' separator + // unescaped, so the encoded prefix stays a leading segment of every encoded full key beneath it and + // NATS subject-wildcard filtering by prefix keeps working. + var token = TestContext.Current.CancellationToken; + var value = new byte[] { 1 }; + + string[] targetKeys = ["tenant#1.one", "tenant#1.two"]; + const string controlKey = "tenant#2.one"; // shares a stem but a different tenant; must survive + + foreach (var key in targetKeys) + { + await DistributedCache.SetAsync(key, value, token); + } + + await DistributedCache.SetAsync(controlKey, value, token); + + var purged = await Maintenance.PurgeByPrefixAsync("tenant#1", token); + + Assert.Equal(targetKeys.Length, (int)purged); + foreach (var key in targetKeys) + { + Assert.Null(await DistributedCache.GetAsync(key, token)); + } + + Assert.NotNull(await DistributedCache.GetAsync(controlKey, token)); + } +} diff --git a/test/UnitTests/Cache/PurgeByPrefixUnitTests.cs b/test/UnitTests/Cache/PurgeByPrefixUnitTests.cs new file mode 100644 index 0000000..da610f0 --- /dev/null +++ b/test/UnitTests/Cache/PurgeByPrefixUnitTests.cs @@ -0,0 +1,32 @@ +using Microsoft.Extensions.Options; +using Moq; +using NATS.Client.Core; + +namespace CodeCargo.Nats.DistributedCache.UnitTests.Cache; + +public class PurgeByPrefixUnitTests +{ + private const string BucketName = "cache"; + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData(" ")] + [InlineData(" ")] + [InlineData(".")] + [InlineData("..")] + public async Task PurgeByPrefixAsync_RejectsEmptyWhitespaceOrDotOnlyPrefix(string? prefix) + { + // The guard runs before the KV store is resolved, so a bare mock connection is enough (no server + // needed). Rejecting these values prevents a scoped purge from collapsing into a purge of the whole + // configured-prefix space (or, with no CacheKeyPrefix, the entire bucket). + var cache = CreateCache(); + + var ex = await Assert.ThrowsAsync(() => cache.PurgeByPrefixAsync(prefix!)); + Assert.Equal("prefix", ex.ParamName); + } + + // The guard never touches the connection, so a bare mock is sufficient and no server is needed. + private static NatsCache CreateCache() => + new(Options.Create(new NatsCacheOptions { BucketName = BucketName }), new Mock().Object); +} diff --git a/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs b/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs index 352ba69..f750835 100644 --- a/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs +++ b/test/UnitTests/Extensions/NatsDistributedCacheExtensionsTests.cs @@ -33,6 +33,41 @@ public void AddNatsCache_RegistersDistributedCacheAsSingleton() Assert.Equal(ServiceLifetime.Singleton, distributedCache.Lifetime); } + [Fact] + public void AddNatsCache_RegistersMaintenanceAsSingleton() + { + // Arrange + var services = new ServiceCollection(); + services.AddSingleton(_mockNatsConnection.Object); + + // Act + services.AddNatsDistributedCache(options => options.BucketName = "cache"); + + // Assert + var maintenance = services.FirstOrDefault(desc => desc.ServiceType == typeof(INatsCacheMaintenance)); + Assert.NotNull(maintenance); + Assert.Equal(ServiceLifetime.Singleton, maintenance.Lifetime); + } + + [Fact] + public void AddNatsCache_ResolvesMaintenanceAsSameInstanceAsDistributedCache() + { + // Arrange + var services = new ServiceCollection(); + services.AddSingleton(_mockNatsConnection.Object); + services.AddNatsDistributedCache(options => options.BucketName = "cache"); + + // Act + var provider = services.BuildServiceProvider(); + var distributedCache = provider.GetRequiredService(); + var maintenance = provider.GetRequiredService(); + + // Assert - both surfaces are backed by the one NatsCache singleton, so a purge and a cache read + // share the same KV store, key prefix, and key encoder. + Assert.Same(distributedCache, maintenance); + Assert.IsType(maintenance); + } + [Fact] public void AddNatsCache_ReplacesPreviouslyUserRegisteredServices() {