From 9f5e1919b9d2ddaddd278eb87449a1271feedabd Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Fri, 26 Sep 2025 13:29:10 +0100 Subject: [PATCH 1/3] Add Prioritized Priority Consumer Policy --- .../Models/ConsumerConfig.cs | 2 +- .../Models/ConsumerGetnextRequest.cs | 13 +++ src/NATS.Client.JetStream/NatsJSConsumer.cs | 3 + .../NatsJSContext.Consumers.cs | 4 +- src/NATS.Client.JetStream/NatsJSOpts.cs | 10 +++ .../PriorityGroupTest.cs | 83 +++++++++++++++++++ 6 files changed, 112 insertions(+), 3 deletions(-) diff --git a/src/NATS.Client.JetStream/Models/ConsumerConfig.cs b/src/NATS.Client.JetStream/Models/ConsumerConfig.cs index 51147bea0..7f7a20212 100644 --- a/src/NATS.Client.JetStream/Models/ConsumerConfig.cs +++ b/src/NATS.Client.JetStream/Models/ConsumerConfig.cs @@ -272,7 +272,7 @@ public ConsumerConfig(string name) #endif /// - /// Specifies the priority policy for consumer message selection, such as prioritizing overflow or pinned_client. + /// Specifies the priority policy for consumer message selection, such as prioritizing prioritized, overflow, or pinned_client. /// [System.Text.Json.Serialization.JsonPropertyName("priority_policy")] [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] diff --git a/src/NATS.Client.JetStream/Models/ConsumerGetnextRequest.cs b/src/NATS.Client.JetStream/Models/ConsumerGetnextRequest.cs index 264e1c62d..d52afbe02 100644 --- a/src/NATS.Client.JetStream/Models/ConsumerGetnextRequest.cs +++ b/src/NATS.Client.JetStream/Models/ConsumerGetnextRequest.cs @@ -82,4 +82,17 @@ public record ConsumerGetnextRequest [System.Text.Json.Serialization.JsonPropertyName("id")] [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] public string? Id { get; set; } + + /// + /// Priority for message delivery when using prioritized priority policy. + /// + /// + /// Lower values indicate higher priority (0 is the highest priority). + /// Maximum priority value is 9. This field is only used when the consumer + /// has PriorityPolicy set to "prioritized". + /// + [System.Text.Json.Serialization.JsonPropertyName("priority")] + [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] + [System.ComponentModel.DataAnnotations.Range(0, 9)] + public byte Priority { get; set; } } diff --git a/src/NATS.Client.JetStream/NatsJSConsumer.cs b/src/NATS.Client.JetStream/NatsJSConsumer.cs index c36c15b4b..8820399ec 100644 --- a/src/NATS.Client.JetStream/NatsJSConsumer.cs +++ b/src/NATS.Client.JetStream/NatsJSConsumer.cs @@ -336,6 +336,7 @@ await sub.CallMsgNextAsync( Group = opts.PriorityGroup?.Group, MinPending = opts.PriorityGroup?.MinPending ?? 0, MinAckPending = opts.PriorityGroup?.MinAckPending ?? 0, + Priority = opts.PriorityGroup?.Priority ?? 0, }, cancellationToken).ConfigureAwait(false); @@ -438,6 +439,7 @@ await sub.CallMsgNextAsync( Group = opts.PriorityGroup?.Group, MinPending = opts.PriorityGroup?.MinPending ?? 0, MinAckPending = opts.PriorityGroup?.MinAckPending ?? 0, + Priority = opts.PriorityGroup?.Priority ?? 0, } : new ConsumerGetnextRequest { @@ -449,6 +451,7 @@ await sub.CallMsgNextAsync( Group = opts.PriorityGroup?.Group, MinPending = opts.PriorityGroup?.MinPending ?? 0, MinAckPending = opts.PriorityGroup?.MinAckPending ?? 0, + Priority = opts.PriorityGroup?.Priority ?? 0, }, cancellationToken).ConfigureAwait(false); diff --git a/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs b/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs index 4f8f7ab99..e7bedbf4c 100644 --- a/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs +++ b/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs @@ -269,9 +269,9 @@ private async ValueTask CreateOrUpdateConsumerInternalAsync( } // TODO: enum these values? - if (config.PriorityPolicy != null && config.PriorityPolicy != "none" && config.PriorityPolicy != "overflow" && config.PriorityPolicy != "pinned_client") + if (config.PriorityPolicy != null && config.PriorityPolicy != "none" && config.PriorityPolicy != "overflow" && config.PriorityPolicy != "pinned_client" && config.PriorityPolicy != "prioritized") { - throw new NatsJSException("Cannot create consumers with priority policy other than 'overflow', 'pinned_client', or 'none'."); + throw new NatsJSException("Cannot create consumers with priority policy other than 'overflow', 'pinned_client', 'prioritized', or 'none'."); } var response = await JSRequestResponseAsync( diff --git a/src/NATS.Client.JetStream/NatsJSOpts.cs b/src/NATS.Client.JetStream/NatsJSOpts.cs index 0611e2c3d..dc40439f5 100644 --- a/src/NATS.Client.JetStream/NatsJSOpts.cs +++ b/src/NATS.Client.JetStream/NatsJSOpts.cs @@ -265,4 +265,14 @@ public record NatsJSPriorityGroupOpts /// When specified, this Pull request will only receive messages when the consumer has at least this many ack pending messages. /// public long MinAckPending { get; set; } + + /// + /// Priority for message delivery when using prioritized priority policy. + /// + /// + /// Lower values indicate higher priority (0 is the highest priority). + /// Maximum priority value is 9. This field is only used when the consumer + /// has PriorityPolicy set to "prioritized". + /// + public byte Priority { get; init; } } diff --git a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs index 5d184d2da..e40ee59d0 100644 --- a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs @@ -151,4 +151,87 @@ public async Task Consume_from_overflow_group() } } } + + [SkipIfNatsServer(versionEarlierThan: "2.12")] + public async Task Fetch_from_prioritized_group_with_priority() + { + await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + var js = new NatsJSContext(nats); + var prefix = _server.GetNextId(); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + await js.CreateStreamAsync($"{prefix}s1", [$"{prefix}s1.>"], cts.Token); + + for (var i = 0; i < 10; i++) + { + var ack = await js.PublishAsync($"{prefix}s1.{i}", i, cancellationToken: cts.Token); + ack.EnsureSuccess(); + } + + var consumerConfig = new ConsumerConfig($"{prefix}c1") + { + PriorityGroups = ["jobs"], + PriorityPolicy = "prioritized", + }; + var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); + + // Test with priority 5 + { + var opts = new NatsJSFetchOpts + { + MaxMsgs = 3, + PriorityGroup = new NatsJSPriorityGroupOpts { Group = "jobs", Priority = 5 }, + }; + var count = 0; + await foreach (var msg in consumer.FetchAsync(opts, cancellationToken: cts.Token)) + { + Assert.Equal(count++, msg.Data); + if (count == 3) break; + } + + Assert.Equal(3, count); + } + + // Test with priority 0 (highest priority) + { + var opts = new NatsJSFetchOpts + { + MaxMsgs = 2, + PriorityGroup = new NatsJSPriorityGroupOpts { Group = "jobs", Priority = 0 }, + }; + var count = 0; + await foreach (var msg in consumer.FetchAsync(opts, cancellationToken: cts.Token)) + { + Assert.Equal(count + 3, msg.Data); // Should continue from where we left off + count++; + if (count == 2) break; + } + + Assert.Equal(2, count); + } + } + + [SkipIfNatsServer(versionEarlierThan: "2.12")] + public async Task Consumer_with_prioritized_policy_validation() + { + await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + var js = new NatsJSContext(nats); + var prefix = _server.GetNextId(); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + await js.CreateStreamAsync($"{prefix}s1", [$"{prefix}s1.>"], cts.Token); + + // Test that prioritized policy is accepted + var consumerConfig = new ConsumerConfig($"{prefix}c1") + { + PriorityGroups = ["jobs"], + PriorityPolicy = "prioritized", + }; + + var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); + Assert.NotNull(consumer); + Assert.Equal("prioritized", consumer.Info.Config.PriorityPolicy); + } } From b8b4520432f6b544e7758d4963029a2fcba6df33 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Fri, 26 Sep 2025 13:47:39 +0100 Subject: [PATCH 2/3] Fix format --- tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs index e40ee59d0..1f13dfb82 100644 --- a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs @@ -187,7 +187,8 @@ public async Task Fetch_from_prioritized_group_with_priority() await foreach (var msg in consumer.FetchAsync(opts, cancellationToken: cts.Token)) { Assert.Equal(count++, msg.Data); - if (count == 3) break; + if (count == 3) + break; } Assert.Equal(3, count); @@ -205,7 +206,8 @@ public async Task Fetch_from_prioritized_group_with_priority() { Assert.Equal(count + 3, msg.Data); // Should continue from where we left off count++; - if (count == 2) break; + if (count == 2) + break; } Assert.Equal(2, count); From 7fc259c78153150deff3ba24ccaaef5eefda8481 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Tue, 30 Sep 2025 13:28:20 +0100 Subject: [PATCH 3/3] Refactor ConsumerConfig to use ConsumerConfigPriorityPolicy enum --- .../Models/ConsumerConfig.cs | 9 +++-- .../Models/ConsumerConfigPriorityPolicy.cs | 22 ++++++++++++ .../NatsJSContext.Consumers.cs | 11 ------ .../NatsJSJsonSerializer.cs | 29 +++++++++++++++ .../EnumJsonTests.cs | 35 +++++++++++++++++++ .../PriorityGroupTest.cs | 12 +++---- 6 files changed, 96 insertions(+), 22 deletions(-) create mode 100644 src/NATS.Client.JetStream/Models/ConsumerConfigPriorityPolicy.cs diff --git a/src/NATS.Client.JetStream/Models/ConsumerConfig.cs b/src/NATS.Client.JetStream/Models/ConsumerConfig.cs index 7f7a20212..3f49fbdaa 100644 --- a/src/NATS.Client.JetStream/Models/ConsumerConfig.cs +++ b/src/NATS.Client.JetStream/Models/ConsumerConfig.cs @@ -272,16 +272,15 @@ public ConsumerConfig(string name) #endif /// - /// Specifies the priority policy for consumer message selection, such as prioritizing prioritized, overflow, or pinned_client. + /// Specifies the priority policy for consumer message selection. /// [System.Text.Json.Serialization.JsonPropertyName("priority_policy")] [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] - [System.ComponentModel.DataAnnotations.StringLength(int.MaxValue, MinimumLength = 1)] - [System.ComponentModel.DataAnnotations.RegularExpression(@"^[^.*>]+$")] #if NET6_0 - public string? PriorityPolicy { get; set; } + [System.Text.Json.Serialization.JsonConverter(typeof(NatsJSJsonStringEnumConverter))] + public ConsumerConfigPriorityPolicy? PriorityPolicy { get; set; } #else - public string? PriorityPolicy { get; init; } + public ConsumerConfigPriorityPolicy? PriorityPolicy { get; init; } #endif /// diff --git a/src/NATS.Client.JetStream/Models/ConsumerConfigPriorityPolicy.cs b/src/NATS.Client.JetStream/Models/ConsumerConfigPriorityPolicy.cs new file mode 100644 index 000000000..ebb2f9826 --- /dev/null +++ b/src/NATS.Client.JetStream/Models/ConsumerConfigPriorityPolicy.cs @@ -0,0 +1,22 @@ +namespace NATS.Client.JetStream.Models; + +/// +/// The priority policy for consumer message selection. +/// +public enum ConsumerConfigPriorityPolicy +{ + /// + /// No priority policy is set. + /// + None = 0, + + /// + /// Messages are delivered based on priority level. + /// + Prioritized = 1, + + /// + /// Messages overflow to the next available consumer. + /// + Overflow = 2, +} diff --git a/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs b/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs index e7bedbf4c..0052898e5 100644 --- a/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs +++ b/src/NATS.Client.JetStream/NatsJSContext.Consumers.cs @@ -263,17 +263,6 @@ private async ValueTask CreateOrUpdateConsumerInternalAsync( throw new NatsJSException("Cannot create consumers with multiple priority groups."); } - if (config.PriorityPolicy is "pinned_client") - { - throw new NotImplementedException("Pinned clients are not supported yet."); - } - - // TODO: enum these values? - if (config.PriorityPolicy != null && config.PriorityPolicy != "none" && config.PriorityPolicy != "overflow" && config.PriorityPolicy != "pinned_client" && config.PriorityPolicy != "prioritized") - { - throw new NatsJSException("Cannot create consumers with priority policy other than 'overflow', 'pinned_client', 'prioritized', or 'none'."); - } - var response = await JSRequestResponseAsync( subject: subject, new ConsumerCreateRequest diff --git a/src/NATS.Client.JetStream/NatsJSJsonSerializer.cs b/src/NATS.Client.JetStream/NatsJSJsonSerializer.cs index 29a8526b7..7d0083aac 100644 --- a/src/NATS.Client.JetStream/NatsJSJsonSerializer.cs +++ b/src/NATS.Client.JetStream/NatsJSJsonSerializer.cs @@ -103,6 +103,7 @@ internal partial class NatsJSJsonSerializerContext : JsonSerializerContext new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), + new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), }, }); #endif @@ -223,6 +224,19 @@ public override TEnum Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSe } } + if (typeToConvert == typeof(ConsumerConfigPriorityPolicy)) + { + switch (stringValue) + { + case "none": + return (TEnum)(object)ConsumerConfigPriorityPolicy.None; + case "prioritized": + return (TEnum)(object)ConsumerConfigPriorityPolicy.Prioritized; + case "overflow": + return (TEnum)(object)ConsumerConfigPriorityPolicy.Overflow; + } + } + throw new InvalidOperationException($"Reading unknown enum type {typeToConvert.Name} or value {stringValue}"); } @@ -345,6 +359,21 @@ public override void Write(Utf8JsonWriter writer, TEnum value, JsonSerializerOpt return; } } + else if (value is ConsumerConfigPriorityPolicy consumerConfigPriorityPolicy) + { + switch (consumerConfigPriorityPolicy) + { + case ConsumerConfigPriorityPolicy.None: + writer.WriteStringValue("none"); + return; + case ConsumerConfigPriorityPolicy.Prioritized: + writer.WriteStringValue("prioritized"); + return; + case ConsumerConfigPriorityPolicy.Overflow: + writer.WriteStringValue("overflow"); + return; + } + } throw new InvalidOperationException($"Writing unknown enum value {value.GetType().Name}.{value}"); } diff --git a/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs b/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs index eb7a653a6..f21dd697a 100644 --- a/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs +++ b/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs @@ -167,4 +167,39 @@ public void ConsumerCreateRequestAction_Test(ConsumerCreateAction value, string Assert.NotNull(result); Assert.Equal(value, result.Action); } + + [Theory] + [InlineData(ConsumerConfigPriorityPolicy.None, "\"priority_policy\":\"none\"")] + [InlineData(ConsumerConfigPriorityPolicy.Prioritized, "\"priority_policy\":\"prioritized\"")] + [InlineData(ConsumerConfigPriorityPolicy.Overflow, "\"priority_policy\":\"overflow\"")] + public void ConsumerConfigPriorityPolicy_test(ConsumerConfigPriorityPolicy value, string expected) + { + var serializer = NatsJSJsonSerializer.Default; + + var bw = new NatsBufferWriter(); + serializer.Serialize(bw, new ConsumerConfig { PriorityPolicy = value }); + + var json = Encoding.UTF8.GetString(bw.WrittenSpan.ToArray()); + Assert.Contains(expected, json); + + var result = serializer.Deserialize(new ReadOnlySequence(bw.WrittenMemory)); + Assert.NotNull(result); + Assert.Equal(value, result.PriorityPolicy); + } + + [Fact] + public void ConsumerConfigPriorityPolicy_null_test() + { + var serializer = NatsJSJsonSerializer.Default; + + var bw = new NatsBufferWriter(); + serializer.Serialize(bw, new ConsumerConfig { PriorityPolicy = null }); + + var json = Encoding.UTF8.GetString(bw.WrittenSpan.ToArray()); + Assert.DoesNotContain("priority_policy", json); + + var result = serializer.Deserialize(new ReadOnlySequence(bw.WrittenMemory)); + Assert.NotNull(result); + Assert.Null(result.PriorityPolicy); + } } diff --git a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs index 1f13dfb82..f07f3cd69 100644 --- a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs @@ -34,7 +34,7 @@ public async Task Next_from_overflow_group() ack.EnsureSuccess(); } - var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], PriorityPolicy = "overflow", }; + var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], PriorityPolicy = ConsumerConfigPriorityPolicy.Overflow, }; var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); // Err on no group @@ -77,7 +77,7 @@ public async Task Fetch_from_overflow_group() ack.EnsureSuccess(); } - var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], PriorityPolicy = "overflow", }; + var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], PriorityPolicy = ConsumerConfigPriorityPolicy.Overflow, }; var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); // Err on no group @@ -123,7 +123,7 @@ public async Task Consume_from_overflow_group() ack.EnsureSuccess(); } - var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], PriorityPolicy = "overflow", }; + var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], PriorityPolicy = ConsumerConfigPriorityPolicy.Overflow, }; var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); // Err on no group @@ -172,7 +172,7 @@ public async Task Fetch_from_prioritized_group_with_priority() var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], - PriorityPolicy = "prioritized", + PriorityPolicy = ConsumerConfigPriorityPolicy.Prioritized, }; var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); @@ -229,11 +229,11 @@ public async Task Consumer_with_prioritized_policy_validation() var consumerConfig = new ConsumerConfig($"{prefix}c1") { PriorityGroups = ["jobs"], - PriorityPolicy = "prioritized", + PriorityPolicy = ConsumerConfigPriorityPolicy.Prioritized, }; var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); Assert.NotNull(consumer); - Assert.Equal("prioritized", consumer.Info.Config.PriorityPolicy); + Assert.Equal(ConsumerConfigPriorityPolicy.Prioritized, consumer.Info.Config.PriorityPolicy); } }