Synchronize Azure Service Bus queue worker based on correlation ID to prevent multiple workflow instances with same correlation ID

This commit is contained in:
Sipke Schoorstra 2021-01-19 18:19:36 +01:00
parent f57b834d36
commit 2a6d8f9bfe

View file

@ -1,8 +1,11 @@
using System;
using System.Diagnostics;
using System.Linq;
using System.Runtime.InteropServices;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Triggers;
using Elsa.DistributedLock;
using Elsa.Models;
using Elsa.Persistence;
using Elsa.Persistence.Specifications;
@ -20,12 +23,14 @@ namespace Elsa.Activities.AzureServiceBus.Services
{
private readonly IMessageReceiver _messageReceiver;
private readonly IServiceProvider _serviceProvider;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly ILogger _logger;
public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, ILogger<QueueWorker> logger)
public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, IDistributedLockProvider distributedLockProvider, ILogger<QueueWorker> logger)
{
_messageReceiver = messageReceiver;
_serviceProvider = serviceProvider;
_distributedLockProvider = distributedLockProvider;
_logger = logger;
_messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler)
@ -47,36 +52,60 @@ namespace Elsa.Activities.AzureServiceBus.Services
using var scope = _serviceProvider.CreateScope();
var workflowRunner = scope.ServiceProvider.GetRequiredService<IWorkflowRunner>();
var queueName = _messageReceiver.Path;
var correlationId = message.CorrelationId;
async Task TriggerNewWorkflowAsync()
{
await workflowRunner!.TriggerWorkflowsAsync<MessageReceivedTrigger>(
x => x.QueueName == queueName && x.CorrelationId == null,
message,
message.CorrelationId,
correlationId,
cancellationToken: cancellationToken);
}
if (string.IsNullOrWhiteSpace(message.CorrelationId))
if (string.IsNullOrWhiteSpace(correlationId))
{
await TriggerNewWorkflowAsync();
return;
}
var workflowSelector = scope.ServiceProvider.GetRequiredService<IWorkflowSelector>();
var workflowInstanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(message.CorrelationId), cancellationToken);
var lockKey = $"azure-service-bus:{queueName}:correlation-{correlationId}";
var stopwatch = new Stopwatch();
if (correlatedWorkflowInstanceCount > 0)
_logger.LogDebug("Acquiring lock {LockKey}", lockKey);
stopwatch.Start();
if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken))
{
// Trigger existing workflows (if blocked on this message).
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);
_logger.LogDebug("Lock {LockKey} already taken", lockKey);
return;
}
else
try
{
// Trigger new workflow.
await TriggerNewWorkflowAsync();
var workflowSelector = scope.ServiceProvider.GetRequiredService<IWorkflowSelector>();
var workflowInstanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(message.CorrelationId), cancellationToken);
if (correlatedWorkflowInstanceCount > 0)
{
// 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);
}
else
{
// Trigger new workflow.
_logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId);
await TriggerNewWorkflowAsync();
}
}
finally
{
await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken);
stopwatch.Stop();
_logger.LogDebug("Lock held for {ElapseTime}", stopwatch.Elapsed);
}
}
@ -85,7 +114,7 @@ namespace Elsa.Activities.AzureServiceBus.Services
switch (e.Exception)
{
case MessageLockLostException:
_logger.LogDebug( e.Exception,"Message lock lost");
_logger.LogDebug(e.Exception, "Message lock lost");
break;
case ServiceBusCommunicationException:
_logger.LogDebug(e.Exception, "Lost service bus communication");