diff --git a/Directory.Packages.props b/Directory.Packages.props index 193cccefae..0af612b056 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -142,7 +142,7 @@ - + diff --git a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs index 1dbe4a6dcc..bdc28c7ca0 100644 --- a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs +++ b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Globalization; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; @@ -120,46 +121,44 @@ public class JobQueue : IJobQueue ChannelAccessor?.Dispose(); } - protected virtual Task EnsureInitializedAsync() + protected virtual async Task EnsureInitializedAsync() { if (ChannelAccessor != null && ChannelAccessor.Channel.IsOpen) { - return Task.CompletedTask; + return; } - ChannelAccessor = ChannelPool.Acquire( + ChannelAccessor = await ChannelPool.AcquireAsync( ChannelPrefix + QueueConfiguration.QueueName, QueueConfiguration.ConnectionName ); - var result = QueueConfiguration.Declare(ChannelAccessor.Channel); + var result = await QueueConfiguration.DeclareAsync(ChannelAccessor.Channel); Logger.LogDebug($"RabbitMQ Queue '{QueueConfiguration.QueueName}' has {result.MessageCount} messages and {result.ConsumerCount} consumers."); // Declare delayed queue - QueueConfiguration.DeclareDelayed(ChannelAccessor.Channel); + await QueueConfiguration.DeclareDelayedAsync(ChannelAccessor.Channel); if (AbpBackgroundJobOptions.IsJobExecutionEnabled) { if (QueueConfiguration.PrefetchCount.HasValue) { - ChannelAccessor.Channel.BasicQos(0, QueueConfiguration.PrefetchCount.Value, false); + await ChannelAccessor.Channel.BasicQosAsync(0, QueueConfiguration.PrefetchCount.Value, false); } - + Consumer = new AsyncEventingBasicConsumer(ChannelAccessor.Channel); - Consumer.Received += MessageReceived; - + Consumer.ReceivedAsync += MessageReceived; + //TODO: What BasicConsume returns? - ChannelAccessor.Channel.BasicConsume( + await ChannelAccessor.Channel.BasicConsumeAsync( queue: QueueConfiguration.QueueName, autoAck: false, consumer: Consumer ); } - - return Task.CompletedTask; } - protected virtual Task PublishAsync( + protected virtual async Task PublishAsync( TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null) @@ -167,29 +166,27 @@ public class JobQueue : IJobQueue //TODO: How to handle priority var routingKey = QueueConfiguration.QueueName; - var basicProperties = CreateBasicPropertiesToPublish(); + var basicProperties = new BasicProperties + { + Persistent = true + }; if (delay.HasValue) { routingKey = QueueConfiguration.DelayedQueueName; - basicProperties.Expiration = delay.Value.TotalMilliseconds.ToString(); + basicProperties.Expiration = delay.Value.TotalMilliseconds.ToString(CultureInfo.InvariantCulture); } - ChannelAccessor!.Channel.BasicPublish( - exchange: "", - routingKey: routingKey, - basicProperties: basicProperties, - body: Serializer.Serialize(args!) - ); - - return Task.CompletedTask; - } - - protected virtual IBasicProperties CreateBasicPropertiesToPublish() - { - var properties = ChannelAccessor!.Channel.CreateBasicProperties(); - properties.Persistent = true; - return properties; + if (ChannelAccessor != null) + { + await ChannelAccessor.Channel.BasicPublishAsync( + exchange: "", + routingKey: routingKey, + mandatory: false, + basicProperties: basicProperties, + body: Serializer.Serialize(args!) + ); + } } protected virtual async Task MessageReceived(object sender, BasicDeliverEventArgs ea) @@ -205,17 +202,17 @@ public class JobQueue : IJobQueue try { await JobExecuter.ExecuteAsync(context); - ChannelAccessor!.Channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false); + await ChannelAccessor!.Channel.BasicAckAsync(deliveryTag: ea.DeliveryTag, multiple: false); } catch (BackgroundJobExecutionException) { //TODO: Reject like that? - ChannelAccessor!.Channel.BasicReject(deliveryTag: ea.DeliveryTag, requeue: true); + await ChannelAccessor!.Channel.BasicRejectAsync(deliveryTag: ea.DeliveryTag, requeue: true); } catch (Exception) { //TODO: Reject like that? - ChannelAccessor!.Channel.BasicReject(deliveryTag: ea.DeliveryTag, requeue: false); + await ChannelAccessor!.Channel.BasicRejectAsync(deliveryTag: ea.DeliveryTag, requeue: false); } } } diff --git a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueueConfiguration.cs b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueueConfiguration.cs index 9425cd0604..c6418f52f7 100644 --- a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueueConfiguration.cs +++ b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueueConfiguration.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Threading.Tasks; using RabbitMQ.Client; using Volo.Abp.RabbitMQ; @@ -34,15 +35,15 @@ public class JobQueueConfiguration : QueueDeclareConfiguration DelayedQueueName = delayedQueueName; } - public virtual QueueDeclareOk DeclareDelayed(IModel channel) + public virtual async Task DeclareDelayedAsync(IChannel channel) { - var delayedArguments = new Dictionary(Arguments) + var delayedArguments = new Dictionary(Arguments) { ["x-dead-letter-routing-key"] = QueueName, ["x-dead-letter-exchange"] = string.Empty }; - return channel.QueueDeclare( + return await channel.QueueDeclareAsync( queue: DelayedQueueName, durable: Durable, exclusive: Exclusive, diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/AbpRabbitMqEventBusOptions.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/AbpRabbitMqEventBusOptions.cs index 130f5b38e2..0b450ca26e 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/AbpRabbitMqEventBusOptions.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/AbpRabbitMqEventBusOptions.cs @@ -17,10 +17,10 @@ public class AbpRabbitMqEventBusOptions public ushort? PrefetchCount { get; set; } - public IDictionary QueueArguments { get; set; } = new Dictionary(); + public IDictionary QueueArguments { get; set; } = new Dictionary(); + + public IDictionary ExchangeArguments { get; set; } = new Dictionary(); - public IDictionary ExchangeArguments { get; set; } = new Dictionary(); - public string GetExchangeTypeOrDefault() { return string.IsNullOrEmpty(ExchangeType) diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs index b0f85afb56..3a647388b5 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs @@ -97,7 +97,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis SubscribeHandlers(AbpDistributedEventBusOptions.Handlers); } - private async Task ProcessEventAsync(IModel channel, BasicDeliverEventArgs ea) + private async Task ProcessEventAsync(IChannel channel, BasicDeliverEventArgs ea) { var eventName = ea.RoutingKey; var eventType = EventTypes.GetOrDefault(eventName); @@ -224,10 +224,10 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis IEnumerable outgoingEvents, OutboxConfig outboxConfig) { - using (var channel = ConnectionPool.Get(AbpRabbitMqEventBusOptions.ConnectionName).CreateModel()) + using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) + .CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true, new ThrottlingRateLimiter(256)))) { var outgoingEventArray = outgoingEvents.ToArray(); - channel.ConfirmSelect(); foreach (var outgoingEvent in outgoingEventArray) { @@ -248,8 +248,6 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis }); } } - - channel.WaitForConfirmsOrDie(); } } @@ -293,31 +291,33 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis return PublishAsync(eventName, body, headersArguments, eventId, correlationId); } - protected virtual Task PublishAsync( + protected virtual async Task PublishAsync( string eventName, byte[] body, Dictionary? headersArguments = null, Guid? eventId = null, string? correlationId = null) { - using (var channel = ConnectionPool.Get(AbpRabbitMqEventBusOptions.ConnectionName).CreateModel()) + using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)).CreateChannelAsync()) { - return PublishAsync(channel, eventName, body, headersArguments, eventId, correlationId); + await PublishAsync(channel, eventName, body, headersArguments, eventId, correlationId); } } - protected virtual Task PublishAsync( - IModel channel, + protected virtual async Task PublishAsync( + IChannel channel, string eventName, byte[] body, Dictionary? headersArguments = null, Guid? eventId = null, string? correlationId = null) { - EnsureExchangeExists(channel); + await EnsureExchangeExistsAsync(channel); - var properties = channel.CreateBasicProperties(); - properties.DeliveryMode = RabbitMqConsts.DeliveryModes.Persistent; + var properties = new BasicProperties + { + DeliveryMode = DeliveryModes.Persistent + }; if (properties.MessageId.IsNullOrEmpty()) { @@ -331,18 +331,16 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis SetEventMessageHeaders(properties, headersArguments); - channel.BasicPublish( + await channel.BasicPublishAsync( exchange: AbpRabbitMqEventBusOptions.ExchangeName, routingKey: eventName, - mandatory: true, + mandatory: false, basicProperties: properties, body: body ); - - return Task.CompletedTask; } - private void EnsureExchangeExists(IModel channel) + protected virtual async Task EnsureExchangeExistsAsync(IChannel channel) { if (_exchangeCreated) { @@ -351,14 +349,14 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis try { - using (var temporaryChannel = ConnectionPool.Get(AbpRabbitMqEventBusOptions.ConnectionName).CreateModel()) + using (var temporaryChannel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)).CreateChannelAsync()) { - temporaryChannel.ExchangeDeclarePassive(AbpRabbitMqEventBusOptions.ExchangeName); + await temporaryChannel.ExchangeDeclarePassiveAsync(AbpRabbitMqEventBusOptions.ExchangeName); } } catch (Exception) { - channel.ExchangeDeclare( + await channel.ExchangeDeclareAsync( AbpRabbitMqEventBusOptions.ExchangeName, AbpRabbitMqEventBusOptions.GetExchangeTypeOrDefault(), durable: true @@ -367,14 +365,14 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis _exchangeCreated = true; } - private void SetEventMessageHeaders(IBasicProperties properties, Dictionary? headersArguments) + protected virtual void SetEventMessageHeaders(IBasicProperties properties, Dictionary? headersArguments) { if (headersArguments == null) { return; } - properties.Headers ??= new Dictionary(); + properties.Headers ??= new Dictionary(); foreach (var header in headersArguments) { diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs index da010b83c9..096ba83a3a 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs @@ -19,7 +19,6 @@ public class AbpRabbitMqModule : AbpModule { foreach (var connectionFactory in options.Connections.Values) { - connectionFactory.DispatchConsumersAsync = true; connectionFactory.AutomaticRecoveryEnabled = false; } }); diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs index b7ed4c80c8..bbb594a5f7 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs @@ -3,10 +3,12 @@ using System.Collections.Concurrent; using System.Diagnostics; using System.Linq; using System.Threading; +using System.Threading.Tasks; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; using RabbitMQ.Client; using Volo.Abp.DependencyInjection; +using Volo.Abp.Threading; namespace Volo.Abp.RabbitMQ; @@ -16,6 +18,8 @@ public class ChannelPool : IChannelPool, ISingletonDependency protected ConcurrentDictionary Channels { get; } + protected SemaphoreSlim Semaphore = new SemaphoreSlim(1, 1); + protected bool IsDisposed { get; private set; } protected TimeSpan TotalDisposeWaitDuration { get; set; } = TimeSpan.FromSeconds(10); @@ -29,16 +33,33 @@ public class ChannelPool : IChannelPool, ISingletonDependency Logger = NullLogger.Instance; } - public virtual IChannelAccessor Acquire(string? channelName = null, string? connectionName = null) + public virtual async Task AcquireAsync(string? channelName = null, string? connectionName = null) { CheckDisposed(); channelName = channelName ?? ""; - var poolItem = Channels.GetOrAdd( - channelName, - _ => new ChannelPoolItem(CreateChannel(channelName, connectionName)) - ); + ChannelPoolItem poolItem; + + if (Channels.TryGetValue(channelName, out var existingChannelPoolItem)) + { + poolItem = existingChannelPoolItem; + } + else + { + using (await Semaphore.LockAsync()) + { + if (Channels.TryGetValue(channelName, out var existingChannelPoolItem2)) + { + poolItem = existingChannelPoolItem2; + } + else + { + poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); + Channels.TryAdd(channelName, poolItem); + } + } + } poolItem.Acquire(); @@ -46,11 +67,20 @@ public class ChannelPool : IChannelPool, ISingletonDependency { poolItem.Dispose(); Channels.TryRemove(channelName, out _); - poolItem = Channels.GetOrAdd( - channelName, - _ => new ChannelPoolItem(CreateChannel(channelName, connectionName)) - ); - + + using (await Semaphore.LockAsync()) + { + if (Channels.TryGetValue(channelName, out var existingChannelPoolItem3)) + { + poolItem = existingChannelPoolItem3; + } + else + { + poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); + Channels.TryAdd(channelName, poolItem); + } + } + poolItem.Acquire(); } @@ -61,14 +91,14 @@ public class ChannelPool : IChannelPool, ISingletonDependency ); } - protected virtual IModel CreateChannel(string channelName, string? connectionName) + protected virtual async Task CreateChannelAsync(string channelName, string? connectionName) { - return ConnectionPool - .Get(connectionName) - .CreateModel(); + return await (await ConnectionPool + .GetAsync(connectionName)) + .CreateChannelAsync(); } - protected void CheckDisposed() + protected virtual void CheckDisposed() { if (IsDisposed) { @@ -130,7 +160,7 @@ public class ChannelPool : IChannelPool, ISingletonDependency protected class ChannelPoolItem : IDisposable { - public IModel Channel { get; } + public IChannel Channel { get; } public bool IsInUse { get => _isInUse; @@ -138,7 +168,7 @@ public class ChannelPool : IChannelPool, ISingletonDependency } private volatile bool _isInUse; - public ChannelPoolItem(IModel channel) + public ChannelPoolItem(IChannel channel) { Channel = channel; } @@ -186,13 +216,13 @@ public class ChannelPool : IChannelPool, ISingletonDependency protected class ChannelAccessor : IChannelAccessor { - public IModel Channel { get; } + public IChannel Channel { get; } public string Name { get; } private readonly Action _disposeAction; - public ChannelAccessor(IModel channel, string name, Action disposeAction) + public ChannelAccessor(IChannel channel, string name, Action disposeAction) { _disposeAction = disposeAction; Name = name; diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs index 5f4556b96a..938b0e4f8c 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs @@ -1,9 +1,11 @@ using System; using System.Collections.Concurrent; -using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; using Microsoft.Extensions.Options; using RabbitMQ.Client; using Volo.Abp.DependencyInjection; +using Volo.Abp.Threading; namespace Volo.Abp.RabbitMQ; @@ -11,52 +13,71 @@ public class ConnectionPool : IConnectionPool, ISingletonDependency { protected AbpRabbitMqOptions Options { get; } - protected ConcurrentDictionary> Connections { get; } + protected ConcurrentDictionary Connections { get; } + + protected SemaphoreSlim Semaphore = new SemaphoreSlim(1, 1); private bool _isDisposed; public ConnectionPool(IOptions options) { Options = options.Value; - Connections = new ConcurrentDictionary>(); + Connections = new ConcurrentDictionary(); } - public virtual IConnection Get(string? connectionName = null) + public virtual async Task GetAsync(string? connectionName = null) { connectionName ??= RabbitMqConnections.DefaultConnectionName; - var connectionFactory = Options.Connections.GetOrDefault(connectionName); - try + + IConnection connection; + + if (Connections.TryGetValue(connectionName, out var existingConnection)) { - var connection = GetConnection(connectionName, connectionFactory); - - if (connection.IsOpen) - { - return connection; - } - - connection.Dispose(); - Connections.TryRemove(connectionName, out _); - return GetConnection(connectionName, connectionFactory); + connection = existingConnection; } - catch (Exception) + else { - Connections.TryRemove(connectionName, out _); - throw; + using (await Semaphore.LockAsync()) + { + try + { + var connectionFactory = Options.Connections.GetOrDefault(connectionName); + if (Connections.TryGetValue(connectionName, out var existingConnection2)) + { + connection = existingConnection2; + } + else + { + connection = await GetConnectionAsync(connectionName, connectionFactory); + Connections.TryAdd(connectionName, connection); + + if (!connection.IsOpen) + { + connection.Dispose(); + Connections.TryRemove(connectionName, out _); + connection = await GetConnectionAsync(connectionName, connectionFactory); + Connections.TryAdd(connectionName, connection); + } + } + } + catch (Exception) + { + Connections.TryRemove(connectionName, out _); + throw; + } + } } + + return connection; } - protected virtual IConnection GetConnection(string connectionName, ConnectionFactory connectionFactory) + protected virtual async Task GetConnectionAsync(string connectionName, ConnectionFactory connectionFactory) { - return Connections.GetOrAdd( - connectionName, () => new Lazy(() => - { - var hostnames = connectionFactory.HostName.TrimEnd(';').Split(';'); - // Handle Rabbit MQ Cluster. - return hostnames.Length == 1 - ? connectionFactory.CreateConnection() - : connectionFactory.CreateConnection(hostnames); - }) - ).Value; + var hostnames = connectionFactory.HostName.TrimEnd(';').Split(';'); + // Handle Rabbit MQ Cluster. + return hostnames.Length == 1 + ? await connectionFactory.CreateConnectionAsync() + : await connectionFactory.CreateConnectionAsync(hostnames); } public void Dispose() @@ -72,11 +93,11 @@ public class ConnectionPool : IConnectionPool, ISingletonDependency { try { - connection.Value.Dispose(); + connection.Dispose(); } catch { - + // ignored } } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs index 5c239893c7..4a5c0f01d0 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs @@ -12,19 +12,19 @@ public class ExchangeDeclareConfiguration public bool AutoDelete { get; set; } - public IDictionary Arguments { get; } + public IDictionary Arguments { get; } public ExchangeDeclareConfiguration( string exchangeName, string type, bool durable = false, bool autoDelete = false, - IDictionary? arguments = null) + IDictionary? arguments = null) { ExchangeName = exchangeName; Type = type; Durable = durable; AutoDelete = autoDelete; - Arguments = arguments?? new Dictionary(); + Arguments = arguments?? new Dictionary(); } } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs index ad8829f8ee..d4eb954a67 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs @@ -10,7 +10,7 @@ public interface IChannelAccessor : IDisposable /// Never dispose the object. /// Instead, dispose the after usage. /// - IModel Channel { get; } + IChannel Channel { get; } /// /// Name of the channel. diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs index 2ba6259fec..06b1cc0ba6 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs @@ -1,8 +1,9 @@ using System; +using System.Threading.Tasks; namespace Volo.Abp.RabbitMQ; public interface IChannelPool : IDisposable { - IChannelAccessor Acquire(string? channelName = null, string? connectionName = null); + Task AcquireAsync(string? channelName = null, string? connectionName = null); } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs index dc97476c84..cca44eea89 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs @@ -1,9 +1,10 @@ using System; +using System.Threading.Tasks; using RabbitMQ.Client; namespace Volo.Abp.RabbitMQ; public interface IConnectionPool : IDisposable { - IConnection Get(string? connectionName = null); + Task GetAsync(string? connectionName = null); } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqMessageConsumer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqMessageConsumer.cs index a24ad69a4a..8501b8958f 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqMessageConsumer.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqMessageConsumer.cs @@ -11,5 +11,5 @@ public interface IRabbitMqMessageConsumer Task UnbindAsync(string routingKey); - void OnMessageReceived(Func callback); + void OnMessageReceived(Func callback); } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs index eb59ec3058..ebb2c3e1cc 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs @@ -1,4 +1,5 @@ using System.Collections.Generic; +using System.Threading.Tasks; using JetBrains.Annotations; using RabbitMQ.Client; @@ -13,10 +14,10 @@ public class QueueDeclareConfiguration public bool Exclusive { get; set; } public bool AutoDelete { get; set; } - - public ushort? PrefetchCount { get; set; } - public IDictionary Arguments { get; } + public ushort? PrefetchCount { get; set; } + + public IDictionary Arguments { get; } public QueueDeclareConfiguration( [NotNull] string queueName, @@ -24,19 +25,19 @@ public class QueueDeclareConfiguration bool exclusive = false, bool autoDelete = false, ushort? prefetchCount = null, - IDictionary? arguments = null) + IDictionary? arguments = null) { QueueName = queueName; Durable = durable; Exclusive = exclusive; AutoDelete = autoDelete; - Arguments = arguments?? new Dictionary(); + Arguments = arguments?? new Dictionary(); PrefetchCount = prefetchCount; } - public virtual QueueDeclareOk Declare(IModel channel) + public virtual async Task DeclareAsync(IChannel channel) { - return channel.QueueDeclare( + return await channel.QueueDeclareAsync( queue: QueueName, durable: Durable, exclusive: Exclusive, diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs index 6f361f7670..74dac68a61 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs @@ -5,6 +5,7 @@ using RabbitMQ.Client; using RabbitMQ.Client.Events; using System; using System.Collections.Concurrent; +using System.Threading; using System.Threading.Tasks; using Volo.Abp.DependencyInjection; using Volo.Abp.ExceptionHandling; @@ -28,13 +29,13 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen protected string? ConnectionName { get; private set; } - protected ConcurrentBag> Callbacks { get; } + protected ConcurrentBag> Callbacks { get; } - protected IModel? Channel { get; private set; } + protected IChannel? Channel { get; private set; } protected ConcurrentQueue QueueBindCommands { get; } - protected object ChannelSendSyncLock { get; } = new object(); + protected SemaphoreSlim Semaphore = new SemaphoreSlim(1, 1); public RabbitMqMessageConsumer( IConnectionPool connectionPool, @@ -47,7 +48,7 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen Logger = NullLogger.Instance; QueueBindCommands = new ConcurrentQueue(); - Callbacks = new ConcurrentBag>(); + Callbacks = new ConcurrentBag>(); Timer.Period = 5000; //5 sec. Timer.Elapsed = Timer_Elapsed; @@ -88,21 +89,21 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen return; } - lock (ChannelSendSyncLock) + using (await Semaphore.LockAsync()) { if (QueueBindCommands.TryPeek(out var command)) { switch (command.Type) { case QueueBindType.Bind: - Channel.QueueBind( + await Channel.QueueBindAsync( queue: Queue.QueueName, exchange: Exchange.ExchangeName, routingKey: command.RoutingKey ); break; case QueueBindType.Unbind: - Channel.QueueUnbind( + await Channel.QueueUnbindAsync( queue: Queue.QueueName, exchange: Exchange.ExchangeName, routingKey: command.RoutingKey @@ -124,7 +125,7 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen } } - public virtual void OnMessageReceived(Func callback) + public virtual void OnMessageReceived(Func callback) { Callbacks.Add(callback); } @@ -144,11 +145,11 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen try { - Channel = ConnectionPool - .Get(ConnectionName) - .CreateModel(); + Channel = await (await ConnectionPool + .GetAsync(ConnectionName)) + .CreateChannelAsync(); - Channel.ExchangeDeclare( + await Channel.ExchangeDeclareAsync( exchange: Exchange.ExchangeName, type: Exchange.Type, durable: Exchange.Durable, @@ -156,7 +157,7 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen arguments: Exchange.Arguments ); - Channel.QueueDeclare( + await Channel.QueueDeclareAsync( queue: Queue.QueueName, durable: Queue.Durable, exclusive: Queue.Exclusive, @@ -166,13 +167,13 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen if (Queue.PrefetchCount.HasValue) { - Channel.BasicQos(0, Queue.PrefetchCount.Value, false); + await Channel.BasicQosAsync(0, Queue.PrefetchCount.Value, false); } - + var consumer = new AsyncEventingBasicConsumer(Channel); - consumer.Received += HandleIncomingMessageAsync; - - Channel.BasicConsume( + consumer.ReceivedAsync += HandleIncomingMessageAsync; + + await Channel.BasicConsumeAsync( queue: Queue.QueueName, autoAck: false, consumer: consumer @@ -194,17 +195,23 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen await callback(Channel!, basicDeliverEventArgs); } - Channel?.BasicAck(basicDeliverEventArgs.DeliveryTag, multiple: false); + if (Channel != null) + { + await Channel.BasicAckAsync(basicDeliverEventArgs.DeliveryTag, multiple: false); + } } catch (Exception ex) { try { - Channel?.BasicNack( - basicDeliverEventArgs.DeliveryTag, - multiple: false, - requeue: true - ); + if (Channel != null) + { + await Channel.BasicNackAsync( + basicDeliverEventArgs.DeliveryTag, + multiple: false, + requeue: true + ); + } } // ReSharper disable once EmptyGeneralCatchClause catch { }