Fix workflow instance locking

This commit is contained in:
Sipke Schoorstra 2021-05-21 14:15:18 +02:00
parent bd012d0965
commit e51e743d78
4 changed files with 20 additions and 19 deletions

View file

@ -7,39 +7,36 @@ using Elsa.Services.Models;
namespace Elsa.Decorators
{
public class LockingWorkflowRunner : IWorkflowRunner
public class LockingWorkflowInstanceExecutor : IWorkflowInstanceExecutor
{
private readonly IWorkflowRunner _workflowRunner;
private readonly IWorkflowInstanceExecutor _workflowInstanceExecutor;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly ElsaOptions _elsaOptions;
public LockingWorkflowRunner(
IWorkflowRunner workflowRunner,
public LockingWorkflowInstanceExecutor(
IWorkflowInstanceExecutor workflowInstanceExecutor,
IDistributedLockProvider distributedLockProvider,
ElsaOptions elsaOptions)
{
_workflowRunner = workflowRunner;
_workflowInstanceExecutor = workflowInstanceExecutor;
_distributedLockProvider = distributedLockProvider;
_elsaOptions = elsaOptions;
}
public async Task<RunWorkflowResult> RunWorkflowAsync(
IWorkflowBlueprint workflowBlueprint,
WorkflowInstance workflowInstance,
string? activityId = default,
object? input = default,
CancellationToken cancellationToken = default)
public async Task<RunWorkflowResult> ExecuteAsync(string workflowInstanceId, string? activityId, object? input = default, CancellationToken cancellationToken = default)
{
var key = $"locking-workflow-runner:{workflowInstance.Id}";
var key = $"workflow-instance:{workflowInstanceId}";
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");
var runWorkflowResult = await _workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, activityId, input, cancellationToken);
var result = await _workflowInstanceExecutor.ExecuteAsync(workflowInstanceId, activityId, input, cancellationToken);
await handle.DisposeAsync();
return runWorkflowResult;
return result;
}
public async Task<RunWorkflowResult> ExecuteAsync(WorkflowInstance workflowInstance, string? activityId, object? input = default, CancellationToken cancellationToken = default) =>
await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken);
}
}

View file

@ -81,7 +81,7 @@ namespace Microsoft.Extensions.DependencyInjection
services.Decorate<IWorkflowDefinitionStore, InitializingWorkflowDefinitionStore>();
services.Decorate<IWorkflowDefinitionStore, EventPublishingWorkflowDefinitionStore>();
services.Decorate<IWorkflowInstanceStore, EventPublishingWorkflowInstanceStore>();
services.Decorate<IWorkflowRunner, LockingWorkflowRunner>();
services.Decorate<IWorkflowInstanceExecutor, LockingWorkflowInstanceExecutor>();
return services;
}

View file

@ -28,7 +28,10 @@ namespace Elsa.Services
if (!ValidatePreconditions(workflowInstanceId, workflowInstance, activityId))
return new RunWorkflowResult(workflowInstance, activityId, false);
return await ExecuteAsync(workflowInstance!, activityId, input, cancellationToken);
return await _workflowRunner.ResumeWorkflowAsync(
workflowInstance!,
activityId,
input, cancellationToken);
}
public async Task<RunWorkflowResult> ExecuteAsync(WorkflowInstance workflowInstance, string? activityId, object? input = default, CancellationToken cancellationToken = default)

View file

@ -210,7 +210,8 @@ namespace Elsa.Services
public Task DispatchPendingWorkflowAsync(string workflowInstanceId, string? activityId, object? input, CancellationToken cancellationToken = default) => DispatchPendingWorkflowAsync(new PendingWorkflow(workflowInstanceId, activityId), input, cancellationToken);
public async Task<RunWorkflowResult> ExecuteStartableWorkflowAsync(StartableWorkflow startableWorkflow, object? input, CancellationToken cancellationToken = default) => await _workflowRunner.RunWorkflowAsync(startableWorkflow.WorkflowBlueprint, startableWorkflow.WorkflowInstance, startableWorkflow.ActivityId, input, cancellationToken);
public async Task<RunWorkflowResult> ExecuteStartableWorkflowAsync(StartableWorkflow startableWorkflow, object? input, CancellationToken cancellationToken = default) =>
await _workflowRunner.RunWorkflowAsync(startableWorkflow.WorkflowBlueprint, startableWorkflow.WorkflowInstance, startableWorkflow.ActivityId, input, cancellationToken);
public async Task<PendingWorkflow> DispatchStartableWorkflowAsync(StartableWorkflow startableWorkflow, object? input, CancellationToken cancellationToken = default)
{