From beedc471dce09590e16c2f1815e4bba2b53c1692 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 19 Apr 2022 20:42:26 +0200 Subject: [PATCH] Incremental work on activity signaling --- src/core/Elsa.Core/Activities/Break.cs | 14 +---- src/core/Elsa.Core/Activities/Container.cs | 40 ++++++++++--- src/core/Elsa.Core/Activities/For.cs | 11 +++- src/core/Elsa.Core/Activities/ForEach.cs | 15 ++++- src/core/Elsa.Core/Activities/Sequence.cs | 15 +++-- src/core/Elsa.Core/Activities/While.cs | 18 ++++-- .../Elsa.Core/Contracts/ILoopingConstruct.cs | 6 -- .../Elsa.Core/Contracts/ISignalHandler.cs | 28 +++++++++ src/core/Elsa.Core/Models/Activity.cs | 53 +++++++++++++++-- .../Models/ActivityExecutionContext.cs | 25 +++++++- .../Components/ActivityInvokerMiddleware.cs | 44 -------------- .../Elsa.Core/Signals/ActivityCompleted.cs | 3 + src/core/Elsa.Core/Signals/BreakSignal.cs | 3 + .../console/Elsa.Samples.Console1/Program.cs | 2 +- .../Workflows/SequentialWorkflowTests.cs | 59 +++++++++++++++++++ 15 files changed, 248 insertions(+), 88 deletions(-) delete mode 100644 src/core/Elsa.Core/Contracts/ILoopingConstruct.cs create mode 100644 src/core/Elsa.Core/Contracts/ISignalHandler.cs create mode 100644 src/core/Elsa.Core/Signals/ActivityCompleted.cs create mode 100644 src/core/Elsa.Core/Signals/BreakSignal.cs create mode 100644 test/Elsa.IntegrationTests/Workflows/SequentialWorkflowTests.cs diff --git a/src/core/Elsa.Core/Activities/Break.cs b/src/core/Elsa.Core/Activities/Break.cs index 3ebccc5fc..07364cfed 100644 --- a/src/core/Elsa.Core/Activities/Break.cs +++ b/src/core/Elsa.Core/Activities/Break.cs @@ -1,22 +1,14 @@ using Elsa.Attributes; -using Elsa.Contracts; using Elsa.Models; -using Microsoft.Extensions.Logging; +using Elsa.Signals; namespace Elsa.Activities; [Activity("Elsa", "Control Flow", "Break out of a loop")] public class Break : Activity { - protected override void Execute(ActivityExecutionContext context) + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { - // Find the first parent looping construct. - var loopingConstructContext = context.GetAncestorActivityExecutionContexts().FirstOrDefault(x => x.Activity is ILoopingConstruct); - var loopingConstructActivity = (ILoopingConstruct?)loopingConstructContext?.Activity; - - if (loopingConstructActivity != null) - { - - } + await context.SignalAsync(new BreakSignal()); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Container.cs b/src/core/Elsa.Core/Activities/Container.cs index 629b9ef90..df84ad7ca 100644 --- a/src/core/Elsa.Core/Activities/Container.cs +++ b/src/core/Elsa.Core/Activities/Container.cs @@ -2,6 +2,7 @@ using System.Collections.ObjectModel; using Elsa.Attributes; using Elsa.Contracts; using Elsa.Models; +using Elsa.Signals; namespace Elsa.Activities; @@ -12,27 +13,52 @@ public abstract class Container : Activity, IContainer { protected Container() { + OnSignalReceived(OnChildActivityCompletedAsync); } - protected Container(params IActivity[] activities) => Activities = activities; - - protected Container(ICollection variables, params IActivity[] activities) + protected Container(params IActivity[] activities) : this() + { + Activities = activities; + } + + protected Container(ICollection variables, params IActivity[] activities) : this(activities) { Variables = variables; - Activities = activities; } [Outbound] public ICollection Activities { get; set; } = new List(); public ICollection Variables { get; set; } = new Collection(); - protected override void Execute(ActivityExecutionContext context) + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { // Register variables. context.ExpressionExecutionContext.Register.Declare(Variables); // Schedule children. - ScheduleChildren(context); + await ScheduleChildrenAsync(context); } - protected abstract void ScheduleChildren(ActivityExecutionContext context); + protected virtual async ValueTask OnChildActivityCompletedAsync(ActivityCompleted signal, SignalContext context) + { + var activityExecutionContext = context.ActivityExecutionContext; + var ownerActivity = activityExecutionContext.Activity; + var childActivityExecutionContext = context.SourceActivityExecutionContext; + var childActivity = childActivityExecutionContext.Activity; + var callbackEntry = activityExecutionContext.WorkflowExecutionContext.CompletionCallbacks.FirstOrDefault(x => x.Owner.Activity == ownerActivity && x.Child == childActivity); + + if (callbackEntry == null) + return; + + await callbackEntry.CompletionCallback(activityExecutionContext, childActivityExecutionContext); + } + + protected virtual ValueTask ScheduleChildrenAsync(ActivityExecutionContext context) + { + ScheduleChildren(context); + return ValueTask.CompletedTask; + } + + protected virtual void ScheduleChildren(ActivityExecutionContext context) + { + } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/For.cs b/src/core/Elsa.Core/Activities/For.cs index 0f6c22709..6b37a264b 100644 --- a/src/core/Elsa.Core/Activities/For.cs +++ b/src/core/Elsa.Core/Activities/For.cs @@ -1,6 +1,8 @@ +using System.Text.Json.Serialization; using Elsa.Attributes; using Elsa.Contracts; using Elsa.Models; +using Elsa.Signals; namespace Elsa.Activities; @@ -15,11 +17,13 @@ public enum ForOperator [Activity("Elsa", "Control Flow", "Iterate over a sequence of steps between a start and an end number.")] public class For : Activity { + [JsonConstructor] public For() { + OnSignalReceived(OnBreak); } - public For(int start, int end, ForOperator forOperator = ForOperator.LessThanOrEqual) + public For(int start, int end, ForOperator forOperator = ForOperator.LessThanOrEqual) : this() { Start = new Input(start); End = new Input(end); @@ -79,4 +83,9 @@ public class For : Activity HandleIteration(ownerActivityExecutionContext); return ValueTask.CompletedTask; } + + private void OnBreak(BreakSignal signal, SignalContext context) + { + context.StopPropagation(); + } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/ForEach.cs b/src/core/Elsa.Core/Activities/ForEach.cs index adc82a415..3b3f1012c 100644 --- a/src/core/Elsa.Core/Activities/ForEach.cs +++ b/src/core/Elsa.Core/Activities/ForEach.cs @@ -2,6 +2,7 @@ using System.Text.Json.Serialization; using Elsa.Attributes; using Elsa.Contracts; using Elsa.Models; +using Elsa.Signals; namespace Elsa.Activities; @@ -9,7 +10,12 @@ namespace Elsa.Activities; public class ForEach : Activity { private const string CurrentIndexProperty = "CurrentIndex"; - + + public ForEach() + { + OnSignalReceived(OnBreak); + } + /// /// The set of values to iterate. /// @@ -57,6 +63,11 @@ public class ForEach : Activity HandleIteration(context); return ValueTask.CompletedTask; } + + private void OnBreak(BreakSignal signal, SignalContext context) + { + context.StopPropagation(); + } } public class ForEach : ForEach @@ -66,7 +77,7 @@ public class ForEach : ForEach { } - public ForEach(Input> items) + public ForEach(Input> items) : this() { Items = items; } diff --git a/src/core/Elsa.Core/Activities/Sequence.cs b/src/core/Elsa.Core/Activities/Sequence.cs index d38621c44..4c39ebc08 100644 --- a/src/core/Elsa.Core/Activities/Sequence.cs +++ b/src/core/Elsa.Core/Activities/Sequence.cs @@ -2,6 +2,7 @@ using System.ComponentModel; using Elsa.Attributes; using Elsa.Contracts; using Elsa.Models; +using Elsa.Signals; namespace Elsa.Activities; @@ -23,27 +24,29 @@ public class Sequence : Container { } - protected override void ScheduleChildren(ActivityExecutionContext context) + protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context) { - HandleItem(context); + await HandleItemAsync(context); } - private void HandleItem(ActivityExecutionContext context) + private async ValueTask HandleItemAsync(ActivityExecutionContext context) { var currentIndex = context.GetProperty(CurrentIndexProperty); var childActivities = Activities.ToList(); if (currentIndex >= childActivities.Count) + { + await context.SignalAsync(new ActivityCompleted()); return; + } var nextActivity = childActivities.ElementAt(currentIndex); context.PostActivity(nextActivity, OnChildCompleted); context.UpdateProperty(CurrentIndexProperty, x => x + 1); } - private ValueTask OnChildCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext) + private async ValueTask OnChildCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext) { - HandleItem(context); - return ValueTask.CompletedTask; + await HandleItemAsync(context); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/While.cs b/src/core/Elsa.Core/Activities/While.cs index d46e8b5f8..2e5be5c7a 100644 --- a/src/core/Elsa.Core/Activities/While.cs +++ b/src/core/Elsa.Core/Activities/While.cs @@ -2,16 +2,24 @@ using System.Text.Json.Serialization; using Elsa.Attributes; using Elsa.Contracts; using Elsa.Models; +using Elsa.Signals; namespace Elsa.Activities; [Activity("Elsa", "Primitives", "Execute an activity while a given condition evaluates to true.")] public class While : Activity { + public static While True(IActivity body) => new(body) + { + Condition = new Input(true) + }; + [JsonConstructor] public While(IActivity? body = default) { Body = body!; + + OnSignalReceived(OnBreak); } public While(Input condition, IActivity? body = default) : this(body) @@ -36,7 +44,7 @@ public class While : Activity } [Input] public Input Condition { get; set; } = new(false); - [Outbound] public IActivity Body { get; set; } = default!; + [Outbound] public IActivity Body { get; set; } protected override void Execute(ActivityExecutionContext context) { @@ -53,9 +61,9 @@ public class While : Activity if (loop) context.PostActivity(Body, OnBodyCompleted); } - - public static While True(IActivity body) => new(body) + + private void OnBreak(BreakSignal signal, SignalContext context) { - Condition = new Input(true) - }; + context.StopPropagation(); + } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Contracts/ILoopingConstruct.cs b/src/core/Elsa.Core/Contracts/ILoopingConstruct.cs deleted file mode 100644 index 9519ccb10..000000000 --- a/src/core/Elsa.Core/Contracts/ILoopingConstruct.cs +++ /dev/null @@ -1,6 +0,0 @@ -namespace Elsa.Contracts; - -public interface ILoopingConstruct : IActivity -{ - -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Contracts/ISignalHandler.cs b/src/core/Elsa.Core/Contracts/ISignalHandler.cs new file mode 100644 index 000000000..5adf7df73 --- /dev/null +++ b/src/core/Elsa.Core/Contracts/ISignalHandler.cs @@ -0,0 +1,28 @@ +using Elsa.Models; + +namespace Elsa.Contracts; + +public interface ISignalHandler : IActivity +{ + ValueTask HandleSignalAsync(object signal, SignalContext context); +} + +public class SignalContext +{ + public SignalContext(ActivityExecutionContext activityExecutionContext, ActivityExecutionContext sourceActivityExecutionContext, CancellationToken cancellationToken) + { + ActivityExecutionContext = activityExecutionContext; + SourceActivityExecutionContext = sourceActivityExecutionContext; + CancellationToken = cancellationToken; + } + + public ActivityExecutionContext ActivityExecutionContext { get; init; } + public ActivityExecutionContext SourceActivityExecutionContext { get; init; } + public CancellationToken CancellationToken { get; init; } + internal bool StopPropagationRequested { get; private set; } + + /// + /// Stops the signal from propagating further up the activity execution context hierarchy. + /// + public void StopPropagation() => StopPropagationRequested = true; +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Models/Activity.cs b/src/core/Elsa.Core/Models/Activity.cs index 7dcc02a55..db3d520c5 100644 --- a/src/core/Elsa.Core/Models/Activity.cs +++ b/src/core/Elsa.Core/Models/Activity.cs @@ -1,11 +1,13 @@ using System.Linq.Expressions; using Elsa.Contracts; using Elsa.Helpers; +using Elsa.Signals; namespace Elsa.Models; -public abstract class Activity : IActivity +public abstract class Activity : ISignalHandler { + private readonly ICollection _signalHandlers = new List(); protected Activity() => TypeName = TypeNameHelper.GenerateTypeName(GetType()); protected Activity(string activityType) => TypeName = activityType; @@ -15,17 +17,58 @@ public abstract class Activity : IActivity public IDictionary ApplicationProperties { get; set; } = new Dictionary(); public IDictionary Metadata { get; set; } = new Dictionary(); - protected virtual ValueTask ExecuteAsync(ActivityExecutionContext context) + protected virtual async ValueTask ExecuteAsync(ActivityExecutionContext context) { Execute(context); + await OnExecutedAsync(context); + } + + protected virtual async ValueTask OnExecutedAsync(ActivityExecutionContext context) + { + // By default, signal that the activity is completed. + await context.SignalAsync(new ActivityCompleted()); + } + + protected virtual ValueTask OnSignalReceivedAsync(object signal, SignalContext context) + { + OnSignalReceived(signal, context); return ValueTask.CompletedTask; } + protected virtual void OnSignalReceived(object signal, SignalContext context) + { + } + protected virtual void Execute(ActivityExecutionContext context) { } + protected void OnSignalReceived(Type signalType, Func handler) => _signalHandlers.Add(new SignalHandlerRegistration(signalType, handler)); + protected void OnSignalReceived(Func handler) => OnSignalReceived(typeof(T), (signal, context) => handler((T)signal, context)); + + protected void OnSignalReceived(Action handler) + { + OnSignalReceived((signal, context) => + { + handler(signal, context); + return ValueTask.CompletedTask; + }); + } + ValueTask IActivity.ExecuteAsync(ActivityExecutionContext context) => ExecuteAsync(context); + + async ValueTask ISignalHandler.HandleSignalAsync(object signal, SignalContext context) + { + // Give derived activity a chance to do something with the signal. + await OnSignalReceivedAsync(signal, context); + + // Invoke registered signal delegates for this particular type of signal. + var signalType = signal.GetType(); + var handlers = _signalHandlers.Where(x => x.SignalType == signalType); + + foreach (var registration in handlers) + await registration.Handler(signal, context); + } } public abstract class ActivityWithResult : Activity @@ -63,7 +106,7 @@ public abstract class Activity : ActivityWithResult public static class ActivityWithResultExtensions { - public static T CaptureOutput(this T activity, Expression> propertyExpression, RegisterLocationReference locationReference) where T:IActivity + public static T CaptureOutput(this T activity, Expression> propertyExpression, RegisterLocationReference locationReference) where T : IActivity { var output = activity.GetPropertyValue(propertyExpression)!; output.Targets.Add(locationReference); @@ -71,4 +114,6 @@ public static class ActivityWithResultExtensions } public static T CaptureOutput(this T activity, RegisterLocationReference locationReference) where T : ActivityWithResult => activity.CaptureOutput(x => x.Result, locationReference); -} \ No newline at end of file +} + +internal record SignalHandlerRegistration(Type SignalType, Func Handler); \ No newline at end of file diff --git a/src/core/Elsa.Core/Models/ActivityExecutionContext.cs b/src/core/Elsa.Core/Models/ActivityExecutionContext.cs index 1960c5b8a..fd558cba5 100644 --- a/src/core/Elsa.Core/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Core/Models/ActivityExecutionContext.cs @@ -1,4 +1,6 @@ using System.Collections.ObjectModel; +using System.Reflection; +using Elsa.Activities; using Elsa.Contracts; namespace Elsa.Models; @@ -167,6 +169,27 @@ public class ActivityExecutionContext /// Stops further execution of the workflow. /// public void PreventContinuation() => Continue = false; + + /// + /// Send a signal up the current branch. + /// + public async ValueTask SignalAsync(object signal) + { + var ancestorContexts = GetAncestorActivityExecutionContexts(); + + foreach (var ancestorContext in ancestorContexts) + { + var signalContext = new SignalContext(ancestorContext, this, CancellationToken); + + if (ancestorContext.Activity is not ISignalHandler handler) + continue; + + await handler.HandleSignalAsync(signal, signalContext); + + if (signalContext.StopPropagationRequested) + return; + } + } /// /// Returns a flattened list of the current context's ancestors. @@ -182,7 +205,7 @@ public class ActivityExecutionContext current = current.ParentActivityExecutionContext; } } - + private RegisterLocation? GetLocation(RegisterLocationReference locationReference) => ExpressionExecutionContext.Register.TryGetLocation(locationReference.Id, out var location) ? location diff --git a/src/core/Elsa.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs b/src/core/Elsa.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs index 508157831..602c2779f 100644 --- a/src/core/Elsa.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs +++ b/src/core/Elsa.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs @@ -1,10 +1,7 @@ -using System.Text.Json; using Elsa.Contracts; using Elsa.Models; using Microsoft.Extensions.Logging; using Delegate = System.Delegate; -using JsonNode = System.Text.Json.Nodes.JsonNode; -using JsonObject = System.Text.Json.Nodes.JsonObject; namespace Elsa.Pipelines.ActivityExecution.Components; @@ -62,9 +59,6 @@ public class ActivityInvokerMiddleware : IActivityExecutionMiddleware // Block current path of execution. return; } - - // Complete parent chain. - await CompleteParentsAsync(context); } private void LogExecutionRecord(ActivityExecutionContext context, string eventName, string? message = default, string? source = default, object? payload = default) => context.AddExecutionLogEntry(eventName, message, source, payload); @@ -84,42 +78,4 @@ public class ActivityInvokerMiddleware : IActivityExecutionMiddleware locationReference.Set(context, value); } } - - private async Task CompleteParentsAsync(ActivityExecutionContext context) - { - var workflowExecutionContext = context.WorkflowExecutionContext; - var currentContext = context; - var currentParentContext = context.ParentActivityExecutionContext; - - while (currentParentContext != null) - { - var scheduledNodes = workflowExecutionContext.Scheduler.List().Select(x => x.ActivityId).ToList(); - var descendantNodes = currentParentContext.ActivityNode.Descendants().Select(x => x.Activity.Id).Distinct().ToList(); - var hasScheduledChildren = scheduledNodes.Intersect(descendantNodes).Any(); - var @continue = currentContext.Continue; - - // Do not continue if the activity instructed not to. - if (!@continue) - return; - - if (!hasScheduledChildren) - { - // Invoke completion callbacks. - var completionCallback = workflowExecutionContext.PopCompletionCallback(currentParentContext, currentContext.Activity); - - if (completionCallback != null) - await completionCallback.Invoke(currentParentContext, currentContext); - - // Remove current activity context. - workflowExecutionContext.ActivityExecutionContexts.Remove(currentContext); - } - - // Do not continue completion callbacks of parents while there are scheduled nodes. - if (workflowExecutionContext.Scheduler.HasAny) - return; - - currentContext = currentParentContext; - currentParentContext = currentContext.ParentActivityExecutionContext; - } - } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Signals/ActivityCompleted.cs b/src/core/Elsa.Core/Signals/ActivityCompleted.cs new file mode 100644 index 000000000..42c2157e1 --- /dev/null +++ b/src/core/Elsa.Core/Signals/ActivityCompleted.cs @@ -0,0 +1,3 @@ +namespace Elsa.Signals; + +public record ActivityCompleted; \ No newline at end of file diff --git a/src/core/Elsa.Core/Signals/BreakSignal.cs b/src/core/Elsa.Core/Signals/BreakSignal.cs new file mode 100644 index 000000000..0dd976510 --- /dev/null +++ b/src/core/Elsa.Core/Signals/BreakSignal.cs @@ -0,0 +1,3 @@ +namespace Elsa.Signals; + +public record BreakSignal; \ No newline at end of file diff --git a/src/samples/console/Elsa.Samples.Console1/Program.cs b/src/samples/console/Elsa.Samples.Console1/Program.cs index c74b41fea..bc911f21a 100644 --- a/src/samples/console/Elsa.Samples.Console1/Program.cs +++ b/src/samples/console/Elsa.Samples.Console1/Program.cs @@ -50,7 +50,7 @@ class Program var workflow14 = new Func(FlowchartWorkflow.Create); var workflow15 = new Func(BreakForWorkflow.Create); - var workflowFactory = workflow15; + var workflowFactory = workflow2; var workflowGraph = workflowFactory(); var workflow = Workflow.FromActivity(workflowGraph); diff --git a/test/Elsa.IntegrationTests/Workflows/SequentialWorkflowTests.cs b/test/Elsa.IntegrationTests/Workflows/SequentialWorkflowTests.cs new file mode 100644 index 000000000..0738a8f01 --- /dev/null +++ b/test/Elsa.IntegrationTests/Workflows/SequentialWorkflowTests.cs @@ -0,0 +1,59 @@ +using System.Linq; +using System.Threading.Tasks; +using Elsa.Activities; +using Elsa.Builders; +using Elsa.Contracts; +using Elsa.Models; +using Elsa.Modules.Activities.Activities.Console; +using Elsa.Testing.Shared; +using Microsoft.Extensions.DependencyInjection; +using Xunit; +using Xunit.Abstractions; + +namespace Elsa.IntegrationTests.Workflows; + +public class SequentialWorkflowTests +{ + private readonly IWorkflowRunner _workflowRunner; + private readonly CapturingTextWriter _capturingTextWriter = new(); + private readonly Workflow _workflow; + + public SequentialWorkflowTests(ITestOutputHelper testOutputHelper) + { + var services = new TestApplicationBuilder(testOutputHelper).WithCapturingTextWriter(_capturingTextWriter).Build(); + _workflowRunner = services.GetRequiredService(); + _workflow = new WorkflowDefinitionBuilder().BuildWorkflow(new SequentialWorkflow()); + } + + [Fact(DisplayName = "Sequence completes only after its child activities complete")] + public async Task Test1() + { + await _workflowRunner.RunAsync(_workflow); + var lines = _capturingTextWriter.Lines.ToList(); + Assert.Equal(new[] { "Start", "Line 1", "Line 2", "Line 3", "End" }, lines); + } + + private class SequentialWorkflow : IWorkflow + { + public void Build(IWorkflowDefinitionBuilder workflow) + { + workflow.WithRoot(new Sequence + { + Activities = + { + new WriteLine("Start"), + new Sequence + { + Activities = + { + new WriteLine("Line 1"), + new WriteLine("Line 2"), + new WriteLine("Line 3") + } + }, + new WriteLine("End"), + } + }); + } + } +} \ No newline at end of file