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 c83c5a74b..e1abb04de 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs @@ -1,5 +1,6 @@ using System.Runtime.CompilerServices; 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; diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.Counters.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.Counters.cs index 7e0194259..06f3ec9df 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.Counters.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.Counters.cs @@ -10,8 +10,10 @@ namespace Elsa.Workflows.Activities.Flowchart.Activities; public partial class Flowchart { private const string ScopeProperty = "FlowScope"; + private const string GraphTransientProperty = "FlowGraph"; private const string BackwardConnectionActivityInput = "BackwardConnection"; - + + private async ValueTask OnChildCompletedCounterBasedLogicAsync(ActivityCompletedContext context) { var flowchartContext = context.TargetContext; @@ -19,9 +21,108 @@ public partial class Flowchart var completedActivity = completedActivityContext.Activity; var result = context.Result; + // Determine the outcomes from the completed activity + var outcomes = result is Outcomes o ? o : Outcomes.Default; + + await ProcessChildCompletedAsync(flowchartContext, completedActivity, completedActivityContext, outcomes); + } + + private IActivity? GetStartActivity(ActivityExecutionContext context) + { + // If there's a trigger that triggered this workflow, use that. + var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId; + var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : null; + + if (triggerActivity != null) + return triggerActivity; + + // If an explicit Start activity was provided, use that. + if (Start != null) + return Start; + + // If there is a Start activity on the flowchart, use that. + var startActivity = Activities.FirstOrDefault(x => x is Start); + + if (startActivity != null) + return startActivity; + + // If there's an activity marked as "Can Start Workflow", use that. + var canStartWorkflowActivity = Activities.FirstOrDefault(x => x.GetCanStartWorkflow()); + + if (canStartWorkflowActivity != null) + return canStartWorkflowActivity; + + // If there is a single activity that has no inbound connections, use that. + var root = GetRootActivity(); + + if (root != null) + return root; + + // If no start activity found, return the first activity. + return Activities.FirstOrDefault(); + } + + /// + /// Checks if there is any pending work for the flowchart. + /// + private bool HasPendingWork(ActivityExecutionContext context) + { + var workflowExecutionContext = context.WorkflowExecutionContext; + + // Use HashSet for O(1) lookups + var activityIds = new HashSet(Activities.Select(x => x.Id)); + + // Short circuit evaluation - check running instances first before more expensive scheduler check + if (context.Children.Any(x => activityIds.Contains(x.Activity.Id) && x.Status == ActivityStatus.Running)) + return true; + + // Scheduler check - optimize to avoid repeated LINQ evaluations + var scheduledItems = workflowExecutionContext.Scheduler.List().ToList(); + + return scheduledItems.Any(workItem => + { + var ownerInstanceId = workItem.Owner?.Id; + + if (ownerInstanceId == null) + return false; + + if (ownerInstanceId == context.Id) + return true; + + var ownerContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.Id == ownerInstanceId); + return ownerContext.GetAncestors().Any(x => x == context); + }); + } + + private IActivity? GetRootActivity() + { + // Get the first activity that has no inbound connections. + var query = + from activity in Activities + let inboundConnections = Connections.Any(x => x.Target.Activity == activity) + where !inboundConnections + select activity; + + 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 ProcessChildCompletedAsync(ActivityExecutionContext flowchartContext, IActivity completedActivity, ActivityExecutionContext completedActivityContext, Outcomes outcomes) + { if (flowchartContext.Activity != this) { - throw new Exception("Target context activity must be this flowchart"); + throw new("Target context activity must be this flowchart"); } // If the completed activity's status is anything but "Completed", do not schedule its outbound activities. @@ -37,14 +138,11 @@ public partial class Flowchart return; } - // Determine the outcomes from the completed activity - var outcomes = result is Outcomes o ? o : Outcomes.Default; - // Schedule the outbound activities - var flowGraph = flowchartContext.GetFlowGraph(); + var flowGraph = GetFlowGraph(flowchartContext); var flowScope = GetFlowScope(flowchartContext); - var completedActivityExecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault(BackwardConnectionActivityInput); - bool hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExecutedByBackwardConnection); + var completedActivityExcecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault(BackwardConnectionActivityInput); + bool hasScheduledActivity = await MaybeScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, OnChildCompletedAsync, completedActivityExcecutedByBackwardConnection); // If there are not any outbound connections, complete the flowchart activity if there is no other pending work if (!hasScheduledActivity) @@ -53,31 +151,30 @@ public partial class Flowchart } } - private FlowScope GetFlowScope(ActivityExecutionContext context) - { - return context.GetProperty(ScopeProperty, () => new FlowScope()); - } - /// /// 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. + /// forward traversal rules, and backward connections. If outcomes is Outcomes.Empty, it indicates + /// that the activity should be skipped - all outbound connections will be visited and treated as + /// not followed. /// /// 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. + /// The callback to invoke upon activity completion. + /// 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) + private static async ValueTask MaybeScheduleOutboundActivitiesAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity activity, Outcomes outcomes, ActivityCompletionCallback completionCallback, bool completedActivityExecutedByBackwardConnection = false) { - var hasScheduledActivity = 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."); + throw new($"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 @@ -93,19 +190,12 @@ public partial class Flowchart flowScope.RegisterConnectionVisit(outboundConnection, connectionFollowed); var outboundActivity = outboundConnection.Target.Activity; - // Determine scheduling strategy based on connection type + // Determine the scheduling strategy based on connection-type. if (flowGraph.IsBackwardConnection(outboundConnection, out var backwardConnectionIsValid)) - { - hasScheduledActivity |= await ScheduleBackwardConnectionActivityAsync(flowGraph, flowchartContext, outboundConnection, outboundActivity, connectionFollowed, backwardConnectionIsValid); - } - else if (outboundActivity is not IJoinNode) - { - hasScheduledActivity |= await ScheduleNonJoinActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity); - } + // Backward connections are scheduled differently + hasScheduledActivity |= await MaybeScheduleBackwardConnectionActivityAsync(flowGraph, flowchartContext, outboundConnection, outboundActivity, connectionFollowed, backwardConnectionIsValid, completionCallback); else - { - hasScheduledActivity |= await ScheduleJoinActivityAsync(flowGraph, flowScope, flowchartContext, outboundConnection, outboundActivity); - } + hasScheduledActivity |= await MaybeScheduleOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundConnection, outboundActivity, completionCallback); } return hasScheduledActivity; @@ -114,7 +204,7 @@ public partial class Flowchart /// /// 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) + private static async ValueTask MaybeScheduleBackwardConnectionActivityAsync(FlowGraph flowGraph, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, bool connectionFollowed, bool backwardConnectionIsValid, ActivityCompletionCallback completionCallback) { if (!connectionFollowed) { @@ -123,18 +213,13 @@ public partial class Flowchart 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."); + throw new($"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 = OnChildCompletedCounterBasedLogicAsync, - Input = new Dictionary() - { - { - BackwardConnectionActivityInput, true - } - } + CompletionCallback = completionCallback, + Input = new Dictionary() { { BackwardConnectionActivityInput, true } } }; await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions); @@ -142,91 +227,136 @@ public partial class Flowchart } /// - /// Schedules a non-join activity if all its forward inbound connections have been visited. + /// Determines the merge mode for a given outbound activity. If the outbound activity is a FlowJoin, it retrieves its configured + /// mode. Otherwise, it defaults to FlowJoinMode.WaitAllActive for implicit joins. /// - private async ValueTask ScheduleNonJoinActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity) + private static async ValueTask GetMergeModeAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity) { - if (!flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) + if (outboundActivity is FlowJoin) { - return false; - } - - if (flowScope.HasFollowedInboundConnection(flowGraph, outboundActivity)) - { - await flowchartContext.ScheduleActivityAsync(outboundActivity, OnChildCompletedCounterBasedLogicAsync); - return true; + var outboundActivityExecutionContext = await flowchartContext.WorkflowExecutionContext.CreateActivityExecutionContextAsync(outboundActivity); + return await outboundActivityExecutionContext.EvaluateInputPropertyAsync(x => x.Mode); } else { - // Propagate skipped connections by scheduling with Outcomes.Empty - return await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty); + // Implicit join case - treat as WaitAllActive + return FlowJoinMode.WaitAllActive; } } /// /// Schedules a join activity based on inbound connection statuses. /// - private async ValueTask ScheduleJoinActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity) + private static async ValueTask MaybeScheduleOutboundActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, ActivityCompletionCallback completionCallback) { - // Ignore the connection if the join activity has already completed (JoinAny scenario) - if (flowScope.ShouldIgnoreConnection(outboundConnection, outboundActivity)) - { - return false; - } + FlowJoinMode mode = await GetMergeModeAsync(flowchartContext, outboundActivity); - // Schedule the join activity only if at least one inbound connection was followed - if (!flowScope.HasFollowedInboundConnection(flowGraph, outboundActivity)) + return mode switch { - 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; - } - } - - if (joinContext is not { Status: ActivityStatus.Running }) - { - var scheduleWorkOptions = new ScheduleWorkOptions - { - CompletionCallback = OnChildCompletedCounterBasedLogicAsync, - ExistingActivityExecutionContext = joinContext - }; - await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions); - return true; - } - else - { - return false; - } + FlowJoinMode.WaitAll => await MaybeScheduleWaitAllActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback), + FlowJoinMode.WaitAllActive => await MaybeScheduleWaitAllActiveActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback), + FlowJoinMode.WaitAny => await MaybeScheduleWaitAnyActivityAsync(flowGraph, flowScope, flowchartContext, outboundConnection, outboundActivity, completionCallback), + _ => throw new($"Unsupported FlowJoinMode: {mode}"), + }; } - public static bool CanWaitAllProceed(ActivityExecutionContext context) + /// + /// Determines whether to schedule an activity based on the FlowJoinMode.WaitAll behavior. + /// If all inbound connections were visited, it checks if they were all followed to decide whether to schedule or skip the activity. + /// + private static async ValueTask MaybeScheduleWaitAllActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) { - var flowchartContext = context.ParentActivityExecutionContext!; - var flowchart = (Flowchart)flowchartContext.Activity; - var flowGraph = flowchartContext.GetFlowGraph(); - var flowScope = flowchart.GetFlowScope(flowchartContext); - var activity = context.Activity; + if (!flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) + // Not all inbound connections have been visited yet; do not schedule anything yet. + return false; - return flowScope.AllInboundConnectionsVisited(flowGraph, activity); + if (flowScope.AllInboundConnectionsFollowed(flowGraph, outboundActivity)) + // All inbound connections were followed; schedule the outbound activity. + return await ScheduleOutboundActivityAsync(flowchartContext, outboundActivity, completionCallback); + else + // No inbound connections were followed; skip the outbound activity. + return await SkipOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback); + } + + /// + /// Determines whether to schedule an activity based on the FlowJoinMode.WaitAllActive behavior. + /// If all inbound connections have been visited, it checks if any were followed to decide whether to schedule or skip the activity. + /// + private static async ValueTask MaybeScheduleWaitAllActiveActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) + { + if (!flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) + // Not all inbound connections have been visited yet; do not schedule anything yet. + return false; + + if (flowScope.AnyInboundConnectionsFollowed(flowGraph, outboundActivity)) + // At least one inbound connection was followed; schedule the outbound activity. + return await ScheduleOutboundActivityAsync(flowchartContext, outboundActivity, completionCallback); + else + // No inbound connections were followed; skip the outbound activity. + return await SkipOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback); + } + + /// + /// Determines whether to schedule an activity based on the FlowJoinMode.WaitAny behavior. + /// If any inbound connection has been followed, it schedules the activity and cancels remaining inbound activities. + /// If a subsequent inbound connection is followed after the activity has been scheduled, it ignores it. + /// + private static async ValueTask MaybeScheduleWaitAnyActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, ActivityCompletionCallback completionCallback) + { + if (flowScope.ShouldIgnoreConnection(outboundConnection, outboundActivity)) + // Ignore the connection if the outbound activity has already completed (JoinAny scenario) + return false; + + if (flowchartContext.WorkflowExecutionContext.Scheduler.List().Any(workItem => workItem.Owner == flowchartContext && workItem.Activity == outboundActivity)) + // Ignore the connection if the outbound activity is already scheduled + return false; + + if (flowScope.AnyInboundConnectionsFollowed(flowGraph, outboundActivity)) + { + // An inbound connection has been followed; cancel remaining inbound activities + await CancelRemainingInboundActivitiesAsync(flowchartContext, outboundActivity); + + // This is the first inbound connection followed; schedule the outbound activity + return await ScheduleOutboundActivityAsync(flowchartContext, outboundActivity, completionCallback); + } + + if (flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) + // All inbound connections have been visited without any being followed; skip the outbound activity + return await SkipOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback); + + // No inbound connections have been followed yet; do not schedule anything yet. + return false; + } + + /// + /// Schedules the outbound activity. + /// + private static async ValueTask ScheduleOutboundActivityAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) + { + await flowchartContext.ScheduleActivityAsync(outboundActivity, completionCallback); + return true; + } + + /// + /// Skips the outbound activity by propagating skipped connections. + /// + private static async ValueTask SkipOutboundActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) + { + return await MaybeScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty, completionCallback); + } + + private static async ValueTask CancelRemainingInboundActivitiesAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity) + { + var flowchart = (Flowchart)flowchartContext.Activity; + var flowGraph = flowchart.GetFlowGraph(flowchartContext); + var ancestorActivities = flowGraph.GetAncestorActivities(outboundActivity); + var inboundActivityExecutionContexts = flowchartContext.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 ValueTask OnScheduleOutcomesAsync(ScheduleActivityOutcomes signal, SignalContext context) @@ -234,15 +364,7 @@ public partial class Flowchart var flowchartContext = context.ReceiverActivityExecutionContext; var schedulingActivityContext = context.SenderActivityExecutionContext; var schedulingActivity = schedulingActivityContext.Activity; - var outcomes = signal.Outcomes; - var outboundConnections = Connections.Where(connection => connection.Source.Activity == schedulingActivity && outcomes.Contains(connection.Source.Port!)).ToList(); - var outboundActivities = outboundConnections.Select(x => x.Target.Activity).ToList(); - - if (outboundActivities.Any()) - { - foreach (var activity in outboundActivities) - await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedCounterBasedLogicAsync); - } + var outcomes = new Outcomes(signal.Outcomes); } private async ValueTask OnCounterFlowActivityCanceledAsync(CancelSignal signal, SignalContext context) @@ -254,6 +376,6 @@ public partial class Flowchart var flowScope = flowchart.GetFlowScope(flowchartContext); // Propagate canceled connections visited count by scheduling with Outcomes.Empty - await flowchart.ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, context.SenderActivityExecutionContext.Activity, Outcomes.Empty); + await MaybeScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, context.SenderActivityExecutionContext.Activity, Outcomes.Empty, OnChildCompletedAsync); } } \ 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 a35408447..8fa1c5d06 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs @@ -1,6 +1,5 @@ using System.ComponentModel; using System.Runtime.CompilerServices; -using Elsa.Extensions; using Elsa.Workflows.Activities.Flowchart.Models; using Elsa.Workflows.Attributes; using Elsa.Workflows.Signals; @@ -40,7 +39,7 @@ public partial class Flowchart : Container /// protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context) { - var startActivity = this.GetStartActivity(context.WorkflowExecutionContext.TriggerActivityId); + var startActivity = GetStartActivity(context); if (startActivity == null) { @@ -52,379 +51,6 @@ public partial class Flowchart : Container await context.ScheduleActivityAsync(startActivity, OnChildCompletedAsync); } -// BEGIN 3.5 - private IActivity? GetStartActivity(ActivityExecutionContext context) - { - // If there's a trigger that triggered this workflow, use that. - var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId; - var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : null; - - if (triggerActivity != null) - return triggerActivity; - - // If an explicit Start activity was provided, use that. - if (Start != null) - return Start; - - // If there is a Start activity on the flowchart, use that. - var startActivity = Activities.FirstOrDefault(x => x is Start); - - if (startActivity != null) - return startActivity; - - // If there's an activity marked as "Can Start Workflow", use that. - var canStartWorkflowActivity = Activities.FirstOrDefault(x => x.GetCanStartWorkflow()); - - if (canStartWorkflowActivity != null) - return canStartWorkflowActivity; - - // If there is a single activity that has no inbound connections, use that. - var root = GetRootActivity(); - - if (root != null) - return root; - - // If no start activity found, return the first activity. - return Activities.FirstOrDefault(); - } - - /// - /// Checks if there is any pending work for the flowchart. - /// - private bool HasPendingWork(ActivityExecutionContext context) - { - var workflowExecutionContext = context.WorkflowExecutionContext; - - // Use HashSet for O(1) lookups - var activityIds = new HashSet(Activities.Select(x => x.Id)); - - // Short circuit evaluation - check running instances first before more expensive scheduler check - if (context.Children.Any(x => activityIds.Contains(x.Activity.Id) && x.Status == ActivityStatus.Running)) - return true; - - // Scheduler check - optimize to avoid repeated LINQ evaluations - var scheduledItems = workflowExecutionContext.Scheduler.List().ToList(); - - return scheduledItems.Any(workItem => - { - var ownerInstanceId = workItem.Owner?.Id; - - if (ownerInstanceId == null) - return false; - - if (ownerInstanceId == context.Id) - return true; - - var ownerContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.Id == ownerInstanceId); - return ownerContext.GetAncestors().Any(x => x == context); - }); - } - - private IActivity? GetRootActivity() - { - // Get the first activity that has no inbound connections. - var query = - from activity in Activities - let inboundConnections = Connections.Any(x => x.Target.Activity == activity) - where !inboundConnections - select activity; - - 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 flowchartContext = context.TargetContext; - var completedActivityContext = context.ChildContext; - var completedActivity = completedActivityContext.Activity; - var result = context.Result; - - // Determine the outcomes from the completed activity - var outcomes = result is Outcomes o ? o : Outcomes.Default; - - await ProcessChildCompletedAsync(flowchartContext, completedActivity, completedActivityContext, outcomes); - } - - private async ValueTask ProcessChildCompletedAsync(ActivityExecutionContext flowchartContext, IActivity completedActivity, ActivityExecutionContext completedActivityContext, Outcomes outcomes) - { - if (flowchartContext.Activity != this) - { - throw new("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; - } - - // Schedule the outbound activities - var flowGraph = GetFlowGraph(flowchartContext); - var flowScope = GetFlowScope(flowchartContext); - var completedActivityExcecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault(BackwardConnectionActivityInput); - bool hasScheduledActivity = await MaybeScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, OnChildCompletedAsync, completedActivityExcecutedByBackwardConnection); - - // If there are not any outbound connections, complete the flowchart activity if there is no other pending work - if (!hasScheduledActivity) - { - await CompleteIfNoPendingWorkAsync(flowchartContext); - } - } - - /// - /// 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. If outcomes is Outcomes.Empty, it indicates - /// that the activity should be skipped - all outbound connections will be visited and treated as - /// not followed. - /// - /// 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. - /// The callback to invoke upon activity completion. - /// Indicates if the completed activity - /// was executed due to a backward connection. - /// True if at least one activity was scheduled; otherwise, false. - private static async ValueTask MaybeScheduleOutboundActivitiesAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity activity, Outcomes outcomes, ActivityCompletionCallback completionCallback, 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($"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)) - { - var connectionFollowed = outcomes.Names.Contains(outboundConnection.Source.Port); - flowScope.RegisterConnectionVisit(outboundConnection, connectionFollowed); - var outboundActivity = outboundConnection.Target.Activity; - - // Determine the scheduling strategy based on connection-type. - if (flowGraph.IsBackwardConnection(outboundConnection, out var backwardConnectionIsValid)) - // Backward connections are scheduled differently - hasScheduledActivity |= await MaybeScheduleBackwardConnectionActivityAsync(flowGraph, flowchartContext, outboundConnection, outboundActivity, connectionFollowed, backwardConnectionIsValid, completionCallback); - else - hasScheduledActivity |= await MaybeScheduleOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundConnection, outboundActivity, completionCallback); - } - - return hasScheduledActivity; - } - - /// - /// Schedules an outbound activity that originates from a backward connection. - /// - private static async ValueTask MaybeScheduleBackwardConnectionActivityAsync(FlowGraph flowGraph, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, bool connectionFollowed, bool backwardConnectionIsValid, ActivityCompletionCallback completionCallback) - { - if (!connectionFollowed) - { - return false; - } - - if (!backwardConnectionIsValid) - { - throw new($"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 = completionCallback, - Input = new Dictionary() { { BackwardConnectionActivityInput, true } } - }; - - await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions); - return true; - } - - /// - /// Determines the merge mode for a given outbound activity. If the outbound activity is a FlowJoin, it retrieves its configured - /// mode. Otherwise, it defaults to FlowJoinMode.WaitAllActive for implicit joins. - /// - private static async ValueTask GetMergeModeAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity) - { - if (outboundActivity is FlowJoin) - { - var outboundActivityExecutionContext = await flowchartContext.WorkflowExecutionContext.CreateActivityExecutionContextAsync(outboundActivity); - return await outboundActivityExecutionContext.EvaluateInputPropertyAsync(x => x.Mode); - } - else - { - // Implicit join case - treat as WaitAllActive - return FlowJoinMode.WaitAllActive; - } - } - - /// - /// Schedules a join activity based on inbound connection statuses. - /// - private static async ValueTask MaybeScheduleOutboundActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, ActivityCompletionCallback completionCallback) - { - FlowJoinMode mode = await GetMergeModeAsync(flowchartContext, outboundActivity); - - return mode switch - { - FlowJoinMode.WaitAll => await MaybeScheduleWaitAllActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback), - FlowJoinMode.WaitAllActive => await MaybeScheduleWaitAllActiveActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback), - FlowJoinMode.WaitAny => await MaybeScheduleWaitAnyActivityAsync(flowGraph, flowScope, flowchartContext, outboundConnection, outboundActivity, completionCallback), - _ => throw new($"Unsupported FlowJoinMode: {mode}"), - }; - } - - /// - /// Determines whether to schedule an activity based on the FlowJoinMode.WaitAll behavior. - /// If all inbound connections were visited, it checks if they were all followed to decide whether to schedule or skip the activity. - /// - private static async ValueTask MaybeScheduleWaitAllActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) - { - if (!flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) - // Not all inbound connections have been visited yet; do not schedule anything yet. - return false; - - if (flowScope.AllInboundConnectionsFollowed(flowGraph, outboundActivity)) - // All inbound connections were followed; schedule the outbound activity. - return await ScheduleOutboundActivityAsync(flowchartContext, outboundActivity, completionCallback); - else - // No inbound connections were followed; skip the outbound activity. - return await SkipOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback); - } - - /// - /// Determines whether to schedule an activity based on the FlowJoinMode.WaitAllActive behavior. - /// If all inbound connections have been visited, it checks if any were followed to decide whether to schedule or skip the activity. - /// - private static async ValueTask MaybeScheduleWaitAllActiveActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) - { - if (!flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) - // Not all inbound connections have been visited yet; do not schedule anything yet. - return false; - - if (flowScope.AnyInboundConnectionsFollowed(flowGraph, outboundActivity)) - // At least one inbound connection was followed; schedule the outbound activity. - return await ScheduleOutboundActivityAsync(flowchartContext, outboundActivity, completionCallback); - else - // No inbound connections were followed; skip the outbound activity. - return await SkipOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback); - } - - /// - /// Determines whether to schedule an activity based on the FlowJoinMode.WaitAny behavior. - /// If any inbound connection has been followed, it schedules the activity and cancels remaining inbound activities. - /// If a subsequent inbound connection is followed after the activity has been scheduled, it ignores it. - /// - private static async ValueTask MaybeScheduleWaitAnyActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, ActivityCompletionCallback completionCallback) - { - if (flowScope.ShouldIgnoreConnection(outboundConnection, outboundActivity)) - // Ignore the connection if the outbound activity has already completed (JoinAny scenario) - return false; - - if (flowchartContext.WorkflowExecutionContext.Scheduler.List().Any(workItem => workItem.Owner == flowchartContext && workItem.Activity == outboundActivity)) - // Ignore the connection if the outbound activity is already scheduled - return false; - - if (flowScope.AnyInboundConnectionsFollowed(flowGraph, outboundActivity)) - { - // An inbound connection has been followed; cancel remaining inbound activities - await CancelRemainingInboundActivitiesAsync(flowchartContext, outboundActivity); - - // This is the first inbound connection followed; schedule the outbound activity - return await ScheduleOutboundActivityAsync(flowchartContext, outboundActivity, completionCallback); - } - - if (flowScope.AllInboundConnectionsVisited(flowGraph, outboundActivity)) - // All inbound connections have been visited without any being followed; skip the outbound activity - return await SkipOutboundActivityAsync(flowGraph, flowScope, flowchartContext, outboundActivity, completionCallback); - - // No inbound connections have been followed yet; do not schedule anything yet. - return false; - } - - /// - /// Schedules the outbound activity. - /// - private static async ValueTask ScheduleOutboundActivityAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) - { - await flowchartContext.ScheduleActivityAsync(outboundActivity, completionCallback); - return true; - } - - /// - /// Skips the outbound activity by propagating skipped connections. - /// - private static async ValueTask SkipOutboundActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback) - { - return await MaybeScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty, completionCallback); - } - - private static async ValueTask CancelRemainingInboundActivitiesAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity) - { - var flowchart = (Flowchart)flowchartContext.Activity; - var flowGraph = flowchart.GetFlowGraph(flowchartContext); - var ancestorActivities = flowGraph.GetAncestorActivities(outboundActivity); - var inboundActivityExecutionContexts = flowchartContext.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) - { - var hasPendingWork = HasPendingWork(context); - - if (!hasPendingWork) - { - var hasFaultedActivities = context.Children.Any(x => x.Status == ActivityStatus.Faulted); - - if (!hasFaultedActivities) - { - await context.CompleteActivityAsync(); - } - } - } - - private async ValueTask OnScheduleOutcomesAsync(ScheduleActivityOutcomes signal, SignalContext context) - { - var flowchartContext = context.ReceiverActivityExecutionContext; - var schedulingActivityContext = context.SenderActivityExecutionContext; - var schedulingActivity = schedulingActivityContext.Activity; - var outcomes = new Outcomes(signal.Outcomes); - - await ProcessChildCompletedAsync(flowchartContext, schedulingActivity, schedulingActivityContext, outcomes); - } - -// END 3.5 - private async ValueTask OnScheduleChildActivityAsync(ScheduleChildActivity signal, SignalContext context) { var flowchartContext = context.ReceiverActivityExecutionContext; @@ -466,9 +92,16 @@ public partial class Flowchart : Container private async Task CompleteIfNoPendingWorkAsync(ActivityExecutionContext context) { - var hasPendingWork = context.HasPendingWork(); + var hasPendingWork = HasPendingWork(context); if (!hasPendingWork) - await context.CompleteActivityAsync(); + { + var hasFaultedActivities = context.Children.Any(x => x.Status == ActivityStatus.Faulted); + + if (!hasFaultedActivities) + { + await context.CompleteActivityAsync(); + } + } } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs index 72f50335a..b0c896410 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs @@ -72,40 +72,30 @@ public static class ExpressionExecutionContextExtensions /// public bool TryGetWorkflowExecutionContext(out WorkflowExecutionContext workflowExecutionContext) => context.TransientProperties.TryGetValue(WorkflowExecutionContextKey, out workflowExecutionContext!); - /// - /// Returns the of the specified - /// - public static WorkflowExecutionContext GetWorkflowExecutionContext(this ExpressionExecutionContext context) - { - return context.TransientProperties.TryGetValue(WorkflowExecutionContextKey, out var value) - ? (WorkflowExecutionContext)value - : throw new InvalidOperationException("WorkflowExecutionContext not found. This value exists only on activity execution contexts."); - } + /// + /// Returns the of the specified + /// + public WorkflowExecutionContext GetWorkflowExecutionContext() + { + return context.TransientProperties.TryGetValue(WorkflowExecutionContextKey, out var value) + ? (WorkflowExecutionContext)value + : throw new InvalidOperationException("WorkflowExecutionContext not found. This value exists only on activity execution contexts."); + } - /// - /// Returns the of the specified - /// - public static ActivityExecutionContext GetActivityExecutionContext(this ExpressionExecutionContext context) - { - return context.TransientProperties.TryGetValue(ActivityExecutionContextKey, out var value) - ? (ActivityExecutionContext)value - : throw new InvalidOperationException("ActivityExecutionContext not found. This value exists only on activity execution contexts."); - } + /// + /// Returns the of the specified + /// + public ActivityExecutionContext GetActivityExecutionContext() + { + return context.TransientProperties.TryGetValue(ActivityExecutionContextKey, out var value) + ? (ActivityExecutionContext)value + : throw new InvalidOperationException("ActivityExecutionContext not found. This value exists only on activity execution contexts."); + } - /// - /// Returns the of the specified - /// - public static bool TryGetActivityExecutionContext(this ExpressionExecutionContext context, out ActivityExecutionContext activityExecutionContext) => context.TransientProperties.TryGetValue(ActivityExecutionContextKey, out activityExecutionContext!); - - /// - /// Returns the of the specified - /// - public static IActivity GetActivity(this ExpressionExecutionContext context) - { - return context.TransientProperties.TryGetValue(ActivityKey, out var value) - ? (IActivity)value - : throw new InvalidOperationException("Activity not found. This value exists only on activity execution contexts."); - } + /// + /// Returns the of the specified + /// + public bool TryGetActivityExecutionContext(out ActivityExecutionContext activityExecutionContext) => context.TransientProperties.TryGetValue(ActivityExecutionContextKey, out activityExecutionContext!); /// /// Returns the of the specified @@ -202,7 +192,7 @@ public static class ExpressionExecutionContextExtensions var variable = context.GetVariable(name); if (variable == null) - return CreateVariable(context, name, value, configure: configure); + return context.CreateVariable(name, value, configure: configure); // Get the context where the variable is defined. var contextWithVariable = context.FindContextContainingBlock(variable.Id) ?? context; @@ -307,7 +297,7 @@ public static class ExpressionExecutionContextExtensions /// Gets all variables names in scope. /// public IEnumerable GetVariableNamesInScope() => - EnumerateVariablesInScope(context) + context.EnumerateVariablesInScope() .Select(x => x.Name) .Where(x => !string.IsNullOrWhiteSpace(x)) .Distinct(); @@ -316,7 +306,7 @@ public static class ExpressionExecutionContextExtensions /// Gets all variables in scope. /// public IEnumerable GetVariablesInScope() => - EnumerateVariablesInScope(context) + context.EnumerateVariablesInScope() .Where(x => !string.IsNullOrWhiteSpace(x.Name)) .DistinctBy(x => x.Name); @@ -325,7 +315,7 @@ public static class ExpressionExecutionContextExtensions /// public void SetVariableInScope(string variableName, object? value) { - var q = from v in EnumerateVariablesInScope(context) + var q = from v in context.EnumerateVariablesInScope() where v.Name == variableName where v.TryGet(context, out _) select v; @@ -336,7 +326,7 @@ public static class ExpressionExecutionContextExtensions variable.Set(context, value); if (variable == null) - CreateVariable(context, variableName, value); + context.CreateVariable(variableName, value); } /// @@ -402,7 +392,6 @@ public static class ExpressionExecutionContextExtensions return serializerOptions; } - /// extension(ExpressionExecutionContext context) { /// diff --git a/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs b/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs index a46e1aa4b..556e89251 100644 --- a/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs +++ b/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs @@ -1,5 +1,5 @@ +using Elsa.Expressions.JavaScript.Contracts; using Elsa.Expressions.Models; -using Elsa.JavaScript.Contracts; using Elsa.Testing.Shared; using Microsoft.Extensions.DependencyInjection; using Xunit.Abstractions; diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs index 5ddb29eb3..1cb751446 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs @@ -369,6 +369,7 @@ public class FlowchartNextActivityTests // F public async Task JoinBehavesCorrectly(bool decisionResult, FlowJoinMode joinMode, string[] expectedLines) { + Flowchart.UseTokenFlow = false; var workflow = new TestWorkflow(workflowBuilder => { var start = new Start() { Id = "Start" }; @@ -423,5 +424,6 @@ public class FlowchartNextActivityTests var lines = _capturingTextWriter.Lines.ToList(); Assert.Equal(WorkflowSubStatus.Finished, result.WorkflowState.SubStatus); Assert.Equal(expectedLines, lines); + Flowchart.UseTokenFlow = true; } } \ No newline at end of file