diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index 67a4c0116..ba74eb7bf 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -1,12 +1,18 @@ using System; +using System.Linq; using System.Threading; using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Triggers; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Specifications; using Elsa.Services; +using Elsa.Triggers; using Microsoft.Azure.ServiceBus; using Microsoft.Azure.ServiceBus.Core; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; +using Open.Linq.AsyncExtensions; namespace Elsa.Activities.AzureServiceBus.Services { @@ -40,16 +46,38 @@ namespace Elsa.Activities.AzureServiceBus.Services { using var scope = _serviceProvider.CreateScope(); var workflowRunner = scope.ServiceProvider.GetRequiredService(); + var queueName = _messageReceiver.Path; + + async Task TriggerNewWorkflowAsync() + { + await workflowRunner!.TriggerWorkflowsAsync( + x => x.QueueName == queueName && x.CorrelationId == null, + message, + message.CorrelationId, + cancellationToken: cancellationToken); + } - Func predicate = string.IsNullOrWhiteSpace(message.CorrelationId) - ? x => x.QueueName == _messageReceiver.Path && x.CorrelationId == null - : x => x.QueueName == _messageReceiver.Path && x.CorrelationId == message.CorrelationId; - - await workflowRunner.TriggerWorkflowsAsync( - predicate, - message, - message.CorrelationId, - cancellationToken: cancellationToken); + if (string.IsNullOrWhiteSpace(message.CorrelationId)) + { + await TriggerNewWorkflowAsync(); + return; + } + + var workflowSelector = scope.ServiceProvider.GetRequiredService(); + var workflowInstanceStore = scope.ServiceProvider.GetRequiredService(); + var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification(message.CorrelationId), cancellationToken); + + if (correlatedWorkflowInstanceCount > 0) + { + // Trigger existing workflows (if blocked on this message). + var existingWorkflows = await workflowSelector.SelectWorkflowsAsync(x => x.QueueName == queueName && x.CorrelationId == message.CorrelationId, cancellationToken).ToList(); + await workflowRunner.TriggerWorkflowsAsync(existingWorkflows, message, message.CorrelationId, cancellationToken: cancellationToken); + } + else + { + // Trigger new workflow. + await TriggerNewWorkflowAsync(); + } } private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e) diff --git a/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs b/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs index 349530e0d..7b9b16dc5 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs @@ -40,7 +40,7 @@ namespace Elsa.Activities.Timers if (ExecuteAt <= now) { - _logger.LogDebug("Scheduled trigger time lies in the past ('{Delta}'). Skipping scheduling.", now - ExecuteAt); + _logger.LogDebug("Scheduled trigger time lies in the past ('{Delta}'). Skipping scheduling", now - ExecuteAt); return Done(); } diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs b/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs index 541f584a7..0d5c7f531 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Builders; @@ -18,6 +19,13 @@ namespace Elsa.Services CancellationToken cancellationToken = default) where TTrigger : ITrigger; + Task TriggerWorkflowsAsync( + IEnumerable results, + object? input = default, + string? correlationId = default, + string? contextId = default, + CancellationToken cancellationToken = default); + ValueTask RunWorkflowAsync( WorkflowInstance workflowInstance, string? activityId = default, diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index 75702731f..b5b99ef0e 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -68,7 +68,16 @@ namespace Elsa.Services where TTrigger : ITrigger { var results = await _workflowSelector.SelectWorkflowsAsync(predicate, cancellationToken).ToList(); - + await TriggerWorkflowsAsync(results, input, correlationId, contextId, cancellationToken); + } + + public async Task TriggerWorkflowsAsync( + IEnumerable results, + object? input = default, + string? correlationId = default, + string? contextId = default, + CancellationToken cancellationToken = default) + { foreach (var result in results) { if (result.WorkflowInstanceId != null)