Cleanup + comments

This commit is contained in:
Sipke Schoorstra 2021-06-02 14:00:56 +02:00
parent 93e80b5248
commit cdaa105f3e
2 changed files with 7 additions and 18 deletions

View file

@ -34,13 +34,11 @@ namespace Elsa.Decorators
public async Task<RunWorkflowResult> 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;
}
}

View file

@ -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<WorkflowInstance>(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<BookmarkFinderResult>();
_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();