From 279cb1b08d62993f4b95a4dcad20c5ce723dd288 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 16 Jan 2021 12:46:45 +0100 Subject: [PATCH] Fix service bus synchronized access --- .../Extensions/ManagementClientExtensions.cs | 17 ----- .../Extensions/ServiceCollectionExtensions.cs | 5 +- .../Services/IMessageReceiverFactory.cs | 11 +++ .../Services/IMessageSenderFactory.cs | 5 -- .../Services/MessageBusFactory.cs | 73 +++++++++++++++++++ .../Services/MessageReceiverFactory.cs | 44 ----------- .../Services/MessageSenderFactory.cs | 44 ----------- 7 files changed, 87 insertions(+), 112 deletions(-) delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/MessageReceiverFactory.cs delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs deleted file mode 100644 index 02f43aae2..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs +++ /dev/null @@ -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); - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs index fdef35981..3fbd9d01b 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -22,8 +22,9 @@ namespace Elsa.Activities.AzureServiceBus.Extensions options.Services .AddSingleton(CreateServiceBusConnection) .AddSingleton(CreateServiceBusManagementClient) - .AddSingleton() - .AddSingleton() + .AddSingleton() + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()) .AddHostedService() .AddTriggerProvider(); diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs new file mode 100644 index 000000000..716a3bfb8 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs @@ -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 GetReceiverAsync(string queueName, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs index c56792c64..250346053 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs @@ -8,9 +8,4 @@ namespace Elsa.Activities.AzureServiceBus.Services { Task GetSenderAsync(string queueName, CancellationToken cancellationToken = default); } - - public interface IMessageReceiverFactory - { - Task GetReceiverAsync(string queueName, CancellationToken cancellationToken = default); - } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs new file mode 100644 index 000000000..50c143c46 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs @@ -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 _senders = new Dictionary(); + private readonly IDictionary _receivers = new Dictionary(); + private readonly SemaphoreSlim _semaphore = new(1); + + public MessageBusFactory(ServiceBusConnection connection, ManagementClient managementClient) + { + _connection = connection; + _managementClient = managementClient; + } + + public async Task 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 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); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageReceiverFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageReceiverFactory.cs deleted file mode 100644 index 054bd654c..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageReceiverFactory.cs +++ /dev/null @@ -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 _receivers = new Dictionary(); - private readonly SemaphoreSlim _semaphore = new(1); - - public MessageReceiverFactory(ServiceBusConnection connection, ManagementClient managementClient) - { - _connection = connection; - _managementClient = managementClient; - } - - public async Task 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(); - } - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs deleted file mode 100644 index 110d5b7c6..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs +++ /dev/null @@ -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 _senders = new Dictionary(); - private readonly SemaphoreSlim _semaphore = new(1); - - public MessageSenderFactory(ServiceBusConnection connection, ManagementClient managementClient) - { - _connection = connection; - _managementClient = managementClient; - } - - public async Task 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(); - } - } - } -} \ No newline at end of file