Update flowchart activity scheduling, implicit join (and FlowJoin/WaitAll) now only waits for followed connections

This commit is contained in:
Bob Hauser 2025-02-12 23:27:40 -05:00
parent 08f7bd7019
commit 4e76a32e4e
14 changed files with 1368 additions and 394 deletions

View file

@ -0,0 +1,23 @@
using Elsa.Workflows;
namespace Elsa.Testing.Shared;
public class TestWorkflow : WorkflowBase
{
private readonly Action<IWorkflowBuilder> _buildWorkflow;
public TestWorkflow(Action<IWorkflowBuilder> buildWorkflow)
{
_buildWorkflow = buildWorkflow;
}
protected override void Build(IWorkflowBuilder workflowBuilder)
{
_buildWorkflow(workflowBuilder);
if (string.IsNullOrEmpty(workflowBuilder.Id))
{
workflowBuilder.Id = Guid.NewGuid().ToString();
}
}
}

View file

@ -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;
/// <summary>
/// Merge multiple branches into a single branch of execution.
/// </summary>
[Activity("Elsa", "Branching", "Merge multiple branches into a single branch of execution.", DisplayName = "Join")]
[PublicAPI]
public class FlowJoin : Activity, IJoinNode
{
/// <inheritdoc />
public FlowJoin([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// 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).
/// </summary>
[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<FlowJoinMode> Mode { get; set; } = new(FlowJoinMode.WaitAny);
/// <inheritdoc />
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;
/// <summary>
/// Merge multiple branches into a single branch of execution.
/// </summary>
[Activity("Elsa", "Branching", "Merge multiple branches into a single branch of execution.", DisplayName = "Join")]
[PublicAPI]
public class FlowJoin : Activity, IJoinNode
{
/// <inheritdoc />
public FlowJoin([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// 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).
/// </summary>
[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<FlowJoinMode> Mode { get; set; } = new(FlowJoinMode.WaitAny);
/// <inheritdoc />
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;
}
}
}
}

View file

@ -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";
/// <inheritdoc />
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();
}
/// <summary>
/// Checks if there is any pending work for the flowchart.
/// </summary>
@ -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<ILogger<Flowchart>>();
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<bool>(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);
}
}
/// <summary>
/// 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.
/// </summary>
/// <param name="flowchart">The flowchart containing the activities.</param>
/// <param name="flowGraph">The graph representation of the flowchart.</param>
/// <param name="flowScope">Tracks activity and connection visits.</param>
/// <param name="flowchartContext">The execution context of the flowchart.</param>
/// <param name="activity">The current activity being processed.</param>
/// <param name="outcomes">The outcomes that determine which connections were followed.</param>
/// <param name="completedActivityExecutedByBackwardConnection">Indicates if the completed activity was executed due to a backward connection.</param>
/// <returns>True if at least one activity was scheduled; otherwise, false.</returns>
private async ValueTask<bool> 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;
}
/// <summary>
/// Schedules an outbound activity that originates from a backward connection.
/// </summary>
private async ValueTask<bool> 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<string, object>() { { BackwardConnectionActivityInput, true } }
};
await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions);
return true;
}
/// <summary>
/// Schedules a non-join activity if all its forward inbound connections have been visited.
/// </summary>
private async ValueTask<bool> 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);
}
}
/// <summary>
/// Schedules a join activity based on inbound connection statuses.
/// </summary>
private async ValueTask<bool> 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)
{

View file

@ -1,107 +1,39 @@
using Elsa.Workflows.Activities.Flowchart.Models;
namespace Elsa.Workflows.Activities.Flowchart.Extensions;
/// <summary>
/// Contains extension methods for <see cref="ICollection{Connection}"/>.
/// </summary>
public static class ConnectionsExtensions
{
/// <summary>
/// Returns all connections that are descendants of the specified parent activity.
/// </summary>
public static IEnumerable<Connection> Descendants(this ICollection<Connection> connections, IActivity parent)
{
var visitedConnections = new HashSet<Connection>();
return connections.Descendants(parent, visitedConnections);
}
/// <summary>
/// Returns all ancestor connections of the specified parent activity.
/// </summary>
public static IEnumerable<Connection> Ancestors(this ICollection<Connection> connections, IActivity activity)
{
var visitedActivities = new HashSet<IActivity>();
return connections.Ancestors(activity, visitedActivities);
}
/// <summary>
/// Returns all inbound connections of the specified activity.
/// </summary>
public static IEnumerable<Connection> InboundConnections(this ICollection<Connection> connections, IActivity activity) => connections.Where(x => x.Target.Activity == activity).ToList();
/// <summary>
/// Returns all "left" inbound connections of the specified activity. "Left" means "not a descendant of the activity".
/// </summary>
public static IEnumerable<Connection> LeftInboundConnections(this ICollection<Connection> 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;
}
/// <summary>
/// Returns all "left" ancestor connections of the specified activity. "Left" means "not a descendant of the activity".
/// </summary>
public static IEnumerable<Connection> LeftAncestorConnections(this ICollection<Connection> 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;
}
/// <summary>
/// Returns all inbound activities of the specified activity.
/// </summary>
public static IEnumerable<IActivity> InboundActivities(this ICollection<Connection> connections, IActivity activity) => connections.InboundConnections(activity).Select(x => x.Source.Activity);
/// <summary>
/// Returns all "left" inbound activities of the specified activity. "Left" means "not a descendant of the activity".
/// </summary>
public static IEnumerable<IActivity> LeftInboundActivities(this ICollection<Connection> connections, IActivity activity) => connections.LeftInboundConnections(activity).Select(x => x.Source.Activity);
/// <summary>
/// Returns all "left" ancestor activities of the specified activity. "Left" means "not a descendant of the activity".
/// </summary>
public static IEnumerable<IActivity> LeftAncestorActivities(this ICollection<Connection> connections, IActivity activity) => connections.LeftAncestorConnections(activity).Select(x => x.Source.Activity);
private static IEnumerable<Connection> Descendants(this ICollection<Connection> connections, IActivity parent, ISet<Connection> 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<Connection> Ancestors(this ICollection<Connection> connections, IActivity activity, ISet<IActivity> 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;
/// <summary>
/// Contains extension methods for <see cref="ICollection{Connection}"/>.
/// </summary>
public static class ConnectionsExtensions
{
/// <summary>
/// Returns all inbound connections of the specified activity.
/// </summary>
public static IEnumerable<Connection> InboundConnections(this ICollection<Connection> connections, IActivity activity) => connections.Where(x => x.Target.Activity == activity).Distinct().ToList();
/// <summary>
/// Returns all inbound activities of the specified activity.
/// </summary>
public static IEnumerable<IActivity> InboundActivities(this ICollection<Connection> connections, IActivity activity) => connections.InboundConnections(activity).Select(x => x.Source.Activity);
/// <summary>
/// Returns all outbound connections of the specified activity.
/// </summary>
public static IEnumerable<Connection> OutboundConnections(this ICollection<Connection> connections, IActivity activity) => connections.Where(x => x.Source.Activity == activity).Distinct().ToList();
/// <summary>
/// Returns all outbound connections of the specified activity matching the specified outcomes.
/// </summary>
public static IEnumerable<Connection> OutboundConnections(this ICollection<Connection> connections, IActivity activity, Outcomes outcomes) => connections.OutboundConnections(activity).Where(c => outcomes.Names.Contains(c.Source.Port));
/// <summary>
/// Returns all outbound activities of the specified activity.
/// </summary>
public static IEnumerable<IActivity> OutboundActivities(this ICollection<Connection> connections, IActivity activity) => connections.OutboundConnections(activity).Select(x => x.Source.Activity);
/// <summary>
/// Returns all outbound activities of the specified activity matching the specified outcomes
/// </summary>
public static IEnumerable<IActivity> OutboundActivities(this ICollection<Connection> connections, IActivity activity, Outcomes outcomes) => connections.OutboundConnections(activity).Where(c => outcomes.Names.Contains(c.Source.Port)).Select(x => x.Source.Activity);
}

View file

@ -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; }
}

View file

@ -5,7 +5,7 @@ namespace Elsa.Workflows.Activities.Flowchart.Models;
/// <summary>
/// A connection between a source and a target endpoint.
/// </summary>
public class Connection
public class Connection : IEquatable<Connection>
{
/// <summary>
/// Initializes a new instance of the <see cref="Connection"/> class.
@ -45,5 +45,36 @@ public class Connection
/// <summary>
/// The target endpoint.
/// </summary>
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);
}
}

View file

@ -0,0 +1,242 @@
using Elsa.Extensions;
using Elsa.Workflows.Activities.Flowchart.Extensions;
namespace Elsa.Workflows.Activities.Flowchart.Models;
/// <summary>
/// Represents a directed graph structure for managing workflow connections.
/// Caches forward and backward connections to optimize graph traversal.
/// </summary>
public class FlowGraph(ICollection<Connection> Connections, IActivity? RootActivity)
{
private List<Connection>? _cachedForwardConnections;
private readonly Dictionary<IActivity, List<Connection>> _cachedInboundForwardConnections = new();
private readonly Dictionary<IActivity, List<Connection>> _cachedOutboundConnections = new();
private readonly Dictionary<Connection, (bool IsBackwardConnection, bool IsValid)> _cachedIsBackwardConnection = new();
private readonly Dictionary<IActivity, bool> _cachedIsDanglingActivity = new();
private readonly Dictionary<IActivity, List<IActivity>> _cachedAncestors = new();
/// <summary>
/// Gets the list of forward connections, computing them if not already cached.
/// </summary>
private List<Connection> ForwardConnections => _cachedForwardConnections ??= RootActivity == null ? new() : GetForwardConnections(Connections, RootActivity);
/// <summary>
/// Retrieves all inbound forward connections for a given activity.
/// </summary>
public List<Connection> GetForwardInboundConnections(IActivity activity) => _cachedInboundForwardConnections.GetOrAdd(activity, () => ForwardConnections.InboundConnections(activity).ToList());
/// <summary>
/// Retrieves all outbound connections for a given activity.
/// </summary>
public List<Connection> GetOutboundConnections(IActivity activity) => _cachedOutboundConnections.GetOrAdd(activity, () => Connections.OutboundConnections(activity).ToList());
/// <summary>
/// Determines if a given activity is "dangling," meaning it does not exist as a target in any forward connection.
/// </summary>
public bool IsDanglingActivity(IActivity activity) => _cachedIsDanglingActivity.GetOrAdd(activity, () => activity != RootActivity && !ForwardConnections.Any(c => c.Target.Activity == activity));
/// <summary>
/// Determines if a given connection is a backward connection (i.e., not part of the forward traversal) and whether it is valid.
/// </summary>
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;
}
/// <summary>
/// Retrieves all ancestor activities for a given activity by traversing ForwardConnections in reverse.
/// </summary>
public List<IActivity> GetAncestorActivities(IActivity activity)
{
return _cachedAncestors.GetOrAdd(activity, () => ComputeAncestors(activity));
}
/// <summary>
/// Computes the list of ancestors by following Source activities in ForwardConnections.
/// </summary>
private List<IActivity> ComputeAncestors(IActivity activity)
{
HashSet<IActivity> ancestors = new();
Queue<IActivity> 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();
}
/// <summary>
/// Computes the list of forward connections in the graph, excluding cyclic connections.
/// </summary>
private static List<Connection> GetForwardConnections(ICollection<Connection> connections, IActivity root)
{
Dictionary<IActivity, List<IActivity>> adjList = new();
foreach (var conn in connections)
{
if (!adjList.ContainsKey(conn.Source.Activity))
adjList[conn.Source.Activity] = new List<IActivity>();
adjList[conn.Source.Activity].Add(conn.Target.Activity);
}
HashSet<IActivity> visited = new();
HashSet<(IActivity, IActivity)> visitedEdges = new();
List<(IActivity Source, IActivity Target)> validEdges = new();
Queue<IActivity> 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();
}
/// <summary>
/// Determines if there is an existing path from the source activity to the target activity.
/// Helps in detecting cyclic connections.
/// </summary>
private static bool HasPathToActivity(ICollection<(IActivity Source, IActivity Target)> edges, IActivity source, IActivity target)
{
if (source == target)
return true;
HashSet<IActivity> visited = new();
Stack<IActivity> 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;
}
/// <summary>
/// Determines whether a backward connection is valid by ensuring all paths from source to root pass through target.
/// </summary>
private static bool IsValidBackwardConnection(List<Connection> 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;
}
/// <summary>
/// Finds all paths from a given start activity to the root using BFS.
/// </summary>
private static List<List<IActivity>> GetPathsToRoot(List<Connection> forwardConnections, IActivity root, IActivity start)
{
List<List<IActivity>> paths = new();
Queue<List<IActivity>> queue = new();
queue.Enqueue(new List<IActivity> { start });
while (queue.Count > 0)
{
var path = queue.Dequeue();
var lastNode = path.Last();
if (lastNode == root)
{
paths.Add(new List<IActivity>(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<IActivity>(path) { prev };
queue.Enqueue(newPath);
}
}
}
return paths;
}
}

View file

@ -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));
}
/// <summary>
/// The activity from which the scope was created.
/// </summary>
public string OwnerActivityId { get; set; } = default!;
/// <summary>
/// A list of scheduled activity IDs and a flag whether they executed or not.
/// </summary>
public IDictionary<string, ActivityFlowState> Activities { get; set; } =
new Dictionary<string, ActivityFlowState>();
public void AddActivities(IEnumerable<IActivity> 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++;
}
/// <summary>
/// Return a list excluding any activities that already executed.
/// </summary>
public IEnumerable<IActivity> ExcludeExecutedActivities(IEnumerable<IActivity> 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);
}
}
}
using System.Text.Json.Serialization;
namespace Elsa.Workflows.Activities.Flowchart.Models;
/// <summary>
/// Represents a scope for tracking activity and connection visits within a flowchart execution.
/// </summary>
public class FlowScope
{
[JsonConstructor]
public FlowScope()
{
}
[JsonInclude]
private Dictionary<string, long> ActivitiesVisitCount { get; init; } = new();
[JsonInclude]
private Dictionary<string, long> ConnectionVisitCount { get; init; } = new();
[JsonInclude]
private Dictionary<string, bool> ConnectionLastVisitFollowed { get; init; } = new();
/// <summary>
/// Registers a visit to the specified activity, incrementing its visit count.
/// </summary>
/// <param name="activity">The activity being visited.</param>
public void RegisterActivityVisit(IActivity activity)
{
string activityId = activity.Id;
ActivitiesVisitCount.TryAdd(activityId, 0);
ActivitiesVisitCount[activityId]++;
}
/// <summary>
/// Gets the number of times the specified activity has been visited.
/// </summary>
/// <param name="activity">The activity to check.</param>
/// <returns>The visit count of the activity.</returns>
private long GetActivityVisitCount(IActivity activity) => ActivitiesVisitCount.TryGetValue(activity.Id, out var count) ? count : 0;
/// <summary>
/// Registers a visit to the specified connection and records whether it was followed.
/// </summary>
/// <param name="connection">The connection being visited.</param>
/// <param name="followed">Indicates whether the connection was followed.</param>
public void RegisterConnectionVisit(Connection connection, bool followed)
{
string connectionId = connection.ToString();
ConnectionVisitCount.TryAdd(connectionId, 0);
ConnectionVisitCount[connectionId]++;
ConnectionLastVisitFollowed[connectionId] = followed;
}
/// <summary>
/// Gets the number of times the specified connection has been visited.
/// </summary>
/// <param name="connection">The connection to check.</param>
/// <returns>The visit count of the connection.</returns>
private long GetConnectionVisitCount(Connection connection) => ConnectionVisitCount.TryGetValue(connection.ToString(), out var count) ? count : 0;
/// <summary>
/// Determines whether the last visit to the specified connection was followed.
/// </summary>
/// <param name="connection">The connection to check.</param>
/// <returns>True if the connection was followed on the last visit, otherwise false.</returns>
private bool GetConnectionLastVisitFollowed(Connection connection) => ConnectionLastVisitFollowed.TryGetValue(connection.ToString(), out var followed) ? followed : false;
/// <summary>
/// Determines whether all inbound connections to the specified activity have been visited.
/// </summary>
/// <param name="flowGraph">The flow graph containing connections.</param>
/// <param name="activity">The activity to check.</param>
/// <returns>True if all inbound connections have been visited, otherwise false.</returns>
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;
}
/// <summary>
/// Determines whether any inbound connection to the specified activity has been followed.
/// </summary>
/// <param name="flowGraph">The flow graph containing connections.</param>
/// <param name="activity">The activity to check.</param>
/// <returns>True if any inbound connection has been followed, otherwise false.</returns>
public bool HasFollowedInboundConnection(FlowGraph flowGraph, IActivity activity)
{
var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity);
return forwardInboundConnections.Any(c => GetConnectionLastVisitFollowed(c));
}
/// <summary>
/// Determines whether a connection should be ignored based on visit counts.
/// </summary>
/// <param name="connection">The connection to check.</param>
/// <param name="activity">The activity associated with the connection.</param>
/// <returns>True if the connection should be ignored, otherwise false.</returns>
public bool ShouldIgnoreConnection(Connection connection, IActivity activity)
{
var connectionVisitCount = GetConnectionVisitCount(connection);
var activityVisitCount = GetActivityVisitCount(activity);
return connectionVisitCount <= activityVisitCount;
}
}

View file

@ -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 <see cref="Activities.Flowchart"/>.
/// </summary>
/// <param name="Names">A list of outcome names.</param>
public record Outcomes(params string[] Names);
public record Outcomes(params string[] Names)
{
public static readonly Outcomes Default = new([null!, "Done"]);
public static readonly Outcomes Empty = new();
}

View file

@ -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 =
{

View file

@ -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<FlowchartWorkflow>();
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<int>("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<object?>(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<int>("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<object?>(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);
}
}

View file

@ -22,7 +22,7 @@ public class JsonObjectJintTests
public async Task Test1()
{
await _services.PopulateRegistriesAsync();
await _workflowRunner.RunAsync<TestWorkflow>();
await _workflowRunner.RunAsync<Workflows.TestWorkflow>();
var lines = _capturingTextWriter.Lines.ToList();
Assert.Equal(new[] { "Baz" }, lines);
}

View file

@ -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<string> expected, IActivity activity)
{
Assert.Equal(expected, flowGraph.GetOutboundConnections(activity).Select(c => c.ToString()));
}
public static void ValidateForwardInboundConnections(this FlowGraph flowGraph, List<string> 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<Activity> expected, Activity activity)
{
Assert.Equal(expected, flowGraph.GetAncestorActivities(activity));
}
}

View file

@ -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<Connection>
{
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<Connection>
{
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<Connection>
{
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<Connection>
{
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<Connection>
{
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<Connection>
{
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;
}
}
}