|
|
|
@ -11,7 +11,11 @@ namespace Volo.Abp.RabbitMQ |
|
|
|
protected IConnectionPool ConnectionPool { get; } |
|
|
|
|
|
|
|
protected ConcurrentDictionary<string, ChannelPoolItem> 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 |
|
|
|
|