From d2d7f23d7368defa743e1783ce84e2cffc656178 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 20 Feb 2023 12:05:57 +0100 Subject: [PATCH] Fix deadlock + multiple trigger support (#3720) * Lock on specific hash for new workflows and split between trigger new /& existing workflows * Update flowchart to take into account the triggered activity to start * Use dispatcher for resuming parent workflows * Update logging middleware (and remove it by default) --- .../Handlers/JobExecutedHandler.cs | 9 ++- .../DispatchWorkflowRequestConsumer.cs | 2 +- .../ProtoActorWorkflowRuntime.cs | 78 +++++++----------- .../TriggerWebhookDrivenActivities.cs | 2 +- .../Flowchart/Activities/Flowchart.cs | 12 ++- .../Activities/LoggingMiddleware.cs | 17 +++- .../Models/WorkflowExecutionContext.cs | 2 +- .../DispatchWorkflowRequestHandler.cs | 2 +- .../ResumeDispatchWorkflowActivityHandler.cs | 8 +- .../Implementations/DefaultWorkflowRuntime.cs | 81 ++++++++++--------- .../Services/IWorkflowRuntime.cs | 12 ++- 11 files changed, 126 insertions(+), 99 deletions(-) diff --git a/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs b/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs index d364e3cbd..459132897 100644 --- a/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs +++ b/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs @@ -7,15 +7,22 @@ using Elsa.Workflows.Runtime.Services; namespace Elsa.Jobs.Activities.Handlers; +/// +/// A handler that resumes workflows waiting for a given job to complete. +/// public class JobExecutedHandler : INotificationHandler { private readonly IWorkflowRuntime _workflowRuntime; + /// + /// Constructor. + /// public JobExecutedHandler(IWorkflowRuntime workflowRuntime) { _workflowRuntime = workflowRuntime; } + /// public async Task HandleAsync(JobExecuted notification, CancellationToken cancellationToken) { var job = notification.Job; @@ -33,7 +40,7 @@ public class JobExecutedHandler : INotificationHandler var jobType = notification.Job.GetType(); var payload = new EnqueuedJobPayload(job.Id); var jobTypeName = JobTypeNameHelper.GenerateTypeName(jobType); - await _workflowRuntime.ResumeWorkflowsAsync(jobTypeName, payload, new ResumeWorkflowRuntimeOptions(), cancellationToken); + await _workflowRuntime.ResumeWorkflowsAsync(jobTypeName, payload, new TriggerWorkflowsRuntimeOptions(), cancellationToken); } } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs index 15cc0adcb..a9112ac33 100644 --- a/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs +++ b/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs @@ -52,7 +52,7 @@ public class DispatchWorkflowRequestConsumer : public async Task Consume(ConsumeContext context) { var message = context.Message; - var options = new ResumeWorkflowRuntimeOptions(CorrelationId: message.CorrelationId, Input: message.Input); + var options = new TriggerWorkflowsRuntimeOptions(CorrelationId: message.CorrelationId, Input: message.Input); await _workflowRuntime.ResumeWorkflowsAsync(message.ActivityTypeName, message.BookmarkPayload, options, context.CancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs index 98f4b583c..605e7391d 100644 --- a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs @@ -89,6 +89,31 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime return new WorkflowExecutionResult(workflowInstanceId, bookmarks); } + /// + public async Task> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) + { + var hash = _hasher.Hash(activityTypeName, bookmarkPayload); + var triggers = await _triggerStore.FindAsync(hash, cancellationToken); + var results = new List(); + + foreach (var trigger in triggers) + { + var definitionId = trigger.WorkflowDefinitionId; + var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId); + var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken); + + // If we can't start the workflow, don't try it. + if(!canStartResult.CanStart) + continue; + + var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken); + + results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks)); + } + + return results; + } + /// public async Task ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) { @@ -109,7 +134,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime } /// - public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) + public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { var hash = _hasher.Hash(activityTypeName, bookmarkPayload); var client = _cluster.GetNamedBookmarkGrain(hash); @@ -122,58 +147,17 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime var bookmarksResponse = await client.Resolve(request, cancellationToken); var bookmarks = bookmarksResponse!.Bookmarks; - return await ResumeWorkflowsAsync(bookmarks, options, cancellationToken); + return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken); } /// public async Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { - var triggeredWorkflows = new List(); - var hash = _hasher.Hash(activityTypeName, bookmarkPayload); + var startedWorkflows = await StartWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken); + var resumedWorkflows = await ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken); + var results = startedWorkflows.Concat(resumedWorkflows).ToList(); - // Start new workflows. - var triggers = await _triggerStore.FindAsync(hash, cancellationToken); - - foreach (var trigger in triggers) - { - var definitionId = trigger.WorkflowDefinitionId; - var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId); - var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken); - - // If we can't start the workflow, don't try it. - if(!canStartResult.CanStart) - continue; - - var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken); - - triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks)); - } - - // Resume existing workflow instances. - var client = _cluster.GetNamedBookmarkGrain(hash); - - var request = new ResolveBookmarksRequest - { - ActivityTypeName = activityTypeName, - CorrelationId = options.CorrelationId.EmptyIfNull() - }; - - var bookmarksResponse = await client.Resolve(request, cancellationToken); - var bookmarks = bookmarksResponse!.Bookmarks; - - foreach (var bookmark in bookmarks) - { - var workflowInstanceId = bookmark.WorkflowInstanceId; - - var resumeResult = await ResumeWorkflowAsync( - workflowInstanceId, - new ResumeWorkflowRuntimeOptions(options.CorrelationId, bookmark.BookmarkId, null, options.Input), - cancellationToken); - - triggeredWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks)); - } - - return new TriggerWorkflowsResult(triggeredWorkflows); + return new TriggerWorkflowsResult(results); } /// diff --git a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs index 72339a8d8..0034329cc 100644 --- a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs +++ b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs @@ -36,7 +36,7 @@ internal class TriggerWebhookDrivenActivities : INotificationHandler FindActivityDescriptors(string eventType) => diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs index 56cab6a03..98949b315 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs @@ -27,17 +27,25 @@ public class Flowchart : Container } [Port] [Browsable(false)] public IActivity? Start { get; set; } + /// + /// A list of connections between activities. + /// public ICollection Connections { get; set; } = new List(); + /// protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context) { - if (Start == null!) + var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId; + var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : default; + var startActivity = triggerActivity ?? Start; + + if (startActivity == null!) { await context.CompleteActivityAsync(); return; } - await context.ScheduleActivityAsync(Start); + await context.ScheduleActivityAsync(startActivity); } private async ValueTask OnDescendantCompletedAsync(ActivityCompleted signal, SignalContext context) diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs index 400619c4c..c30490bca 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs @@ -6,12 +6,18 @@ using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Core.Middleware.Activities; +/// +/// An activity execution middleware component that logs information about the activity being executed. +/// public class LoggingMiddleware : IActivityExecutionMiddleware { private readonly ActivityMiddlewareDelegate _next; private readonly ILogger _logger; private readonly Stopwatch _stopwatch; + /// + /// Constructor. + /// public LoggingMiddleware(ActivityMiddlewareDelegate next, ILogger logger) { _next = next; @@ -19,18 +25,25 @@ public class LoggingMiddleware : IActivityExecutionMiddleware _stopwatch = new Stopwatch(); } + /// public async ValueTask InvokeAsync(ActivityExecutionContext context) { var activity = context.Activity; - _logger.LogDebug("Executing activity {ActivityType}", activity.Type); + _logger.LogInformation("Executing activity {ActivityId}", activity.Id); _stopwatch.Restart(); await _next(context); _stopwatch.Stop(); - _logger.LogDebug("Executed activity {ActivityType} in {Elapsed}", activity.Type, _stopwatch.Elapsed); + _logger.LogInformation("Executed activity {ActivityId} in {Elapsed}", activity.Id, _stopwatch.Elapsed); } } +/// +/// Extends to install the component. +/// public static class LoggingMiddlewareExtensions { + /// + /// Installs the component. + /// public static IActivityExecutionPipelineBuilder UseLogging(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs index 3ea5fdeac..a255f2314 100644 --- a/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs @@ -266,7 +266,7 @@ public class WorkflowExecutionContext /// Returns the with the specified ID from the workflow graph. /// public IActivity FindActivityByNodeId(string nodeId) => FindNodeById(nodeId).Activity; - + /// /// Returns a custom property with the specified key from the dictionary. /// diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs index 354e7af85..f5a5646fa 100644 --- a/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs @@ -46,7 +46,7 @@ internal class DispatchWorkflowRequestHandler : public async Task HandleAsync(DispatchResumeWorkflowsCommand command, CancellationToken cancellationToken) { - var options = new ResumeWorkflowRuntimeOptions(CorrelationId: command.CorrelationId, Input: command.Input); + var options = new TriggerWorkflowsRuntimeOptions(CorrelationId: command.CorrelationId, Input: command.Input); await _workflowRuntime.ResumeWorkflowsAsync(command.ActivityTypeName, command.BookmarkPayload, options, cancellationToken); return Unit.Instance; diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs index e8f96472f..6ca90d31b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs @@ -4,6 +4,7 @@ using Elsa.Workflows.Core.Models; using Elsa.Workflows.Core.Notifications; using Elsa.Workflows.Runtime.Activities; using Elsa.Workflows.Runtime.Bookmarks; +using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Services; using JetBrains.Annotations; @@ -15,9 +16,9 @@ namespace Elsa.Workflows.Runtime.Handlers; [PublicAPI] internal class ResumeDispatchWorkflowActivityHandler : INotificationHandler { - private readonly IWorkflowRuntime _workflowRuntime; + private readonly IWorkflowDispatcher _workflowRuntime; - public ResumeDispatchWorkflowActivityHandler(IWorkflowRuntime workflowRuntime) + public ResumeDispatchWorkflowActivityHandler(IWorkflowDispatcher workflowRuntime) { _workflowRuntime = workflowRuntime; } @@ -29,6 +30,7 @@ internal class ResumeDispatchWorkflowActivityHandler : INotificationHandler(); - await _workflowRuntime.TriggerWorkflowsAsync(activityTypeName, bookmark, new TriggerWorkflowsRuntimeOptions(), cancellationToken); + var request = new DispatchResumeWorkflowsRequest(activityTypeName, bookmark); + await _workflowRuntime.DispatchAsync(request, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs index 477a078f8..6592d0416 100644 --- a/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs @@ -68,6 +68,40 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks); } + /// + public async Task> StartWorkflowsAsync( + string activityTypeName, + object bookmarkPayload, + TriggerWorkflowsRuntimeOptions options, + CancellationToken cancellationToken = default) + { + var results = new List(); + var hash = _hasher.Hash(activityTypeName, bookmarkPayload); + + // Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances from being created. + var sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows__{hash}"; + await using (await _distributedLockProvider.AcquireLockAsync(sharedResource, TimeSpan.FromMinutes(10), cancellationToken)) + { + var triggers = await _triggerStore.FindAsync(hash, cancellationToken); + + foreach (var trigger in triggers) + { + var definitionId = trigger.WorkflowDefinitionId; + var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId); + var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken); + + // If we can't start the workflow, don't try it. + if (!canStartResult.CanStart) + continue; + + var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken); + results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks)); + } + } + + return results; + } + /// public async Task ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) { @@ -104,52 +138,21 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime } /// - public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) + public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { var hash = _hasher.Hash(activityTypeName, bookmarkPayload); var correlationId = options.CorrelationId; - var bookmarks = correlationId == null ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken); - return await ResumeWorkflowsAsync(bookmarks, options, cancellationToken); + var bookmarks = string.IsNullOrWhiteSpace(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken); + return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken); } /// - public async Task TriggerWorkflowsAsync( - string activityTypeName, - object bookmarkPayload, - TriggerWorkflowsRuntimeOptions options, - CancellationToken cancellationToken = default) + public async Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { - var triggeredWorkflows = new List(); - var hash = _hasher.Hash(activityTypeName, bookmarkPayload); - - // Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances being created. - const string sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows"; - await using (await _distributedLockProvider.AcquireLockAsync(sharedResource, TimeSpan.FromMinutes(10), cancellationToken)) - { - var triggers = await _triggerStore.FindAsync(hash, cancellationToken); - - foreach (var trigger in triggers) - { - var definitionId = trigger.WorkflowDefinitionId; - var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId); - var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken); - - // If we can't start the workflow, don't try it. - if (!canStartResult.CanStart) - continue; - - var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken); - triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks)); - } - } - - // Resume bookmarks. - var correlationId = options.CorrelationId; - var bookmarks = (string.IsNullOrEmpty(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken)).ToList(); - var resumedWorkflows = await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken); - - triggeredWorkflows.AddRange(resumedWorkflows.Select(x => new WorkflowExecutionResult(x.InstanceId, x.Bookmarks))); - return new TriggerWorkflowsResult(triggeredWorkflows); + var startedWorkflows = await StartWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken); + var resumedWorkflows = await ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken); + var results = startedWorkflows.Concat(resumedWorkflows).ToList(); + return new TriggerWorkflowsResult(results); } /// diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs index 712d4d2a9..d95454d4d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs @@ -22,6 +22,16 @@ public interface IWorkflowRuntime /// /// Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default); + + /// + /// Starts all workflows with triggers matching the specified activity type and bookmark payload. + /// + /// + Task> StartWorkflowsAsync( + string activityTypeName, + object bookmarkPayload, + TriggerWorkflowsRuntimeOptions options, + CancellationToken cancellationToken = default); /// /// Resumes an existing workflow instance. @@ -34,7 +44,7 @@ public interface IWorkflowRuntime /// /// Resumes all workflows that are bookmarked on the specified activity type. /// - Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default); + Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default); /// /// Starts all workflows and resumes existing workflow instances based on the specified activity type and bookmark payload.