diff --git a/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs b/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs index 2a644b9d3..d9291c8ba 100644 --- a/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs +++ b/src/core/Elsa.Core/Decorators/LockingWorkflowRunner.cs @@ -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; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs index df371908e..3aa163650 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs @@ -34,6 +34,7 @@ namespace Elsa.Dispatch.Handlers IDistributedLockProvider distributedLockProvider, IWorkflowDefinitionDispatcher workflowDefinitionDispatcher, IWorkflowInstanceDispatcher workflowInstanceDispatcher, + IMediator mediator, ElsaOptions elsaOptions, ILogger 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(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 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); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowStarter.cs b/src/core/Elsa.Core/Services/WorkflowStarter.cs index 8bde14e9c..01fb89cfa 100644 --- a/src/core/Elsa.Core/Services/WorkflowStarter.cs +++ b/src/core/Elsa.Core/Services/WorkflowStarter.cs @@ -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 _workflowBuilderFactory; private readonly IWorkflowRunner _workflowRunner; + private readonly IWorkflowInstanceStore _workflowInstanceStore; - public WorkflowStarter(ITriggerFinder triggerFinder, IWorkflowFactory workflowFactory, Func workflowBuilderFactory, IWorkflowRunner workflowRunner) + public WorkflowStarter(ITriggerFinder triggerFinder, IWorkflowFactory workflowFactory, Func 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); }