From 9defb83ba8d3a7a4f2785a45b9eccd5e87656f69 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 9 Apr 2021 17:46:31 +0200 Subject: [PATCH] Lock on correlation ID to prevent race condition between incoming events and newly started workflows not yet persisted --- .../Handlers/ExecuteWorkflowDefinition.cs | 23 ++++++++++++++++++- .../Dispatch/Handlers/TriggerWorkflows.cs | 5 +++- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs b/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs index c6cc55032..88e218554 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs @@ -1,3 +1,4 @@ +using System; using System.Threading; using System.Threading.Tasks; using Elsa.Exceptions; @@ -7,8 +8,10 @@ using Elsa.Persistence.Specifications; using Elsa.Persistence.Specifications.WorkflowInstances; using Elsa.Services; using Elsa.Services.Models; +using Medallion.Threading; using MediatR; using Microsoft.Extensions.Logging; +using IDistributedLockProvider = Elsa.Services.IDistributedLockProvider; namespace Elsa.Dispatch.Handlers { @@ -47,14 +50,32 @@ namespace Elsa.Dispatch.Handlers return Unit.Value; var lockKey = $"execute-workflow-definition:tenant:{tenantId}:workflow-definition:{workflowDefinitionId}"; + var correlationId = request.CorrelationId; + var correlationLockHandle = default(IDistributedSynchronizationHandle?); - await using var handle = await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken); + if (!string.IsNullOrWhiteSpace(correlationId)) + { + _logger.LogDebug("Acquiring lock on correlation ID {CorrelationId}", correlationId); + correlationLockHandle = await _distributedLockProvider.AcquireLockAsync(correlationId, _elsaOptions.DistributedLockTimeout, cancellationToken); + } + + if(correlationLockHandle == null) + throw new LockAcquisitionException($"Failed to acquire a lock on correlation ID {correlationId}"); + + try + { + await using var handle = await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken); if(handle == null) throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}"); if (!workflowBlueprint!.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false) await _startsWorkflow.StartWorkflowAsync(workflowBlueprint, request.ActivityId, request.Input, request.CorrelationId, request.ContextId, cancellationToken); + } + finally + { + await correlationLockHandle.DisposeAsync(); + } return Unit.Value; } diff --git a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs index e0d071c50..1968a0efe 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs @@ -56,8 +56,9 @@ namespace Elsa.Dispatch.Handlers if (!string.IsNullOrWhiteSpace(correlationId)) { - var lockKey = $"trigger-workflows:correlation:{correlationId}"; + var lockKey = correlationId; + _logger.LogDebug("Acquiring lock on correlation ID {CorrelationId}", correlationId); await using var handle = await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken); if (handle == null) @@ -67,6 +68,8 @@ namespace Elsa.Dispatch.Handlers ? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId).WithStatus(WorkflowStatus.Suspended), cancellationToken) : 0; + _logger.LogDebug("Found {CorrelatedWorkflowCount} correlated workflows,", correlatedWorkflowInstanceCount); + if (correlatedWorkflowInstanceCount > 0) { _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId);