diff --git a/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs b/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs index 945dc56bf..9d1f18753 100644 --- a/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs +++ b/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs @@ -34,13 +34,11 @@ namespace Elsa.Decorators public async Task ExecuteAsync(string workflowInstanceId, string? activityId, object? input = default, CancellationToken cancellationToken = default) { var workflowInstanceLockKey = $"workflow-instance:{workflowInstanceId}"; - _logger.LogDebug("Acquiring lock on {LockKey}", workflowInstanceLockKey); await using var workflowInstanceLockHandle = await _distributedLockProvider.AcquireLockAsync(workflowInstanceLockKey, _elsaOptions.DistributedLockTimeout, cancellationToken); if (workflowInstanceLockHandle == null) throw new LockAcquisitionException("Could not acquire a lock within the configured amount of time"); - - _logger.LogDebug("Lock acquired on {LockKey}", workflowInstanceLockKey); + var workflowInstance = await _workflowInstanceStore.FindByIdAsync(workflowInstanceId, cancellationToken); if (workflowInstance == null) @@ -53,22 +51,20 @@ namespace Elsa.Decorators if (!string.IsNullOrWhiteSpace(correlationId)) { - _logger.LogDebug("Acquiring lock on correlation {CorrelationId}", correlationId); + // We need to lock on correlation ID to prevent a race condition with WorkflowLaunchpad that is used to find workflows by correlation ID to execute. + // The race condition is: when a workflow instance is done executing, the BookmarkIndexer will collect bookmarks. + // But if in the meantime an event comes in that triggers correlated workflows, the bookmarks may not have been created yet. await using var correlationLockHandle = await _distributedLockProvider.AcquireLockAsync(correlationId, _elsaOptions.DistributedLockTimeout, cancellationToken); if(correlationLockHandle == null) throw new LockAcquisitionException($"Could not acquire a lock on correlation {correlationId} within the configured amount of time"); - _logger.LogDebug("Lock acquired on correlation {CorrelationId}", correlationId); var result = await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); - _logger.LogDebug("Released lock on correlation {CorrelationId}", correlationId); - _logger.LogDebug("Released lock on {LockKey}", workflowInstanceLockKey); return result; } else { var result = await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); - _logger.LogDebug("Released lock on {LockKey}", workflowInstanceLockKey); return result; } } diff --git a/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs b/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs index 6daa6af24..90a7bff60 100644 --- a/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs +++ b/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs @@ -163,8 +163,6 @@ namespace Elsa.Services if (correlationLockHandle != null) await correlationLockHandle.DisposeAsync(); } - - return null; } public async Task CollectAndExecuteStartableWorkflowAsync(string workflowDefinitionId, string? activityId, string? correlationId = default, string? contextId = default, object? input = default, string? tenantId = default, CancellationToken cancellationToken = default) @@ -265,30 +263,25 @@ namespace Elsa.Services { var correlationId = context.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) throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}"); - - _logger.LogDebug("Acquired lock on correlation ID {CorrelationId}", correlationId); - + var correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId) ? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId).WithStatus(WorkflowStatus.Suspended), cancellationToken) : 0; - _logger.LogDebug("Found {CorrelatedWorkflowCount} correlated workflows,", correlatedWorkflowInstanceCount); + _logger.LogDebug("Found {{CorrelatedWorkflowCount}} workflows with correlation ID {CorrelationId}", correlatedWorkflowInstanceCount, correlationId); if (correlatedWorkflowInstanceCount > 0) { - _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId); var bookmarkResults = context.Bookmark != null ? await _bookmarkFinder.FindBookmarksAsync(context.ActivityType, context.Bookmark, correlationId, context.TenantId, cancellationToken).ToList() : new List(); _logger.LogDebug("Found {BookmarkCount} bookmarks for activity type {ActivityType}", bookmarkResults.Count, context.ActivityType); return bookmarkResults.Select(x => new PendingWorkflow(x.WorkflowInstanceId, x.ActivityId)).ToList(); } } - _logger.LogDebug("Released lock on correlation ID {CorrelationId}", correlationId); var startableWorkflows = await CollectStartableWorkflowsAsync(context, cancellationToken); return startableWorkflows.Select(x => new PendingWorkflow(x.WorkflowInstance.Id, x.ActivityId)).ToList();