From 25cf309cbad8a27218ec153c8238d20c1d2b5614 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 21 Jan 2021 19:30:13 +0100 Subject: [PATCH] Update RunWorkflowDefinitionConsumer to avoid creating multiple workflows with same correlation --- .../RunWorkflowDefinitionConsumer.cs | 27 +++++++++++++++++-- 1 file changed, 25 insertions(+), 2 deletions(-) diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs index 9d650a01d..3b6f60bda 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs @@ -1,8 +1,11 @@ using System.Threading.Tasks; using Elsa.Messages; using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Specifications; using Elsa.Services; using Elsa.Services.Models; +using Elsa.Triggers; using Microsoft.Extensions.Logging; using Rebus.Handlers; @@ -12,12 +15,16 @@ namespace Elsa.Consumers { private readonly IWorkflowRunner _workflowRunner; private readonly IWorkflowRegistry _workflowRegistry; + private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly IWorkflowSelector _workflowSelector; private readonly ILogger _logger; - public RunWorkflowDefinitionConsumer(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, ILogger logger) + public RunWorkflowDefinitionConsumer(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, IWorkflowSelector workflowSelector, ILogger logger) { _workflowRunner = workflowRunner; _workflowRegistry = workflowRegistry; + _workflowInstanceStore = workflowInstanceStore; + _workflowSelector = workflowSelector; _logger = logger; } @@ -29,6 +36,22 @@ namespace Elsa.Consumers if (!ValidatePreconditions(workflowDefinitionId, workflowBlueprint)) return; + + var correlationId = message.CorrelationId; + + if (!string.IsNullOrWhiteSpace(message.CorrelationId)) + { + var correlatedWorkflowInstanceCount = await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId)); + + if (correlatedWorkflowInstanceCount > 0) + { + // Do not create a new workflow instance. + _logger.LogWarning("There's already a workflow with correlation ID '{CorrelationId}'", correlationId); + return; + } + + _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); + } await _workflowRunner.RunWorkflowAsync(workflowBlueprint!, message.ActivityId, message.Input, message.CorrelationId, message.ContextId); } @@ -37,7 +60,7 @@ namespace Elsa.Consumers { if (workflowBlueprint == null) { - _logger.LogError("Could not run workflow with ID {WorkflowDefinitionId} because it does not exist.", workflowDefinitionId); + _logger.LogError("Could not run workflow with ID {WorkflowDefinitionId} because it does not exist", workflowDefinitionId); return false; }