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 b05fbca525..d34e131187 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 @@ -111,7 +111,7 @@ namespace Volo.Abp.BackgroundJobs.RabbitMQ return Task.CompletedTask; } - ChannelAccessor = ChannelPool.Acquire(QueueName); + ChannelAccessor = ChannelPool.Acquire(QueueName + ".JobQueue"); var queueOptions = RabbitMqOptions.Queues.GetOrDefault(QueueName) ?? new QueueOptions(QueueName); 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 cf3bf40455..4c2ae3d5c0 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs @@ -16,7 +16,10 @@ namespace Volo.Abp.RabbitMQ public override void OnApplicationShutdown(ApplicationShutdownContext context) { - context.ServiceProvider.GetRequiredService().Dispose(); + context.ServiceProvider + .GetRequiredService() + .Dispose(); + //TODO: Dispose channel pool when it's implemented! } } 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 3e3d7ab6ad..94fb359cf9 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs @@ -1,4 +1,7 @@ -using RabbitMQ.Client; +using System; +using System.Collections.Concurrent; +using System.Threading; +using RabbitMQ.Client; using Volo.Abp.DependencyInjection; namespace Volo.Abp.RabbitMQ @@ -7,16 +10,43 @@ namespace Volo.Abp.RabbitMQ { protected IConnectionPool ConnectionPool { get; } + protected ConcurrentDictionary Channels { get; } + public ChannelPool(IConnectionPool connectionPool) { ConnectionPool = connectionPool; + Channels = new ConcurrentDictionary(); } public virtual IChannelAccessor Acquire(string channelName = null) { - //TODO: Pool channels! + channelName = channelName ?? ""; + + var poolItem = Channels.GetOrAdd(channelName, _ => new ChannelPoolItem + { + Channel = CreateChannel(channelName) + }); + + lock (poolItem) + { + while (poolItem.IsInUse) + { + Monitor.Wait(poolItem); + } + + poolItem.IsInUse = true; + } + return new ChannelAccessor( - CreateChannel(channelName) + poolItem.Channel, + () => + { + lock (poolItem) + { + poolItem.IsInUse = false; + Monitor.PulseAll(poolItem); + } + } ); } @@ -29,15 +59,17 @@ namespace Volo.Abp.RabbitMQ protected class ChannelAccessor : IChannelAccessor { public IModel Channel { get; } + private readonly Action _disposeAction; - public ChannelAccessor(IModel channel) + public ChannelAccessor(IModel channel, Action disposeAction) { + _disposeAction = disposeAction; Channel = channel; } public void Dispose() { - Channel.Dispose(); + _disposeAction.Invoke(); } } } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs new file mode 100644 index 0000000000..2fac563c54 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs @@ -0,0 +1,16 @@ +using RabbitMQ.Client; + +namespace Volo.Abp.RabbitMQ +{ + public class ChannelPoolItem + { + public IModel Channel { get; set; } + + public bool IsInUse + { + get => _isInUse; + set => _isInUse = value; + } + private volatile bool _isInUse; + } +} \ No newline at end of file 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 e354ca90ec..9121bc1994 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs @@ -1,5 +1,4 @@ -using System; -using System.Collections.Concurrent; +using System.Collections.Concurrent; using System.Collections.Generic; using Microsoft.Extensions.Options; using RabbitMQ.Client; @@ -7,7 +6,7 @@ using Volo.Abp.DependencyInjection; namespace Volo.Abp.RabbitMQ { - public class ConnectionPool : IConnectionPool, IDisposable, ISingletonDependency + public class ConnectionPool : IConnectionPool, ISingletonDependency { protected AbpRabbitMqOptions Options { get; } 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 32a6584d62..43b6f0023d 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs @@ -5,6 +5,11 @@ namespace Volo.Abp.RabbitMQ { public interface IChannelAccessor : IDisposable { + /// + /// Reference to the channel. + /// Never dispose the object. + /// Instead, dispose the after usage. + /// IModel Channel { get; } } } \ No newline at end of file diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/Jobs/WriteToConsoleGreenJob.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/Jobs/WriteToConsoleGreenJob.cs index f2d2854f1f..ed89d648e8 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/Jobs/WriteToConsoleGreenJob.cs +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/Jobs/WriteToConsoleGreenJob.cs @@ -9,7 +9,7 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Shared.Jobs { if (RandomHelper.GetRandom(0, 100) < 70) { - throw new ApplicationException("A sample exception from the WriteToConsoleGreenJob!"); + //throw new ApplicationException("A sample exception from the WriteToConsoleGreenJob!"); } var oldColor = Console.ForegroundColor;