From bd87f5b2df4252823b55b9a9a5bc98c40714d2d4 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 29 Mar 2021 22:11:37 +0200 Subject: [PATCH] Incremental work on default dispatchers in preparation for actor model implementation --- .../Services/QueueWorker.cs | 75 ++----------- .../Services/TopicWorker.cs | 1 + .../Webhooks/Consumers/TriggerWorkflows.cs | 1 + .../Jobs/RunQuartzWorkflowJob.cs | 80 ++----------- .../ICorrelatingWorkflowDispatcher.cs | 14 +++ .../Dispatch/IWorkflowDefinitionDispatcher.cs | 13 +++ .../Dispatch/IWorkflowInstanceDispatcher.cs | 13 +++ src/core/Elsa.Abstractions/Dispatch/Models.cs | 9 ++ .../IDistributedLockProvider.cs | 2 +- .../Extensions/WorkflowQueueExtensions.cs | 40 +++---- .../RunWorkflowDefinitionConsumer.cs | 3 + .../Consumers/RunWorkflowInstanceConsumer.cs | 1 + ...xecuteCorrelatedWorkflowRequestConsumer.cs | 106 ++++++++++++++++++ ...xecuteWorkflowDefinitionRequestConsumer.cs | 62 ++++++++++ .../ExecuteWorkflowRequestConsumer.cs | 100 +++++++++++++++++ .../Dispatch/QueuingWorkflowDispatcher.cs | 18 +++ .../DistributedLock/DefaultLockProvider.cs | 1 + src/core/Elsa.Core/ElsaOptions.cs | 1 + .../StartupTasks/ContinueRunningWorkflows.cs | 1 + .../AzureBlobLockProvider.cs | 1 + .../RedisLockProvider.cs | 1 + .../SqlLockProvider.cs | 1 + 22 files changed, 390 insertions(+), 154 deletions(-) create mode 100644 src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs create mode 100644 src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs create mode 100644 src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs create mode 100644 src/core/Elsa.Abstractions/Dispatch/Models.cs rename src/core/Elsa.Abstractions/{DistributedLock => DistributedLocking}/IDistributedLockProvider.cs (94%) create mode 100644 src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs create mode 100644 src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowDefinitionRequestConsumer.cs create mode 100644 src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowRequestConsumer.cs create mode 100644 src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index d17067792..6aaf3759b 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -6,11 +6,14 @@ using Elsa.Activities.AzureServiceBus.Bookmarks; using Elsa.Activities.AzureServiceBus.Models; using Elsa.Activities.AzureServiceBus.Options; using Elsa.Bookmarks; +using Elsa.Dispatch; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; using Elsa.Services; +using Elsa.Services.Models; using Elsa.Triggers; using Microsoft.Azure.ServiceBus; using Microsoft.Azure.ServiceBus.Core; @@ -27,20 +30,18 @@ namespace Elsa.Activities.AzureServiceBus.Services private const string TenantId = default; private readonly IMessageReceiver _messageReceiver; - private readonly IServiceScopeFactory _serviceScopeFactory; - private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ICorrelatingWorkflowDispatcher _workflowDispatcher; private readonly ILogger _logger; public QueueWorker( IMessageReceiver messageReceiver, + ICorrelatingWorkflowDispatcher workflowDispatcher, IServiceScopeFactory serviceScopeFactory, - IDistributedLockProvider distributedLockProvider, IOptions options, ILogger logger) { _messageReceiver = messageReceiver; - _serviceScopeFactory = serviceScopeFactory; - _distributedLockProvider = distributedLockProvider; + _workflowDispatcher = workflowDispatcher; _logger = logger; _messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler) @@ -61,11 +62,8 @@ namespace Elsa.Activities.AzureServiceBus.Services private async Task TriggerWorkflowsAsync(Message message, CancellationToken cancellationToken) { - using var scope = _serviceScopeFactory.CreateScope(); - var workflowQueue = scope.ServiceProvider.GetRequiredService(); var queueName = _messageReceiver.Path; var correlationId = message.CorrelationId; - var triggerFinder = scope.ServiceProvider.GetRequiredService(); var model = new MessageModel { @@ -85,63 +83,10 @@ namespace Elsa.Activities.AzureServiceBus.Services ScheduledEnqueueTimeUtc = message.ScheduledEnqueueTimeUtc }; - async Task TriggerNewWorkflowAsync() - { - var bookmark = new QueueMessageReceivedBookmark(queueName); - var triggers = await triggerFinder.FindTriggersAsync(bookmark, TenantId, cancellationToken); - - foreach (var trigger in triggers) - { - var workflowBlueprint = trigger.WorkflowBlueprint; - await workflowQueue.EnqueueWorkflowDefinition(workflowBlueprint.Id, workflowBlueprint.TenantId, trigger.ActivityId, model, correlationId, null, cancellationToken); - } - } - - if (string.IsNullOrWhiteSpace(correlationId)) - { - await TriggerNewWorkflowAsync(); - return; - } - - var lockKey = $"azure-service-bus:{queueName}:correlation-{correlationId}"; - var stopwatch = new Stopwatch(); - - _logger.LogDebug("Acquiring lock {LockKey}", lockKey); - stopwatch.Start(); - - if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken)) - { - _logger.LogDebug("Lock {LockKey} already taken", lockKey); - return; - } - - try - { - var bookmarkFinder = scope.ServiceProvider.GetRequiredService(); - var workflowInstanceStore = scope.ServiceProvider.GetRequiredService(); - var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification(model.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 bookmark = new QueueMessageReceivedBookmark(queueName, correlationId); - var existingWorkflows = await bookmarkFinder.FindBookmarksAsync(bookmark, TenantId, cancellationToken).ToList(); - await workflowQueue.EnqueueWorkflowsAsync(existingWorkflows, model, model.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); - } + var bookmark = new QueueMessageReceivedBookmark(queueName, correlationId); + var trigger = new QueueMessageReceivedBookmark(queueName); + var activityType = nameof(AzureServiceBusQueueMessageReceived); + await _workflowDispatcher.DispatchAsync(new ExecuteCorrelatedWorkflowRequest(correlationId, bookmark, trigger, activityType, model, TenantId: TenantId), cancellationToken); } private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs index f3620fd3e..58b433474 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs @@ -7,6 +7,7 @@ using Elsa.Activities.AzureServiceBus.Models; using Elsa.Activities.AzureServiceBus.Options; using Elsa.Bookmarks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; diff --git a/src/activities/Elsa.Activities.Telnyx/Webhooks/Consumers/TriggerWorkflows.cs b/src/activities/Elsa.Activities.Telnyx/Webhooks/Consumers/TriggerWorkflows.cs index 551936aff..0ae5a8739 100644 --- a/src/activities/Elsa.Activities.Telnyx/Webhooks/Consumers/TriggerWorkflows.cs +++ b/src/activities/Elsa.Activities.Telnyx/Webhooks/Consumers/TriggerWorkflows.cs @@ -7,6 +7,7 @@ using Elsa.Activities.Telnyx.Webhooks.Payloads.Abstract; using Elsa.Activities.Telnyx.Webhooks.Payloads.Call; using Elsa.Activities.Telnyx.Webhooks.Services; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Jobs/RunQuartzWorkflowJob.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Jobs/RunQuartzWorkflowJob.cs index c7aabb3b5..ac2ad315f 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Jobs/RunQuartzWorkflowJob.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Jobs/RunQuartzWorkflowJob.cs @@ -1,37 +1,20 @@ -using System.Diagnostics; -using System.Threading.Tasks; -using Elsa.DistributedLock; -using Elsa.Models; -using Elsa.Persistence; -using Elsa.Persistence.Specifications; -using Elsa.Persistence.Specifications.WorkflowInstances; -using Elsa.Services; -using Microsoft.Extensions.Logging; +using System.Threading.Tasks; +using Elsa.Dispatch; using Quartz; namespace Elsa.Activities.Temporal.Quartz.Jobs { public class RunQuartzWorkflowJob : IJob { - private readonly IWorkflowRegistry _workflowRegistry; - private readonly IWorkflowInstanceStore _workflowInstanceStore; - private readonly IWorkflowQueue _workflowQueue; - private readonly IDistributedLockProvider _distributedLockProvider; - private readonly ILogger _logger; - private readonly Stopwatch _stopwatch = new(); + private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher; + private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher; public RunQuartzWorkflowJob( - IWorkflowRegistry workflowRegistry, - IWorkflowInstanceStore workflowInstanceStore, - IWorkflowQueue workflowQueue, - IDistributedLockProvider distributedLockProvider, - ILogger logger) + IWorkflowDefinitionDispatcher workflowDefinitionDispatcher, + IWorkflowInstanceDispatcher workflowInstanceDispatcher) { - _workflowRegistry = workflowRegistry; - _workflowInstanceStore = workflowInstanceStore; - _workflowQueue = workflowQueue; - _distributedLockProvider = distributedLockProvider; - _logger = logger; + _workflowDefinitionDispatcher = workflowDefinitionDispatcher; + _workflowInstanceDispatcher = workflowInstanceDispatcher; } public async Task Execute(IJobExecutionContext context) @@ -42,50 +25,11 @@ namespace Elsa.Activities.Temporal.Quartz.Jobs var tenantId = dataMap.GetString("TenantId"); var workflowDefinitionId = dataMap.GetString("WorkflowDefinitionId")!; var activityId = dataMap.GetString("ActivityId")!; - var lockKey = (workflowInstanceId, workflowDefinitionId, activityId).GetHashCode().ToString(); - _logger.LogDebug("Acquiring lock on {WorkflowInstanceId} / {WorkflowDefinitionId} / {ActivityId}", workflowInstanceId, workflowDefinitionId, activityId); - - if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken)) - { - _logger.LogDebug("Failed to acquire lock on {WorkflowInstanceId} / {WorkflowDefinitionId} / {ActivityId}", workflowInstanceId, workflowDefinitionId, activityId); - return; - } - - _stopwatch.Restart(); - - try - { - if (workflowInstanceId == null) - { - var workflowBlueprint = (await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.Published, cancellationToken)); - - if (workflowBlueprint == null) - { - _logger.LogWarning("No workflow definition {WorkflowDefinitionId} found. Make sure the scheduled workflow definition is published and enabled", workflowDefinitionId); - return; - } - - if (!workflowBlueprint.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false) - await _workflowQueue.EnqueueWorkflowDefinition(workflowDefinitionId, tenantId, activityId, null, null, null, cancellationToken); - } - else - { - await _workflowQueue.EnqueueWorkflowInstance(workflowInstanceId, activityId, null, cancellationToken); - } - } - finally - { - _stopwatch.Stop(); - await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken); - _logger.LogDebug("Held lock on {WorkflowInstanceId} / {WorkflowDefinitionId} / {ActivityId} for {LockTime}", workflowInstanceId, workflowDefinitionId, activityId, _stopwatch.Elapsed); - } - } - - private async Task GetWorkflowIsAlreadyExecutingAsync(string? tenantId, string workflowDefinitionId) - { - var specification = new TenantSpecification(tenantId).WithWorkflowDefinition(workflowDefinitionId).And(new WorkflowIsAlreadyExecutingSpecification()); - return await _workflowInstanceStore.FindAsync(specification) != null; + if (workflowInstanceId == null) + await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId, TenantId: tenantId), cancellationToken); + else + await _workflowInstanceDispatcher.DispatchAsync(new ExecuteWorkflowInstanceRequest(workflowInstanceId, activityId), cancellationToken); } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs b/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs new file mode 100644 index 000000000..252f4ae1c --- /dev/null +++ b/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs @@ -0,0 +1,14 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Dispatch +{ + /// + /// The correlating dispatcher is responsible for finding workflows correlated by the specified correlation ID. + /// If no correlated workflows are found, a new one is started. + /// + public interface ICorrelatingWorkflowDispatcher + { + Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken? cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs b/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs new file mode 100644 index 000000000..4a95cac58 --- /dev/null +++ b/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs @@ -0,0 +1,13 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Dispatch +{ + /// + /// Dispatches requests for executing workflow definitions. + /// + public interface IWorkflowDefinitionDispatcher + { + Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken? cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs b/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs new file mode 100644 index 000000000..806b5aee1 --- /dev/null +++ b/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs @@ -0,0 +1,13 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Dispatch +{ + /// + /// Dispatches requests for executing workflow instances. + /// + public interface IWorkflowInstanceDispatcher + { + Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken? cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/Models.cs b/src/core/Elsa.Abstractions/Dispatch/Models.cs new file mode 100644 index 000000000..4a756b68e --- /dev/null +++ b/src/core/Elsa.Abstractions/Dispatch/Models.cs @@ -0,0 +1,9 @@ +using Elsa.Bookmarks; + +namespace Elsa.Dispatch +{ + public record ExecuteCorrelatedWorkflowRequest(string CorrelationId, IBookmark Bookmark, IBookmark Trigger, string ActivityType, object? Input = default, string? ContextId = default, string? TenantId = default); + public record ExecuteWorkflowDefinitionRequest(string WorkflowDefinitionId, string? ActivityId = default, object? Input = default, string? CorrelationId = default, string? ContextId = default, string? TenantId = default); + public record ExecuteWorkflowInstanceRequest(string WorkflowInstanceId, string ActivityId, object? Input = default); + +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/DistributedLock/IDistributedLockProvider.cs b/src/core/Elsa.Abstractions/DistributedLocking/IDistributedLockProvider.cs similarity index 94% rename from src/core/Elsa.Abstractions/DistributedLock/IDistributedLockProvider.cs rename to src/core/Elsa.Abstractions/DistributedLocking/IDistributedLockProvider.cs index 37920dde0..07e0fed26 100644 --- a/src/core/Elsa.Abstractions/DistributedLock/IDistributedLockProvider.cs +++ b/src/core/Elsa.Abstractions/DistributedLocking/IDistributedLockProvider.cs @@ -1,7 +1,7 @@ using System.Threading; using System.Threading.Tasks; -namespace Elsa.DistributedLock +namespace Elsa.DistributedLocking { /// /// Provides functionality to acquire and release locks which are distributed across all instances in a web farm. diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowQueueExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowQueueExtensions.cs index 0a0ed54e9..b5dd07270 100644 --- a/src/core/Elsa.Abstractions/Extensions/WorkflowQueueExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowQueueExtensions.cs @@ -1,20 +1,20 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Bookmarks; -using Elsa.Services; - -namespace Elsa -{ - public static class WorkflowQueueExtensions - { - public static Task EnqueueWorkflowsAsync( - this IWorkflowQueue workflowQueue, - IBookmark bookmark, - string? tenantId, - object? input = default, - string? correlationId = default, - string? contextId = default, - CancellationToken cancellationToken = default) where T : IActivity => - workflowQueue.EnqueueWorkflowsAsync(typeof(T).Name, bookmark, tenantId, input, correlationId, contextId, cancellationToken); - } -} \ No newline at end of file +// using System.Threading; +// using System.Threading.Tasks; +// using Elsa.Bookmarks; +// using Elsa.Services; +// +// namespace Elsa +// { +// public static class WorkflowQueueExtensions +// { +// public static Task EnqueueWorkflowsAsync( +// this IWorkflowQueue workflowQueue, +// IBookmark bookmark, +// string? tenantId, +// object? input = default, +// string? correlationId = default, +// string? contextId = default, +// CancellationToken cancellationToken = default) where T : IActivity => +// workflowQueue.EnqueueWorkflowsAsync(typeof(T).Name, bookmark, tenantId, input, correlationId, contextId, cancellationToken); +// } +// } \ No newline at end of file diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs index c9876159f..17ed39416 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowDefinitionConsumer.cs @@ -1,5 +1,7 @@ +using System; using System.Threading.Tasks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Messages; using Elsa.Models; using Elsa.Persistence; @@ -12,6 +14,7 @@ using Rebus.Handlers; namespace Elsa.Consumers { + [Obsolete] public class RunWorkflowDefinitionConsumer : IHandleMessages { private readonly IWorkflowRunner _workflowRunner; diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs index fee014c37..d2538ebd4 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs @@ -2,6 +2,7 @@ using System.Diagnostics; using System.Linq; using System.Threading.Tasks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Messages; using Elsa.Models; using Elsa.Persistence; diff --git a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs new file mode 100644 index 000000000..9a5d71182 --- /dev/null +++ b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs @@ -0,0 +1,106 @@ +using System.Collections.Generic; +using System.Diagnostics; +using System.Threading.Tasks; +using Elsa.Bookmarks; +using Elsa.DistributedLocking; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Specifications; +using Elsa.Services; +using Elsa.Triggers; +using Microsoft.Extensions.Logging; +using Open.Linq.AsyncExtensions; +using Rebus.Handlers; + +namespace Elsa.Dispatch.Consumers +{ + public class ExecuteCorrelatedWorkflowRequestConsumer : IHandleMessages + { + private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly IBookmarkFinder _bookmarkFinder; + private readonly ITriggerFinder _triggerFinder; + private readonly ICommandSender _commandSender; + private readonly ILogger _logger; + private readonly Stopwatch _stopwatch = new(); + + public ExecuteCorrelatedWorkflowRequestConsumer( + IWorkflowInstanceStore workflowInstanceStore, + IDistributedLockProvider distributedLockProvider, + IBookmarkFinder bookmarkFinder, + ITriggerFinder triggerFinder, + ICommandSender commandSender, + ILogger logger) + { + _workflowInstanceStore = workflowInstanceStore; + _distributedLockProvider = distributedLockProvider; + _bookmarkFinder = bookmarkFinder; + _triggerFinder = triggerFinder; + _commandSender = commandSender; + _logger = logger; + } + + public async Task Handle(ExecuteCorrelatedWorkflowRequest message) + { + var correlationId = message.CorrelationId; + var lockKey = $"correlated-workflow-request:correlation-{correlationId}"; + + _logger.LogDebug("Acquiring lock {LockKey}", lockKey); + _stopwatch.Restart(); + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey)) + { + _logger.LogDebug("Lock {LockKey} already taken", lockKey); + await _commandSender.SendAsync(message); + return; + } + + try + { + var correlatedWorkflowInstanceCount = await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId)); + + if (correlatedWorkflowInstanceCount > 0) + { + _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId); + var existingWorkflows = await _bookmarkFinder.FindBookmarksAsync(message.ActivityType, message.Bookmark, message.TenantId).ToList(); + await EnqueueWorkflowsAsync(existingWorkflows, message.Input); + } + else + { + // Trigger new workflow. + _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); + await TriggerNewWorkflowAsync(message); + } + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey); + _stopwatch.Stop(); + _logger.LogDebug("Lock held for {ElapseTime}", _stopwatch.Elapsed); + } + } + + async Task TriggerNewWorkflowAsync(ExecuteCorrelatedWorkflowRequest message) + { + var filter = message.Trigger; + var triggers = await _triggerFinder.FindTriggersAsync(message.ActivityType, filter, message.TenantId); + + foreach (var trigger in triggers) + { + var workflowBlueprint = trigger.WorkflowBlueprint; + await EnqueueWorkflowDefinition(workflowBlueprint.Id, workflowBlueprint.TenantId, trigger.ActivityId, message.Input, message.CorrelationId, message.ContextId); + } + } + + private async Task EnqueueWorkflowsAsync(IEnumerable results, object? input) + { + foreach (var result in results) + await EnqueueWorkflowInstance(result.WorkflowInstanceId, result.ActivityId, input); + } + + public async Task EnqueueWorkflowInstance(string workflowInstanceId, string activityId, object? input) => await _commandSender.SendAsync(new ExecuteWorkflowInstanceRequest(workflowInstanceId, activityId, input)); + + public async Task EnqueueWorkflowDefinition(string workflowDefinitionId, string? tenantId, string activityId, object? input, string? correlationId, string? contextId) => + await _commandSender.SendAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId, input, correlationId, contextId, tenantId)); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowDefinitionRequestConsumer.cs b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowDefinitionRequestConsumer.cs new file mode 100644 index 000000000..36bda865e --- /dev/null +++ b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowDefinitionRequestConsumer.cs @@ -0,0 +1,62 @@ +using System.Threading.Tasks; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Specifications; +using Elsa.Persistence.Specifications.WorkflowInstances; +using Elsa.Services; +using Elsa.Services.Models; +using Microsoft.Extensions.Logging; +using Rebus.Handlers; + +namespace Elsa.Dispatch.Consumers +{ + public class ExecuteWorkflowDefinitionRequestConsumer : IHandleMessages + { + private readonly IWorkflowRunner _workflowRunner; + private readonly IWorkflowRegistry _workflowRegistry; + private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly ILogger _logger; + + public ExecuteWorkflowDefinitionRequestConsumer( + IWorkflowRunner workflowRunner, + IWorkflowRegistry workflowRegistry, + IWorkflowInstanceStore workflowInstanceStore, + ILogger logger) + { + _workflowRunner = workflowRunner; + _workflowRegistry = workflowRegistry; + _workflowInstanceStore = workflowInstanceStore; + _logger = logger; + } + + public async Task Handle(ExecuteWorkflowDefinitionRequest message) + { + var workflowDefinitionId = message.WorkflowDefinitionId; + var tenantId = message.TenantId; + var workflowBlueprint = await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.Published); + + if (!ValidatePreconditions(workflowDefinitionId, workflowBlueprint)) + return; + + if (!workflowBlueprint!.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false) + await _workflowRunner.RunWorkflowAsync(workflowBlueprint, message.ActivityId, message.Input, message.CorrelationId, message.ContextId); + } + + private bool ValidatePreconditions(string? workflowDefinitionId, IWorkflowBlueprint? workflowBlueprint) + { + if (workflowBlueprint == null) + { + _logger.LogWarning("No workflow definition {WorkflowDefinitionId} found. Make sure the scheduled workflow definition is published and enabled", workflowDefinitionId); + return false; + } + + return true; + } + + private async Task GetWorkflowIsAlreadyExecutingAsync(string? tenantId, string workflowDefinitionId) + { + var specification = new TenantSpecification(tenantId).WithWorkflowDefinition(workflowDefinitionId).And(new WorkflowIsAlreadyExecutingSpecification()); + return await _workflowInstanceStore.FindAsync(specification) != null; + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowRequestConsumer.cs b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowRequestConsumer.cs new file mode 100644 index 000000000..7dfe83570 --- /dev/null +++ b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowRequestConsumer.cs @@ -0,0 +1,100 @@ +using System.Diagnostics; +using System.Linq; +using System.Threading.Tasks; +using Elsa.DistributedLocking; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Services; +using Microsoft.Extensions.Logging; +using Rebus.Handlers; + +namespace Elsa.Dispatch.Consumers +{ + public class ExecuteWorkflowRequestConsumer : IHandleMessages + { + private readonly IWorkflowRunner _workflowRunner; + private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ICommandSender _commandSender; + private readonly ILogger _logger; + private readonly Stopwatch _stopwatch = new(); + + public ExecuteWorkflowRequestConsumer( + IWorkflowRunner workflowRunner, + IWorkflowInstanceStore workflowInstanceStore, + IDistributedLockProvider distributedLockProvider, + ICommandSender commandSender, + ILogger logger) + { + _workflowRunner = workflowRunner; + _workflowInstanceStore = workflowInstanceStore; + _distributedLockProvider = distributedLockProvider; + _commandSender = commandSender; + _logger = logger; + } + + public async Task Handle(ExecuteWorkflowInstanceRequest message) + { + var workflowInstanceId = message.WorkflowInstanceId; + var lockKey = workflowInstanceId; + + _logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}", workflowInstanceId); + _stopwatch.Restart(); + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey)) + { + _logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Re-queueing message", workflowInstanceId); + await _commandSender.SendAsync(message); + return; + } + + try + { + var workflowInstance = await _workflowInstanceStore.FindByIdAsync(message.WorkflowInstanceId); + + if (!ValidatePreconditions(workflowInstanceId, workflowInstance, message.ActivityId)) + return; + + await _workflowRunner.RunWorkflowAsync( + workflowInstance!, + message.ActivityId, + message.Input); + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey); + _stopwatch.Stop(); + _logger.LogDebug("Held lock on workflow instance {WorkflowInstanceId} for {ElapsedTime}", workflowInstanceId, _stopwatch.Elapsed); + } + } + + private bool ValidatePreconditions(string? workflowInstanceId, WorkflowInstance? workflowInstance, string? activityId) + { + if (workflowInstance == null) + { + _logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist", workflowInstanceId); + return false; + } + + if (workflowInstance.WorkflowStatus != WorkflowStatus.Suspended && workflowInstance.WorkflowStatus != WorkflowStatus.Running) + { + _logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it has a status other than Suspended or Running. Its actual status is {WorkflowStatus}", workflowInstanceId, workflowInstance.WorkflowStatus); + return false; + } + + if (activityId != null) + { + var activityIsBlocking = workflowInstance.BlockingActivities.Any(x => x.ActivityId == activityId); + var activityIsScheduled = workflowInstance.ScheduledActivities.Any(x => x.ActivityId == activityId) || workflowInstance.CurrentActivity?.ActivityId == activityId; + + if (!activityIsBlocking && !activityIsScheduled) + { + _logger.LogWarning("Did not run workflow {WorkflowInstanceId} for activity {ActivityId} because the workflow is not blocked on that activity nor is that activity scheduled for execution", workflowInstanceId, activityId); + return false; + } + } + + return true; + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs b/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs new file mode 100644 index 000000000..4ac0d4fc3 --- /dev/null +++ b/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs @@ -0,0 +1,18 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services; + +namespace Elsa.Dispatch +{ + /// + /// The default strategy that process workflow execution requests by sending them to a queue. + /// + public class QueuingWorkflowDispatcher : IWorkflowDefinitionDispatcher, IWorkflowInstanceDispatcher, ICorrelatingWorkflowDispatcher + { + private readonly ICommandSender _commandSender; + public QueuingWorkflowDispatcher(ICommandSender commandSender) => _commandSender = commandSender; + public async Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request); + public async Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request); + public async Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/DistributedLock/DefaultLockProvider.cs b/src/core/Elsa.Core/DistributedLock/DefaultLockProvider.cs index 66d9a9036..6e7212744 100644 --- a/src/core/Elsa.Core/DistributedLock/DefaultLockProvider.cs +++ b/src/core/Elsa.Core/DistributedLock/DefaultLockProvider.cs @@ -1,6 +1,7 @@ using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; +using Elsa.DistributedLocking; namespace Elsa.DistributedLock { diff --git a/src/core/Elsa.Core/ElsaOptions.cs b/src/core/Elsa.Core/ElsaOptions.cs index 0dd17f870..42a0b853c 100644 --- a/src/core/Elsa.Core/ElsaOptions.cs +++ b/src/core/Elsa.Core/ElsaOptions.cs @@ -6,6 +6,7 @@ using AutoMapper; using Elsa.Builders; using Elsa.Caching; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.InMemory; diff --git a/src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflows.cs b/src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflows.cs index 188a08567..b31439f37 100644 --- a/src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflows.cs +++ b/src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflows.cs @@ -2,6 +2,7 @@ using System.Linq; using System.Threading; using System.Threading.Tasks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications.WorkflowInstances; diff --git a/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs b/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs index 86469f8cf..c56b35787 100644 --- a/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs +++ b/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs @@ -5,6 +5,7 @@ using System.Linq; using System.Threading; using System.Threading.Tasks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Microsoft.Azure.Storage; using Microsoft.Azure.Storage.Blob; using Microsoft.Azure.Storage.RetryPolicies; diff --git a/src/locking/Elsa.DistributedLocking.Redis/RedisLockProvider.cs b/src/locking/Elsa.DistributedLocking.Redis/RedisLockProvider.cs index b2385af61..c93d210f4 100644 --- a/src/locking/Elsa.DistributedLocking.Redis/RedisLockProvider.cs +++ b/src/locking/Elsa.DistributedLocking.Redis/RedisLockProvider.cs @@ -3,6 +3,7 @@ using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Microsoft.Extensions.Logging; using RedLockNet; diff --git a/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs b/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs index fe41221c7..59af52bcf 100644 --- a/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs +++ b/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs @@ -5,6 +5,7 @@ using System.Data.SqlClient; using System.Threading; using System.Threading.Tasks; using Elsa.DistributedLock; +using Elsa.DistributedLocking; using Microsoft.Extensions.Logging; namespace Elsa