From 2d2fa4451fc76134c60d98fa744d75710874d4df Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 28 May 2021 19:36:16 +0200 Subject: [PATCH] Fix concurrent access to Azure Service Bus workers --- .../Services/ServiceBusQueuesStarter.cs | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs index f972a228a..3894484db 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs @@ -19,6 +19,7 @@ namespace Elsa.Activities.AzureServiceBus.Services private readonly IServiceProvider _serviceProvider; private readonly ILogger _logger; private readonly ICollection _workers; + private readonly SemaphoreSlim _semaphore = new(1); public ServiceBusQueuesStarter( IQueueMessageReceiverClientFactory messageReceiverClientFactory, @@ -35,11 +36,20 @@ namespace Elsa.Activities.AzureServiceBus.Services public async Task CreateWorkersAsync(CancellationToken cancellationToken = default) { - await DisposeExistingWorkersAsync(); - var queueNames = (await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); + await _semaphore.WaitAsync(cancellationToken); - foreach (var queueName in queueNames) - await CreateAndAddWorkerAsync(queueName, cancellationToken); + try + { + await DisposeExistingWorkersAsync(); + var queueNames = (await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); + + foreach (var queueName in queueNames) + await CreateAndAddWorkerAsync(queueName, cancellationToken); + } + finally + { + _semaphore.Release(); + } } private async Task CreateAndAddWorkerAsync(string queueName, CancellationToken cancellationToken)