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