From e39d5cdb681de80c90b985e4f35f5283cff683e6 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 21 Jan 2021 12:11:25 +0100 Subject: [PATCH] Ensure workflow instance locking when receiving Azure Service Bus messages --- .../Services/QueueWorker.cs | 10 ++++------ src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs | 11 +++++++++++ src/core/Elsa.Core/Services/WorkflowQueue.cs | 7 ++++++- 3 files changed, 21 insertions(+), 7 deletions(-) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index e7a369a16..a01ed4200 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -48,18 +48,16 @@ namespace Elsa.Activities.AzureServiceBus.Services private async Task TriggerWorkflowsAsync(Message message, CancellationToken cancellationToken) { using var scope = _serviceProvider.CreateScope(); - var workflowRunner = scope.ServiceProvider.GetRequiredService(); + var workflowRunner = scope.ServiceProvider.GetRequiredService(); var queueName = _messageReceiver.Path; var correlationId = message.CorrelationId; - async Task TriggerNewWorkflowAsync() - { - await workflowRunner!.TriggerWorkflowsAsync( + async Task TriggerNewWorkflowAsync() => + await workflowRunner.EnqueueWorkflowsAsync( x => x.QueueName == queueName && x.CorrelationId == null, message, correlationId, cancellationToken: cancellationToken); - } if (string.IsNullOrWhiteSpace(correlationId)) { @@ -90,7 +88,7 @@ namespace Elsa.Activities.AzureServiceBus.Services // Trigger existing workflows (if blocked on this message). _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}'. Resuming them", correlatedWorkflowInstanceCount, correlationId); var existingWorkflows = await workflowSelector.SelectWorkflowsAsync(x => x.QueueName == queueName && x.CorrelationId == message.CorrelationId, cancellationToken).ToList(); - await workflowRunner.TriggerWorkflowsAsync(existingWorkflows, message, message.CorrelationId, cancellationToken: cancellationToken); + await workflowRunner.EnqueueWorkflowsAsync(existingWorkflows, message, message.CorrelationId, cancellationToken: cancellationToken); } else { diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs b/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs index 8d9cbdf08..69213ff0b 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Triggers; @@ -18,6 +19,16 @@ namespace Elsa.Services CancellationToken cancellationToken = default) where TTrigger : ITrigger; + /// + /// Enqueues the specified workflows for execution. + /// + Task EnqueueWorkflowsAsync( + IEnumerable results, + object? input = default, + string? correlationId = default, + string? contextId = default, + CancellationToken cancellationToken = default); + /// /// Enqueues the specified workflow instance and activity for execution. /// diff --git a/src/core/Elsa.Core/Services/WorkflowQueue.cs b/src/core/Elsa.Core/Services/WorkflowQueue.cs index 97fb90c51..8ad2b9e3d 100644 --- a/src/core/Elsa.Core/Services/WorkflowQueue.cs +++ b/src/core/Elsa.Core/Services/WorkflowQueue.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Messages; @@ -19,7 +20,7 @@ namespace Elsa.Services _bus = bus; _commandSender = commandSender; } - + public async Task EnqueueWorkflowsAsync( Func predicate, object? input = default, @@ -29,7 +30,11 @@ namespace Elsa.Services where TTrigger : ITrigger { var results = await _workflowSelector.SelectWorkflowsAsync(predicate, cancellationToken).ToList(); + await EnqueueWorkflowsAsync(results, input, correlationId, contextId, cancellationToken); + } + public async Task EnqueueWorkflowsAsync(IEnumerable results, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default) + { foreach (var result in results) { if (result.WorkflowInstanceId != null)