From 9e1bf9ebad9f731ccaba722a7ec721b77c4750e3 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Sun, 23 Mar 2025 07:32:58 +0000 Subject: [PATCH 1/3] Add KV marker reason check --- src/NATS.Client.KeyValueStore/NatsKVStore.cs | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/src/NATS.Client.KeyValueStore/NatsKVStore.cs b/src/NATS.Client.KeyValueStore/NatsKVStore.cs index 7be4cba92..9c85a5583 100644 --- a/src/NATS.Client.KeyValueStore/NatsKVStore.cs +++ b/src/NATS.Client.KeyValueStore/NatsKVStore.cs @@ -44,6 +44,7 @@ public class NatsKVStore : INatsKVStore private const string NatsSequence = "Nats-Sequence"; private const string NatsTimeStamp = "Nats-Time-Stamp"; private const string NatsTtl = "Nats-TTL"; + private const string NatsMarkerReason = "Nats-Marker-Reason"; private static readonly Regex ValidKeyRegex = new(pattern: @"\A[-/_=\.a-zA-Z0-9]+\z", RegexOptions.Compiled); private static readonly NatsKVException MissingSequenceHeaderException = new("Missing sequence header"); private static readonly NatsKVException MissingTimestampHeaderException = new("Missing timestamp header"); @@ -421,6 +422,22 @@ public async ValueTask>> TryGetEntryAsync(string ke if (!Enum.TryParse(operationValues[0], ignoreCase: true, out operation)) return InvalidOperationException; } + else if (headers.TryGetValue(NatsMarkerReason, out var markerReasonValues)) + { + var reason = markerReasonValues.Last(); + if (reason is "MaxAge" or "Purge") + { + operation = NatsKVOperation.Purge; + } + else if (reason is "Remove") + { + operation = NatsKVOperation.Del; + } + else + { + return InvalidOperationException; + } + } if (operation is NatsKVOperation.Del or NatsKVOperation.Purge) { From 41e6f8d169d508c7e1091bb89b4e98c8a45d6e69 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Sun, 23 Mar 2025 12:32:23 +0000 Subject: [PATCH 2/3] Add KV watcher support for marker reason --- .../Internal/NatsKVWatcher.cs | 12 ++++++++++++ src/NATS.Client.KeyValueStore/NatsKVStore.cs | 2 +- 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/src/NATS.Client.KeyValueStore/Internal/NatsKVWatcher.cs b/src/NATS.Client.KeyValueStore/Internal/NatsKVWatcher.cs index ad2812eeb..7c1475a2c 100644 --- a/src/NATS.Client.KeyValueStore/Internal/NatsKVWatcher.cs +++ b/src/NATS.Client.KeyValueStore/Internal/NatsKVWatcher.cs @@ -207,6 +207,18 @@ private async Task CommandLoop() _ => operation, }; } + else if (headers.TryGetValue(NatsKVStore.NatsMarkerReason, out var markerReasonValues)) + { + var reason = markerReasonValues.Last(); + if (reason is "MaxAge" or "Purge") + { + operation = NatsKVOperation.Purge; + } + else if (reason is "Remove") + { + operation = NatsKVOperation.Del; + } + } if (headers is { Code: 100, MessageText: "FlowControl Request" }) { diff --git a/src/NATS.Client.KeyValueStore/NatsKVStore.cs b/src/NATS.Client.KeyValueStore/NatsKVStore.cs index 9c85a5583..cf1cd0618 100644 --- a/src/NATS.Client.KeyValueStore/NatsKVStore.cs +++ b/src/NATS.Client.KeyValueStore/NatsKVStore.cs @@ -34,6 +34,7 @@ public enum NatsKVOperation /// public class NatsKVStore : INatsKVStore { + internal const string NatsMarkerReason = "Nats-Marker-Reason"; private const string NatsExpectedLastSubjectSequence = "Nats-Expected-Last-Subject-Sequence"; private const string KVOperation = "KV-Operation"; private const string NatsRollup = "Nats-Rollup"; @@ -44,7 +45,6 @@ public class NatsKVStore : INatsKVStore private const string NatsSequence = "Nats-Sequence"; private const string NatsTimeStamp = "Nats-Time-Stamp"; private const string NatsTtl = "Nats-TTL"; - private const string NatsMarkerReason = "Nats-Marker-Reason"; private static readonly Regex ValidKeyRegex = new(pattern: @"\A[-/_=\.a-zA-Z0-9]+\z", RegexOptions.Compiled); private static readonly NatsKVException MissingSequenceHeaderException = new("Missing sequence header"); private static readonly NatsKVException MissingTimestampHeaderException = new("Missing timestamp header"); From e220cd8b6c8a75640ac5611d8757ebbde46ba280 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Sun, 23 Mar 2025 17:39:51 +0000 Subject: [PATCH 3/3] Add test for SubjectDeleteMarkerTTL --- .../KeyValueStoreTest.cs | 50 +++++++++++++++++++ 1 file changed, 50 insertions(+) diff --git a/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs b/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs index f9a308c43..76a63f2c9 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs @@ -525,6 +525,56 @@ public async Task SetsSubjectDeleteMarkerTTL() Assert.Equal(TimeSpan.FromSeconds(2), info.Info.Config.SubjectDeleteMarkerTTL); } + [SkipIfNatsServer(versionEarlierThan: "2.11")] + public async Task SubjectDeleteMarkerTTL_enabled_removals_should_be_interpreted_as_Operation_Purge() + { + await using var server = await NatsServer.StartJSAsync(); + await using var nats = await server.CreateClientConnectionAsync(); + + var js = new NatsJSContext(nats); + var kv = new NatsKVContext(js); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + var cancellationToken = cts.Token; + + var store = await kv.CreateStoreAsync( + new NatsKVConfig("kv1") + { + AllowMsgTTL = true, + SubjectDeleteMarkerTTL = TimeSpan.FromHours(1), + MaxAge = TimeSpan.FromSeconds(4), + }, + cancellationToken: cancellationToken); + + var r1 = await store.CreateAsync("foo", "LOCKED", cancellationToken: cancellationToken); + Assert.Equal(1ul, r1); + + var create = Task.Run( + async () => + { + await Task.Delay(6000, cancellationToken); + Console.WriteLine("6 seconds passed — creating..."); + return await store.CreateAsync("foo", "LOCKED", ttl: TimeSpan.FromSeconds(1), cancellationToken: cancellationToken); + }, + cancellationToken); + + var checkOps = new List(); + await foreach (var entry in store.WatchAsync("foo", opts: new() { IncludeHistory = true }, cancellationToken: cancellationToken)) + { + checkOps.Add(entry.Operation); + if (entry.Revision == 3) + break; + } + + Assert.Equal(3, checkOps.Count); + Assert.Equal(NatsKVOperation.Put, checkOps[0]); + Assert.Equal(NatsKVOperation.Purge, checkOps[1]); + Assert.Equal(NatsKVOperation.Put, checkOps[2]); + + var r2 = await create; + Assert.Equal(3ul, r2); + } + [SkipIfNatsServer(versionEarlierThan: "2.11")] public async Task TestMessageNeverExpire() {