From e51e743d78fcdbd58da4afa584a80da3646ff7c4 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 21 May 2021 14:15:18 +0200 Subject: [PATCH] Fix workflow instance locking --- .../Decorators/LockingWorkflowRunner.cs | 29 +++++++++---------- .../ElsaServiceCollectionExtensions.cs | 2 +- .../Services/WorkflowInstanceExecutor.cs | 5 +++- .../Elsa.Core/Services/WorkflowLaunchpad.cs | 3 +- 4 files changed, 20 insertions(+), 19 deletions(-) diff --git a/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs b/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs index 468983cbd..04b82f070 100644 --- a/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs +++ b/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs @@ -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 RunWorkflowAsync( - IWorkflowBlueprint workflowBlueprint, - WorkflowInstance workflowInstance, - string? activityId = default, - object? input = default, - CancellationToken cancellationToken = default) + public async Task 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 ExecuteAsync(WorkflowInstance workflowInstance, string? activityId, object? input = default, CancellationToken cancellationToken = default) => + await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index a35e13bba..42275dfe7 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -81,7 +81,7 @@ namespace Microsoft.Extensions.DependencyInjection services.Decorate(); services.Decorate(); services.Decorate(); - services.Decorate(); + services.Decorate(); return services; } diff --git a/src/core/Elsa.Core/Services/WorkflowInstanceExecutor.cs b/src/core/Elsa.Core/Services/WorkflowInstanceExecutor.cs index f3e68036a..03b426d4a 100644 --- a/src/core/Elsa.Core/Services/WorkflowInstanceExecutor.cs +++ b/src/core/Elsa.Core/Services/WorkflowInstanceExecutor.cs @@ -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 ExecuteAsync(WorkflowInstance workflowInstance, string? activityId, object? input = default, CancellationToken cancellationToken = default) diff --git a/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs b/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs index e06a680c6..e489ddd59 100644 --- a/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs +++ b/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs @@ -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 ExecuteStartableWorkflowAsync(StartableWorkflow startableWorkflow, object? input, CancellationToken cancellationToken = default) => await _workflowRunner.RunWorkflowAsync(startableWorkflow.WorkflowBlueprint, startableWorkflow.WorkflowInstance, startableWorkflow.ActivityId, input, cancellationToken); + public async Task ExecuteStartableWorkflowAsync(StartableWorkflow startableWorkflow, object? input, CancellationToken cancellationToken = default) => + await _workflowRunner.RunWorkflowAsync(startableWorkflow.WorkflowBlueprint, startableWorkflow.WorkflowInstance, startableWorkflow.ActivityId, input, cancellationToken); public async Task DispatchStartableWorkflowAsync(StartableWorkflow startableWorkflow, object? input, CancellationToken cancellationToken = default) {