From 71db4f6e6197758a882542bceba1458051812dab Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Thu, 24 Sep 2026 11:18:04 -0500 Subject: [PATCH 1/5] Keep classic priority arguments off quorum queues --- .../Messaging/RabbitMQMessageBus.cs | 2 +- .../Messaging/RabbitMQMessageBusOptions.cs | 10 +- .../Messaging/RabbitMqMessageBusTestBase.cs | 133 ++++++++++++------ 3 files changed, 97 insertions(+), 48 deletions(-) diff --git a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs index f4678c50..3ac45184 100644 --- a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs +++ b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs @@ -817,7 +817,7 @@ private async Task CreateQueueAsync(IChannel channel) if (_options.SingleActiveConsumer) arguments["x-single-active-consumer"] = true; - if (_options.MaxPriority.HasValue) + if (_options.MaxPriority.HasValue && !_isQuorumQueue) arguments["x-max-priority"] = (int)_options.MaxPriority.Value; if (_options.DelayedRetryType.HasValue) diff --git a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs index f87bc379..9c584c6f 100644 --- a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs +++ b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs @@ -171,10 +171,10 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions public bool SingleActiveConsumer { get; set; } /// - /// Maximum number of priority levels for the queue (1-32). + /// Maximum number of priority levels for classic queues (1-32). /// Messages published with a higher priority value are delivered to consumers before lower-priority messages. - /// RabbitMQ 4.3+ quorum queues support 32 strict priority levels. - /// Set via the x-max-priority queue argument. + /// Set via the x-max-priority queue argument for classic queues only. + /// Quorum queues use their broker-defined priority behavior without this argument. /// See: https://www.rabbitmq.com/docs/priority /// public byte? MaxPriority { get; set; } @@ -429,9 +429,9 @@ public RabbitMQMessageBusOptionsBuilder UseSingleActiveConsumer(bool enabled = t } /// - /// Enables message priority on the queue. Messages published with higher priority are + /// Configures classic queue priority. Messages published with higher priority are /// delivered to consumers before lower-priority messages. - /// RabbitMQ 4.3+ quorum queues support up to 32 strict priority levels. + /// Does not configure or cap quorum queue priority levels. /// /// Maximum priority levels (1-32). Default: 32. public RabbitMQMessageBusOptionsBuilder UseMessagePriority(byte maxPriority = 32) diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs index 1b1a98b1..769986b0 100644 --- a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs @@ -1,5 +1,7 @@ using System; using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Globalization; using System.Linq; using System.Threading; using System.Threading.Tasks; @@ -8,6 +10,7 @@ using Foundatio.Tests.Extensions; using Foundatio.Tests.Messaging; using Microsoft.Extensions.Logging; +using RabbitMQ.Client; using Xunit; namespace Foundatio.RabbitMQ.Tests.Messaging; @@ -186,6 +189,7 @@ public override Task PublishAsync_WithDeliveryDelayExtension_DelaysDeliveryAsync return base.PublishAsync_WithDeliveryDelayExtension_DelaysDeliveryAsync(); } + [Fact] public override Task PublishAsync_WithDelayedMessageAndDisposeBeforeDelivery_DiscardsMessageAsync() { @@ -283,50 +287,28 @@ await messageBus.SubscribeAsync(_ => } [Fact] - public async Task PublishAsync_WithPriority_DeliversHighPriorityFirst() + public Task PublishAsync_WithPriority_DeliversHighPriorityFirst() { - Assert.SkipWhen(string.IsNullOrEmpty(ConnectionString), "RabbitMQ infrastructure not available"); - // Arrange - string topic = $"test_topic_priority_{DateTime.UtcNow.Ticks}"; - string queueName = $"{topic}_{Guid.NewGuid():N}"; - - await using var publisher = new RabbitMQMessageBus(o => o - .ConnectionString(ConnectionString) - .SubscriptionQueueName(queueName) - .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) - .UseQuorumQueues() - .UseMessagePriority() - .PrefetchCount(1) - .LoggerFactory(Log)); - - await publisher.PublishAsync(new SimpleMessageA { Data = "low" }, - new MessageOptions { Properties = { ["Priority"] = "1" } }, TestCancellationToken); - await publisher.PublishAsync(new SimpleMessageA { Data = "high" }, - new MessageOptions { Properties = { ["Priority"] = "10" } }, TestCancellationToken); - await publisher.PublishAsync(new SimpleMessageA { Data = "medium" }, - new MessageOptions { Properties = { ["Priority"] = "5" } }, TestCancellationToken); - - await Task.Delay(TimeSpan.FromMilliseconds(500), TestCancellationToken); - - var received = new ConcurrentQueue(); - var countdownEvent = new AsyncCountdownEvent(3); - - // Act - await publisher.SubscribeAsync(msg => - { - received.Enqueue(msg.Data!); - countdownEvent.Signal(); - }, TestCancellationToken); - - await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10)); - - // Assert - var messages = received.ToArray(); - Assert.Equal(3, messages.Length); - Assert.Equal("high", messages[0]); - Assert.Equal("medium", messages[1]); - Assert.Equal("low", messages[2]); + Assert.SkipWhen(string.IsNullOrEmpty(ConnectionString), "RabbitMQ infrastructure not available"); + // Act and Assert: the shared check verifies broker ordering. + return VerifyPriorityAsync(ConnectionString, (topic, queueName) => + Assert.IsType(GetMessageBus(options => + { + var rabbitOptions = Assert.IsType(options); + rabbitOptions.Topic = topic; + rabbitOptions.SubscriptionQueueName = queueName; + rabbitOptions.IsDurable = true; + rabbitOptions.IsSubscriptionQueueExclusive = false; + rabbitOptions.SubscriptionQueueAutoDelete = false; + rabbitOptions.AcknowledgementStrategy = AcknowledgementStrategy.Automatic; + rabbitOptions.PrefetchCount = 1; + rabbitOptions.PublisherConfirmsEnabled = true; + rabbitOptions.MaxPriority = 32; + rabbitOptions.Arguments ??= new Dictionary(); + rabbitOptions.Arguments.TryAdd("x-queue-type", "classic"); + return rabbitOptions; + })), fullRange: false, TestCancellationToken); } [Fact] @@ -395,4 +377,71 @@ await messageBus2.SubscribeAsync(msg => await messageBus2.DisposeAsync(); } + + internal static async Task VerifyPriorityAsync(string connectionString, + Func createMessageBus, bool fullRange, CancellationToken cancellationToken) + { + // Arrange + using var timeout = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + timeout.CancelAfter(TimeSpan.FromSeconds(60)); + var token = timeout.Token; + string topic = $"priority-{Guid.NewGuid():N}"; + string queueName = $"{topic}-subscription"; + var factory = new ConnectionFactory { Uri = new Uri(connectionString), AutomaticRecoveryEnabled = false }; + await using var adminConnection = await factory.CreateConnectionAsync(token); + await using var admin = await adminConnection.CreateChannelAsync(cancellationToken: token); + Assert.Equal(new Version(4, 2, 5), RabbitMQMessageBus.ParseServerVersion(adminConnection.ServerProperties)); + + (string Id, byte Priority)[] publications = fullRange + ? Enumerable.Range(0, 32).Select(priority => ($"priority-{priority}", (byte)priority)).ToArray() + : [("low", 1), ("high", 10), ("medium", 5)]; + // The client omits zero; 4.2.5 uses normal priority for omitted quorum + // priorities and level zero for classic queues. + Assert.False(new BasicProperties { Priority = 0 }.IsPriorityPresent()); + string[] expected = publications.OrderByDescending(message => message.Priority) + .Select(message => message.Id).ToArray(); + + try + { + // Provision through the provider, then remove the setup consumer so the + // ordering assertion observes a complete backlog rather than live arrivals. + await using (var setup = createMessageBus(topic, queueName)) + { + await setup.SubscribeAsync((_, _) => Task.CompletedTask, cancellationToken: token); + } + + await using var publisher = createMessageBus(topic, queueName); + foreach (var message in publications) + { + await publisher.PublishAsync(new SimpleMessageA { Data = message.Id }, + new MessageOptions { Properties = { ["Priority"] = message.Priority.ToString(CultureInfo.InvariantCulture) } }, token); + } + + var queued = await admin.QueueDeclarePassiveAsync(queueName, token); + Assert.Equal(0u, queued.ConsumerCount); + Assert.Equal((uint)expected.Length, queued.MessageCount); + + var received = new ConcurrentQueue(); + var completed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + await using var subscriber = createMessageBus(topic, queueName); + // Act + await subscriber.SubscribeAsync(message => + { + received.Enqueue(message.Data!); + if (received.Count == expected.Length) + completed.TrySetResult(); + }, token); + + await completed.Task.WaitAsync(token); + // Assert + Assert.Equal(expected, received.ToArray()); + Assert.Equal(0u, (await admin.QueueDeclarePassiveAsync(queueName, token)).MessageCount); + } + finally + { + using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + await admin.QueueDeleteAsync(queueName, cancellationToken: cleanup.Token); + await admin.ExchangeDeleteAsync(topic, cancellationToken: cleanup.Token); + } + } } From 6d5694b1708c13afe9187751942164593fbc4722 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Thu, 24 Sep 2026 12:31:39 -0500 Subject: [PATCH 2/5] Document classic and quorum priority by broker version --- .../Messaging/RabbitMQMessageBusOptions.cs | 8 ++++---- .../Messaging/RabbitMqMessageBusTestBase.cs | 1 - 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs index 9c584c6f..3f2d6c08 100644 --- a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs +++ b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs @@ -429,11 +429,11 @@ public RabbitMQMessageBusOptionsBuilder UseSingleActiveConsumer(bool enabled = t } /// - /// Configures classic queue priority. Messages published with higher priority are - /// delivered to consumers before lower-priority messages. - /// Does not configure or cap quorum queue priority levels. + /// Sets x-max-priority for classic queues only. RabbitMQ 4.2 quorum queues use + /// normal/high tiers; RabbitMQ 4.3+ quorum queues have 32 strict levels without + /// configuration. This option does not configure or cap quorum priorities. /// - /// Maximum priority levels (1-32). Default: 32. + /// Classic queue maximum priority (1-32). Default: 32. public RabbitMQMessageBusOptionsBuilder UseMessagePriority(byte maxPriority = 32) { ArgumentOutOfRangeException.ThrowIfZero(maxPriority); diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs index 769986b0..8a8646b2 100644 --- a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs @@ -189,7 +189,6 @@ public override Task PublishAsync_WithDeliveryDelayExtension_DelaysDeliveryAsync return base.PublishAsync_WithDeliveryDelayExtension_DelaysDeliveryAsync(); } - [Fact] public override Task PublishAsync_WithDelayedMessageAndDisposeBeforeDelivery_DiscardsMessageAsync() { From 73064e373cbb54fa861ab3582271dc7b1bfecc31 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Fri, 25 Sep 2026 16:18:09 -0500 Subject: [PATCH 3/5] Validate classic-only priority configuration --- README.md | 2 +- .../Messaging/RabbitMQMessageBus.cs | 6 +- .../Messaging/RabbitMQMessageBusOptions.cs | 41 ++++- .../RabbitMqMessageBusClassicTestBase.cs | 64 ++++++++ .../Messaging/RabbitMqMessageBusTestBase.cs | 149 +++++++----------- .../Messaging/RabbitMqPriorityOptionTests.cs | 116 ++++++++++++++ 6 files changed, 281 insertions(+), 97 deletions(-) create mode 100644 tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityOptionTests.cs diff --git a/README.md b/README.md index bcf720e3..03d2d8e5 100644 --- a/README.md +++ b/README.md @@ -82,7 +82,7 @@ The `rabbitmq_delayed_message_exchange` plugin is [archived and no longer mainta **Supported (AMQP 0.9.1 compatible):** -- 32 strict message priority levels on quorum queues (via `UseMessagePriority()`) +- 32 strict message priority levels on quorum queues automatically; `UseMessagePriority()` configures classic queues only - Delayed retries with linear backoff (via `UseDelayedRetries()`) - Per-queue consumer timeouts (via `ConsumerTimeout()`) - Single active consumer (via `UseSingleActiveConsumer()`) diff --git a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs index 3ac45184..10aaf311 100644 --- a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs +++ b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs @@ -50,7 +50,9 @@ public RabbitMQMessageBus(RabbitMQMessageBusOptions options) : base(options) !primaryUri.Scheme.Equals("amqps", StringComparison.OrdinalIgnoreCase)) throw new ArgumentException($"ConnectionString must use amqp:// or amqps:// scheme: {SanitizeUri(primaryUri)}"); - _isQuorumQueue = options.Arguments is not null && options.Arguments.TryGetValue("x-queue-type", out object? queueType) && queueType is string type && String.Equals(type, "quorum", StringComparison.OrdinalIgnoreCase); + _isQuorumQueue = RabbitMQMessageBusOptions.IsQuorumQueue(options.Arguments); + if (_isQuorumQueue && options.MaxPriority.HasValue) + throw new InvalidOperationException("MaxPriority applies only to classic queues and cannot be used with quorum queues."); // Initialize the connection factory with credentials/vhost from connection string // Automatic recovery will allow the connections to be restored in case the server is @@ -817,7 +819,7 @@ private async Task CreateQueueAsync(IChannel channel) if (_options.SingleActiveConsumer) arguments["x-single-active-consumer"] = true; - if (_options.MaxPriority.HasValue && !_isQuorumQueue) + if (_options.MaxPriority.HasValue) arguments["x-max-priority"] = (int)_options.MaxPriority.Value; if (_options.DelayedRetryType.HasValue) diff --git a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs index 3f2d6c08..d293c987 100644 --- a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs +++ b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs @@ -24,7 +24,8 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions public TimeSpan? DefaultMessageTimeToLive { get; set; } /// - /// Arguments passed to QueueDeclare. Some brokers use it to implement additional features like message TTL. + /// Arguments passed to QueueDeclare. Configure this mutable dictionary before constructing + /// the bus; changing it while the bus is running is unsupported. /// public IDictionary? Arguments { get; set; } @@ -171,13 +172,36 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions public bool SingleActiveConsumer { get; set; } /// - /// Maximum number of priority levels for classic queues (1-32). + /// Maximum number of priority levels for classic queues (1-255). Higher limits cost + /// more broker CPU and memory; UseMessagePriority() limits its convenience API to 32. /// Messages published with a higher priority value are delivered to consumers before lower-priority messages. /// Set via the x-max-priority queue argument for classic queues only. - /// Quorum queues use their broker-defined priority behavior without this argument. + /// RabbitMQ 4.2 quorum queues use normal/high tiers; RabbitMQ 4.3+ quorum queues support + /// 32 strict priority levels automatically and cannot use this setting. /// See: https://www.rabbitmq.com/docs/priority /// - public byte? MaxPriority { get; set; } + public byte? MaxPriority + { + get; + set + { + if (value is 0) + throw new ArgumentOutOfRangeException(nameof(MaxPriority), value, "Classic queue maximum priority must be positive."); + + field = value; + } + } + + /// + /// Checks whether the supplied declaration arguments explicitly request a quorum queue. + /// This does not discover a broker or virtual-host default, inspect an existing queue, + /// or validate other argument types and values. + /// + internal static bool IsQuorumQueue(IDictionary? arguments) + { + return arguments is not null && arguments.TryGetValue("x-queue-type", out object? queueType) + && queueType is string type && String.Equals(type, "quorum", StringComparison.OrdinalIgnoreCase); + } /// /// Configures native delayed retry for quorum queues (RabbitMQ 4.3+). @@ -343,6 +367,9 @@ public RabbitMQMessageBusOptionsBuilder PublishRecoveryTimeout(TimeSpan timeout) /// The builder instance for method chaining. public RabbitMQMessageBusOptionsBuilder UseQuorumQueues() { + if (Target.MaxPriority.HasValue) + throw new InvalidOperationException("MaxPriority applies only to classic queues and cannot be used with quorum queues."); + Target.SubscriptionQueueAutoDelete = false; Target.IsSubscriptionQueueExclusive = false; @@ -430,14 +457,16 @@ public RabbitMQMessageBusOptionsBuilder UseSingleActiveConsumer(bool enabled = t /// /// Sets x-max-priority for classic queues only. RabbitMQ 4.2 quorum queues use - /// normal/high tiers; RabbitMQ 4.3+ quorum queues have 32 strict levels without - /// configuration. This option does not configure or cap quorum priorities. + /// normal/high tiers; RabbitMQ 4.3+ quorum queues have 32 strict levels automatically. + /// Cannot be combined with UseQuorumQueues(). /// /// Classic queue maximum priority (1-32). Default: 32. public RabbitMQMessageBusOptionsBuilder UseMessagePriority(byte maxPriority = 32) { ArgumentOutOfRangeException.ThrowIfZero(maxPriority); ArgumentOutOfRangeException.ThrowIfGreaterThan(maxPriority, (byte)32); + if (RabbitMQMessageBusOptions.IsQuorumQueue(Target.Arguments)) + throw new InvalidOperationException("MaxPriority applies only to classic queues and cannot be used with quorum queues."); Target.MaxPriority = maxPriority; return this; diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs index e9ba7e4c..9858a587 100644 --- a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs @@ -1,9 +1,14 @@ using System; +using System.Collections.Concurrent; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; +using Foundatio.AsyncEx; using Foundatio.Messaging; +using Foundatio.Tests.Extensions; using Foundatio.Tests.Messaging; using Microsoft.Extensions.Logging; +using RabbitMQ.Client; using Xunit; namespace Foundatio.RabbitMQ.Tests.Messaging; @@ -63,4 +68,63 @@ await messageBus.SubscribeAsync(_ => await CleanupMessageBusAsync(messageBus); } } + + [Fact] + public async Task PublishAsync_WithClassicPriority_DeliversHighPriorityFirst() + { + Assert.SkipWhen(string.IsNullOrEmpty(ConnectionString), "RabbitMQ infrastructure not available"); + + // Arrange + string topic = $"test_topic_classic_priority_{Guid.NewGuid():N}"; + string queueName = $"{topic}_queue"; + var factory = new ConnectionFactory { Uri = new Uri(ConnectionString) }; + await using var connection = await factory.CreateConnectionAsync(TestCancellationToken); + await using var channel = await connection.CreateChannelAsync(cancellationToken: TestCancellationToken); + try + { + await channel.ExchangeDeclareAsync(topic, "fanout", durable: true, cancellationToken: TestCancellationToken); + await channel.QueueDeclareAsync(queueName, durable: true, exclusive: false, autoDelete: false, + arguments: new Dictionary { ["x-queue-type"] = "classic", ["x-max-priority"] = 10 }, + cancellationToken: TestCancellationToken); + await channel.QueueBindAsync(queueName, topic, String.Empty, cancellationToken: TestCancellationToken); + + await using var bus = new RabbitMQMessageBus(o => o + .ConnectionString(ConnectionString) + .Topic(topic) + .SubscriptionQueueName(queueName) + .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) + .IsSubscriptionQueueExclusive(false) + .SubscriptionQueueAutoDelete(false) + .Arguments(new Dictionary { ["x-queue-type"] = "classic" }) + .UseMessagePriority(10) + .PrefetchCount(1) + .PublisherConfirmsEnabled() + .LoggerFactory(Log)); + await bus.PublishAsync(new SimpleMessageA { Data = "low" }, + new MessageOptions { Properties = { ["Priority"] = "1" } }, TestCancellationToken); + await bus.PublishAsync(new SimpleMessageA { Data = "high" }, + new MessageOptions { Properties = { ["Priority"] = "10" } }, TestCancellationToken); + await bus.PublishAsync(new SimpleMessageA { Data = "medium" }, + new MessageOptions { Properties = { ["Priority"] = "5" } }, TestCancellationToken); + var received = new ConcurrentQueue(); + var countdownEvent = new AsyncCountdownEvent(3); + + // Act + await bus.SubscribeAsync(msg => + { + received.Enqueue(msg.Data!); + countdownEvent.Signal(); + }, TestCancellationToken); + await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10)); + + // Assert + Assert.Equal(["high", "medium", "low"], received.ToArray()); + } + finally + { + using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + await channel.QueueDeleteAsync(queueName, cancellationToken: cleanup.Token); + await channel.ExchangeDeleteAsync(topic, cancellationToken: cleanup.Token); + } + } } diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs index 8a8646b2..0bcb8e70 100644 --- a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs @@ -1,8 +1,6 @@ using System; using System.Collections.Concurrent; using System.Collections.Generic; -using System.Globalization; -using System.Linq; using System.Threading; using System.Threading.Tasks; using Foundatio.AsyncEx; @@ -286,28 +284,70 @@ await messageBus.SubscribeAsync(_ => } [Fact] - public Task PublishAsync_WithPriority_DeliversHighPriorityFirst() + public async Task PublishAsync_WithPriority_DeliversHighPriorityFirst() { - // Arrange Assert.SkipWhen(string.IsNullOrEmpty(ConnectionString), "RabbitMQ infrastructure not available"); - // Act and Assert: the shared check verifies broker ordering. - return VerifyPriorityAsync(ConnectionString, (topic, queueName) => - Assert.IsType(GetMessageBus(options => + + // Arrange + string topic = $"test_topic_priority_{DateTime.UtcNow.Ticks}"; + string queueName = $"{topic}_{Guid.NewGuid():N}"; + var factory = new ConnectionFactory { Uri = new Uri(ConnectionString) }; + await using var connection = await factory.CreateConnectionAsync(TestCancellationToken); + // Strict quorum priority ordering requires RabbitMQ 4.3+; 4.2 uses normal/high tiers. + Assert.SkipWhen(RabbitMQMessageBus.ParseServerVersion(connection.ServerProperties) is not { } version + || version < new Version(4, 3), "Strict quorum priority ordering requires RabbitMQ 4.3+"); + await using var channel = await connection.CreateChannelAsync(cancellationToken: TestCancellationToken); + try + { + // A bound queue must exist before publishing the backlog, without an active consumer. + await channel.ExchangeDeclareAsync(topic, "fanout", durable: true, cancellationToken: TestCancellationToken); + await channel.QueueDeclareAsync(queueName, durable: true, exclusive: false, autoDelete: false, + arguments: new Dictionary { ["x-queue-type"] = "quorum", ["x-delivery-limit"] = 2L }, + cancellationToken: TestCancellationToken); + await channel.QueueBindAsync(queueName, topic, String.Empty, cancellationToken: TestCancellationToken); + + await using var publisher = new RabbitMQMessageBus(o => o + .ConnectionString(ConnectionString) + .Topic(topic) + .SubscriptionQueueName(queueName) + .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) + .UseQuorumQueues() + .PrefetchCount(1) + .PublisherConfirmsEnabled() + .LoggerFactory(Log)); + + await publisher.PublishAsync(new SimpleMessageA { Data = "low" }, + new MessageOptions { Properties = { ["Priority"] = "1" } }, TestCancellationToken); + await publisher.PublishAsync(new SimpleMessageA { Data = "high" }, + new MessageOptions { Properties = { ["Priority"] = "10" } }, TestCancellationToken); + await publisher.PublishAsync(new SimpleMessageA { Data = "medium" }, + new MessageOptions { Properties = { ["Priority"] = "5" } }, TestCancellationToken); + + var received = new ConcurrentQueue(); + var countdownEvent = new AsyncCountdownEvent(3); + + // Act + await publisher.SubscribeAsync(msg => { - var rabbitOptions = Assert.IsType(options); - rabbitOptions.Topic = topic; - rabbitOptions.SubscriptionQueueName = queueName; - rabbitOptions.IsDurable = true; - rabbitOptions.IsSubscriptionQueueExclusive = false; - rabbitOptions.SubscriptionQueueAutoDelete = false; - rabbitOptions.AcknowledgementStrategy = AcknowledgementStrategy.Automatic; - rabbitOptions.PrefetchCount = 1; - rabbitOptions.PublisherConfirmsEnabled = true; - rabbitOptions.MaxPriority = 32; - rabbitOptions.Arguments ??= new Dictionary(); - rabbitOptions.Arguments.TryAdd("x-queue-type", "classic"); - return rabbitOptions; - })), fullRange: false, TestCancellationToken); + received.Enqueue(msg.Data!); + countdownEvent.Signal(); + }, TestCancellationToken); + + await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10)); + + // Assert + var messages = received.ToArray(); + Assert.Equal(3, messages.Length); + Assert.Equal("high", messages[0]); + Assert.Equal("medium", messages[1]); + Assert.Equal("low", messages[2]); + } + finally + { + using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + await channel.QueueDeleteAsync(queueName, cancellationToken: cleanup.Token); + await channel.ExchangeDeleteAsync(topic, cancellationToken: cleanup.Token); + } } [Fact] @@ -376,71 +416,4 @@ await messageBus2.SubscribeAsync(msg => await messageBus2.DisposeAsync(); } - - internal static async Task VerifyPriorityAsync(string connectionString, - Func createMessageBus, bool fullRange, CancellationToken cancellationToken) - { - // Arrange - using var timeout = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - timeout.CancelAfter(TimeSpan.FromSeconds(60)); - var token = timeout.Token; - string topic = $"priority-{Guid.NewGuid():N}"; - string queueName = $"{topic}-subscription"; - var factory = new ConnectionFactory { Uri = new Uri(connectionString), AutomaticRecoveryEnabled = false }; - await using var adminConnection = await factory.CreateConnectionAsync(token); - await using var admin = await adminConnection.CreateChannelAsync(cancellationToken: token); - Assert.Equal(new Version(4, 2, 5), RabbitMQMessageBus.ParseServerVersion(adminConnection.ServerProperties)); - - (string Id, byte Priority)[] publications = fullRange - ? Enumerable.Range(0, 32).Select(priority => ($"priority-{priority}", (byte)priority)).ToArray() - : [("low", 1), ("high", 10), ("medium", 5)]; - // The client omits zero; 4.2.5 uses normal priority for omitted quorum - // priorities and level zero for classic queues. - Assert.False(new BasicProperties { Priority = 0 }.IsPriorityPresent()); - string[] expected = publications.OrderByDescending(message => message.Priority) - .Select(message => message.Id).ToArray(); - - try - { - // Provision through the provider, then remove the setup consumer so the - // ordering assertion observes a complete backlog rather than live arrivals. - await using (var setup = createMessageBus(topic, queueName)) - { - await setup.SubscribeAsync((_, _) => Task.CompletedTask, cancellationToken: token); - } - - await using var publisher = createMessageBus(topic, queueName); - foreach (var message in publications) - { - await publisher.PublishAsync(new SimpleMessageA { Data = message.Id }, - new MessageOptions { Properties = { ["Priority"] = message.Priority.ToString(CultureInfo.InvariantCulture) } }, token); - } - - var queued = await admin.QueueDeclarePassiveAsync(queueName, token); - Assert.Equal(0u, queued.ConsumerCount); - Assert.Equal((uint)expected.Length, queued.MessageCount); - - var received = new ConcurrentQueue(); - var completed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - await using var subscriber = createMessageBus(topic, queueName); - // Act - await subscriber.SubscribeAsync(message => - { - received.Enqueue(message.Data!); - if (received.Count == expected.Length) - completed.TrySetResult(); - }, token); - - await completed.Task.WaitAsync(token); - // Assert - Assert.Equal(expected, received.ToArray()); - Assert.Equal(0u, (await admin.QueueDeclarePassiveAsync(queueName, token)).MessageCount); - } - finally - { - using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(10)); - await admin.QueueDeleteAsync(queueName, cancellationToken: cleanup.Token); - await admin.ExchangeDeleteAsync(topic, cancellationToken: cleanup.Token); - } - } } diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityOptionTests.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityOptionTests.cs new file mode 100644 index 00000000..03eb8884 --- /dev/null +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityOptionTests.cs @@ -0,0 +1,116 @@ +using System; +using System.Collections.Generic; +using Foundatio.Messaging; +using Xunit; + +namespace Foundatio.RabbitMQ.Tests.Messaging; + +public class RabbitMqPriorityOptionTests +{ + [Fact] + public void Constructor_WithDirectQuorumMaxPriority_Throws() + { + // Arrange + var options = new RabbitMQMessageBusOptions + { + ConnectionString = "amqp://localhost", + MaxPriority = 3, + Arguments = new Dictionary { ["x-queue-type"] = "quorum" } + }; + + // Act + Action create = () => new RabbitMQMessageBus(options); + + // Assert + var exception = Assert.Throws(create); + Assert.Contains("classic", exception.Message, StringComparison.OrdinalIgnoreCase); + } + + [Fact] + public void Constructor_WithMutatedArguments_Throws() + { + // Arrange + var arguments = new Dictionary(); + var options = new RabbitMQMessageBusOptions { ConnectionString = "amqp://localhost", Arguments = arguments, MaxPriority = 3 }; + arguments["x-queue-type"] = "quorum"; + + // Act + Action create = () => new RabbitMQMessageBus(options); + + // Assert + Assert.Throws(create); + } + + [Fact] + public void MaxPriority_WithClassicQueue_AllowsBrokerSupportedValue() + { + // Arrange + var options = new RabbitMQMessageBusOptions(); + + // Act + options.MaxPriority = byte.MaxValue; + + // Assert + Assert.Equal(byte.MaxValue, options.MaxPriority); + } + + [Fact] + public void MaxPriority_WithZero_ThrowsWithoutChangingOptions() + { + // Arrange + var options = new RabbitMQMessageBusOptions(); + + // Act + Action configure = () => options.MaxPriority = 0; + + // Assert + Assert.Throws(configure); + Assert.Null(options.MaxPriority); + } + + [Fact] + public void UseMessagePriority_AfterUseQuorumQueues_Throws() + { + // Arrange + var builder = new RabbitMQMessageBusOptionsBuilder().UseQuorumQueues(); + + // Act + Action configure = () => builder.UseMessagePriority(); + + // Assert + var exception = Assert.Throws(configure); + Assert.Contains("classic", exception.Message, StringComparison.OrdinalIgnoreCase); + Assert.Null(builder.Build().MaxPriority); + } + + [Fact] + public void UseMessagePriority_WithClassicQueue_SetsMaximum() + { + // Arrange + var builder = new RabbitMQMessageBusOptionsBuilder(); + + // Act + builder.UseMessagePriority(4); + + // Assert + Assert.Equal((byte)4, builder.Build().MaxPriority); + } + + [Fact] + public void UseQuorumQueues_AfterUseMessagePriority_ThrowsWithoutChangingOptions() + { + // Arrange + var builder = new RabbitMQMessageBusOptionsBuilder().UseMessagePriority(); + + // Act + Action configure = () => builder.UseQuorumQueues(); + + // Assert + var exception = Assert.Throws(configure); + Assert.Contains("classic", exception.Message, StringComparison.OrdinalIgnoreCase); + var options = builder.Build(); + Assert.Null(options.Arguments); + Assert.True(options.IsSubscriptionQueueExclusive); + Assert.True(options.SubscriptionQueueAutoDelete); + } +} From 260c8dfdd5efb36ebc158bfda3e1f41a6187f3a7 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Mon, 28 Sep 2026 21:01:32 -0500 Subject: [PATCH 4/5] Exercise priority ordering through message bus subscriptions --- .../Messaging/RabbitMQMessageBusOptions.cs | 14 ++-- .../RabbitMqMessageBusClassicTestBase.cs | 67 +++++++++---------- .../Messaging/RabbitMqMessageBusTestBase.cs | 62 ++++++++--------- 3 files changed, 75 insertions(+), 68 deletions(-) diff --git a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs index d293c987..69c96dae 100644 --- a/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs +++ b/src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs @@ -24,8 +24,13 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions public TimeSpan? DefaultMessageTimeToLive { get; set; } /// - /// Arguments passed to QueueDeclare. Configure this mutable dictionary before constructing - /// the bus; changing it while the bus is running is unsupported. + /// Optional queue declaration arguments for broker features such as message TTL. + /// Prefer the fluent builder methods for supported settings, such as UseQuorumQueues() + /// and UseMessagePriority(). Use Arguments for settings without a dedicated builder method; + /// for example, Arguments(new Dictionary<string, object?> { ["x-message-ttl"] = 60000 }). + /// Configure this mutable dictionary before constructing the bus; changing it while running is unsupported. + /// See: https://www.rabbitmq.com/docs/queues#optional-arguments + /// See: https://www.rabbitmq.com/docs/ttl#per-queue-message-ttl-in-queues /// public IDictionary? Arguments { get; set; } @@ -172,13 +177,14 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions public bool SingleActiveConsumer { get; set; } /// - /// Maximum number of priority levels for classic queues (1-255). Higher limits cost + /// Maximum message priority for classic queues (1-255). Higher limits cost /// more broker CPU and memory; UseMessagePriority() limits its convenience API to 32. - /// Messages published with a higher priority value are delivered to consumers before lower-priority messages. + /// Prioritizes queued messages; deliveries already sent to consumers are not reordered. /// Set via the x-max-priority queue argument for classic queues only. /// RabbitMQ 4.2 quorum queues use normal/high tiers; RabbitMQ 4.3+ quorum queues support /// 32 strict priority levels automatically and cannot use this setting. /// See: https://www.rabbitmq.com/docs/priority + /// See: https://www.rabbitmq.com/docs/4.2/priority#quorum-queues /// public byte? MaxPriority { diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs index 9858a587..591c34c5 100644 --- a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusClassicTestBase.cs @@ -1,6 +1,5 @@ using System; using System.Collections.Concurrent; -using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Foundatio.AsyncEx; @@ -8,7 +7,6 @@ using Foundatio.Tests.Extensions; using Foundatio.Tests.Messaging; using Microsoft.Extensions.Logging; -using RabbitMQ.Client; using Xunit; namespace Foundatio.RabbitMQ.Tests.Messaging; @@ -19,7 +17,6 @@ public RabbitMqMessageBusClassicTestBase(string connectionString, ITestOutputHel { } - protected override IMessageBus? GetMessageBus(Func? config = null) { if (string.IsNullOrEmpty(ConnectionString)) @@ -77,44 +74,47 @@ public async Task PublishAsync_WithClassicPriority_DeliversHighPriorityFirst() // Arrange string topic = $"test_topic_classic_priority_{Guid.NewGuid():N}"; string queueName = $"{topic}_queue"; - var factory = new ConnectionFactory { Uri = new Uri(ConnectionString) }; - await using var connection = await factory.CreateConnectionAsync(TestCancellationToken); - await using var channel = await connection.CreateChannelAsync(cancellationToken: TestCancellationToken); + await using var bus = new RabbitMQMessageBus(o => o + .ConnectionString(ConnectionString) + .Topic(topic) + .SubscriptionQueueName(queueName) + .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) + .UseMessagePriority(10) + .PrefetchCount(1) + .PublisherConfirmsEnabled() + .LoggerFactory(Log)); + var received = new ConcurrentQueue(); + var countdownEvent = new AsyncCountdownEvent(3); + var warmupReceived = new AsyncManualResetEvent(); + var releaseWarmup = new AsyncManualResetEvent(); + try { - await channel.ExchangeDeclareAsync(topic, "fanout", durable: true, cancellationToken: TestCancellationToken); - await channel.QueueDeclareAsync(queueName, durable: true, exclusive: false, autoDelete: false, - arguments: new Dictionary { ["x-queue-type"] = "classic", ["x-max-priority"] = 10 }, - cancellationToken: TestCancellationToken); - await channel.QueueBindAsync(queueName, topic, String.Empty, cancellationToken: TestCancellationToken); - - await using var bus = new RabbitMQMessageBus(o => o - .ConnectionString(ConnectionString) - .Topic(topic) - .SubscriptionQueueName(queueName) - .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) - .IsSubscriptionQueueExclusive(false) - .SubscriptionQueueAutoDelete(false) - .Arguments(new Dictionary { ["x-queue-type"] = "classic" }) - .UseMessagePriority(10) - .PrefetchCount(1) - .PublisherConfirmsEnabled() - .LoggerFactory(Log)); + await bus.SubscribeAsync(async msg => + { + if (msg.Data == "warmup") + { + warmupReceived.Set(); + await releaseWarmup.WaitAsync(TestCancellationToken); + return; + } + + received.Enqueue(msg.Data!); + countdownEvent.Signal(); + }, TestCancellationToken); + + // Hold the only prefetched delivery so the priority messages wait in the queue. + await bus.PublishAsync(new SimpleMessageA { Data = "warmup" }, cancellationToken: TestCancellationToken); + await warmupReceived.WaitAsync(TestCancellationToken).WaitAsync(TimeSpan.FromSeconds(10), TestCancellationToken); await bus.PublishAsync(new SimpleMessageA { Data = "low" }, new MessageOptions { Properties = { ["Priority"] = "1" } }, TestCancellationToken); await bus.PublishAsync(new SimpleMessageA { Data = "high" }, new MessageOptions { Properties = { ["Priority"] = "10" } }, TestCancellationToken); await bus.PublishAsync(new SimpleMessageA { Data = "medium" }, new MessageOptions { Properties = { ["Priority"] = "5" } }, TestCancellationToken); - var received = new ConcurrentQueue(); - var countdownEvent = new AsyncCountdownEvent(3); // Act - await bus.SubscribeAsync(msg => - { - received.Enqueue(msg.Data!); - countdownEvent.Signal(); - }, TestCancellationToken); + releaseWarmup.Set(); await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10)); // Assert @@ -122,9 +122,8 @@ await bus.SubscribeAsync(msg => } finally { - using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(10)); - await channel.QueueDeleteAsync(queueName, cancellationToken: cleanup.Token); - await channel.ExchangeDeleteAsync(topic, cancellationToken: cleanup.Token); + releaseWarmup.Set(); + await CleanupMessageBusAsync(bus); } } } diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs index 0bcb8e70..4148b31e 100644 --- a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqMessageBusTestBase.cs @@ -1,6 +1,5 @@ using System; using System.Collections.Concurrent; -using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Foundatio.AsyncEx; @@ -296,26 +295,38 @@ public async Task PublishAsync_WithPriority_DeliversHighPriorityFirst() // Strict quorum priority ordering requires RabbitMQ 4.3+; 4.2 uses normal/high tiers. Assert.SkipWhen(RabbitMQMessageBus.ParseServerVersion(connection.ServerProperties) is not { } version || version < new Version(4, 3), "Strict quorum priority ordering requires RabbitMQ 4.3+"); - await using var channel = await connection.CreateChannelAsync(cancellationToken: TestCancellationToken); + await using var publisher = new RabbitMQMessageBus(o => o + .ConnectionString(ConnectionString) + .Topic(topic) + .SubscriptionQueueName(queueName) + .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) + .UseQuorumQueues() + .PrefetchCount(1) + .PublisherConfirmsEnabled() + .LoggerFactory(Log)); + var received = new ConcurrentQueue(); + var countdownEvent = new AsyncCountdownEvent(3); + var warmupReceived = new AsyncManualResetEvent(); + var releaseWarmup = new AsyncManualResetEvent(); + try { - // A bound queue must exist before publishing the backlog, without an active consumer. - await channel.ExchangeDeclareAsync(topic, "fanout", durable: true, cancellationToken: TestCancellationToken); - await channel.QueueDeclareAsync(queueName, durable: true, exclusive: false, autoDelete: false, - arguments: new Dictionary { ["x-queue-type"] = "quorum", ["x-delivery-limit"] = 2L }, - cancellationToken: TestCancellationToken); - await channel.QueueBindAsync(queueName, topic, String.Empty, cancellationToken: TestCancellationToken); - - await using var publisher = new RabbitMQMessageBus(o => o - .ConnectionString(ConnectionString) - .Topic(topic) - .SubscriptionQueueName(queueName) - .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) - .UseQuorumQueues() - .PrefetchCount(1) - .PublisherConfirmsEnabled() - .LoggerFactory(Log)); + await publisher.SubscribeAsync(async msg => + { + if (msg.Data == "warmup") + { + warmupReceived.Set(); + await releaseWarmup.WaitAsync(TestCancellationToken); + return; + } + + received.Enqueue(msg.Data!); + countdownEvent.Signal(); + }, TestCancellationToken); + // Hold the only prefetched delivery so the priority messages wait in the queue. + await publisher.PublishAsync(new SimpleMessageA { Data = "warmup" }, cancellationToken: TestCancellationToken); + await warmupReceived.WaitAsync(TestCancellationToken).WaitAsync(TimeSpan.FromSeconds(10), TestCancellationToken); await publisher.PublishAsync(new SimpleMessageA { Data = "low" }, new MessageOptions { Properties = { ["Priority"] = "1" } }, TestCancellationToken); await publisher.PublishAsync(new SimpleMessageA { Data = "high" }, @@ -323,16 +334,8 @@ await publisher.PublishAsync(new SimpleMessageA { Data = "high" }, await publisher.PublishAsync(new SimpleMessageA { Data = "medium" }, new MessageOptions { Properties = { ["Priority"] = "5" } }, TestCancellationToken); - var received = new ConcurrentQueue(); - var countdownEvent = new AsyncCountdownEvent(3); - // Act - await publisher.SubscribeAsync(msg => - { - received.Enqueue(msg.Data!); - countdownEvent.Signal(); - }, TestCancellationToken); - + releaseWarmup.Set(); await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10)); // Assert @@ -344,9 +347,8 @@ await publisher.SubscribeAsync(msg => } finally { - using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(10)); - await channel.QueueDeleteAsync(queueName, cancellationToken: cleanup.Token); - await channel.ExchangeDeleteAsync(topic, cancellationToken: cleanup.Token); + releaseWarmup.Set(); + await CleanupMessageBusAsync(publisher); } } From ae09eef6bd247633ae15545e6b1f4d4e61816dad Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Thu, 1 Oct 2026 22:22:40 -0500 Subject: [PATCH 5/5] Add independent RabbitMQ priority upgrade and routing regressions --- .../RabbitMqPriorityBehaviorTestBase.cs | 229 ++++++++++++++++++ .../RabbitMqPriorityBehaviorTests.cs | 6 + 2 files changed, 235 insertions(+) create mode 100644 tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTestBase.cs create mode 100644 tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTests.cs diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTestBase.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTestBase.cs new file mode 100644 index 00000000..f10fb0a9 --- /dev/null +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTestBase.cs @@ -0,0 +1,229 @@ +using System; +using System.Collections.Concurrent; +using System.Threading.Tasks; +using Foundatio.AsyncEx; +using Foundatio.Messaging; +using Foundatio.Tests.Extensions; +using Foundatio.Tests.Messaging; +using RabbitMQ.Client; +using Xunit; + +namespace Foundatio.RabbitMQ.Tests.Messaging; + +public abstract class RabbitMqPriorityBehaviorTestBase(string connectionString, ITestOutputHelper output) : MessageBusTestBase(output) +{ + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task PublishAsync_BeforeSubscription_DoesNotRetainUnroutableMessages(bool quorum) + { + Assert.SkipWhen(string.IsNullOrEmpty(connectionString), "RabbitMQ infrastructure not available"); + + // Arrange + await using var bus = CreateBus(quorum); + var received = new ConcurrentQueue(); + var delivered = new AsyncManualResetEvent(); + try + { + await bus.PublishAsync(new SimpleMessageA { Data = "before" }, cancellationToken: TestCancellationToken); + + // Act + await bus.SubscribeAsync(message => + { + received.Enqueue(message.Data!); + delivered.Set(); + }, TestCancellationToken); + await bus.PublishAsync(new SimpleMessageA { Data = "after" }, cancellationToken: TestCancellationToken); + await delivered.WaitAsync(TestCancellationToken).WaitAsync(TimeSpan.FromSeconds(10), TestCancellationToken); + + // Assert + Assert.Equal(["after"], received.ToArray()); + } + finally + { + await CleanupMessageBusAsync(bus); + } + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task PublishAsync_WithAnInFlightLowPriorityMessage_DoesNotPreemptDelivery(bool quorum) + { + Assert.SkipWhen(string.IsNullOrEmpty(connectionString), "RabbitMQ infrastructure not available"); + + // Arrange + await using var bus = CreateBus(quorum); + var received = new ConcurrentQueue(); + var lowReceived = new AsyncManualResetEvent(); + var releaseLow = new AsyncManualResetEvent(); + var countdown = new AsyncCountdownEvent(2); + try + { + await bus.SubscribeAsync(async message => + { + received.Enqueue(message.Data!); + if (message.Data == "low") + { + lowReceived.Set(); + await releaseLow.WaitAsync(TestCancellationToken); + } + countdown.Signal(); + }, TestCancellationToken); + await bus.PublishAsync(new SimpleMessageA { Data = "low" }, + new MessageOptions { Properties = { ["Priority"] = "1" } }, TestCancellationToken); + await lowReceived.WaitAsync(TestCancellationToken).WaitAsync(TimeSpan.FromSeconds(10), TestCancellationToken); + + // Act + await bus.PublishAsync(new SimpleMessageA { Data = "high" }, + new MessageOptions { Properties = { ["Priority"] = "10" } }, TestCancellationToken); + releaseLow.Set(); + await countdown.WaitAsync(TimeSpan.FromSeconds(10)); + + // Assert + Assert.Equal(["low", "high"], received.ToArray()); + } + finally + { + releaseLow.Set(); + await CleanupMessageBusAsync(bus); + } + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task PublishAsync_WithOmittedOrZeroPriority_UsesQueueSpecificDefault(bool quorum) + { + Assert.SkipWhen(string.IsNullOrEmpty(connectionString), "RabbitMQ infrastructure not available"); + + // Arrange + var version = await GetBrokerVersionAsync(); + + // Act + var received = await PublishQueuedAsync(quorum, ("one", "1"), ("omitted", null), ("zero", "0")); + + // Assert + // RabbitMQ.Client 7.2.2 omits priority zero; quorum 4.3 defaults an absent priority to 4. + Assert.Equal(quorum && version >= new Version(4, 3) + ? ["omitted", "zero", "one"] + : ["one", "omitted", "zero"], received); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task PublishAsync_WithQueuedPriorities_UsesQueueSpecificPriorityLevels(bool quorum) + { + Assert.SkipWhen(string.IsNullOrEmpty(connectionString), "RabbitMQ infrastructure not available"); + + // Arrange + var version = await GetBrokerVersionAsync(); + + // Act + var received = await PublishQueuedAsync(quorum, + ("five-a", "5"), ("ten-a", "10"), ("five-b", "5"), ("ten-b", "10")); + + // Assert + // Before 4.3 these quorum messages share one priority group (or FIFO before 4.0). + Assert.Equal(!quorum || version >= new Version(4, 3) + ? ["ten-a", "ten-b", "five-a", "five-b"] + : ["five-a", "ten-a", "five-b", "ten-b"], received); + } + + [Fact] + public async Task PublishAsync_WithQueuedQuorumPriorities_UsesVersionSpecificFairness() + { + Assert.SkipWhen(string.IsNullOrEmpty(connectionString), "RabbitMQ infrastructure not available"); + + // Arrange + var version = await GetBrokerVersionAsync(); + + // Act + var received = await PublishQueuedAsync(true, ("normal", "1"), + ("high-a", "10"), ("high-b", "10"), ("high-c", "10"), + ("high-d", "10"), ("high-e", "10"), ("high-f", "10")); + + // Assert + Assert.Equal(7, received.Length); + Assert.Equal(["high-a", "high-b", "high-c", "high-d", "high-e", "high-f"], + Array.FindAll(received, message => message != "normal")); + if (version >= new Version(4, 3)) + Assert.Equal("normal", received[^1]); + else if (version >= new Version(4, 0)) + { + Assert.StartsWith("high-", received[0]); + Assert.InRange(Array.IndexOf(received, "normal"), 1, 5); + } + else + Assert.Equal("normal", received[0]); + } + + private RabbitMQMessageBus CreateBus(bool quorum) + { + return new RabbitMQMessageBus(options => + { + options.ConnectionString(connectionString) + .Topic($"priority-behavior-{Guid.NewGuid():N}") + .SubscriptionQueueName($"priority-behavior-{Guid.NewGuid():N}") + .AcknowledgementStrategy(AcknowledgementStrategy.Automatic) + .PrefetchCount(1) + .PublisherConfirmsEnabled() + .LoggerFactory(Log); + return quorum ? options.UseQuorumQueues() : options.UseMessagePriority(10); + }); + } + + private async Task GetBrokerVersionAsync() + { + var factory = new ConnectionFactory { Uri = new Uri(connectionString) }; + await using var connection = await factory.CreateConnectionAsync(TestCancellationToken); + var version = RabbitMQMessageBus.ParseServerVersion(connection.ServerProperties); + Assert.NotNull(version); + return version; + } + + private async Task PublishQueuedAsync(bool quorum, params (string Data, string? Priority)[] messages) + { + await using var bus = CreateBus(quorum); + var received = new ConcurrentQueue(); + var countdown = new AsyncCountdownEvent(messages.Length); + var warmupReceived = new AsyncManualResetEvent(); + var releaseWarmup = new AsyncManualResetEvent(); + try + { + await bus.SubscribeAsync(async message => + { + if (message.Data == "warmup") + { + warmupReceived.Set(); + await releaseWarmup.WaitAsync(TestCancellationToken); + return; + } + + received.Enqueue(message.Data!); + countdown.Signal(); + }, TestCancellationToken); + + // Fill the single delivery slot so priorities are compared while messages are queued. + await bus.PublishAsync(new SimpleMessageA { Data = "warmup" }, cancellationToken: TestCancellationToken); + await warmupReceived.WaitAsync(TestCancellationToken).WaitAsync(TimeSpan.FromSeconds(10), TestCancellationToken); + foreach (var message in messages) + { + var options = new MessageOptions(); + if (message.Priority is not null) + options.Properties["Priority"] = message.Priority; + await bus.PublishAsync(new SimpleMessageA { Data = message.Data }, options, TestCancellationToken); + } + + releaseWarmup.Set(); + await countdown.WaitAsync(TimeSpan.FromSeconds(10)); + return received.ToArray(); + } + finally + { + releaseWarmup.Set(); + await CleanupMessageBusAsync(bus); + } + } +} diff --git a/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTests.cs b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTests.cs new file mode 100644 index 00000000..ab87dafa --- /dev/null +++ b/tests/Foundatio.RabbitMQ.Tests/Messaging/RabbitMqPriorityBehaviorTests.cs @@ -0,0 +1,6 @@ +using Xunit; + +namespace Foundatio.RabbitMQ.Tests.Messaging; + +public class RabbitMqPriorityBehaviorTests(AspireFixture fixture, ITestOutputHelper output) + : RabbitMqPriorityBehaviorTestBase(fixture.MessagingConnectionString!, output), IClassFixture;