Lock on correlation ID to prevent race condition between incoming events and newly started workflows not yet persisted
This commit is contained in:
parent
ed7eb136f5
commit
9defb83ba8
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<WorkflowInstance>(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);
|
||||
|
|
|
|||
Loading…
Reference in a new issue