From feff07bda59d6c3ef2a3efaf681a072629ba1735 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 2 Nov 2021 20:57:42 +0100 Subject: [PATCH] Make ServiceBusFactory thread-safe --- .../Services/Messaging/ServiceBusFactory.cs | 25 +++++++++++++------ 1 file changed, 18 insertions(+), 7 deletions(-) diff --git a/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs b/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs index 4da28f9c9..c626768c0 100644 --- a/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs +++ b/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs @@ -1,6 +1,7 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Threading; using Elsa.Options; using Elsa.Serialization; using Microsoft.Extensions.Logging; @@ -20,6 +21,7 @@ namespace Elsa.Services.Messaging private readonly IDictionary _serviceBuses = new Dictionary(); private readonly IDictionary _messageTypeQueueDictionary = new Dictionary(); private readonly DependencyInjectionHandlerActivator _handlerActivator; + private readonly SemaphoreSlim _semaphore = new(1); public ServiceBusFactory(ElsaOptions elsaOptions, ILoggerFactory loggerFactory, IServiceProvider serviceProvider) { @@ -67,15 +69,24 @@ namespace Elsa.Services.Messaging private IBus GetOrCreateServiceBus(Type messageType, string? queueName) { - queueName ??= _messageTypeQueueDictionary[messageType]; - - if (!_serviceBuses.TryGetValue(queueName, out var bus)) + _semaphore.Wait(); + + try { - bus = ConfigureServiceBus(new[] { messageType }, queueName); - _serviceBuses[queueName] = bus; - } + queueName ??= _messageTypeQueueDictionary[messageType]; - return bus; + if (!_serviceBuses.TryGetValue(queueName, out var bus)) + { + bus = ConfigureServiceBus(new[] { messageType }, queueName); + _serviceBuses[queueName] = bus; + } + + return bus; + } + finally + { + _semaphore.Release(); + } } private string PrefixQueueName(string name) => $"{_elsaOptions.ServiceBusOptions.QueuePrefix}{name}";