From d70875311f282ddff98a52b45c4ada82c0e96c2d Mon Sep 17 00:00:00 2001 From: Halil ibrahim Kalkan Date: Thu, 26 Jul 2018 16:07:41 +0300 Subject: [PATCH] Dispose channel pool on application shutdown. --- .../AbpBackgroundJobsRabbitMqModule.cs | 7 ++ .../Volo/Abp/Modularity/ModuleManager.cs | 4 +- .../Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs | 6 +- .../Volo/Abp/RabbitMQ/ChannelPool.cs | 119 +++++++++++++++--- .../Volo/Abp/RabbitMQ/ChannelPoolItem.cs | 16 --- .../Volo/Abp/RabbitMQ/ConnectionPool.cs | 2 + .../Volo/Abp/RabbitMQ/IChannelPool.cs | 6 +- 7 files changed, 119 insertions(+), 41 deletions(-) delete mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs diff --git a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs index 5f79068ae2..d57300bdd7 100644 --- a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs +++ b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs @@ -25,5 +25,12 @@ namespace Volo.Abp.BackgroundJobs.RabbitMQ .GetRequiredService() .Start(); } + + public override void OnApplicationShutdown(ApplicationShutdownContext context) + { + context.ServiceProvider + .GetRequiredService() + .Stop(); + } } } diff --git a/framework/src/Volo.Abp.Core/Volo/Abp/Modularity/ModuleManager.cs b/framework/src/Volo.Abp.Core/Volo/Abp/Modularity/ModuleManager.cs index 005fe09cff..57488692e3 100644 --- a/framework/src/Volo.Abp.Core/Volo/Abp/Modularity/ModuleManager.cs +++ b/framework/src/Volo.Abp.Core/Volo/Abp/Modularity/ModuleManager.cs @@ -57,9 +57,11 @@ namespace Volo.Abp.Modularity public void ShutdownModules(ApplicationShutdownContext context) { + var modules = _moduleContainer.Modules.Reverse().ToList(); + foreach (var contributer in _lifecycleContributers) { - foreach (var module in _moduleContainer.Modules) + foreach (var module in modules) { contributer.Shutdown(context, module.Instance); } 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 4c2ae3d5c0..21c31a76bc 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs @@ -17,10 +17,12 @@ namespace Volo.Abp.RabbitMQ public override void OnApplicationShutdown(ApplicationShutdownContext context) { context.ServiceProvider - .GetRequiredService() + .GetRequiredService() .Dispose(); - //TODO: Dispose channel pool when it's implemented! + context.ServiceProvider + .GetRequiredService() + .Dispose(); } } } 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 94fb359cf9..fcab957132 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs @@ -11,7 +11,11 @@ namespace Volo.Abp.RabbitMQ protected IConnectionPool ConnectionPool { get; } protected ConcurrentDictionary Channels { get; } + + protected bool IsDisposed { get; private set; } + protected TimeSpan ChannelDisposeWaitDuration { get; set; } = TimeSpan.FromSeconds(15); + public ChannelPool(IConnectionPool connectionPool) { ConnectionPool = connectionPool; @@ -20,40 +24,115 @@ namespace Volo.Abp.RabbitMQ public virtual IChannelAccessor Acquire(string channelName = null) { + CheckDisposed(); + channelName = channelName ?? ""; - var poolItem = Channels.GetOrAdd(channelName, _ => new ChannelPoolItem + var poolItem = Channels.GetOrAdd( + channelName, + _ => new ChannelPoolItem(CreateChannel(channelName)) + ); + + poolItem.Acquire(); + + return new ChannelAccessor( + poolItem.Channel, + () => poolItem.Release() + ); + } + + protected virtual IModel CreateChannel(string channelName) + { + //TODO: How to determine the right connection name? + return ConnectionPool.Get().CreateModel(); + } + + protected void CheckDisposed() + { + if (IsDisposed) + { + throw new ObjectDisposedException(nameof(ChannelPool)); + } + } + + public void Dispose() + { + if (IsDisposed) { - Channel = CreateChannel(channelName) - }); + return; + } - lock (poolItem) + IsDisposed = true; + + foreach (var poolItem in Channels.Values) { - while (poolItem.IsInUse) + try { - Monitor.Wait(poolItem); + poolItem.WaitIfInUse(ChannelDisposeWaitDuration); + poolItem.Channel.Dispose(); } - - poolItem.IsInUse = true; + catch + { } } - return new ChannelAccessor( - poolItem.Channel, - () => + Channels.Clear(); + } + + protected class ChannelPoolItem : IDisposable + { + public IModel Channel { get; } + + public bool IsInUse + { + get => _isInUse; + private set => _isInUse = value; + } + private volatile bool _isInUse; + + public ChannelPoolItem(IModel channel) + { + Channel = channel; + } + + public void Acquire() + { + lock (this) { - lock (poolItem) + while (IsInUse) { - poolItem.IsInUse = false; - Monitor.PulseAll(poolItem); + Monitor.Wait(this); } + + IsInUse = true; } - ); - } + } - protected virtual IModel CreateChannel(string channelName) - { - //TODO: How to determine the right connection name? - return ConnectionPool.Get().CreateModel(); + public void WaitIfInUse(TimeSpan timeout) + { + lock (this) + { + if (!IsInUse) + { + return; + } + + Monitor.Wait(this, timeout); + } + } + + public void Release() + { + lock (this) + { + IsInUse = false; + Monitor.PulseAll(this); + } + } + + public void Dispose() + { + Channel.Dispose(); + } } protected class ChannelAccessor : IChannelAccessor diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs deleted file mode 100644 index 2fac563c54..0000000000 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPoolItem.cs +++ /dev/null @@ -1,16 +0,0 @@ -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 9121bc1994..51ebb1a047 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs @@ -53,6 +53,8 @@ namespace Volo.Abp.RabbitMQ } } + + Connections.Clear(); } } } \ No newline at end of file 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 36c011d098..6c3251bd16 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs @@ -1,6 +1,8 @@ -namespace Volo.Abp.RabbitMQ +using System; + +namespace Volo.Abp.RabbitMQ { - public interface IChannelPool + public interface IChannelPool : IDisposable { IChannelAccessor Acquire(string channelName = null); }