From 2a6d8f9bfecaa182b30ae53a7b7d011ce89dadeb Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 19 Jan 2021 18:19:36 +0100 Subject: [PATCH] Synchronize Azure Service Bus queue worker based on correlation ID to prevent multiple workflow instances with same correlation ID --- .../Services/QueueWorker.cs | 59 ++++++++++++++----- 1 file changed, 44 insertions(+), 15 deletions(-) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index ba74eb7bf..0b389757b 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -1,8 +1,11 @@ using System; +using System.Diagnostics; using System.Linq; +using System.Runtime.InteropServices; using System.Threading; using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Triggers; +using Elsa.DistributedLock; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; @@ -20,12 +23,14 @@ namespace Elsa.Activities.AzureServiceBus.Services { private readonly IMessageReceiver _messageReceiver; private readonly IServiceProvider _serviceProvider; + private readonly IDistributedLockProvider _distributedLockProvider; private readonly ILogger _logger; - public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, ILogger logger) + public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, IDistributedLockProvider distributedLockProvider, ILogger logger) { _messageReceiver = messageReceiver; _serviceProvider = serviceProvider; + _distributedLockProvider = distributedLockProvider; _logger = logger; _messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler) @@ -47,36 +52,60 @@ namespace Elsa.Activities.AzureServiceBus.Services using var scope = _serviceProvider.CreateScope(); var workflowRunner = scope.ServiceProvider.GetRequiredService(); var queueName = _messageReceiver.Path; + var correlationId = message.CorrelationId; async Task TriggerNewWorkflowAsync() { await workflowRunner!.TriggerWorkflowsAsync( x => x.QueueName == queueName && x.CorrelationId == null, message, - message.CorrelationId, + correlationId, cancellationToken: cancellationToken); } - - if (string.IsNullOrWhiteSpace(message.CorrelationId)) + + if (string.IsNullOrWhiteSpace(correlationId)) { await TriggerNewWorkflowAsync(); return; } - var workflowSelector = scope.ServiceProvider.GetRequiredService(); - var workflowInstanceStore = scope.ServiceProvider.GetRequiredService(); - var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification(message.CorrelationId), cancellationToken); + var lockKey = $"azure-service-bus:{queueName}:correlation-{correlationId}"; + var stopwatch = new Stopwatch(); - if (correlatedWorkflowInstanceCount > 0) + _logger.LogDebug("Acquiring lock {LockKey}", lockKey); + stopwatch.Start(); + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken)) { - // Trigger existing workflows (if blocked on this message). - var existingWorkflows = await workflowSelector.SelectWorkflowsAsync(x => x.QueueName == queueName && x.CorrelationId == message.CorrelationId, cancellationToken).ToList(); - await workflowRunner.TriggerWorkflowsAsync(existingWorkflows, message, message.CorrelationId, cancellationToken: cancellationToken); + _logger.LogDebug("Lock {LockKey} already taken", lockKey); + return; } - else + + try { - // Trigger new workflow. - await TriggerNewWorkflowAsync(); + var workflowSelector = scope.ServiceProvider.GetRequiredService(); + var workflowInstanceStore = scope.ServiceProvider.GetRequiredService(); + var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification(message.CorrelationId), cancellationToken); + + if (correlatedWorkflowInstanceCount > 0) + { + // Trigger existing workflows (if blocked on this message). + _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}'. Resuming them", correlatedWorkflowInstanceCount, correlationId); + var existingWorkflows = await workflowSelector.SelectWorkflowsAsync(x => x.QueueName == queueName && x.CorrelationId == message.CorrelationId, cancellationToken).ToList(); + await workflowRunner.TriggerWorkflowsAsync(existingWorkflows, message, message.CorrelationId, cancellationToken: cancellationToken); + } + else + { + // Trigger new workflow. + _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); + await TriggerNewWorkflowAsync(); + } + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken); + stopwatch.Stop(); + _logger.LogDebug("Lock held for {ElapseTime}", stopwatch.Elapsed); } } @@ -85,7 +114,7 @@ namespace Elsa.Activities.AzureServiceBus.Services switch (e.Exception) { case MessageLockLostException: - _logger.LogDebug( e.Exception,"Message lock lost"); + _logger.LogDebug(e.Exception, "Message lock lost"); break; case ServiceBusCommunicationException: _logger.LogDebug(e.Exception, "Lost service bus communication");