Fix service bus synchronized access

This commit is contained in:
Sipke Schoorstra 2021-01-16 12:46:45 +01:00
parent 7d576a7ddc
commit 279cb1b08d
7 changed files with 87 additions and 112 deletions

View file

@ -1,17 +0,0 @@
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Azure.ServiceBus.Management;
namespace Elsa.Activities.AzureServiceBus.Extensions
{
public static class ManagementClientExtensions
{
public static async Task EnsureQueueExistsAsync(this ManagementClient managementClient, string queueName, CancellationToken cancellationToken)
{
if (await managementClient.QueueExistsAsync(queueName, cancellationToken))
return;
await managementClient.CreateQueueAsync(queueName, cancellationToken);
}
}
}

View file

@ -22,8 +22,9 @@ namespace Elsa.Activities.AzureServiceBus.Extensions
options.Services
.AddSingleton(CreateServiceBusConnection)
.AddSingleton(CreateServiceBusManagementClient)
.AddSingleton<IMessageSenderFactory, MessageSenderFactory>()
.AddSingleton<IMessageReceiverFactory, MessageReceiverFactory>()
.AddSingleton<MessageBusFactory>()
.AddSingleton<IMessageSenderFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddSingleton<IMessageReceiverFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddHostedService<StartServiceBusQueues>()
.AddTriggerProvider<MessageReceivedTriggerProvider>();

View file

@ -0,0 +1,11 @@
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Azure.ServiceBus.Core;
namespace Elsa.Activities.AzureServiceBus.Services
{
public interface IMessageReceiverFactory
{
Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken = default);
}
}

View file

@ -8,9 +8,4 @@ namespace Elsa.Activities.AzureServiceBus.Services
{
Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken = default);
}
public interface IMessageReceiverFactory
{
Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,73 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Extensions;
using Microsoft.Azure.ServiceBus;
using Microsoft.Azure.ServiceBus.Core;
using Microsoft.Azure.ServiceBus.Management;
namespace Elsa.Activities.AzureServiceBus.Services
{
public class MessageBusFactory : IMessageSenderFactory, IMessageReceiverFactory
{
private readonly ServiceBusConnection _connection;
private readonly ManagementClient _managementClient;
private readonly IDictionary<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
private readonly IDictionary<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
private readonly SemaphoreSlim _semaphore = new(1);
public MessageBusFactory(ServiceBusConnection connection, ManagementClient managementClient)
{
_connection = connection;
_managementClient = managementClient;
}
public async Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken)
{
await _semaphore.WaitAsync(cancellationToken);
try
{
if (_senders.TryGetValue(queueName, out var messageSender))
return messageSender;
await EnsureQueueExistsAsync(queueName, cancellationToken);
var newMessageSender = new MessageSender(_connection, queueName);
_senders.Add(queueName, newMessageSender);
return newMessageSender;
}
finally
{
_semaphore.Release();
}
}
public async Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken)
{
await _semaphore.WaitAsync(cancellationToken);
if (_receivers.TryGetValue(queueName, out var messageReceiver))
return messageReceiver;
try
{
await EnsureQueueExistsAsync(queueName, cancellationToken);
var newMessageReceiver = new MessageReceiver(_connection, queueName);
_receivers.Add(queueName, newMessageReceiver);
return newMessageReceiver;
}
finally
{
_semaphore.Release();
}
}
private async Task EnsureQueueExistsAsync(string queueName, CancellationToken cancellationToken)
{
if (await _managementClient.QueueExistsAsync(queueName, cancellationToken))
return;
await _managementClient.CreateQueueAsync(queueName, cancellationToken);
}
}
}

View file

@ -1,44 +0,0 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Extensions;
using Microsoft.Azure.ServiceBus;
using Microsoft.Azure.ServiceBus.Core;
using Microsoft.Azure.ServiceBus.Management;
namespace Elsa.Activities.AzureServiceBus.Services
{
public class MessageReceiverFactory : IMessageReceiverFactory
{
private readonly ServiceBusConnection _connection;
private readonly ManagementClient _managementClient;
private readonly IDictionary<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
private readonly SemaphoreSlim _semaphore = new(1);
public MessageReceiverFactory(ServiceBusConnection connection, ManagementClient managementClient)
{
_connection = connection;
_managementClient = managementClient;
}
public async Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken)
{
if (_receivers.TryGetValue(queueName, out var messageReceiver))
return messageReceiver;
await _semaphore.WaitAsync(cancellationToken);
try
{
await _managementClient.EnsureQueueExistsAsync(queueName, cancellationToken);
var newMessageReceiver = new MessageReceiver(_connection, queueName);
_receivers.Add(queueName, newMessageReceiver);
return newMessageReceiver;
}
finally
{
_semaphore.Release();
}
}
}
}

View file

@ -1,44 +0,0 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Extensions;
using Microsoft.Azure.ServiceBus;
using Microsoft.Azure.ServiceBus.Core;
using Microsoft.Azure.ServiceBus.Management;
namespace Elsa.Activities.AzureServiceBus.Services
{
public class MessageSenderFactory : IMessageSenderFactory
{
private readonly ServiceBusConnection _connection;
private readonly ManagementClient _managementClient;
private readonly IDictionary<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
private readonly SemaphoreSlim _semaphore = new(1);
public MessageSenderFactory(ServiceBusConnection connection, ManagementClient managementClient)
{
_connection = connection;
_managementClient = managementClient;
}
public async Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken)
{
await _semaphore.WaitAsync(cancellationToken);
try
{
if (_senders.TryGetValue(queueName, out var messageSender))
return messageSender;
await _managementClient.EnsureQueueExistsAsync(queueName, cancellationToken);
var newMessageSender = new MessageSender(_connection, queueName);
_senders.Add(queueName, newMessageSender);
return newMessageSender;
}
finally
{
_semaphore.Release();
}
}
}
}