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();