From 356e0ba0bac72f27df83276cd8b62baf06ca19ec Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 30 Dec 2020 12:02:13 +0100 Subject: [PATCH] Update Azure Service Bus queue worker to schedule workflow execution in the background This to prevent scenarios where workflows take a long time to execute and the message gets re-delivered. This change ACKS the message right after the workflow task has been queued internally. --- .../Services/QueueWorker.cs | 33 ++++++++++--------- 1 file changed, 18 insertions(+), 15 deletions(-) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index 86d71018a..1ccdfd811 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -14,14 +14,16 @@ namespace Elsa.Activities.AzureServiceBus.Services { private readonly IMessageReceiver _messageReceiver; private readonly IServiceProvider _serviceProvider; + private readonly IBackgroundWorker _backgroundWorker; private readonly ILogger _logger; - public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, ILogger logger) + public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, IBackgroundWorker backgroundWorker, ILogger logger) { _messageReceiver = messageReceiver; _serviceProvider = serviceProvider; + _backgroundWorker = backgroundWorker; _logger = logger; - + _messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler) { AutoComplete = false, @@ -31,17 +33,21 @@ namespace Elsa.Activities.AzureServiceBus.Services private async Task OnMessageReceived(Message message, CancellationToken cancellationToken) { - using (var scope = _serviceProvider.CreateScope()) - { - var workflowRunner = scope.ServiceProvider.GetRequiredService(); - - await workflowRunner.TriggerWorkflowsAsync(x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), message, - message.CorrelationId, cancellationToken: cancellationToken); - } - + await _backgroundWorker.ScheduleTask(async () => await TriggerWorkflowsAsync(message, cancellationToken), cancellationToken); await _messageReceiver.CompleteAsync(message.SystemProperties.LockToken); } + private async Task TriggerWorkflowsAsync(Message message, CancellationToken cancellationToken) + { + using var scope = _serviceProvider.CreateScope(); + var workflowRunner = scope.ServiceProvider.GetRequiredService(); + + await workflowRunner.TriggerWorkflowsAsync( + x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), + message, message.CorrelationId, + cancellationToken: cancellationToken); + } + private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e) { var context = e.ExceptionReceivedContext; @@ -60,13 +66,10 @@ namespace Elsa.Activities.AzureServiceBus.Services _logger.LogError("- Executing Action: {Action}", context.Action); break; } - + return Task.CompletedTask; } - public async ValueTask DisposeAsync() - { - await _messageReceiver.CloseAsync(); - } + public async ValueTask DisposeAsync() => await _messageReceiver.CloseAsync(); } } \ No newline at end of file