diff --git a/src/NATS.Client.JetStream/Models/ConsumerConfig.cs b/src/NATS.Client.JetStream/Models/ConsumerConfig.cs index 51147bea0..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 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/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..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") - { - throw new NatsJSException("Cannot create consumers with priority policy other than 'overflow', 'pinned_client', 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 c00104b21..fcec79d28 100644 --- a/src/NATS.Client.JetStream/NatsJSJsonSerializer.cs +++ b/src/NATS.Client.JetStream/NatsJSJsonSerializer.cs @@ -104,6 +104,7 @@ internal partial class NatsJSJsonSerializerContext : JsonSerializerContext new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), + new JsonStringEnumConverter(JsonNamingPolicy.SnakeCaseLower), }, }); #endif @@ -235,6 +236,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}"); } @@ -369,6 +383,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/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/EnumJsonTests.cs b/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs index 7c2a82332..4eef1a5d4 100644 --- a/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs +++ b/tests/NATS.Client.JetStream.Tests/EnumJsonTests.cs @@ -252,4 +252,39 @@ public void StreamConfigPersistMode_roundtrip_preserves_server_values() var jsonOut3 = Encoding.UTF8.GetString(bw3.WrittenSpan.ToArray()); Assert.DoesNotContain("persist_mode", jsonOut3); } + + [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 5d184d2da..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 @@ -151,4 +151,89 @@ 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 = ConsumerConfigPriorityPolicy.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 = ConsumerConfigPriorityPolicy.Prioritized, + }; + + var consumer = await js.CreateOrUpdateConsumerAsync($"{prefix}s1", consumerConfig, cancellationToken: cts.Token); + Assert.NotNull(consumer); + Assert.Equal(ConsumerConfigPriorityPolicy.Prioritized, consumer.Info.Config.PriorityPolicy); + } }