Ensure workflow instance locking when receiving Azure Service Bus messages

This commit is contained in:
Sipke Schoorstra 2021-01-21 12:11:25 +01:00
parent 2a64cc1473
commit e39d5cdb68
3 changed files with 21 additions and 7 deletions

View file

@ -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<IWorkflowRunner>();
var workflowRunner = scope.ServiceProvider.GetRequiredService<IWorkflowQueue>();
var queueName = _messageReceiver.Path;
var correlationId = message.CorrelationId;
async Task TriggerNewWorkflowAsync()
{
await workflowRunner!.TriggerWorkflowsAsync<MessageReceivedTrigger>(
async Task TriggerNewWorkflowAsync() =>
await workflowRunner.EnqueueWorkflowsAsync<MessageReceivedTrigger>(
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<MessageReceivedTrigger>(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
{

View file

@ -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;
/// <summary>
/// Enqueues the specified workflows for execution.
/// </summary>
Task EnqueueWorkflowsAsync(
IEnumerable<WorkflowSelectorResult> results,
object? input = default,
string? correlationId = default,
string? contextId = default,
CancellationToken cancellationToken = default);
/// <summary>
/// Enqueues the specified workflow instance and activity for execution.
/// </summary>

View file

@ -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<TTrigger>(
Func<TTrigger, bool> 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<WorkflowSelectorResult> results, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default)
{
foreach (var result in results)
{
if (result.WorkflowInstanceId != null)