Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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()`)
Expand Down
4 changes: 3 additions & 1 deletion src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBus.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 45 additions & 10 deletions src/Foundatio.RabbitMQ/Messaging/RabbitMQMessageBusOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,13 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions
public TimeSpan? DefaultMessageTimeToLive { get; set; }

/// <summary>
/// Arguments passed to QueueDeclare. Some brokers use it to implement additional features like message TTL.
/// 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&lt;string, object?&gt; { ["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
/// </summary>
public IDictionary<string, object?>? Arguments { get; set; }

Expand Down Expand Up @@ -171,13 +177,37 @@ public class RabbitMQMessageBusOptions : SharedMessageBusOptions
public bool SingleActiveConsumer { get; set; }

/// <summary>
/// Maximum number of priority levels for the queue (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.
/// Maximum message priority for classic queues (1-255). Higher limits cost
/// more broker CPU and memory; UseMessagePriority() limits its convenience API to 32.
/// 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
/// </summary>
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;
}
}

/// <summary>
/// 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.
/// </summary>
internal static bool IsQuorumQueue(IDictionary<string, object?>? arguments)
{
return arguments is not null && arguments.TryGetValue("x-queue-type", out object? queueType)
&& queueType is string type && String.Equals(type, "quorum", StringComparison.OrdinalIgnoreCase);
}

/// <summary>
/// Configures native delayed retry for quorum queues (RabbitMQ 4.3+).
Expand Down Expand Up @@ -343,6 +373,9 @@ public RabbitMQMessageBusOptionsBuilder PublishRecoveryTimeout(TimeSpan timeout)
/// <returns>The builder instance for method chaining.</returns>
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;

Expand Down Expand Up @@ -429,15 +462,17 @@ public RabbitMQMessageBusOptionsBuilder UseSingleActiveConsumer(bool enabled = t
}

/// <summary>
/// Enables message priority on the queue. 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.
/// 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 automatically.
/// Cannot be combined with UseQuorumQueues().
/// </summary>
/// <param name="maxPriority">Maximum priority levels (1-32). Default: 32.</param>
/// <param name="maxPriority">Classic queue maximum priority (1-32). Default: 32.</param>
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;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
using System;
using System.Collections.Concurrent;
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 Xunit;
Expand All @@ -14,7 +17,6 @@ public RabbitMqMessageBusClassicTestBase(string connectionString, ITestOutputHel
{
}


protected override IMessageBus? GetMessageBus(Func<SharedMessageBusOptions, SharedMessageBusOptions>? config = null)
{
if (string.IsNullOrEmpty(ConnectionString))
Expand Down Expand Up @@ -63,4 +65,65 @@ await messageBus.SubscribeAsync<SimpleMessageA>(_ =>
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";
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<string>();
var countdownEvent = new AsyncCountdownEvent(3);
var warmupReceived = new AsyncManualResetEvent();
var releaseWarmup = new AsyncManualResetEvent();

try
{
await bus.SubscribeAsync<SimpleMessageA>(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);

// Act
releaseWarmup.Set();
await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10));

// Assert
Assert.Equal(["high", "medium", "low"], received.ToArray());
}
finally
{
releaseWarmup.Set();
await CleanupMessageBusAsync(bus);
}
}
}
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
using System;
using System.Collections.Concurrent;
using System.Linq;
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;
Expand Down Expand Up @@ -290,43 +290,66 @@ public async Task PublishAsync_WithPriority_DeliversHighPriorityFirst()
// 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 publisher = new RabbitMQMessageBus(o => o
.ConnectionString(ConnectionString)
.Topic(topic)
.SubscriptionQueueName(queueName)
.AcknowledgementStrategy(AcknowledgementStrategy.Automatic)
.UseQuorumQueues()
.UseMessagePriority()
.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);

await Task.Delay(TimeSpan.FromMilliseconds(500), TestCancellationToken);

var received = new ConcurrentQueue<string>();
var countdownEvent = new AsyncCountdownEvent(3);
var warmupReceived = new AsyncManualResetEvent();
var releaseWarmup = new AsyncManualResetEvent();

// Act
await publisher.SubscribeAsync<SimpleMessageA>(msg =>
try
{
received.Enqueue(msg.Data!);
countdownEvent.Signal();
}, TestCancellationToken);

await countdownEvent.WaitAsync(TimeSpan.FromSeconds(10));
await publisher.SubscribeAsync<SimpleMessageA>(async msg =>
{
if (msg.Data == "warmup")
{
warmupReceived.Set();
await releaseWarmup.WaitAsync(TestCancellationToken);
return;
}

received.Enqueue(msg.Data!);
countdownEvent.Signal();
}, TestCancellationToken);

// 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]);
// 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" },
new MessageOptions { Properties = { ["Priority"] = "10" } }, TestCancellationToken);
await publisher.PublishAsync(new SimpleMessageA { Data = "medium" },
new MessageOptions { Properties = { ["Priority"] = "5" } }, TestCancellationToken);

// Act
releaseWarmup.Set();
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
{
releaseWarmup.Set();
await CleanupMessageBusAsync(publisher);
}
}

[Fact]
Expand Down
Loading
Loading