From 01dcf2707f02873b5a76b4b637d99aebadbb7ca7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Benoi=CC=82t=20Belz?= Date: Tue, 9 Jun 2026 15:17:07 +0200 Subject: [PATCH] Use AsyncEventHandler for consumer message handling Convert the consumer's ConsumeAsync callback from the synchronous EventHandler to RabbitMQ's AsyncEventHandler, awaited inside ReceivedAsync. MessagingEvent moves into the MerQure project and derives from AsyncEventArgs so it satisfies the delegate constraint. Serialization switches from a lock to a SemaphoreSlim to allow awaiting. --- samples/MerQure.Samples/DeadLetterExample.cs | 4 +-- samples/MerQure.Samples/SimpleExample.cs | 6 ++-- samples/MerQure.Samples/StopExample.cs | 2 +- src/MerQure.RbMQ/Clients/Consumer.cs | 33 +++++++++++-------- src/MerQure.Tools/Buses/Consumer.cs | 1 + src/MerQure/Clients/IConsumer.cs | 2 +- .../Events/MessagingEvent.cs | 7 ++-- .../Buses/ConsumerTests.cs | 8 +++-- 8 files changed, 37 insertions(+), 26 deletions(-) rename src/{MerQure.RbMQ => MerQure}/Events/MessagingEvent.cs (66%) diff --git a/samples/MerQure.Samples/DeadLetterExample.cs b/samples/MerQure.Samples/DeadLetterExample.cs index 47faeb8..80d9670 100644 --- a/samples/MerQure.Samples/DeadLetterExample.cs +++ b/samples/MerQure.Samples/DeadLetterExample.cs @@ -51,12 +51,12 @@ public async Task RunAsync() // Get the consumer on the existing queue and consume its messages var consumer = await _messagingService.GetConsumerAsync("deadletter.queue"); - await consumer.ConsumeAsync((object sender, IMessagingEvent args) => + await consumer.ConsumeAsync((object sender, MessagingEvent args) => { var realDelay = DateTime.Now.Subtract(dateStart).TotalSeconds; Console.WriteLine(string.Format("{0} received after {1:#.##}s.", args.Message.GetRoutingKey(), realDelay)); // send ACK: acknowlegdment to the queue - consumer.AcknowlegdeDeliveredMessageAsync(args); + return consumer.AcknowlegdeDeliveredMessageAsync(args).AsTask(); }); } } \ No newline at end of file diff --git a/samples/MerQure.Samples/SimpleExample.cs b/samples/MerQure.Samples/SimpleExample.cs index d994612..e5c87c1 100644 --- a/samples/MerQure.Samples/SimpleExample.cs +++ b/samples/MerQure.Samples/SimpleExample.cs @@ -32,20 +32,20 @@ public async Task RunAsync() // Get the consumer on the existing queue and consume its messages var consumer = await _messagingService.GetConsumerAsync("simple.queue"); var random = new Random(); - await consumer.ConsumeAsync((object sender, IMessagingEvent args) => + await consumer.ConsumeAsync((object sender, MessagingEvent args) => { // we simulate the delivery success if (random.Next() % 2 == 0) { Console.WriteLine("Retry " + args.Message.GetRoutingKey()); // send NACK: negative acknowlegdment to the queue - consumer.RejectDeliveredMessageAsync(args); + return consumer.RejectDeliveredMessageAsync(args).AsTask(); } else { Console.WriteLine(args.Message.GetBody()); // send ACK: acknowlegdment to the queue - consumer.AcknowlegdeDeliveredMessageAsync(args); + return consumer.AcknowlegdeDeliveredMessageAsync(args).AsTask(); } }); } diff --git a/samples/MerQure.Samples/StopExample.cs b/samples/MerQure.Samples/StopExample.cs index 9e28d6a..c1ec97b 100644 --- a/samples/MerQure.Samples/StopExample.cs +++ b/samples/MerQure.Samples/StopExample.cs @@ -39,7 +39,7 @@ await consumer.ConsumeAsync((_, args) => Thread.Sleep(10); Console.WriteLine(args.Message.GetBody()); // send ACK: acknowlegdment to the queue - consumer.AcknowlegdeDeliveredMessageAsync(args); + return consumer.AcknowlegdeDeliveredMessageAsync(args).AsTask(); }); // Stop Consuming after 100 ms ~ 10 messages diff --git a/src/MerQure.RbMQ/Clients/Consumer.cs b/src/MerQure.RbMQ/Clients/Consumer.cs index d328436..ef29e14 100644 --- a/src/MerQure.RbMQ/Clients/Consumer.cs +++ b/src/MerQure.RbMQ/Clients/Consumer.cs @@ -1,12 +1,12 @@ using MerQure.Messages; using MerQure.RbMQ.Content; -using MerQure.RbMQ.Events; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System; using System.Collections.Generic; using System.Linq; using System.Text; +using System.Threading; using System.Threading.Tasks; namespace MerQure.RbMQ.Clients; @@ -17,33 +17,37 @@ class Consumer : RabbitMqClient, IConsumer private AsyncEventingBasicConsumer _consumer; private readonly ushort _prefetchCount; - private readonly object _consumingLock; + private readonly SemaphoreSlim _consumingLock; public Consumer(IChannel channel, string queueName, ushort prefetchCount) : base(channel) { QueueName = queueName.ToLowerInvariant(); _prefetchCount = prefetchCount; - _consumingLock = new object(); + _consumingLock = new SemaphoreSlim(1, 1); } - public async Task ConsumeAsync(EventHandler onMessageReceived) + public async Task ConsumeAsync(AsyncEventHandler onMessageReceived) { await Channel.BasicQosAsync(0, _prefetchCount, false); _consumer = new AsyncEventingBasicConsumer(Channel); - _consumer.ReceivedAsync += (sender, args) => + _consumer.ReceivedAsync += async (sender, args) => { if (onMessageReceived != null) { - lock (_consumingLock) + await _consumingLock.WaitAsync(); + try { var message = ParseDeliveredMessage(args); var messageEventArgs = new MessagingEvent(message, args.DeliveryTag.ToString()); - onMessageReceived(sender, messageEventArgs); + await onMessageReceived(sender, messageEventArgs); + } + finally + { + _consumingLock.Release(); } } - return Task.CompletedTask; }; await Channel.BasicConsumeAsync(QueueName, false, _consumer); @@ -89,17 +93,18 @@ public async Task StopConsuming(AsyncEventHandler onConsumerS { if (IsConsuming()) { - lock (_consumingLock) + await _consumingLock.WaitAsync(); + try { if (onConsumerStopped != null) { - _consumer.UnregisteredAsync += (sender, e) => - { - onConsumerStopped(sender, e); - return Task.CompletedTask; - }; + _consumer.UnregisteredAsync += onConsumerStopped; } } + finally + { + _consumingLock.Release(); + } // Must be outside the lock to avoid deadlock foreach (var tag in _consumer.ConsumerTags) diff --git a/src/MerQure.Tools/Buses/Consumer.cs b/src/MerQure.Tools/Buses/Consumer.cs index 0672e82..97ac6ca 100644 --- a/src/MerQure.Tools/Buses/Consumer.cs +++ b/src/MerQure.Tools/Buses/Consumer.cs @@ -24,6 +24,7 @@ public async Task ConsumeAsync(Channel channel, EventHandler callback) await consumer.ConsumeAsync((_, messagingEvent) => { OnMessageReceived(callback, messagingEvent); + return Task.CompletedTask; }); } diff --git a/src/MerQure/Clients/IConsumer.cs b/src/MerQure/Clients/IConsumer.cs index 7b07de2..bdb2ef0 100644 --- a/src/MerQure/Clients/IConsumer.cs +++ b/src/MerQure/Clients/IConsumer.cs @@ -19,7 +19,7 @@ public interface IConsumer : IAsyncDisposable /// Start listening on the queue /// /// Handler called each time a message arrives for this consumer. - Task ConsumeAsync(EventHandler onMessageReceived); + Task ConsumeAsync(AsyncEventHandler onMessageReceived); /// /// Indicates if the Consumer is registred on the queue and waiting for messages diff --git a/src/MerQure.RbMQ/Events/MessagingEvent.cs b/src/MerQure/Events/MessagingEvent.cs similarity index 66% rename from src/MerQure.RbMQ/Events/MessagingEvent.cs rename to src/MerQure/Events/MessagingEvent.cs index 9161e60..39f219c 100644 --- a/src/MerQure.RbMQ/Events/MessagingEvent.cs +++ b/src/MerQure/Events/MessagingEvent.cs @@ -1,6 +1,9 @@ -namespace MerQure.RbMQ.Events +using MerQure.Messages; +using RabbitMQ.Client.Events; + +namespace MerQure { - class MessagingEvent : IMessagingEvent + public class MessagingEvent : AsyncEventArgs, IMessagingEvent { public IMessage Message { get; set; } diff --git a/tests/MerQure.Tools.Tests/Buses/ConsumerTests.cs b/tests/MerQure.Tools.Tests/Buses/ConsumerTests.cs index 494810c..0caa89b 100644 --- a/tests/MerQure.Tools.Tests/Buses/ConsumerTests.cs +++ b/tests/MerQure.Tools.Tests/Buses/ConsumerTests.cs @@ -1,8 +1,10 @@ -using MerQure.Tools.Buses; +using MerQure.Messages; +using MerQure.Tools.Buses; using MerQure.Tools.Configurations; using MerQure.Tools.Messages; using Moq; using Newtonsoft.Json; +using RabbitMQ.Client.Events; using System; using System.Threading.Tasks; using Xunit; @@ -18,9 +20,9 @@ public class ConsumerTests : IDisposable public ConsumerTests() { _mockMerQureConsumer = new Mock(); - _mockMerQureConsumer.Setup(m => m.ConsumeAsync(It.IsAny>())).Callback((EventHandler action) => + _mockMerQureConsumer.Setup(m => m.ConsumeAsync(It.IsAny>())).Callback((AsyncEventHandler action) => { - action(this, new Mock().Object); + action(this, new MessagingEvent(new Mock().Object, "1")); }); _mockMessagingService = new Mock();