From 0993b38b7affa8e814928c28afd99306fb7b19a6 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 29 Mar 2024 19:58:35 +0100 Subject: [PATCH] Fix Alteration Runner (#5142) * Refactor alteration runner logic and simplify CancelActivityHandler The DefaultAlterationRunner.cs file was refactored to include a workflow middleware pipeline, replacing the previous method of updating workflow state. Comments were added to discuss potential solutions for reusing the same pipeline in the workflow runtime. The CancelActivityHandler was simplified by removing multiple functions and replacing them with a single CancelAsync method. * Add RunAlterationsMiddleware to handle workflow alterations The commit introduces a new middleware, RunAlterationsMiddleware, that is designed to process workflow alterations. The middleware is in charge of executing alteration handlers and taking care of any required commit actions. The original code for handling alterations in DefaultAlterationRunner has been significantly reduced as it now delegates most of its responsibility to this new middleware. * Add comments to IAlterationPlanManager interface methods In the IAlterationPlanManager interface, explanatory comments were added to each method. These include methods for getting a plan by ID, checking if all jobs in the plan have been completed, and completing an alteration plan. The changes made will greatly aid in understanding the purpose and functionality of each method. * Refactor pipeline execution methods and replace middleware The workflow execution pipeline has been refactored to allow for dynamic configuration. The middleware components are now retrievable properties and can be replaced individually. This change provides more extensibility with modifying the pipeline execution and replacing the DefaultActivitySchedulerMiddleware with desired middleware. * Refactor pipeline alteration and addition methods Simplified the pipeline alteration process in Elsa.Alterations. Instead of manually handling middleware delegates, added a new extension method, ReplaceTerminal, in WorkflowExecutionMiddlewareExtensions.cs to replace the terminal middleware component. This approach improves code readability and maintenance. * Refactor RunAlterationsMiddleware class Removed an unnecessary extension class and updated the RunAlterationsMiddleware class to streamline its structure. The handlers are now initialized directly in the constructor, eliminating the need for an additional field. Renamed local variable for better code clarity. * Refactor variable name in RunAlterationsMiddleware The variable name 'handlers1' has been renamed to 'supportedHandlers' in the RunAlterationsMiddleware.cs file. This change improves code readability and makes it clear that the list contains only the handlers that can handle the given alteration. * Change RunAsync method to return void The RunAsync function in the RunAlterationsMiddleware class no longer returns a boolean value. The returned 'false' has been replaced with a 'return' statement, and the 'return true' statement has been completely removed. This refactor simplifies the control flow when executing alterations. --- .../Contracts/IAlterationPlanManager.cs | 14 +++++ .../CancelActivityHandler.cs | 42 +++---------- .../Workflows/RunAlterationsMiddleware.cs | 54 ++++++++++++++++ .../Services/DefaultAlterationRunner.cs | 63 +++++++------------ .../Contracts/IWorkflowExecutionPipeline.cs | 21 +++++++ .../IWorkflowExecutionPipelineBuilder.cs | 13 ++++ .../WorkflowExecutionMiddlewareExtensions.cs | 43 +++++++++++-- .../WorkflowExecutionPipeline.cs | 12 ++-- .../WorkflowExecutionPipelineBuilder.cs | 9 +++ 9 files changed, 186 insertions(+), 85 deletions(-) create mode 100644 src/modules/Elsa.Alterations/Middleware/Workflows/RunAlterationsMiddleware.cs diff --git a/src/modules/Elsa.Alterations.Core/Contracts/IAlterationPlanManager.cs b/src/modules/Elsa.Alterations.Core/Contracts/IAlterationPlanManager.cs index 79b72c6d9..617cf2824 100644 --- a/src/modules/Elsa.Alterations.Core/Contracts/IAlterationPlanManager.cs +++ b/src/modules/Elsa.Alterations.Core/Contracts/IAlterationPlanManager.cs @@ -2,9 +2,23 @@ using Elsa.Alterations.Core.Entities; namespace Elsa.Alterations.Core.Contracts; +/// +/// Represents a manager for alteration plans. +/// public interface IAlterationPlanManager { + /// + /// Gets an alteration plan by ID. + /// Task GetPlanAsync(string planId, CancellationToken cancellationToken = default); + + /// + /// Gets a value indicating whether all jobs in the plan have been completed. + /// Task GetIsAllJobsCompletedAsync(string planId, CancellationToken cancellationToken = default); + + /// + /// Completes an alteration plan. + /// Task CompletePlanAsync(AlterationPlan plan, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Alterations/AlterationHandlers/CancelActivityHandler.cs b/src/modules/Elsa.Alterations/AlterationHandlers/CancelActivityHandler.cs index 05f58d16b..95ffa2db7 100644 --- a/src/modules/Elsa.Alterations/AlterationHandlers/CancelActivityHandler.cs +++ b/src/modules/Elsa.Alterations/AlterationHandlers/CancelActivityHandler.cs @@ -1,9 +1,8 @@ using Elsa.Alterations.AlterationTypes; using Elsa.Alterations.Core.Abstractions; using Elsa.Alterations.Core.Contexts; +using Elsa.Extensions; using Elsa.Workflows; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Filters; using JetBrains.Annotations; namespace Elsa.Alterations.AlterationHandlers; @@ -14,20 +13,6 @@ namespace Elsa.Alterations.AlterationHandlers; [UsedImplicitly] public class CancelActivityHandler : AlterationHandlerBase { - private readonly IBookmarkManager _bookmarkManager; - private readonly IActivityExecutionManager _activityExecutionManager; - private readonly IActivityExecutionMapper _activityExecutionMapper; - - /// - /// Initializes a new instance of the class. - /// - public CancelActivityHandler(IBookmarkManager bookmarkManager, IActivityExecutionManager activityExecutionManager, IActivityExecutionMapper activityExecutionMapper) - { - _bookmarkManager = bookmarkManager; - _activityExecutionManager = activityExecutionManager; - _activityExecutionMapper = activityExecutionMapper; - } - /// protected override ValueTask HandleAsync(AlterationContext context, CancelActivity alteration) { @@ -49,32 +34,19 @@ public class CancelActivityHandler : AlterationHandlerBase return ValueTask.CompletedTask; } - context.Succeed(() => CleanupAsync(activityExecutionContexts)); + context.Succeed(() => CancelAsync(activityExecutionContexts)); return ValueTask.CompletedTask; } - private async Task CleanupAsync(IEnumerable activityExecutionContexts) + private async Task CancelAsync(IEnumerable activityExecutionContexts) { - foreach (var activityExecutionContext in activityExecutionContexts) - await CleanupAsync(activityExecutionContext); - } - - private async Task CleanupAsync(ActivityExecutionContext activityExecutionContext) - { - await RemoveBookmarksAsync(activityExecutionContext); - await SaveActivityExecutionRecordAsync(activityExecutionContext); + foreach (var activityExecutionContext in activityExecutionContexts) + await CancelAsync(activityExecutionContext); } - private async Task RemoveBookmarksAsync(ActivityExecutionContext activityExecutionContext) + private async Task CancelAsync(ActivityExecutionContext activityExecutionContext) { - var filter = new BookmarkFilter { ActivityInstanceId = activityExecutionContext.Id }; - await _bookmarkManager.DeleteManyAsync(filter, activityExecutionContext.CancellationToken); - } - - private async Task SaveActivityExecutionRecordAsync(ActivityExecutionContext activityExecutionContext) - { - var activityExecutionRecord = _activityExecutionMapper.Map(activityExecutionContext); - await _activityExecutionManager.SaveAsync(activityExecutionRecord, CancellationToken.None); + await activityExecutionContext.CancelActivityAsync(); } private static IEnumerable GetActivityExecutionContexts(AlterationContext context, CancelActivity alteration) diff --git a/src/modules/Elsa.Alterations/Middleware/Workflows/RunAlterationsMiddleware.cs b/src/modules/Elsa.Alterations/Middleware/Workflows/RunAlterationsMiddleware.cs new file mode 100644 index 000000000..140c598fb --- /dev/null +++ b/src/modules/Elsa.Alterations/Middleware/Workflows/RunAlterationsMiddleware.cs @@ -0,0 +1,54 @@ +using Elsa.Alterations.Core.Contexts; +using Elsa.Alterations.Core.Contracts; +using Elsa.Alterations.Core.Models; +using Elsa.Extensions; +using Elsa.Workflows; +using Elsa.Workflows.Pipelines.WorkflowExecution; + +namespace Elsa.Alterations.Middleware.Workflows; + +/// +/// Middleware that runs alterations. +/// +internal class RunAlterationsMiddleware(WorkflowMiddlewareDelegate next, IEnumerable handlers) : WorkflowExecutionMiddleware(next) +{ + public static readonly object AlterationsPropertyKey = new(); + public static readonly object AlterationsLogPropertyKey = new(); + + public override async ValueTask InvokeAsync(WorkflowExecutionContext context) + { + var alterations = (IEnumerable)(context.TransientProperties.GetValue(AlterationsPropertyKey) ?? throw new InvalidOperationException("No alterations found in the transient properties.")); + var log = (AlterationLog)(context.TransientProperties.GetValue(AlterationsLogPropertyKey) ?? throw new InvalidOperationException("No alteration log found in the transient properties.")); + await RunAsync(context, alterations, log, context.CancellationTokens.ApplicationCancellationToken); + } + + private async Task RunAsync(WorkflowExecutionContext workflowExecutionContext, IEnumerable alterations, AlterationLog log, CancellationToken cancellationToken = default) + { + var commitActions = new List>(); + + foreach (var alteration in alterations) + { + // Find handlers. + var supportedHandlers = handlers.Where(x => x.CanHandle(alteration)).ToList(); + + foreach (var handler in supportedHandlers) + { + // Execute handler. + var alterationContext = new AlterationContext(alteration, workflowExecutionContext, log, cancellationToken); + await handler.HandleAsync(alterationContext); + + // If the handler has failed, exit. + if (alterationContext.HasFailed) + return; + + // Collect the commit handler, if any. + if (alterationContext.CommitAction != null) + commitActions.Add(alterationContext.CommitAction); + } + } + + // Execute commit handlers. + foreach (var commitAction in commitActions) + await commitAction(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs b/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs index 5cf2df277..dfddcd7fa 100644 --- a/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs +++ b/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs @@ -2,10 +2,13 @@ using Elsa.Alterations.Core.Contexts; using Elsa.Alterations.Core.Contracts; using Elsa.Alterations.Core.Models; using Elsa.Alterations.Core.Results; +using Elsa.Alterations.Middleware.Workflows; using Elsa.Common.Contracts; +using Elsa.Extensions; using Elsa.Workflows; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Pipelines.WorkflowExecution; using Elsa.Workflows.Runtime.Contracts; using Microsoft.Extensions.Logging; @@ -14,8 +17,8 @@ namespace Elsa.Alterations.Services; /// public class DefaultAlterationRunner : IAlterationRunner { - private readonly IEnumerable _handlers; private readonly IWorkflowRuntime _workflowRuntime; + private readonly IWorkflowExecutionPipeline _workflowExecutionPipeline; private readonly IWorkflowDefinitionService _workflowDefinitionService; private readonly IWorkflowStateExtractor _workflowStateExtractor; private readonly ISystemClock _systemClock; @@ -25,15 +28,15 @@ public class DefaultAlterationRunner : IAlterationRunner /// Initializes a new instance of the class. /// public DefaultAlterationRunner( - IEnumerable handlers, IWorkflowRuntime workflowRuntime, + IWorkflowExecutionPipeline workflowExecutionPipeline, IWorkflowDefinitionService workflowDefinitionService, IWorkflowStateExtractor workflowStateExtractor, ISystemClock systemClock, IServiceProvider serviceProvider) { - _handlers = handlers; _workflowRuntime = workflowRuntime; + _workflowExecutionPipeline = workflowExecutionPipeline; _workflowDefinitionService = workflowDefinitionService; _workflowStateExtractor = workflowStateExtractor; _systemClock = systemClock; @@ -86,15 +89,23 @@ public class DefaultAlterationRunner : IAlterationRunner // Create workflow execution context. var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflow, workflowState, cancellationTokens: cancellationToken); + workflowExecutionContext.TransientProperties.Add(RunAlterationsMiddleware.AlterationsPropertyKey, alterations); + workflowExecutionContext.TransientProperties.Add(RunAlterationsMiddleware.AlterationsLogPropertyKey, log); - // Execute alterations. - var success = await RunAsync(workflowExecutionContext, alterations, log, cancellationToken); - - // If the alterations have failed, exit. - if (!success) - return result; + // Build a new workflow execution pipeline. + var pipelineBuilder = new WorkflowExecutionPipelineBuilder(_serviceProvider); + _workflowExecutionPipeline.ConfigurePipelineBuilder(pipelineBuilder); - // Update workflow state. + // Replace the terminal DefaultActivitySchedulerMiddleware with the RunAlterationsMiddleware terminal. + pipelineBuilder.ReplaceTerminal(); + + // Build modified pipeline. + var pipeline = pipelineBuilder.Build(); + + // Execute the pipeline. + await pipeline(workflowExecutionContext); + + // Extract workflow state. workflowState = _workflowStateExtractor.Extract(workflowExecutionContext); // Apply updated workflow state. @@ -105,36 +116,4 @@ public class DefaultAlterationRunner : IAlterationRunner return result; } - - private async Task RunAsync(WorkflowExecutionContext workflowExecutionContext, IEnumerable alterations, AlterationLog log, CancellationToken cancellationToken = default) - { - var commitActions = new List>(); - - foreach (var alteration in alterations) - { - // Find handlers. - var handlers = _handlers.Where(x => x.CanHandle(alteration)).ToList(); - - foreach (var handler in handlers) - { - // Execute handler. - var alterationContext = new AlterationContext(alteration, workflowExecutionContext, log, cancellationToken); - await handler.HandleAsync(alterationContext); - - // If the handler has failed, exit. - if (alterationContext.HasFailed) - return false; - - // Collect the commit handler, if any. - if (alterationContext.CommitAction != null) - commitActions.Add(alterationContext.CommitAction); - } - } - - // Execute commit handlers. - foreach (var commitAction in commitActions) - await commitAction(); - - return true; - } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipeline.cs b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipeline.cs index eea9b5f8f..4bf7429f4 100644 --- a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipeline.cs +++ b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipeline.cs @@ -2,9 +2,30 @@ using Elsa.Workflows.Pipelines.WorkflowExecution; namespace Elsa.Workflows.Contracts; +/// +/// Represents a workflow execution pipeline. +/// public interface IWorkflowExecutionPipeline { + /// + /// The pipeline builder factory. + /// + Action ConfigurePipelineBuilder { get; } + + /// + /// Sets up the pipeline using the specified pipeline builder factory. + /// WorkflowMiddlewareDelegate Setup(Action setup); + + /// + /// The constructed pipeline delegate. + /// WorkflowMiddlewareDelegate Pipeline { get; } + + /// + /// Executes the pipeline with the specified workflow execution context. + /// + /// + /// Task ExecuteAsync(WorkflowExecutionContext context); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipelineBuilder.cs b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipelineBuilder.cs index 1a895b4d5..a2f37c9b0 100644 --- a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipelineBuilder.cs +++ b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowExecutionPipelineBuilder.cs @@ -12,6 +12,11 @@ public interface IWorkflowExecutionPipelineBuilder /// public IDictionary Properties { get; } + /// + /// The middleware components that have been installed. + /// + public IEnumerable> Components { get; } + /// /// The current service provider to resolve services from. /// @@ -31,4 +36,12 @@ public interface IWorkflowExecutionPipelineBuilder /// Clears the current pipeline. /// IWorkflowExecutionPipelineBuilder Reset(); + + /// + /// Replaces the middleware component at the specified index with the specified delegate. + /// + /// The index of the middleware component to replace. + /// The delegate to use as the new middleware component. + /// A new instance of with the middleware component replaced. + IWorkflowExecutionPipelineBuilder Replace(int index, Func middleware); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionMiddlewareExtensions.cs b/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionMiddlewareExtensions.cs index f220f8c46..4e7221c11 100644 --- a/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionMiddlewareExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionMiddlewareExtensions.cs @@ -1,3 +1,4 @@ +using System.Diagnostics.CodeAnalysis; using Elsa.Workflows.Contracts; using Microsoft.Extensions.DependencyInjection; @@ -11,16 +12,50 @@ public static class WorkflowExecutionMiddlewareExtensions /// /// Installs the specified middleware component into the pipeline being built. /// - public static IWorkflowExecutionPipelineBuilder UseMiddleware(this IWorkflowExecutionPipelineBuilder pipelineBuilder, params object[] args) where TMiddleware: IWorkflowExecutionMiddleware + public static IWorkflowExecutionPipelineBuilder UseMiddleware<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicConstructors)] TMiddleware>( + this IWorkflowExecutionPipelineBuilder pipelineBuilder, params object[] args) where TMiddleware : IWorkflowExecutionMiddleware + { + var delegateFactory = CreateMiddlewareDelegateFactory(pipelineBuilder, args); + return pipelineBuilder.Use(delegateFactory); + } + + /// + /// Replaces the terminal middleware component with the specified middleware component. + /// + public static IWorkflowExecutionPipelineBuilder ReplaceTerminal<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicConstructors)] TMiddleware>( + this IWorkflowExecutionPipelineBuilder pipelineBuilder, params object[] args) where TMiddleware : IWorkflowExecutionMiddleware + { + var index = pipelineBuilder.Components.Count() - 1; + return pipelineBuilder.Replace(index, args); + } + + /// + /// Replaces the middleware component at the specified index with the specified middleware component. + /// + public static IWorkflowExecutionPipelineBuilder Replace<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicConstructors)] TMiddleware>( + this IWorkflowExecutionPipelineBuilder pipelineBuilder, int index, params object[] args) where TMiddleware : IWorkflowExecutionMiddleware + { + var delegateFactory = CreateMiddlewareDelegateFactory(pipelineBuilder, args); + return pipelineBuilder.Replace(index, delegateFactory); + } + + /// + /// Creates a middleware delegate for the specified middleware component. + /// + public static Func CreateMiddlewareDelegateFactory<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicConstructors)] TMiddleware>( + this IWorkflowExecutionPipelineBuilder pipelineBuilder, params object[] args) where TMiddleware : IWorkflowExecutionMiddleware { var middleware = typeof(TMiddleware); - return pipelineBuilder.Use(next => + return next => { var invokeMethod = MiddlewareHelpers.GetInvokeMethod(middleware); - var ctorParams = new[] { next }.Concat(args).Select(x => x).ToArray(); + var ctorParams = new[] + { + next + }.Concat(args).Select(x => x).ToArray(); var instance = ActivatorUtilities.CreateInstance(pipelineBuilder.ServiceProvider, middleware, ctorParams); return (WorkflowMiddlewareDelegate)invokeMethod.CreateDelegate(typeof(WorkflowMiddlewareDelegate), instance); - }); + }; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipeline.cs b/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipeline.cs index 9b1023dd6..a53439dcd 100644 --- a/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipeline.cs +++ b/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipeline.cs @@ -10,17 +10,21 @@ public class WorkflowExecutionPipeline : IWorkflowExecutionPipeline private WorkflowMiddlewareDelegate? _pipeline; /// - /// Constructor. + /// Initializes a new instance of the class. /// - public WorkflowExecutionPipeline(IServiceProvider serviceProvider, Action pipelineBuilder) + public WorkflowExecutionPipeline(IServiceProvider serviceProvider, Action configurePipelineBuilder) { _serviceProvider = serviceProvider; - Setup(pipelineBuilder); + ConfigurePipelineBuilder = configurePipelineBuilder; + Setup(configurePipelineBuilder); } - + /// public WorkflowMiddlewareDelegate Pipeline => _pipeline ??= CreateDefaultPipeline(); + /// + public Action ConfigurePipelineBuilder { get; } + /// public WorkflowMiddlewareDelegate Setup(Action setup) { diff --git a/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipelineBuilder.cs b/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipelineBuilder.cs index 3e46b827f..c61c6c4e5 100644 --- a/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipelineBuilder.cs +++ b/src/modules/Elsa.Workflows.Core/Pipelines/WorkflowExecution/WorkflowExecutionPipelineBuilder.cs @@ -19,6 +19,9 @@ public class WorkflowExecutionPipelineBuilder : IWorkflowExecutionPipelineBuilde /// public IDictionary Properties { get; } = new Dictionary(); + /// + public IEnumerable> Components => _components.ToList(); + /// public IServiceProvider ServiceProvider { @@ -51,6 +54,12 @@ public class WorkflowExecutionPipelineBuilder : IWorkflowExecutionPipelineBuilde return this; } + public IWorkflowExecutionPipelineBuilder Replace(int index, Func middleware) + { + _components[index] = middleware; + return this; + } + private T? GetProperty(string key) => Properties.TryGetValue(key, out var value) ? (T?)value : default(T); private void SetProperty(string key, T value) => Properties[key] = value; } \ No newline at end of file