Fix timing issue of workflow resumption by always storing new workflow instances before executing

This commit is contained in:
Sipke Schoorstra 2021-04-07 21:07:00 +02:00
parent 52cd21d0ff
commit a0d345bebd
3 changed files with 21 additions and 15 deletions

View file

@ -32,15 +32,14 @@ namespace Elsa.Decorators
{
var key = $"locking-workflow-runner:{workflowInstance.Id}";
await using (var handle = await _distributedLockProvider.AcquireLockAsync(key, _elsaOptions.DistributedLockTimeout, cancellationToken))
{
if (handle == null)
throw new LockAcquisitionException("Could not acquire a lock within the configured amount of time");
await using var handle = await _distributedLockProvider.AcquireLockAsync(key, _elsaOptions.DistributedLockTimeout, cancellationToken);
if (handle == null)
throw new LockAcquisitionException("Could not acquire a lock within the configured amount of time");
workflowInstance = await _workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, activityId, input, cancellationToken);
await handle.DisposeAsync();
return workflowInstance;
}
workflowInstance = await _workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, activityId, input, cancellationToken);
await handle.DisposeAsync();
return workflowInstance;
}
}
}

View file

@ -34,6 +34,7 @@ namespace Elsa.Dispatch.Handlers
IDistributedLockProvider distributedLockProvider,
IWorkflowDefinitionDispatcher workflowDefinitionDispatcher,
IWorkflowInstanceDispatcher workflowInstanceDispatcher,
IMediator mediator,
ElsaOptions elsaOptions,
ILogger<TriggerWorkflows> logger)
{
@ -43,6 +44,7 @@ namespace Elsa.Dispatch.Handlers
_distributedLockProvider = distributedLockProvider;
_workflowDefinitionDispatcher = workflowDefinitionDispatcher;
_workflowInstanceDispatcher = workflowInstanceDispatcher;
_mediator = mediator;
_elsaOptions = elsaOptions;
_logger = logger;
}
@ -56,15 +58,15 @@ namespace Elsa.Dispatch.Handlers
var lockKey = $"trigger-workflows:correlation:{correlationId}";
await using var handle = await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken);
if(handle == null)
if (handle == null)
throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}");
var correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId) ? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(correlationId), cancellationToken) : 0;
// Release lock before executing workflows. If we don't, we potentially enter a deadlock.
await handle.DisposeAsync();
if (correlatedWorkflowInstanceCount > 0)
{
_logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId);
@ -87,8 +89,9 @@ namespace Elsa.Dispatch.Handlers
foreach (var trigger in triggers)
{
var workflowBlueprint = trigger.WorkflowBlueprint;
await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowBlueprint.Id, trigger.ActivityId, request.Input, request.CorrelationId, request.ContextId, workflowBlueprint.TenantId),
await _mediator.Send(
new ExecuteWorkflowDefinitionRequest(workflowBlueprint.Id, trigger.ActivityId, request.Input, request.CorrelationId, request.ContextId, workflowBlueprint.TenantId),
cancellationToken);
}
@ -98,7 +101,7 @@ namespace Elsa.Dispatch.Handlers
private async Task ResumeWorkflowsAsync(IEnumerable<BookmarkFinderResult> results, object? input, CancellationToken cancellationToken)
{
foreach (var result in results)
await _workflowInstanceDispatcher.DispatchAsync(new ExecuteWorkflowInstanceRequest(result.WorkflowInstanceId, result.ActivityId, input), cancellationToken);
await _mediator.Send(new ExecuteWorkflowInstanceRequest(result.WorkflowInstanceId, result.ActivityId, input), cancellationToken);
}
}
}

View file

@ -5,6 +5,7 @@ using System.Threading.Tasks;
using Elsa.Bookmarks;
using Elsa.Builders;
using Elsa.Models;
using Elsa.Persistence;
using Elsa.Services.Models;
using Elsa.Triggers;
using Open.Linq.AsyncExtensions;
@ -17,13 +18,15 @@ namespace Elsa.Services
private readonly IWorkflowFactory _workflowFactory;
private readonly Func<IWorkflowBuilder> _workflowBuilderFactory;
private readonly IWorkflowRunner _workflowRunner;
private readonly IWorkflowInstanceStore _workflowInstanceStore;
public WorkflowStarter(ITriggerFinder triggerFinder, IWorkflowFactory workflowFactory, Func<IWorkflowBuilder> workflowBuilderFactory, IWorkflowRunner workflowRunner)
public WorkflowStarter(ITriggerFinder triggerFinder, IWorkflowFactory workflowFactory, Func<IWorkflowBuilder> workflowBuilderFactory, IWorkflowRunner workflowRunner, IWorkflowInstanceStore workflowInstanceStore)
{
_triggerFinder = triggerFinder;
_workflowFactory = workflowFactory;
_workflowBuilderFactory = workflowBuilderFactory;
_workflowRunner = workflowRunner;
_workflowInstanceStore = workflowInstanceStore;
}
public async Task FindAndStartWorkflowsAsync(string activityType, IBookmark bookmark, string? tenantId, object? input = default, string? contextId = default, CancellationToken cancellationToken = default)
@ -71,6 +74,7 @@ namespace Elsa.Services
contextId,
cancellationToken);
await _workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken);
return await _workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, activityId, input, cancellationToken);
}