diff --git a/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs b/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs new file mode 100644 index 000000000..3a74edb44 --- /dev/null +++ b/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs @@ -0,0 +1,23 @@ +using Elsa.Workflows; + +namespace Elsa.Testing.Shared; + +public class TestWorkflow : WorkflowBase +{ + private readonly Action _buildWorkflow; + + public TestWorkflow(Action buildWorkflow) + { + _buildWorkflow = buildWorkflow; + } + + protected override void Build(IWorkflowBuilder workflowBuilder) + { + _buildWorkflow(workflowBuilder); + + if (string.IsNullOrEmpty(workflowBuilder.Id)) + { + workflowBuilder.Id = Guid.NewGuid().ToString(); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs index 28b154495..1e93ca546 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs @@ -1,81 +1,60 @@ -using System.Runtime.CompilerServices; -using Elsa.Extensions; -using Elsa.Workflows.Activities.Flowchart.Contracts; -using Elsa.Workflows.Activities.Flowchart.Extensions; -using Elsa.Workflows.Activities.Flowchart.Models; -using Elsa.Workflows.Attributes; -using Elsa.Workflows.Models; -using Elsa.Workflows.UIHints; -using JetBrains.Annotations; - -namespace Elsa.Workflows.Activities.Flowchart.Activities; - -/// -/// Merge multiple branches into a single branch of execution. -/// -[Activity("Elsa", "Branching", "Merge multiple branches into a single branch of execution.", DisplayName = "Join")] -[PublicAPI] -public class FlowJoin : Activity, IJoinNode -{ - /// - public FlowJoin([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) - { - } - - /// - /// The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All). - /// - [Input( - Description = "The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All).", - DefaultValue = FlowJoinMode.WaitAny, - UIHint = InputUIHints.DropDown - )] - public Input Mode { get; set; } = new(FlowJoinMode.WaitAny); - - /// - protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) - { - var flowchartContext = context.ParentActivityExecutionContext!; - var flowchart = (Flowchart)flowchartContext.Activity; - var inboundActivities = flowchart.Connections.LeftInboundActivities(this).ToList(); - var flowScope = flowchartContext.GetProperty(Flowchart.ScopeProperty, () => new FlowScope()); - var executionCount = flowScope.GetExecutionCount(this); - var mode = context.Get(Mode); - - switch (mode) - { - case FlowJoinMode.WaitAll: - { - // If all left-inbound activities have executed, complete & continue. - var haveAllInboundActivitiesExecuted = inboundActivities.All(x => flowScope.GetExecutionCount(x) > executionCount); - - if (haveAllInboundActivitiesExecuted) - { - await CancelActivitiesInInboundPathAsync(flowchart, flowchartContext, context); - await context.CompleteActivityAsync(); - } - - break; - } - case FlowJoinMode.WaitAny: - { - await CancelActivitiesInInboundPathAsync(flowchart, flowchartContext, context); - await context.CompleteActivityAsync(); - break; - } - } - } - - private async Task CancelActivitiesInInboundPathAsync(Flowchart flowchart, ActivityExecutionContext flowchartContext, ActivityExecutionContext joinContext) - { - // Cancel all activities between this join activity and its most recent fork. - var connections = flowchart.Connections; - var workflowExecutionContext = joinContext.WorkflowExecutionContext; - var inboundActivities = connections.LeftAncestorActivities(this).Select(x => workflowExecutionContext.FindNodeByActivity(x)).Select(x => x!.Activity).ToList(); - var inboundActivityExecutionContexts = workflowExecutionContext.ActivityExecutionContexts.Where(x => inboundActivities.Contains(x.Activity) && x.ParentActivityExecutionContext == flowchartContext).ToList(); - - // Cancel each inbound activity. - foreach (var activityExecutionContext in inboundActivityExecutionContexts) - await activityExecutionContext.CancelActivityAsync(); - } +using System.Runtime.CompilerServices; +using Elsa.Extensions; +using Elsa.Workflows.Activities.Flowchart.Contracts; +using Elsa.Workflows.Activities.Flowchart.Extensions; +using Elsa.Workflows.Activities.Flowchart.Models; +using Elsa.Workflows.Attributes; +using Elsa.Workflows.Models; +using Elsa.Workflows.UIHints; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Activities.Flowchart.Activities; + +/// +/// Merge multiple branches into a single branch of execution. +/// +[Activity("Elsa", "Branching", "Merge multiple branches into a single branch of execution.", DisplayName = "Join")] +[PublicAPI] +public class FlowJoin : Activity, IJoinNode +{ + /// + public FlowJoin([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) + { + } + + /// + /// The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All). + /// + [Input( + Description = "The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All).", + DefaultValue = FlowJoinMode.WaitAny, + UIHint = InputUIHints.DropDown + )] + public Input Mode { get; set; } = new(FlowJoinMode.WaitAny); + + /// + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) + { + var mode = context.Get(Mode); + + switch (mode) + { + case FlowJoinMode.WaitAll: + { + if (Flowchart.CanWaitAllProceed(context)) + { + Flowchart.CancelAncestorActivatesAsync(context); + await context.CompleteActivityAsync(); + } + + break; + } + case FlowJoinMode.WaitAny: + { + Flowchart.CancelAncestorActivatesAsync(context); + await context.CompleteActivityAsync(); + break; + } + } + } } \ No newline at end of file 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 3fa0921ea..5fbe321b1 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs @@ -2,12 +2,10 @@ using System.ComponentModel; using System.Runtime.CompilerServices; using Elsa.Extensions; using Elsa.Workflows.Activities.Flowchart.Contracts; -using Elsa.Workflows.Activities.Flowchart.Extensions; using Elsa.Workflows.Activities.Flowchart.Models; using Elsa.Workflows.Attributes; using Elsa.Workflows.Options; using Elsa.Workflows.Signals; -using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Activities.Flowchart.Activities; @@ -18,7 +16,9 @@ namespace Elsa.Workflows.Activities.Flowchart.Activities; [Browsable(false)] public class Flowchart : Container { - public const string ScopeProperty = "Scope"; + private const string ScopeProperty = "FlowScope"; + private const string GraphTransientProperty = "FlowGraph"; + private const string BackwardConnectionActivityInput = "BackwardConnection"; /// public Flowchart([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) @@ -88,7 +88,7 @@ public class Flowchart : Container // If no start activity found, return the first activity. return Activities.FirstOrDefault(); } - + /// /// Checks if there is any pending work for the flowchart. /// @@ -129,97 +129,232 @@ public class Flowchart : Container var rootActivity = query.FirstOrDefault(); return rootActivity; - } - + } + + private FlowGraph GetFlowGraph(ActivityExecutionContext context) + { + // Store in TransientProperties so FlowChart is not persisted in WorkflowState + return context.TransientProperties.GetOrAdd(GraphTransientProperty, () => new FlowGraph(Connections, GetStartActivity(context))); + } + + private FlowScope GetFlowScope(ActivityExecutionContext context) + { + return context.GetProperty(ScopeProperty, () => new FlowScope()); + } + private async ValueTask OnChildCompletedAsync(ActivityCompletedContext context) - { - var logger = context.GetRequiredService>(); - var flowchartContext = context.TargetContext; - var completedActivityContext = context.ChildContext; - var completedActivity = completedActivityContext.Activity; - var result = context.Result; - - // If the complete activity's status is anything but "Completed", do not schedule its outbound activities. - var scheduleChildren = completedActivityContext.Status == ActivityStatus.Completed; - var outcomeNames = result is Outcomes outcomes - ? outcomes.Names - : [null!, "Done"]; - - // Only query the outbound connections if the completed activity wasn't already completed. - var outboundConnections = Connections.Where(connection => connection.Source.Activity == completedActivity && outcomeNames.Contains(connection.Source.Port)).ToList(); - var children = outboundConnections.Select(x => x.Target.Activity).ToList(); - var scope = flowchartContext.GetProperty(ScopeProperty, () => new FlowScope()); - - scope.RegisterActivityExecution(completedActivity); - - // If the complete activity is a terminal node, complete the flowchart immediately. - if (completedActivity is ITerminalNode) + { + var flowchartContext = context.TargetContext; + var completedActivityContext = context.ChildContext; + var completedActivity = completedActivityContext.Activity; + var result = context.Result; + + if (flowchartContext.Activity != this) { - await flowchartContext.CompleteActivityAsync(); + throw new Exception("Target context activity must be this flowchart"); + } + + // If the completed activity's status is anything but "Completed", do not schedule its outbound activities. + if (completedActivityContext.Status != ActivityStatus.Completed) + { + return; + } + + // If the complete activity is a terminal node, complete the flowchart immediately. + if (completedActivity is ITerminalNode) + { + await flowchartContext.CompleteActivityAsync(); + return; + } + + // Determine the outcomes from the completed activity + var outcomes = result is Outcomes o ? o : Outcomes.Default; + + // Schedule the outbound activities + var flowGraph = GetFlowGraph(flowchartContext); + var flowScope = GetFlowScope(flowchartContext); + var completedActivityExcecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault(BackwardConnectionActivityInput); + bool hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExcecutedByBackwardConnection); + + // If there are not any outbound connections, complete the flowchart activity if there is no other pending work + if (!hasScheduledActivity) + { + await CompleteIfNoPendingWorkAsync(flowchartContext); } - else if (scheduleChildren) - { - if (children.Any()) - { - // Schedule each child, but only if all of its left inbound activities have already executed. - foreach (var activity in children) - { - var existingActivity = scope.ContainsActivity(activity); - scope.AddActivity(activity); - - var inboundActivities = Connections.LeftInboundActivities(activity).ToList(); - - // If the completed activity is not part of the left inbound path, always allow its children to be scheduled. - if (!inboundActivities.Contains(completedActivity)) - { - await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync); - continue; - } - - // If the activity is anything but a join activity, only schedule it if all of its left-inbound activities have executed, effectively implementing a "wait all" join. - if (activity is not IJoinNode) - { - var executionCount = scope.GetExecutionCount(activity); - var haveInboundActivitiesExecuted = inboundActivities.All(x => scope.GetExecutionCount(x) > executionCount); - - if (haveInboundActivitiesExecuted) - await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync); - } - else - { - // Select an existing activity execution context for this activity, if any. - var joinContext = flowchartContext.WorkflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => - x.ParentActivityExecutionContext == flowchartContext && x.Activity == activity); - var scheduleWorkOptions = new ScheduleWorkOptions - { - CompletionCallback = OnChildCompletedAsync, - ExistingActivityExecutionContext = joinContext, - PreventDuplicateScheduling = true - }; - - if (joinContext != null) - logger.LogDebug("Next activity {ChildActivityId} is a join activity. Attaching to existing join context {JoinContext}", activity.Id, joinContext.Id); - else if (!existingActivity) - logger.LogDebug("Next activity {ChildActivityId} is a join activity. Creating new join context", activity.Id); - else - { - logger.LogDebug("Next activity {ChildActivityId} is a join activity. Join context was not found, but activity is already being created", activity.Id); - continue; - } - - await flowchartContext.ScheduleActivityAsync(activity, scheduleWorkOptions); - } - } - } - - if (!children.Any()) - { - await CompleteIfNoPendingWorkAsync(flowchartContext); - } - } - - flowchartContext.SetProperty(ScopeProperty, scope); - } + } + + /// + /// Schedules outbound activities based on the flowchart's structure and execution state. + /// This method determines whether an activity should be scheduled based on visited connections, + /// forward traversal rules, and backward connections. + /// + /// The flowchart containing the activities. + /// The graph representation of the flowchart. + /// Tracks activity and connection visits. + /// The execution context of the flowchart. + /// The current activity being processed. + /// The outcomes that determine which connections were followed. + /// Indicates if the completed activity was executed due to a backward connection. + /// True if at least one activity was scheduled; otherwise, false. + private async ValueTask ScheduleOutboundActivitiesAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity activity, Outcomes outcomes, bool completedActivityExecutedByBackwardConnection = false) + { + bool hasScheduledActivity = false; + + // Check if the activity is dangling (i.e., it is not reachable from the flowchart graph) + if (flowGraph.IsDanglingActivity(activity)) + { + throw new Exception($"Activity {activity.Id} is not reachable from the flowchart graph. Unable to schedule it's outbound activities."); + } + + // Register the activity as visited unless it was executed due to a backward connection + if (!completedActivityExecutedByBackwardConnection) + { + flowScope.RegisterActivityVisit(activity); + } + + // Process each outbound connection from the current activity + foreach (var outboundConnection in flowGraph.GetOutboundConnections(activity)) + { + bool connectionFollowed = outcomes.Names.Contains(outboundConnection.Source.Port); + flowScope.RegisterConnectionVisit(outboundConnection, connectionFollowed); + var outboundActivity = outboundConnection.Target.Activity; + + // Determine scheduling strategy based on connection type + if (flowGraph.IsBackwardConnection(outboundConnection, out bool backwardConnectionIsValid)) + { + hasScheduledActivity |= await ScheduleBackwardConnectionActivityAsync(flowGraph, flowchartContext, outboundConnection, outboundActivity, connectionFollowed, backwardConnectionIsValid); + } + else if (outboundActivity is not IJoinNode) + { + hasScheduledActivity |= await ScheduleNonJoinActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity); + } + else + { + hasScheduledActivity |= await ScheduleJoinActivityAsync(flowGraph, flowScope, flowchartContext, outboundConnection, outboundActivity); + } + } + return hasScheduledActivity; + } + + /// + /// Schedules an outbound activity that originates from a backward connection. + /// + private async ValueTask ScheduleBackwardConnectionActivityAsync(FlowGraph flowGraph, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, bool connectionFollowed, bool backwardConnectionIsValid) + { + if (!connectionFollowed) + { + return false; + } + + if (!backwardConnectionIsValid) + { + throw new Exception($"Invalid backward connection: Every path from the source ('{outboundConnection.Source.Activity.Id}') must go through the target ('{outboundConnection.Target.Activity.Id}') when tracing back to the start."); + } + + var scheduleWorkOptions = new ScheduleWorkOptions + { + CompletionCallback = OnChildCompletedAsync, + Input = new Dictionary() { { BackwardConnectionActivityInput, true } } + }; + + await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions); + return true; + } + + /// + /// Schedules a non-join activity if all its forward inbound connections have been visited. + /// + private async ValueTask ScheduleNonJoinActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity) + { + if (!flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) + { + return false; + } + + if (flowScope.HasFollowedInboundConnection(flowGraph, outboundActivity)) + { + await flowchartContext.ScheduleActivityAsync(outboundActivity, OnChildCompletedAsync); + return true; + } + else + { + // Propagate skipped connections by scheduling with Outcomes.Empty + return await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty); + } + } + + /// + /// Schedules a join activity based on inbound connection statuses. + /// + private async ValueTask ScheduleJoinActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity) + { + // Ignore the connection if the join activity has already completed (JoinAny scenario) + if (flowScope.ShouldIgnoreConnection(outboundConnection, outboundActivity)) + { + return false; + } + + // Schedule the join activity only if at least one inbound connection was followed + if (!flowScope.HasFollowedInboundConnection(flowGraph, outboundActivity)) + { + if (flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) + { + // Propagate skipped connections by scheduling with Outcomes.Empty + return await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty); + } + return false; + } + + // Check for an existing execution context for the join activity + var joinContext = flowchartContext.WorkflowExecutionContext.ActivityExecutionContexts.LastOrDefault(x => + x.ParentActivityExecutionContext == flowchartContext && + x.Activity == outboundActivity && + x.Status is ActivityStatus.Pending or ActivityStatus.Running); + + // If the join activity was already scheduled, do not schedule it again + if (joinContext == null) + { + var activityScheduled = flowchartContext.WorkflowExecutionContext.Scheduler.List().Any(workItem => workItem.Owner == flowchartContext && workItem.Activity == outboundActivity); + if (activityScheduled) + { + return true; + } + } + + var scheduleWorkOptions = new ScheduleWorkOptions + { + CompletionCallback = OnChildCompletedAsync, + ExistingActivityExecutionContext = joinContext + }; + await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions); + return true; + } + + public static bool CanWaitAllProceed(ActivityExecutionContext context) + { + var flowchartContext = context.ParentActivityExecutionContext!; + var flowchart = (Flowchart)flowchartContext.Activity; + var flowGraph = flowchart.GetFlowGraph(flowchartContext); + var flowScope = flowchart.GetFlowScope(flowchartContext); + var activity = context.Activity; + + return flowScope.AllInboundConnectionsVisited(flowGraph, activity); + } + + public static async void CancelAncestorActivatesAsync(ActivityExecutionContext context) + { + var flowchartContext = context.ParentActivityExecutionContext!; + var flowchart = (Flowchart)flowchartContext.Activity; + var flowGraph = flowchart.GetFlowGraph(flowchartContext); + var ancestorActivities = flowGraph.GetAncestorActivities(context.Activity); + var inboundActivityExecutionContexts = context.WorkflowExecutionContext.ActivityExecutionContexts.Where(x => ancestorActivities.Contains(x.Activity) && x.ParentActivityExecutionContext == flowchartContext).ToList(); + + // Cancel each ancestor activity. + foreach (var activityExecutionContext in inboundActivityExecutionContexts) + { + await activityExecutionContext.CancelActivityAsync(); + } + } private async Task CompleteIfNoPendingWorkAsync(ActivityExecutionContext context) { diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Extensions/ConnectionsExtensions.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Extensions/ConnectionsExtensions.cs index cd653a2ba..46ff61b7f 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Extensions/ConnectionsExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Extensions/ConnectionsExtensions.cs @@ -1,107 +1,39 @@ -using Elsa.Workflows.Activities.Flowchart.Models; - -namespace Elsa.Workflows.Activities.Flowchart.Extensions; - -/// -/// Contains extension methods for . -/// -public static class ConnectionsExtensions -{ - /// - /// Returns all connections that are descendants of the specified parent activity. - /// - public static IEnumerable Descendants(this ICollection connections, IActivity parent) - { - var visitedConnections = new HashSet(); - return connections.Descendants(parent, visitedConnections); - } - - /// - /// Returns all ancestor connections of the specified parent activity. - /// - public static IEnumerable Ancestors(this ICollection connections, IActivity activity) - { - var visitedActivities = new HashSet(); - return connections.Ancestors(activity, visitedActivities); - } - - /// - /// Returns all inbound connections of the specified activity. - /// - public static IEnumerable InboundConnections(this ICollection connections, IActivity activity) => connections.Where(x => x.Target.Activity == activity).ToList(); - - /// - /// Returns all "left" inbound connections of the specified activity. "Left" means "not a descendant of the activity". - /// - public static IEnumerable LeftInboundConnections(this ICollection connections, IActivity activity) - { - // We only take "left" inbound connections, which means we exclude descendent connections looping back. - var descendantConnections = connections.Descendants(activity).ToList(); - var filteredConnections = connections.InboundConnections(activity).Except(descendantConnections).ToList(); - - return filteredConnections; - } - - /// - /// Returns all "left" ancestor connections of the specified activity. "Left" means "not a descendant of the activity". - /// - public static IEnumerable LeftAncestorConnections(this ICollection connections, IActivity activity) - { - // We only take "left" inbound connections, which means we exclude descendent connections looping back. - var descendantConnections = connections.Descendants(activity).ToList(); - var filteredConnections = connections.Ancestors(activity).Except(descendantConnections).ToList(); - - return filteredConnections; - } - - /// - /// Returns all inbound activities of the specified activity. - /// - public static IEnumerable InboundActivities(this ICollection connections, IActivity activity) => connections.InboundConnections(activity).Select(x => x.Source.Activity); - - /// - /// Returns all "left" inbound activities of the specified activity. "Left" means "not a descendant of the activity". - /// - public static IEnumerable LeftInboundActivities(this ICollection connections, IActivity activity) => connections.LeftInboundConnections(activity).Select(x => x.Source.Activity); - - /// - /// Returns all "left" ancestor activities of the specified activity. "Left" means "not a descendant of the activity". - /// - public static IEnumerable LeftAncestorActivities(this ICollection connections, IActivity activity) => connections.LeftAncestorConnections(activity).Select(x => x.Source.Activity); - - private static IEnumerable Descendants(this ICollection connections, IActivity parent, ISet visitedConnections) - { - var children = connections.Where(x => parent == x.Source.Activity && !visitedConnections.Contains(x)).ToList(); - - foreach (var child in children) - { - visitedConnections.Add(child); - yield return child; - - var descendants = connections.Descendants(child.Target.Activity, visitedConnections).ToList(); - - foreach (var descendant in descendants) - { - yield return descendant; - } - } - } - - private static IEnumerable Ancestors(this ICollection connections, IActivity activity, ISet visitedActivities) - { - var parents = connections.Where(x => activity == x.Target.Activity && !visitedActivities.Contains(x.Source.Activity)).ToList(); - - foreach (var parent in parents) - { - visitedActivities.Add(parent.Source.Activity); - yield return parent; - - var ancestors = connections.Ancestors(parent.Source.Activity, visitedActivities).ToList(); - - foreach (var ancestor in ancestors) - { - yield return ancestor; - } - } - } +using Elsa.Workflows.Activities.Flowchart.Models; + +namespace Elsa.Workflows.Activities.Flowchart.Extensions; + +/// +/// Contains extension methods for . +/// +public static class ConnectionsExtensions +{ + /// + /// Returns all inbound connections of the specified activity. + /// + public static IEnumerable InboundConnections(this ICollection connections, IActivity activity) => connections.Where(x => x.Target.Activity == activity).Distinct().ToList(); + + /// + /// Returns all inbound activities of the specified activity. + /// + public static IEnumerable InboundActivities(this ICollection connections, IActivity activity) => connections.InboundConnections(activity).Select(x => x.Source.Activity); + + /// + /// Returns all outbound connections of the specified activity. + /// + public static IEnumerable OutboundConnections(this ICollection connections, IActivity activity) => connections.Where(x => x.Source.Activity == activity).Distinct().ToList(); + + /// + /// Returns all outbound connections of the specified activity matching the specified outcomes. + /// + public static IEnumerable OutboundConnections(this ICollection connections, IActivity activity, Outcomes outcomes) => connections.OutboundConnections(activity).Where(c => outcomes.Names.Contains(c.Source.Port)); + + /// + /// Returns all outbound activities of the specified activity. + /// + public static IEnumerable OutboundActivities(this ICollection connections, IActivity activity) => connections.OutboundConnections(activity).Select(x => x.Source.Activity); + + /// + /// Returns all outbound activities of the specified activity matching the specified outcomes + /// + public static IEnumerable OutboundActivities(this ICollection connections, IActivity activity, Outcomes outcomes) => connections.OutboundConnections(activity).Where(c => outcomes.Names.Contains(c.Source.Port)).Select(x => x.Source.Activity); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/ActivityFlowState.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/ActivityFlowState.cs deleted file mode 100644 index 3d5569cec..000000000 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/ActivityFlowState.cs +++ /dev/null @@ -1,22 +0,0 @@ -using System.Diagnostics; -using System.Text.Json.Serialization; - -namespace Elsa.Workflows.Activities.Flowchart.Models; - -[DebuggerDisplay("ActivityId = {ActivityId}, ExecutionCount = {ExecutionCount}")] -public class ActivityFlowState -{ - [JsonConstructor] - public ActivityFlowState() - { - } - - public ActivityFlowState(string activityId, long executionCount = 0) - { - ActivityId = activityId; - ExecutionCount = executionCount; - } - - public string ActivityId { get; set; } = default!; - public long ExecutionCount { get; set; } -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Connection.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Connection.cs index ba22c4f15..698e3bf3d 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Connection.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Connection.cs @@ -5,7 +5,7 @@ namespace Elsa.Workflows.Activities.Flowchart.Models; /// /// A connection between a source and a target endpoint. /// -public class Connection +public class Connection : IEquatable { /// /// Initializes a new instance of the class. @@ -45,5 +45,36 @@ public class Connection /// /// The target endpoint. /// - public Endpoint Target { get; set; } = default!; + public Endpoint Target { get; set; } = default!; + + public override string ToString() => + $"{Source.Activity.Id}{(string.IsNullOrEmpty(Source.Port) ? "" : $":{Source.Port}")}->" + + $"{Target.Activity.Id}{(string.IsNullOrEmpty(Target.Port) ? "" : $":{Target.Port}")}"; + + // Implement equality logic + public bool Equals(Connection? other) + { + if (other == null) return false; + return AreEndpointsEqual(Source, other.Source) && AreEndpointsEqual(Target, other.Target); + } + + public override bool Equals(object? obj) + { + return obj is Connection other && Equals(other); + } + + public override int GetHashCode() + { + return HashCode.Combine(GetEndpointHashCode(Source), GetEndpointHashCode(Target)); + } + + private static bool AreEndpointsEqual(Endpoint e1, Endpoint e2) + { + return e1.Activity.Equals(e2.Activity) && e1.Port == e2.Port; + } + + private static int GetEndpointHashCode(Endpoint endpoint) + { + return HashCode.Combine(endpoint.Activity.GetHashCode(), endpoint.Port?.GetHashCode() ?? 0); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowGraph.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowGraph.cs new file mode 100644 index 000000000..a73e0ca32 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowGraph.cs @@ -0,0 +1,242 @@ +using Elsa.Extensions; +using Elsa.Workflows.Activities.Flowchart.Extensions; + +namespace Elsa.Workflows.Activities.Flowchart.Models; + +/// +/// Represents a directed graph structure for managing workflow connections. +/// Caches forward and backward connections to optimize graph traversal. +/// +public class FlowGraph(ICollection Connections, IActivity? RootActivity) +{ + private List? _cachedForwardConnections; + private readonly Dictionary> _cachedInboundForwardConnections = new(); + private readonly Dictionary> _cachedOutboundConnections = new(); + private readonly Dictionary _cachedIsBackwardConnection = new(); + private readonly Dictionary _cachedIsDanglingActivity = new(); + private readonly Dictionary> _cachedAncestors = new(); + + /// + /// Gets the list of forward connections, computing them if not already cached. + /// + private List ForwardConnections => _cachedForwardConnections ??= RootActivity == null ? new() : GetForwardConnections(Connections, RootActivity); + + /// + /// Retrieves all inbound forward connections for a given activity. + /// + public List GetForwardInboundConnections(IActivity activity) => _cachedInboundForwardConnections.GetOrAdd(activity, () => ForwardConnections.InboundConnections(activity).ToList()); + + /// + /// Retrieves all outbound connections for a given activity. + /// + public List GetOutboundConnections(IActivity activity) => _cachedOutboundConnections.GetOrAdd(activity, () => Connections.OutboundConnections(activity).ToList()); + + /// + /// Determines if a given activity is "dangling," meaning it does not exist as a target in any forward connection. + /// + public bool IsDanglingActivity(IActivity activity) => _cachedIsDanglingActivity.GetOrAdd(activity, () => activity != RootActivity && !ForwardConnections.Any(c => c.Target.Activity == activity)); + + /// + /// Determines if a given connection is a backward connection (i.e., not part of the forward traversal) and whether it is valid. + /// + public bool IsBackwardConnection(Connection connection, out bool isValid) + { + // Check if result is already cached + if (_cachedIsBackwardConnection.TryGetValue(connection, out var result)) + { + isValid = result.IsValid; + return result.IsBackwardConnection; + } + + // Compute if the connection is backward + bool isBackwardConnection = !GetForwardInboundConnections(connection.Target.Activity).Contains(connection); + + // Compute if the backward connection is valid + isValid = isBackwardConnection ? IsValidBackwardConnection(ForwardConnections, RootActivity, connection) : false; + + // Cache the result + _cachedIsBackwardConnection[connection] = (isBackwardConnection, isValid); + + return isBackwardConnection; + } + + /// + /// Retrieves all ancestor activities for a given activity by traversing ForwardConnections in reverse. + /// + public List GetAncestorActivities(IActivity activity) + { + return _cachedAncestors.GetOrAdd(activity, () => ComputeAncestors(activity)); + } + + /// + /// Computes the list of ancestors by following Source activities in ForwardConnections. + /// + private List ComputeAncestors(IActivity activity) + { + HashSet ancestors = new(); + Queue queue = new(); + + // Find all connections where this activity is the target + foreach (var connection in ForwardConnections.Where(c => c.Target.Activity == activity)) + { + if (ancestors.Add(connection.Source.Activity)) + queue.Enqueue(connection.Source.Activity); + } + + // Traverse upwards through the graph + while (queue.Count > 0) + { + var current = queue.Dequeue(); + + foreach (var connection in ForwardConnections.Where(c => c.Target.Activity == current)) + { + if (ancestors.Add(connection.Source.Activity)) + queue.Enqueue(connection.Source.Activity); + } + } + + return ancestors.ToList(); + } + + /// + /// Computes the list of forward connections in the graph, excluding cyclic connections. + /// + private static List GetForwardConnections(ICollection connections, IActivity root) + { + Dictionary> adjList = new(); + + foreach (var conn in connections) + { + if (!adjList.ContainsKey(conn.Source.Activity)) + adjList[conn.Source.Activity] = new List(); + + adjList[conn.Source.Activity].Add(conn.Target.Activity); + } + + HashSet visited = new(); + HashSet<(IActivity, IActivity)> visitedEdges = new(); + List<(IActivity Source, IActivity Target)> validEdges = new(); + Queue queue = new(); + + queue.Enqueue(root); + + while (queue.Count > 0) + { + var source = queue.Dequeue(); + visited.Add(source); + + if (!adjList.ContainsKey(source)) continue; + + foreach (var target in adjList[source]) + { + var edge = (source, target); + if (visitedEdges.Contains(edge)) + continue; + + if (HasPathToActivity(validEdges, target, source)) + continue; + + visitedEdges.Add(edge); + validEdges.Add((source, target)); + + if (!visited.Contains(target)) + queue.Enqueue(target); + } + } + + return validEdges + .SelectMany(e => connections.Where(c => c.Source.Activity == e.Source && c.Target.Activity == e.Target)) + .Distinct() + .ToList(); + } + + /// + /// Determines if there is an existing path from the source activity to the target activity. + /// Helps in detecting cyclic connections. + /// + private static bool HasPathToActivity(ICollection<(IActivity Source, IActivity Target)> edges, IActivity source, IActivity target) + { + if (source == target) + return true; + + HashSet visited = new(); + Stack stack = new(); + stack.Push(source); + + while (stack.Count > 0) + { + var current = stack.Pop(); + + if (current == target) + return true; + + if (visited.Contains(current)) + continue; + + visited.Add(current); + + foreach (var next in edges.Where(x => x.Source == current).Select(e => e.Target)) + { + if (!visited.Contains(next)) + stack.Push(next); + } + } + + return false; + } + + /// + /// Determines whether a backward connection is valid by ensuring all paths from source to root pass through target. + /// + private static bool IsValidBackwardConnection(List forwardConnections, IActivity? root, Connection connection) + { + if (root == null) return false; + + var pathsToRoot = GetPathsToRoot(forwardConnections, root, connection.Source.Activity); + + foreach (var path in pathsToRoot) + { + if (!path.Contains(connection.Target.Activity)) + return false; + } + + return true; + } + + /// + /// Finds all paths from a given start activity to the root using BFS. + /// + private static List> GetPathsToRoot(List forwardConnections, IActivity root, IActivity start) + { + List> paths = new(); + Queue> queue = new(); + queue.Enqueue(new List { start }); + + while (queue.Count > 0) + { + var path = queue.Dequeue(); + var lastNode = path.Last(); + + if (lastNode == root) + { + paths.Add(new List(path)); + continue; + } + + var previousNodes = forwardConnections + .Where(c => c.Target.Activity == lastNode) + .Select(c => c.Source.Activity); + + foreach (var prev in previousNodes) + { + if (!path.Contains(prev)) + { + var newPath = new List(path) { prev }; + queue.Enqueue(newPath); + } + } + } + + return paths; + } +} diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs index 1e55ae2e0..d0f74dd42 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs @@ -1,89 +1,107 @@ -using System.Text.Json.Serialization; - -namespace Elsa.Workflows.Activities.Flowchart.Models; - -public class FlowScope -{ - [JsonConstructor] - public FlowScope() - { - } - - public FlowScope(string ownerActivityId) - { - OwnerActivityId = ownerActivityId; - Activities.Add(ownerActivityId, new ActivityFlowState(ownerActivityId, 1)); - } - - /// - /// The activity from which the scope was created. - /// - public string OwnerActivityId { get; set; } = default!; - - /// - /// A list of scheduled activity IDs and a flag whether they executed or not. - /// - public IDictionary Activities { get; set; } = - new Dictionary(); - - public void AddActivities(IEnumerable activities, long executionCount = 0) - { - foreach (var activity in activities) - AddActivity(activity, executionCount); - } - - public void AddActivity(IActivity activity, long executionCount = 0) => EnsureActivity(activity, executionCount); - - public ActivityFlowState EnsureActivity(IActivity activity, long executionCount = 0) - { - if (Activities.ContainsKey(activity.Id)) - return Activities[activity.Id]; - - var state = new ActivityFlowState(activity.Id, executionCount); - Activities.Add(activity.Id, state); - return state; - - } - - public bool ContainsActivity(IActivity activity) - { - return Activities.ContainsKey(activity.Id); - } - - public void RegisterActivityExecution(IActivity activity) - { - var state = Activities.TryGetValue(activity.Id, out var s) ? s : default; - - if (state == null) - { - state = new ActivityFlowState(activity.Id); - Activities[activity.Id] = state; - } - - state.ExecutionCount++; - } - - /// - /// Return a list excluding any activities that already executed. - /// - public IEnumerable ExcludeExecutedActivities(IEnumerable activities) => - activities.Where(x => !Activities.ContainsKey(x.Id) || Activities[x.Id].ExecutionCount == 0); - - public bool HasPendingActivities() - { - var sample = Activities.Values.First().ExecutionCount; - return Activities.Values.Any(x => x.ExecutionCount != sample); - } - - public long GetExecutionCount(IActivity activity) => Activities.ContainsKey(activity.Id) ? Activities[activity.Id].ExecutionCount : 0; - - public void Clear() => Activities.Clear(); - - public void Remove(IActivity activity) - { - if (Activities.ContainsKey(activity.Id)) - { - Activities.Remove(activity.Id); - } - } -} \ No newline at end of file +using System.Text.Json.Serialization; + +namespace Elsa.Workflows.Activities.Flowchart.Models; + +/// +/// Represents a scope for tracking activity and connection visits within a flowchart execution. +/// +public class FlowScope +{ + [JsonConstructor] + public FlowScope() + { + } + + [JsonInclude] + private Dictionary ActivitiesVisitCount { get; init; } = new(); + + [JsonInclude] + private Dictionary ConnectionVisitCount { get; init; } = new(); + + [JsonInclude] + private Dictionary ConnectionLastVisitFollowed { get; init; } = new(); + + /// + /// Registers a visit to the specified activity, incrementing its visit count. + /// + /// The activity being visited. + public void RegisterActivityVisit(IActivity activity) + { + string activityId = activity.Id; + ActivitiesVisitCount.TryAdd(activityId, 0); + ActivitiesVisitCount[activityId]++; + } + + /// + /// Gets the number of times the specified activity has been visited. + /// + /// The activity to check. + /// The visit count of the activity. + private long GetActivityVisitCount(IActivity activity) => ActivitiesVisitCount.TryGetValue(activity.Id, out var count) ? count : 0; + + /// + /// Registers a visit to the specified connection and records whether it was followed. + /// + /// The connection being visited. + /// Indicates whether the connection was followed. + public void RegisterConnectionVisit(Connection connection, bool followed) + { + string connectionId = connection.ToString(); + ConnectionVisitCount.TryAdd(connectionId, 0); + ConnectionVisitCount[connectionId]++; + ConnectionLastVisitFollowed[connectionId] = followed; + } + + /// + /// Gets the number of times the specified connection has been visited. + /// + /// The connection to check. + /// The visit count of the connection. + private long GetConnectionVisitCount(Connection connection) => ConnectionVisitCount.TryGetValue(connection.ToString(), out var count) ? count : 0; + + /// + /// Determines whether the last visit to the specified connection was followed. + /// + /// The connection to check. + /// True if the connection was followed on the last visit, otherwise false. + private bool GetConnectionLastVisitFollowed(Connection connection) => ConnectionLastVisitFollowed.TryGetValue(connection.ToString(), out var followed) ? followed : false; + + /// + /// Determines whether all inbound connections to the specified activity have been visited. + /// + /// The flow graph containing connections. + /// The activity to check. + /// True if all inbound connections have been visited, otherwise false. + public bool AllInboundConnectionsVisited(FlowGraph flowGraph, IActivity activity) + { + var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity); + var outboundActivityVisitCount = GetActivityVisitCount(activity); + var minConnectionVisitCount = forwardInboundConnections.Min(c => GetConnectionVisitCount(c)); + return minConnectionVisitCount > outboundActivityVisitCount; + } + + /// + /// Determines whether any inbound connection to the specified activity has been followed. + /// + /// The flow graph containing connections. + /// The activity to check. + /// True if any inbound connection has been followed, otherwise false. + public bool HasFollowedInboundConnection(FlowGraph flowGraph, IActivity activity) + { + var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity); + return forwardInboundConnections.Any(c => GetConnectionLastVisitFollowed(c)); + } + + /// + /// Determines whether a connection should be ignored based on visit counts. + /// + /// The connection to check. + /// The activity associated with the connection. + /// True if the connection should be ignored, otherwise false. + public bool ShouldIgnoreConnection(Connection connection, IActivity activity) + { + var connectionVisitCount = GetConnectionVisitCount(connection); + var activityVisitCount = GetActivityVisitCount(activity); + return connectionVisitCount <= activityVisitCount; + } +} diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs index f036f0da0..9cfa3cfde 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs @@ -4,4 +4,8 @@ namespace Elsa.Workflows.Activities.Flowchart.Models; /// Represents a list of outcomes that can be send when completing an activity. This information is used by . /// /// A list of outcome names. -public record Outcomes(params string[] Names); \ No newline at end of file +public record Outcomes(params string[] Names) +{ + public static readonly Outcomes Default = new([null!, "Done"]); + public static readonly Outcomes Empty = new(); +} diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/FlowJoins/Workflows.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/FlowJoins/Workflows.cs index 198d1b3e5..a3750bb0d 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/FlowJoins/Workflows.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/FlowJoins/Workflows.cs @@ -6,7 +6,7 @@ public class SingleJoinWorkflow : WorkflowBase { protected override void Build(IWorkflowBuilder builder) { - builder.Root = new Flowchart + builder.Root = new Elsa.Workflows.Activities.Flowchart.Activities.Flowchart { Activities = { diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs index 7276c344f..bfad64033 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs @@ -1,5 +1,10 @@ +using Elsa.Expressions.Models; using Elsa.Testing.Shared; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Activities.Flowchart.Activities; +using Elsa.Workflows.Activities.Flowchart.Models; using Elsa.Workflows.IntegrationTests.Scenarios.FlowchartNextActivity.Workflows; +using Elsa.Workflows.Memory; using Microsoft.Extensions.DependencyInjection; using Xunit.Abstractions; @@ -28,5 +33,214 @@ public class FlowchartNextActivityTests await _workflowRunner.RunAsync(); var lines = _capturingTextWriter.Lines.ToList(); Assert.Equal(new[] { "Line 1" }, lines); + } + + [Fact(DisplayName = "Flowchart with backward connections and a dangling activity")] + public async Task BackwardConnectionTest() + { + var workflow = new TestWorkflow(workflowBuilder => + { + var loopVariable = new Variable("LoopCount", 0); + + var start = new Start(); + var dangling = new WriteLine("dangling"); + var writeLineDecision = new FlowSwitch() + { + Cases = { + new FlowSwitchCase("LessThanThree", new Expression("JavaScript", "getVariable('LoopCount') < 3")), + new FlowSwitchCase("LessThanOne", new Expression("JavaScript", "getVariable('LoopCount') < 1")), + }, + Mode = new(SwitchMode.MatchAny) + }; + var a = new WriteLine("A"); + var b = new WriteLine("B"); + var c = new WriteLine("C"); + var incrementLoop = new SetVariable() + { + Variable = loopVariable, + Value = new Models.Input(new Expression("JavaScript", "getVariable('LoopCount') + 1")) + }; + var loopbackDecision = new FlowSwitch() + { + Cases = { + new FlowSwitchCase("EqualOne", new Expression("JavaScript", "getVariable('LoopCount') == 1")), + new FlowSwitchCase("LessThanFour", new Expression("JavaScript", "getVariable('LoopCount') < 4")), + }, + Mode = new(SwitchMode.MatchFirst) + }; + var d = new WriteLine("D"); + var e = new WriteLine("E"); + var f = new WriteLine("F"); + var end = new End(); + + + workflowBuilder.Root = new Flowchart + { + Variables = + { + loopVariable + }, + Activities = + { + start, + dangling, + writeLineDecision, + a, + b, + c, + incrementLoop, + loopbackDecision, + d, + e, + f, + end + }, + Connections = + { + new(start, writeLineDecision), + new(dangling, writeLineDecision), + new(new Endpoint(writeLineDecision, "LessThanThree"), new Endpoint(a)), + new(new Endpoint(writeLineDecision, "LessThanThree"), new Endpoint(b)), + new(new Endpoint(writeLineDecision, "LessThanOne"), new Endpoint(c)), + new(new Endpoint(writeLineDecision, "Default"), new Endpoint(incrementLoop)), + new(a, incrementLoop), + new(b, incrementLoop), + new(c, incrementLoop), + new(incrementLoop, loopbackDecision), + new(new Endpoint(loopbackDecision, "EqualOne"), new Endpoint(d)), + new(d, incrementLoop), + new(new Endpoint(loopbackDecision, "LessThanFour"), new Endpoint(e)), + new(e, writeLineDecision), + new(new Endpoint(loopbackDecision, "Default"), new Endpoint(f)), + new(f, end), + } + }; + }); + + await _services.PopulateRegistriesAsync(); + var result = await _workflowRunner.RunAsync(workflow); + var lines = _capturingTextWriter.Lines.ToList(); + Assert.Equal(WorkflowSubStatus.Finished, result.WorkflowState.SubStatus); + Assert.Equal(new[] { "A", "B", "C", "D", "E", "A", "B", "E", "F" }, lines); + } + + [Fact(DisplayName = "Flowchart with an invalid backward connection")] + public async Task InvalidBackwardConnectionTest() + { + var workflow = new TestWorkflow(workflowBuilder => + { + var start = new Start() { Id = "Start" }; + var a = new WriteLine("A") { Id = "WriteLineA" }; + var b = new WriteLine("B") { Id = "WriteLineB" }; + var c = new WriteLine("C") { Id = "WriteLineC" }; + var d = new WriteLine("D") { Id = "WriteLineD" }; + var e = new WriteLine("E") { Id = "WriteLineE" }; + + workflowBuilder.Root = new Flowchart + { + Activities = + { + start, + a, + b, + c, + d, + e, + }, + Connections = + { + new(start, a), + new(a, b), + new(b, c), + new(b, d), + new(c, e), + new(d, e), + new(e, c), + } + }; + }); + + await _services.PopulateRegistriesAsync(); + var result = await _workflowRunner.RunAsync(workflow); + var lines = _capturingTextWriter.Lines.ToList(); + Assert.Equal(WorkflowSubStatus.Faulted, result.WorkflowState.SubStatus); + Assert.Equal(1, result.WorkflowState.Incidents.Count()); + Assert.Equal("Invalid backward connection: Every path from the source ('WriteLineE') must go through the target ('WriteLineC') when tracing back to the start.", result.WorkflowState.Incidents.First().Message); + Assert.Equal(new[] { "A", "B", "C", "D", "E" }, lines); + } + + [Theory(DisplayName = "Flowchart with a Join activity executed multiple times")] + [InlineData(FlowJoinMode.WaitAll)] + [InlineData(FlowJoinMode.WaitAny)] + public async Task WaitAnyLoopTest(FlowJoinMode joinMode) + { + var workflow = new TestWorkflow(workflowBuilder => + { + var loopVariable = new Variable("LoopCount", 0); + + var start = new Start(); + var a = new WriteLine("A"); + var b = new WriteLine("B"); + var c = new WriteLine("C"); + var d = new WriteLine("D"); + var join = new FlowJoin() + { + Mode = new(joinMode) + }; + var incrementLoop = new SetVariable() + { + Variable = loopVariable, + Value = new Models.Input(new Expression("JavaScript", "getVariable('LoopCount') + 1")) + }; + var loopbackDecision = new FlowSwitch() + { + Cases = { + new FlowSwitchCase("LessThanThree", new Expression("JavaScript", "getVariable('LoopCount') < 3")), + }, + Mode = new(SwitchMode.MatchFirst) + }; + var end = new End(); + + + workflowBuilder.Root = new Flowchart + { + Variables = + { + loopVariable + }, + Activities = + { + start, + a, + b, + c, + d, + join, + incrementLoop, + loopbackDecision, + end + }, + Connections = + { + new(start, a), + new(a, b), + new(a, c), + new(a, d), + new(b, join), + new(c, join), + new(d, join), + new(join, incrementLoop), + new(incrementLoop,loopbackDecision), + new(new Endpoint(loopbackDecision, "LessThanThree"), new Endpoint(a)), + new(new Endpoint(loopbackDecision, "Default"), new Endpoint(end)), + } + }; + }); + + await _services.PopulateRegistriesAsync(); + var result = await _workflowRunner.RunAsync(workflow); + var lines = _capturingTextWriter.Lines.ToList(); + Assert.Equal(WorkflowSubStatus.Finished, result.WorkflowState.SubStatus); + Assert.Equal(new[] { "A", "B", "C", "D", "A", "B", "C", "D", "A", "B", "C", "D"}, lines); } } \ No newline at end of file diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs index c89a2db5c..b25390c8a 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs @@ -22,7 +22,7 @@ public class JsonObjectJintTests public async Task Test1() { await _services.PopulateRegistriesAsync(); - await _workflowRunner.RunAsync(); + await _workflowRunner.RunAsync(); var lines = _capturingTextWriter.Lines.ToList(); Assert.Equal(new[] { "Baz" }, lines); } diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/Flowchart/FlowGraphExtensions.cs b/test/unit/Elsa.Workflows.Core.UnitTests/Flowchart/FlowGraphExtensions.cs new file mode 100644 index 000000000..fe431998c --- /dev/null +++ b/test/unit/Elsa.Workflows.Core.UnitTests/Flowchart/FlowGraphExtensions.cs @@ -0,0 +1,32 @@ +using Elsa.Workflows.Activities.Flowchart.Models; + +namespace Elsa.Workflows.Core.UnitTests.Flowchart; + +public static class FlowGraphExtensions +{ + public static void ValidateOutboundConnections(this FlowGraph flowGraph, List expected, IActivity activity) + { + Assert.Equal(expected, flowGraph.GetOutboundConnections(activity).Select(c => c.ToString())); + } + + public static void ValidateForwardInboundConnections(this FlowGraph flowGraph, List expected, IActivity activity) + { + Assert.Equal(expected, flowGraph.GetForwardInboundConnections(activity).Select(c => c.ToString())); + } + public static void ValidateBackwardConnection(this FlowGraph flowGraph, bool expectedBackward, bool expectedValid, Connection connection) + { + var actualBackward = flowGraph.IsBackwardConnection(connection, out var actualValid); + Assert.Equal(expectedBackward, actualBackward); + Assert.Equal(expectedValid, actualValid); + } + + public static void ValidateDanglingActivity(this FlowGraph flowGraph, bool expected, Activity activity) + { + Assert.Equal(expected, flowGraph.IsDanglingActivity(activity)); + } + + public static void ValidateAncestorActivities(this FlowGraph flowGraph, List expected, Activity activity) + { + Assert.Equal(expected, flowGraph.GetAncestorActivities(activity)); + } +} diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/Flowchart/FlowGraphTests.cs b/test/unit/Elsa.Workflows.Core.UnitTests/Flowchart/FlowGraphTests.cs new file mode 100644 index 000000000..b090e18ac --- /dev/null +++ b/test/unit/Elsa.Workflows.Core.UnitTests/Flowchart/FlowGraphTests.cs @@ -0,0 +1,386 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.Activities.Flowchart.Models; + +namespace Elsa.Workflows.Core.UnitTests.Flowchart; +public class FlowGraphTests +{ + // Start + // ↓ + // A ← W Y + // ↙ ↘ ↙ + // B C ← X ← Z + // ↘ ↙ + // D + // ↓ + // End + [Fact] + public void InvalidDanglingActivitiesTest() + { + var start = new TestActivity("start"); + var a = new TestActivity("a"); + var b = new TestActivity("b"); + var c = new TestActivity("c"); + var d = new TestActivity("d"); + var end = new TestActivity("end"); + var w = new TestActivity("w"); + var x = new TestActivity("x"); + var y = new TestActivity("y"); + var z = new TestActivity("z"); + + var connections = new List + { + new(start, a), + new(a, b), + new(a, c), + new(b, d), + new(c, d), + new(d, end), + new(w, a), // invalid dangling connector + new(x, c), // invalid dangling connector + new(y, x), // invalid dangling connector + new(z, x), // invalid dangling connector + }; + + var flowGraph = new FlowGraph(connections, start); + + // FlowGraph.GetForwardInboundConnections + flowGraph.ValidateForwardInboundConnections(["start->a"], a); + flowGraph.ValidateForwardInboundConnections(["a->b"], b); + flowGraph.ValidateForwardInboundConnections(["a->c"], c); + flowGraph.ValidateForwardInboundConnections(["b->d", "c->d"], d); + flowGraph.ValidateForwardInboundConnections(["d->end"], end); + flowGraph.ValidateForwardInboundConnections([], x); + + // FlowGraph.GetOutboundConnections + flowGraph.ValidateOutboundConnections(["start->a"], start); + flowGraph.ValidateOutboundConnections(["a->b", "a->c"], a); + flowGraph.ValidateOutboundConnections(["b->d"], b); + flowGraph.ValidateOutboundConnections(["c->d"], c); + flowGraph.ValidateOutboundConnections(["d->end"], d); + flowGraph.ValidateOutboundConnections([], end); + flowGraph.ValidateOutboundConnections(["x->c"], x); + + // FlowGraph.GetAncestorActivities + flowGraph.ValidateAncestorActivities([], start); + flowGraph.ValidateAncestorActivities([start], a); + flowGraph.ValidateAncestorActivities([a, start], b); + flowGraph.ValidateAncestorActivities([a, start], c); + flowGraph.ValidateAncestorActivities([b, c, a, start], d); + flowGraph.ValidateAncestorActivities([], x); + + // FlowGraph IsDanglingActivity + flowGraph.ValidateDanglingActivity(false, start); + flowGraph.ValidateDanglingActivity(false, a); + flowGraph.ValidateDanglingActivity(false, b); + flowGraph.ValidateDanglingActivity(false, c); + flowGraph.ValidateDanglingActivity(false, d); + flowGraph.ValidateDanglingActivity(false, end); + flowGraph.ValidateDanglingActivity(true, w); + flowGraph.ValidateDanglingActivity(true, x); + flowGraph.ValidateDanglingActivity(true, y); + flowGraph.ValidateDanglingActivity(true, z); + } + + // Start + // ↓ + // A ← ← + // ↙ ↘ ↖ + // B C ↑ + // ↘ ↙ ↗ + // D → → + [Fact] + public void ValidBackwardConnectionTest() + { + var start = new TestActivity("start"); + var a = new TestActivity("a"); + var b = new TestActivity("b"); + var c = new TestActivity("c"); + var d = new TestActivity("d"); + + var connections = new List + { + new(start, a), + new(a, b), + new(a, c), + new(b, d), + new(c, d), + new(d, a), // loopback + }; + + var flowGraph = new FlowGraph(connections, start); + + // FlowGraph.GetForwardInboundConnections + flowGraph.ValidateForwardInboundConnections([], start); + flowGraph.ValidateForwardInboundConnections(["start->a"], a); + flowGraph.ValidateForwardInboundConnections(["a->b"], b); + flowGraph.ValidateForwardInboundConnections(["a->c"], c); + flowGraph.ValidateForwardInboundConnections(["b->d", "c->d"], d); + + // FlowGraph.GetOutboundConnections + flowGraph.ValidateOutboundConnections(["start->a"], start); + flowGraph.ValidateOutboundConnections(["a->b", "a->c"], a); + flowGraph.ValidateOutboundConnections(["b->d"], b); + flowGraph.ValidateOutboundConnections(["c->d"], c); + flowGraph.ValidateOutboundConnections(["d->a"], d); + + // FlowGraph.GetAncestorActivities + flowGraph.ValidateAncestorActivities([], start); + flowGraph.ValidateAncestorActivities([start], a); + flowGraph.ValidateAncestorActivities([a, start], b); + flowGraph.ValidateAncestorActivities([a, start], c); + flowGraph.ValidateAncestorActivities([b, c, a, start], d); + + // FlowGraph.IsBackwardConnection + flowGraph.ValidateBackwardConnection(false, false, new(a, c)); + flowGraph.ValidateBackwardConnection(true, true, new(d, a)); + } + + + // Start + // ↙ ↘ + // A B ← ← + // ↘ ↙ ↘ ↖ + // C D ↰ ↑ + // ↘ ↙ ↗ ↗ + // E → → + [Fact] + public void InvalidLoopbackTest() + { + var start = new TestActivity("start"); + var a = new TestActivity("a"); + var b = new TestActivity("b"); + var c = new TestActivity("c"); + var d = new TestActivity("d"); + var e = new TestActivity("e"); + + var connections = new List + { + new(start, a), + new(start, b), + new(a, c), + new(b, c), + new(b, d), + new(d, e), + new(c, e), + new(e, b), // invalid loopback + new(e, d), // invalid loopback + }; + + var flowGraph = new FlowGraph(connections, start); + + // FlowGraph.GetForwardInboundConnections + flowGraph.ValidateForwardInboundConnections([], start); + flowGraph.ValidateForwardInboundConnections(["start->a"], a); + flowGraph.ValidateForwardInboundConnections(["start->b"], b); + flowGraph.ValidateForwardInboundConnections(["a->c", "b->c"], c); + flowGraph.ValidateForwardInboundConnections(["b->d"], d); + flowGraph.ValidateForwardInboundConnections(["c->e", "d->e"], e); + + // FlowGraph.GetOutboundConnections + flowGraph.ValidateOutboundConnections(["start->a", "start->b"], start); + flowGraph.ValidateOutboundConnections(["a->c"], a); + flowGraph.ValidateOutboundConnections(["b->c", "b->d"], b); + flowGraph.ValidateOutboundConnections(["c->e"], c); + flowGraph.ValidateOutboundConnections(["d->e"], d); + flowGraph.ValidateOutboundConnections(["e->b", "e->d"], e); + + // FlowGraph.GetAncestorActivities + flowGraph.ValidateAncestorActivities([], start); + flowGraph.ValidateAncestorActivities([start], a); + flowGraph.ValidateAncestorActivities([start], b); + flowGraph.ValidateAncestorActivities([a, b, start], c); + flowGraph.ValidateAncestorActivities([b, start], d); + flowGraph.ValidateAncestorActivities([c, d, a, b, start], e); + + // FlowGraph.IsBackwardConnection + flowGraph.ValidateBackwardConnection(true, false, new(e, b)); + flowGraph.ValidateBackwardConnection(true, false, new(e, d)); + } + + // Start + // ↙ ↘ + // ↓ A + // ↓ ↓ + // ↓ B + // ↓ ↙ ↓ + // C D + // ↘ ↙ + // End + [Fact] + public void LongLegTest() + { + var start = new TestActivity("start"); + var a = new TestActivity("a"); + var b = new TestActivity("b"); + var c = new TestActivity("c"); + var d = new TestActivity("d"); + var end = new TestActivity("end"); + + var connections = new List + { + new(start, a), + new(start, c), + new(a, b), + new(b, c), + new(b, d), + new(c, end), + new(d, end), + }; + + var flowGraph = new FlowGraph(connections, start); + + // FlowGraph.GetForwardInboundConnections + flowGraph.ValidateForwardInboundConnections([], start); + flowGraph.ValidateForwardInboundConnections(["start->a"], a); + flowGraph.ValidateForwardInboundConnections(["a->b"], b); + flowGraph.ValidateForwardInboundConnections(["start->c", "b->c"], c); + flowGraph.ValidateForwardInboundConnections(["b->d"], d); + flowGraph.ValidateForwardInboundConnections(["c->end", "d->end"], end); + + // FlowGraph.GetOutboundConnections + flowGraph.ValidateOutboundConnections(["start->a", "start->c"], start); + flowGraph.ValidateOutboundConnections(["a->b"], a); + flowGraph.ValidateOutboundConnections(["b->c", "b->d"], b); + flowGraph.ValidateOutboundConnections(["c->end"], c); + flowGraph.ValidateOutboundConnections(["d->end"], d); + flowGraph.ValidateOutboundConnections([], end); + + // FlowGraph.GetAncestorActivities + flowGraph.ValidateAncestorActivities([], start); + flowGraph.ValidateAncestorActivities([start], a); + flowGraph.ValidateAncestorActivities([a, start], b); + flowGraph.ValidateAncestorActivities([start, b, a], c); + flowGraph.ValidateAncestorActivities([b, a, start], d); + flowGraph.ValidateAncestorActivities([c, d, start, b, a], end); + } + + // Start + // ↙ ↘ + // A B + // ↓↘ ↙ ↘ + // ↳→C D + // ↓↘ ↙ + // ↳→E + [Fact] + public void SameEdgeDuplicateTest() + { + var start = new TestActivity("start"); + var a = new TestActivity("a"); + var b = new TestActivity("b"); + var c = new TestActivity("c"); + var d = new TestActivity("d"); + var e = new TestActivity("e"); + + var connections = new List + { + new(start, a), + new(start, b), + new(a, c), + new(a, c), // duplicate + new(b, c), + new(b, d), + new(d, e), + new(c, e), + new(c, e), // duplicate + }; + + var flowGraph = new FlowGraph(connections, start); + + // FlowGraph.GetForwardInboundConnections + flowGraph.ValidateForwardInboundConnections([], start); + flowGraph.ValidateForwardInboundConnections(["start->a"], a); + flowGraph.ValidateForwardInboundConnections(["start->b"], b); + flowGraph.ValidateForwardInboundConnections(["a->c", "b->c"], c); + flowGraph.ValidateForwardInboundConnections(["b->d"], d); + flowGraph.ValidateForwardInboundConnections(["c->e", "d->e"], e); + + // FlowGraph.GetOutboundConnections + flowGraph.ValidateOutboundConnections(["start->a", "start->b"], start); + flowGraph.ValidateOutboundConnections(["a->c"], a); + flowGraph.ValidateOutboundConnections(["b->c", "b->d"], b); + flowGraph.ValidateOutboundConnections(["c->e"], c); + flowGraph.ValidateOutboundConnections(["d->e"], d); + flowGraph.ValidateOutboundConnections([], e); + + // FlowGraph.GetAncestorActivities + flowGraph.ValidateAncestorActivities([], start); + flowGraph.ValidateAncestorActivities([start], a); + flowGraph.ValidateAncestorActivities([start], b); + flowGraph.ValidateAncestorActivities([a, b, start], c); + flowGraph.ValidateAncestorActivities([b, start], d); + flowGraph.ValidateAncestorActivities([c, d, a, b, start], e); + } + + // Start + // ↙ ↘ + // A B + // ↓↘ ↙ ↘ + // ↳→C D + // ↓↘ ↙ + // ↳→E + [Fact] + public void SameEdgeDifferentPortDuplicateTest() + { + var start = new TestActivity("start"); + var a = new TestActivity("a"); + var b = new TestActivity("b"); + var c = new TestActivity("c"); + var d = new TestActivity("d"); + var e = new TestActivity("e"); + + var connections = new List + { + new(start, a), + new(start, b), + new(start, b), // duplicate + new(a, c), + new(new Endpoint(a, "Yes"), new Endpoint(c)), + new(new Endpoint(a, "Yes"), new Endpoint(c)), // duplicate + new(b, c), + new(b, d), + new(c, e), + new(new Endpoint(c, "Yes"), new Endpoint(e)), + new(d, e), + }; + + var flowGraph = new FlowGraph(connections, start); + + // FlowGraph.GetForwardInboundConnections + flowGraph.ValidateForwardInboundConnections([], start); + flowGraph.ValidateForwardInboundConnections(["start->a"], a); + flowGraph.ValidateForwardInboundConnections(["start->b"], b); + flowGraph.ValidateForwardInboundConnections(["a->c", "a:Yes->c", "b->c"], c); + flowGraph.ValidateForwardInboundConnections(["b->d"], d); + flowGraph.ValidateForwardInboundConnections(["c->e", "c:Yes->e", "d->e"], e); + + // FlowGraph.GetOutboundConnections + flowGraph.ValidateOutboundConnections(["start->a", "start->b"], start); + flowGraph.ValidateOutboundConnections(["a->c", "a:Yes->c"], a); + flowGraph.ValidateOutboundConnections(["b->c", "b->d"], b); + flowGraph.ValidateOutboundConnections(["c->e", "c:Yes->e"], c); + flowGraph.ValidateOutboundConnections(["d->e"], d); + flowGraph.ValidateOutboundConnections([], e); + + // FlowGraph.GetAncestorActivities + flowGraph.ValidateAncestorActivities([], start); + flowGraph.ValidateAncestorActivities([start], a); + flowGraph.ValidateAncestorActivities([start], b); + flowGraph.ValidateAncestorActivities([a, b, start], c); + flowGraph.ValidateAncestorActivities([b, start], d); + flowGraph.ValidateAncestorActivities([c, d, a, b, start], e); + } + + class TestActivity : Activity + { + public TestActivity(string id) + { + Id = id; + Name = id; + } + + public override string ToString() + { + return Id; + } + } +}