-
Notifications
You must be signed in to change notification settings - Fork 2
#45 Add purge-by-prefix cache maintenance helper #60
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
4 commits
Select commit
Hold shift + click to select a range
5f2ba82
#45 Add purge-by-prefix cache maintenance helper
matthewdevenny 79a59e2
#45 Remove unused using in NatsCache.Maintenance
matthewdevenny ba21828
#45 Purge by prefix via a single JetStream stream purge
matthewdevenny b52c9c3
#45 Clarify PurgeByPrefixAsync return-count semantics in docs
matthewdevenny File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,42 @@ | ||
| namespace CodeCargo.Nats.DistributedCache; | ||
|
|
||
| /// <summary> | ||
| /// Maintenance operations for a NATS-backed distributed cache that fall outside the | ||
| /// <see cref="Microsoft.Extensions.Caching.Distributed.IDistributedCache"/> surface. | ||
| /// </summary> | ||
| /// <remarks> | ||
| /// Resolve this from dependency injection to reach the same cache instance registered by | ||
| /// <see cref="NatsDistributedCacheExtensions.AddNatsDistributedCache"/>. | ||
| /// </remarks> | ||
| public interface INatsCacheMaintenance | ||
| { | ||
| /// <summary> | ||
| /// Purges every cache entry whose key begins with <paramref name="prefix"/>, treated as a sub-prefix | ||
| /// beneath the configured <see cref="NatsCacheOptions.CacheKeyPrefix"/> (entries stored as | ||
| /// <c>{prefix}.{key}</c> are matched). This is a bulk, irreversible maintenance operation intended for | ||
| /// scenarios such as evicting every entry belonging to a single tenant. | ||
| /// </summary> | ||
| /// <remarks> | ||
| /// The matching entries are removed in a single JetStream stream purge (a subject-filtered purge of the | ||
| /// bucket's backing <c>KV_<bucket></c> stream), so the messages are deleted outright rather than | ||
| /// left as purge-marker tombstones. This requires JetStream stream-purge permission on the | ||
| /// <c>KV_<bucket></c> stream, in addition to the ordinary KV access the cache uses. | ||
| /// </remarks> | ||
| /// <param name="prefix"> | ||
| /// The key sub-prefix to purge, relative to the configured cache key prefix. Must be non-empty and not | ||
| /// consist solely of whitespace or <c>'.'</c> characters, so a scoped purge cannot collapse into a purge | ||
| /// of the entire prefix space (or, when no <see cref="NatsCacheOptions.CacheKeyPrefix"/> is configured, | ||
| /// the whole bucket). | ||
| /// </param> | ||
| /// <param name="cancellationToken">A token used to cancel the operation.</param> | ||
| /// <returns> | ||
| /// The number of stream messages purged beneath the prefix. For the cache's single-revision | ||
| /// (<c>History = 1</c>) 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. | ||
| /// </returns> | ||
| /// <exception cref="System.ArgumentException"> | ||
| /// <paramref name="prefix"/> is null, empty, whitespace, or consists solely of <c>'.'</c> characters. | ||
| /// </exception> | ||
| Task<long> PurgeByPrefixAsync(string prefix, CancellationToken cancellationToken = default); | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,60 @@ | ||
| using NATS.Client.JetStream.Models; | ||
|
|
||
| namespace CodeCargo.Nats.DistributedCache; | ||
|
|
||
| public partial class NatsCache : INatsCacheMaintenance | ||
| { | ||
| /// <inheritdoc /> | ||
| public async Task<long> 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.<bucket>.<encodedKey>', 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.<bucket>.{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_<bucket>') 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; | ||
| } | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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<INatsCacheMaintenance>(); | ||
|
|
||
| [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)); | ||
| } | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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<ArgumentException>(() => 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<INatsConnection>().Object); | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.