diff --git a/.github/copilot-instructions.md b/.github/copilot-instructions.md index 41f39a82b..2cb307856 100644 --- a/.github/copilot-instructions.md +++ b/.github/copilot-instructions.md @@ -2,15 +2,15 @@ ## Repository Overview -**Elsa Workflows** is a powerful .NET workflow library that enables workflow execution within any .NET application. This is version 3.0, supporting .NET 9.0 and providing both a visual designer and programmatic workflow definition capabilities. +**Elsa Workflows** is a powerful .NET workflow library that enables workflow execution within any .NET application. This is version 3.0, supporting .NET 8.0, .NET 9.0 and .NET 10.0 and providing both a visual designer (from a different repository, elsa-studio) and programmatic workflow definition capabilities. ### Key Statistics -- **Language**: C# (.NET 9.0) +- **Language**: C# (.NET 10.0) - **Architecture**: Modular library with 104+ projects - **Code Size**: ~3,500 C# files across modules - **License**: MIT - **Build System**: NUKE build automation -- **Target Frameworks**: .NET 9.0 (primary) +- **Target Frameworks**: .NET 10.0 (primary) ## High-Level Architecture @@ -19,9 +19,6 @@ src/ ├── apps/ # Reference applications (5 projects) │ ├── Elsa.Server.Web # Workflow server only -│ ├── Elsa.ServerAndStudio.Web # Combined server + studio -│ ├── Elsa.Studio.Web # Studio web interface -│ ├── ElsaStudioWebAssembly # Studio WebAssembly app │ └── Elsa.Server.LoadBalancer # Load balancer ├── common/ # Shared libraries (8 projects) ├── modules/ # Core functionality modules (70+ projects) @@ -46,7 +43,7 @@ docker/ # Docker configurations - **Elsa.Workflows.Runtime**: Workflow execution runtime - **Elsa.Workflows.Api**: RESTful API for workflow management - **Elsa.Workflows.Management**: Workflow definition management -- **Elsa modules**: Specialized functionality (HTTP, email, scheduling, etc.) +- **Elsa modules**: Specialized functionality (HTTP, persistence, scheduling, etc.) ## Build Instructions @@ -54,20 +51,6 @@ docker/ # Docker configurations - **.NET 10.0 SDK** - **Build time**: Initial restore ~1-2 minutes, full compile ~5-10 minutes -### Critical Build Information - -⚠️ **IMPORTANT**: The repository has external dependencies that may cause build failures: - -1. **External NuGet Feeds**: Some projects depend on packages from: - - `https://f.feedz.io/elsa-workflows/elsa-3/nuget/index.json` (Elsa Studio packages) - - `https://f.feedz.io/sfmskywalker/webhooks-core/nuget/index.json` (Webhooks packages) - -2. **Build Failure Workarounds**: - - Studio apps (`Elsa.Studio.Web`, `ElsaStudioWebAssembly`, `Elsa.ServerAndStudio.Web`) depend on prebuilt studio packages that may not be accessible - - Server app (`Elsa.Server.Web`) depends on WebhooksCore package that may not be accessible - - Core workflow functionality can be built independently - - Some test projects may fail due to missing external packages - ### Build Commands **Primary build script**: `./build.sh` (Linux/macOS) or `.\build.cmd` (Windows) diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index 02d18e5db..14458a322 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -70,6 +70,17 @@ jobs: dotnet test "$project" --configuration Release --no-build --logger "GitHubActions;report-warnings=false" /p:CollectCoverage=true done + - name: Dump docker logs on failure + if: failure() + run: | + echo "=== Docker containers ===" + docker ps -a || true + echo "=== Docker logs ===" + for container in $(docker ps -aq); do + echo "--- Logs for container $container ---" + docker logs "$container" 2>&1 | tail -100 || true + done + - name: Install ReportGenerator run: dotnet tool install -g dotnet-reportgenerator-globaltool diff --git a/Directory.Packages.props b/Directory.Packages.props index ecf434e00..ce19a6446 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -43,38 +43,38 @@ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + @@ -126,7 +126,7 @@ - + @@ -151,20 +151,21 @@ - + - + - + + diff --git a/NuGet.Config b/NuGet.Config index 31400bd46..54f245bd3 100644 --- a/NuGet.Config +++ b/NuGet.Config @@ -4,21 +4,10 @@ - - - - - - - - - - - \ No newline at end of file diff --git a/docker/ElsaServer-Datadog.Dockerfile b/docker/ElsaServer-Datadog.Dockerfile index b7f1a3f4b..9305f71e8 100644 --- a/docker/ElsaServer-Datadog.Dockerfile +++ b/docker/ElsaServer-Datadog.Dockerfile @@ -1,7 +1,7 @@ # Version: 1 # Description: Dockerfile for building and running Elsa Server with Datadog and OpenTelemetry auto-instrumentation -FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:9.0-bookworm-slim AS build +FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:10.0 AS build WORKDIR /source # Copy sources. @@ -15,10 +15,10 @@ RUN dotnet restore "./src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj" # Build and publish (UseAppHost=false creates platform independent binaries). WORKDIR /source/src/bundles/Elsa.Server.Web RUN dotnet build "Elsa.Server.Web.csproj" -c Release -o /app/build -RUN dotnet publish "Elsa.Server.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net9.0 +RUN dotnet publish "Elsa.Server.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net10.0 # Move binaries into smaller base image. -FROM mcr.microsoft.com/dotnet/aspnet:9.0-bookworm-slim AS base +FROM mcr.microsoft.com/dotnet/aspnet:10.0 AS base WORKDIR /app COPY --from=build /app/publish ./ diff --git a/docker/ElsaServer.Dockerfile b/docker/ElsaServer.Dockerfile index 2695ca4de..c194e4253 100644 --- a/docker/ElsaServer.Dockerfile +++ b/docker/ElsaServer.Dockerfile @@ -1,4 +1,4 @@ -FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:9.0-bookworm-slim AS build +FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:10.0 AS build WORKDIR /source # copy sources. @@ -12,10 +12,10 @@ RUN dotnet restore "./src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj" # build and publish (UseAppHost=false creates platform independent binaries). WORKDIR /source/src/apps/Elsa.Server.Web RUN dotnet build "Elsa.Server.Web.csproj" -c Release -o /app/build -RUN dotnet publish "Elsa.Server.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net9.0 +RUN dotnet publish "Elsa.Server.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net10.0 # move binaries into smaller base image. -FROM mcr.microsoft.com/dotnet/aspnet:9.0-bookworm-slim AS base +FROM mcr.microsoft.com/dotnet/aspnet:10.0 AS base WORKDIR /app COPY --from=build /app/publish ./ diff --git a/docker/ElsaServerAndStudio.Dockerfile b/docker/ElsaServerAndStudio.Dockerfile index aaadba910..726904cd0 100644 --- a/docker/ElsaServerAndStudio.Dockerfile +++ b/docker/ElsaServerAndStudio.Dockerfile @@ -1,4 +1,4 @@ -FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:9.0-bookworm-slim AS build +FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:10.0 AS build WORKDIR /source # copy sources. @@ -13,10 +13,10 @@ RUN dotnet restore "./src/apps/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web # build and publish (UseAppHost=false creates platform independent binaries). WORKDIR /source/src/apps/Elsa.ServerAndStudio.Web RUN dotnet build "Elsa.ServerAndStudio.Web.csproj" -c Release -o /app/build -RUN dotnet publish "Elsa.ServerAndStudio.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net9.0 +RUN dotnet publish "Elsa.ServerAndStudio.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net10.0 # move binaries into smaller base image. -FROM mcr.microsoft.com/dotnet/aspnet:9.0-bookworm-slim AS base +FROM mcr.microsoft.com/dotnet/aspnet:10.0 AS base WORKDIR /app COPY --from=build /app/publish ./ diff --git a/docker/ElsaStudio.Dockerfile b/docker/ElsaStudio.Dockerfile index 6e76acb3d..34b69c0c9 100644 --- a/docker/ElsaStudio.Dockerfile +++ b/docker/ElsaStudio.Dockerfile @@ -1,4 +1,4 @@ -FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:9.0-bookworm-slim AS build +FROM --platform=$BUILDPLATFORM mcr.microsoft.com/dotnet/sdk:10.0 AS build WORKDIR /source # copy sources. @@ -13,10 +13,10 @@ RUN dotnet restore "./src/apps/Elsa.Studio.Web/Elsa.Studio.Web.csproj" # build and publish (UseAppHost=false creates platform independent binaries). WORKDIR /source/src/apps/Elsa.Studio.Web RUN dotnet build "Elsa.Studio.Web.csproj" -c Release -o /app/build -RUN dotnet publish "Elsa.Studio.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net9.0 +RUN dotnet publish "Elsa.Studio.Web.csproj" -c Release -o /app/publish /p:UseAppHost=false --no-restore -f net10.0 # move binaries into smaller base image. -FROM mcr.microsoft.com/dotnet/aspnet:9.0-bookworm-slim AS base +FROM mcr.microsoft.com/dotnet/aspnet:10.0 AS base WORKDIR /app COPY --from=build /app/publish ./ diff --git a/src/modules/Elsa.Http/Activities/HttpEndpoint.cs b/src/modules/Elsa.Http/Activities/HttpEndpoint.cs index 7b550d6fb..e62d553ab 100644 --- a/src/modules/Elsa.Http/Activities/HttpEndpoint.cs +++ b/src/modules/Elsa.Http/Activities/HttpEndpoint.cs @@ -1,4 +1,4 @@ -using System.Runtime.CompilerServices; +using System.Runtime.CompilerServices; using System.Text.Json; using Elsa.Expressions.Models; using Elsa.Extensions; @@ -182,7 +182,7 @@ public class HttpEndpoint : Trigger { var path = Path.Get(context); var methods = SupportedMethods.GetOrDefault(context) ?? new List { HttpMethods.Get }; - await context.WaitForHttpRequestAsync(path, methods, OnResumeAsync); + await context.WaitForHttpRequestAsync(path, methods, OnResumeAsync, Elsa.Http.HttpStimulusNames.HttpEndpoint); } private async ValueTask OnResumeAsync(ActivityExecutionContext context) diff --git a/src/modules/Elsa.Http/Activities/HttpEndpointBase.cs b/src/modules/Elsa.Http/Activities/HttpEndpointBase.cs index 9257a73c1..8d6a708ce 100644 --- a/src/modules/Elsa.Http/Activities/HttpEndpointBase.cs +++ b/src/modules/Elsa.Http/Activities/HttpEndpointBase.cs @@ -22,8 +22,8 @@ public abstract class HttpEndpointBase : Trigger protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { - var options = GetOptions(); - await context.WaitForHttpRequestAsync(options, HttpRequestReceivedAsync); + var options = GetOptions(); + await context.WaitForHttpRequestAsync(options, HttpRequestReceivedAsync, Elsa.Http.HttpStimulusNames.HttpEndpoint); } protected override IEnumerable GetTriggerPayloads(TriggerIndexingContext context) diff --git a/src/modules/Elsa.Http/Extensions/HttpEndpointActivityExecutionContextExtensions.cs b/src/modules/Elsa.Http/Extensions/HttpEndpointActivityExecutionContextExtensions.cs index ca62a0782..4c56c3cd3 100644 --- a/src/modules/Elsa.Http/Extensions/HttpEndpointActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Http/Extensions/HttpEndpointActivityExecutionContextExtensions.cs @@ -9,27 +9,27 @@ namespace Elsa.Http.Extensions; public static class HttpEndpointActivityExecutionContextExtensions { -public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, string path, string method, ExecuteActivityDelegate? callback = null) +public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, string path, string method, ExecuteActivityDelegate? callback = null, string? bookmarkName = null) { var options = new HttpEndpointOptions { Path = path, Methods = [method] }; - await WaitForHttpRequestAsync(context, options, callback); + await WaitForHttpRequestAsync(context, options, callback, bookmarkName); } -public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, string path, IEnumerable methods, ExecuteActivityDelegate? callback = null) +public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, string path, IEnumerable methods, ExecuteActivityDelegate? callback = null, string? bookmarkName = null) { var options = new HttpEndpointOptions { Path = path, Methods = methods.ToList() }; - await WaitForHttpRequestAsync(context, options, callback); + await WaitForHttpRequestAsync(context, options, callback, bookmarkName); } - public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, HttpEndpointOptions options, ExecuteActivityDelegate? callback = null) + public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, HttpEndpointOptions options, ExecuteActivityDelegate? callback = null, string? bookmarkName = null) { var path = options.Path; if (path.Contains("//")) @@ -38,7 +38,8 @@ public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionCont var expressionExecutionContext = context.ExpressionExecutionContext; if (!context.IsTriggerOfWorkflow()) { - context.CreateBookmarks(expressionExecutionContext.GetHttpEndpointStimuli(options), includeActivityInstanceId: false, callback: callback); + var name = bookmarkName ?? Elsa.Http.HttpStimulusNames.HttpEndpoint; + context.CreateBookmarks(expressionExecutionContext.GetHttpEndpointStimuli(options), includeActivityInstanceId: false, bookmarkName: name, callback: callback); return; } 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 c41419322..e1abb04de 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs @@ -27,7 +27,7 @@ public class FlowJoin : Activity, IJoinNode /// The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All). /// [Input( - Description = "The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All).", + Description = "The join mode determines whether this activity should continue as soon as one inbound path comes in (WaitAny), or once all inbound paths have executed (WaitAll). To wait for all activated inbound paths, set the mode to WaitAllActive.", DefaultValue = FlowJoinMode.WaitAny, UIHint = InputUIHints.DropDown )] @@ -40,19 +40,7 @@ public class FlowJoin : Activity, IJoinNode if (context.ParentActivityExecutionContext != null) await context.ParentActivityExecutionContext.CancelInboundAncestorsAsync(this); + // The join behavior is handled by Flowchart, so we can simply complete the activity here. await context.CompleteActivityAsync(); } - - protected override bool CanExecute(ActivityExecutionContext context) - { - if(Flowchart.UseTokenFlow) - return true; - - return context.Get(Mode) switch - { - FlowJoinMode.WaitAny => true, - FlowJoinMode.WaitAll => Flowchart.CanWaitAllProceed(context), - _ => true - }; - } } \ No newline at end of file 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 6c72282e7..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.Workflows.Activities.Flowchart.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) { @@ -93,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/Activities/Flowchart/Models/FlowJoinMode.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowJoinMode.cs index cd23b523b..7372af7c8 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowJoinMode.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowJoinMode.cs @@ -3,5 +3,6 @@ namespace Elsa.Workflows.Activities.Flowchart.Models; public enum FlowJoinMode { WaitAll, - WaitAny + WaitAllActive, + WaitAny, } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs index 742fd3222..8d0155fc1 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs @@ -86,7 +86,7 @@ public class FlowScope /// The flow graph containing connections. /// The activity to check. /// True if any inbound connection has been followed, otherwise false. - public bool HasFollowedInboundConnection(FlowGraph flowGraph, IActivity activity) + public bool AnyInboundConnectionsFollowed(FlowGraph flowGraph, IActivity activity) { var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity); var outboundActivityVisitCount = GetActivityVisitCount(activity); @@ -95,6 +95,21 @@ public class FlowScope && forwardInboundConnections.Any(c => GetConnectionVisitCount(c) == maxConnectionVisitCount && GetConnectionLastVisitFollowed(c)); } + /// + /// Determines whether all inbound connection to the specified activity has been followed. + /// + /// The flow graph containing connections. + /// The activity to check. + /// True if all inbound connection has been followed, otherwise false. + public bool AllInboundConnectionsFollowed(FlowGraph flowGraph, IActivity activity) + { + var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity); + var outboundActivityVisitCount = GetActivityVisitCount(activity); + var maxConnectionVisitCount = forwardInboundConnections.Max(GetConnectionVisitCount); + return maxConnectionVisitCount > outboundActivityVisitCount + && forwardInboundConnections.All(c => GetConnectionVisitCount(c) == maxConnectionVisitCount && GetConnectionLastVisitFollowed(c)); + } + /// /// Determines whether a connection should be ignored based on visit counts. /// diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs index 0b2f642ad..a1bb1883d 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs @@ -371,13 +371,15 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable /// The payloads to create bookmarks for. /// An optional callback that is invoked when the bookmark is resumed. /// Whether or not the activity instance ID should be included in the bookmark payload. - public void CreateBookmarks(IEnumerable payloads, ExecuteActivityDelegate? callback = null, bool includeActivityInstanceId = true) + /// An optional name to use for the bookmark. Defaults to the activity type. + public void CreateBookmarks(IEnumerable payloads, ExecuteActivityDelegate? callback = null, bool includeActivityInstanceId = true, string? bookmarkName = null) { foreach (var payload in payloads) CreateBookmark(new() { Stimulus = payload, Callback = callback, + BookmarkName = bookmarkName, IncludeActivityInstanceId = includeActivityInstanceId }); } diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs index b3f149017..b0c896410 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs @@ -8,9 +8,7 @@ using Elsa.Workflows; using Elsa.Workflows.Activities; using Elsa.Workflows.Memory; using Elsa.Workflows.Models; -using Elsa.Workflows.Options; using Humanizer; -using Microsoft.Extensions.Options; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; @@ -77,15 +75,25 @@ public static class ExpressionExecutionContextExtensions /// /// Returns the of the specified /// - public WorkflowExecutionContext GetWorkflowExecutionContext() => (WorkflowExecutionContext)context.TransientProperties[WorkflowExecutionContextKey]; + 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 ActivityExecutionContext GetActivityExecutionContext() => (ActivityExecutionContext)context.TransientProperties[ActivityExecutionContextKey]; + 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 + /// Returns the of the specified /// public bool TryGetActivityExecutionContext(out ActivityExecutionContext activityExecutionContext) => context.TransientProperties.TryGetValue(ActivityExecutionContextKey, out activityExecutionContext!); @@ -184,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; @@ -289,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(); @@ -298,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); @@ -307,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; @@ -318,7 +326,7 @@ public static class ExpressionExecutionContextExtensions variable.Set(context, value); if (variable == null) - CreateVariable(context, variableName, value); + context.CreateVariable(variableName, value); } /// @@ -384,7 +392,6 @@ public static class ExpressionExecutionContextExtensions return serializerOptions; } - /// extension(ExpressionExecutionContext context) { /// diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 501a67276..5835d852a 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -118,6 +118,14 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module) /// public Func DispatchWorkflowCommandHandler { get; set; } = sp => sp.GetRequiredService(); + /// + /// A factory that instantiates an . + /// + public Func WorkflowResumer { get; set; } = sp => sp.GetRequiredService(); + + /// + /// A factory that instantiates an . + /// public Func BookmarkQueueWorker { get; set; } = sp => sp.GetRequiredService(); /// @@ -249,7 +257,7 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module) .AddSingleton(BackgroundActivityScheduler) .AddSingleton() .AddSingleton() - .AddScoped() + .AddScoped(BookmarkQueueWorker) .AddScoped() .AddScoped() .AddScoped() @@ -275,7 +283,10 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module) .AddScoped() .AddScoped() .AddScoped() - .AddScoped() + .AddScoped(WorkflowResumer) + .AddScoped() + .AddScoped(BookmarkQueueWorker) + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/test/component/Directory.Build.props b/test/component/Directory.Build.props index e6ff014e4..fb63b368e 100644 --- a/test/component/Directory.Build.props +++ b/test/component/Directory.Build.props @@ -1,14 +1,15 @@ - + - - 40 - + + 40 + - - - - + + + + + \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index 7d4751127..dbfd8dcec 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -22,6 +22,7 @@ + diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Abstractions/AppComponentTest.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Abstractions/AppComponentTest.cs index b0db161d7..d69ed62f1 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Abstractions/AppComponentTest.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Abstractions/AppComponentTest.cs @@ -13,9 +13,7 @@ public abstract class AppComponentTest(App app) : IDisposable void IDisposable.Dispose() { - // Disposing the Scope here and in other places where it is created somehow seems to cause the test runner to hang when running other test projects. - // Let's comment it out for the time being. - //Scope.Dispose(); + Scope.Dispose(); OnDispose(); } diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/Infrastructure.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/Infrastructure.cs index 497951b60..a254bbe89 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/Infrastructure.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/Infrastructure.cs @@ -1,3 +1,4 @@ +using Testcontainers.MsSql; using Testcontainers.PostgreSql; using Testcontainers.RabbitMq; @@ -5,11 +6,10 @@ namespace Elsa.Workflows.ComponentTests.Fixtures; public class Infrastructure : IAsyncLifetime { - public readonly PostgreSqlContainer DbContainer = new PostgreSqlBuilder() - .WithImage("postgres:latest") - .WithDatabase("elsa") - .WithUsername("postgres") - .WithPassword("postgres") + //public readonly PostgreSqlContainer DbContainer = new PostgreSqlBuilder().Build(); + + public readonly MsSqlContainer DbContainer = new MsSqlBuilder() + //.WithImage("mcr.microsoft.com/mssql/server:2025-GA-ubuntu") .Build(); public readonly RabbitMqContainer RabbitMqContainer = new RabbitMqBuilder() diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs index aeae8967e..1a10ecba1 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -78,15 +78,27 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl var workflowsDirectory = Path.Join(workflowsDirectorySegments); return StorageFactory.Blobs.DirectoryFiles(workflowsDirectory); }); - elsa.UseIdentity(identity => identity.UseEntityFrameworkCore(ef => ef.UsePostgreSql(dbConnectionString))); + elsa.UseIdentity(identity => identity.UseEntityFrameworkCore(ef => + { + //ef.UsePostgreSql(dbConnectionString); + ef.UseSqlServer(dbConnectionString); + })); elsa.UseWorkflowManagement(management => { - management.UseEntityFrameworkCore(ef => ef.UsePostgreSql(dbConnectionString)); + management.UseEntityFrameworkCore(ef => + { + //ef.UsePostgreSql(dbConnectionString); + ef.UseSqlServer(dbConnectionString); + }); management.UseCache(); }); elsa.UseWorkflowRuntime(runtime => { - runtime.UseEntityFrameworkCore(ef => ef.UsePostgreSql(dbConnectionString)); + runtime.UseEntityFrameworkCore(ef => + { + //ef.UsePostgreSql(dbConnectionString); + ef.UseSqlServer(dbConnectionString); + }); runtime.UseCache(); runtime.UseDistributedRuntime(); }); @@ -100,7 +112,11 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl }); elsa.UseAlterations(alterations => { - alterations.UseEntityFrameworkCore(ef => ef.UsePostgreSql(dbConnectionString)); + alterations.UseEntityFrameworkCore(ef => + { + //ef.UsePostgreSql(dbConnectionString); + ef.UseSqlServer(dbConnectionString); + }); }); elsa.UseHttp(http => { diff --git a/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs b/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs new file mode 100644 index 000000000..556e89251 --- /dev/null +++ b/test/integration/Elsa.JavaScript.IntegrationTests/GuidTests.cs @@ -0,0 +1,33 @@ +using Elsa.Expressions.JavaScript.Contracts; +using Elsa.Expressions.Models; +using Elsa.Testing.Shared; +using Microsoft.Extensions.DependencyInjection; +using Xunit.Abstractions; + +namespace Elsa.JavaScript.IntegrationTests; + +public class GuidTests +{ + private readonly IJavaScriptEvaluator _evaluator; + private readonly IServiceProvider _serviceProvider; + + public GuidTests(ITestOutputHelper testOutputHelper) + { + _serviceProvider = new TestApplicationBuilder(testOutputHelper).Build(); + _evaluator = _serviceProvider.GetRequiredService(); + } + + [Fact] + public async Task NewGuidReturnsGuid() + { + //Setup + var script = "newGuid()"; + var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, new()); + + //Act + var result = (Guid)(await _evaluator.EvaluateAsync(script, typeof(Guid), expressionExecutionContext))!; + + //Assert + Assert.IsType(result); + } +} \ No newline at end of file diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs index dd25a4f73..1cb751446 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FlowchartNextActivity/Tests.cs @@ -198,8 +198,9 @@ public class FlowchartNextActivityTests [Theory(DisplayName = "Flowchart with a Join activity executed multiple times")] [InlineData(FlowJoinMode.WaitAll)] + [InlineData(FlowJoinMode.WaitAllActive)] [InlineData(FlowJoinMode.WaitAny)] - public async Task WaitAnyLoopTest(FlowJoinMode joinMode) + public async Task JoinLoopTest(FlowJoinMode joinMode) { var workflow = new TestWorkflow(workflowBuilder => { @@ -272,8 +273,9 @@ public class FlowchartNextActivityTests [Theory(DisplayName = "Flowchart with a Join activity executed multiple times, bug 6479")] [InlineData(FlowJoinMode.WaitAll)] + [InlineData(FlowJoinMode.WaitAllActive)] [InlineData(FlowJoinMode.WaitAny)] - public async Task WaitLoopBug6479Test(FlowJoinMode joinMode) + public async Task JoinLoopBug6479Test(FlowJoinMode joinMode) { var workflow = new TestWorkflow(workflowBuilder => { @@ -339,4 +341,89 @@ public class FlowchartNextActivityTests "A", "A", "A", "B" }, lines); } + + [Theory(DisplayName = "Flowchart Join behaves correctly")] + [InlineData(false, FlowJoinMode.WaitAll, new[] { "A", "B", "C", "D", "F" })] // "E" is not scheduled because join has an unfollowed inbound connection + [InlineData(false, FlowJoinMode.WaitAllActive, new[] { "A", "B", "C", "D", "E", "F" })] // "E" gets scheduled by join with an unfollowed inbound connection + [InlineData(false, FlowJoinMode.WaitAny, new[] { "A", "B", "C", "D", "E", "F" })] // "E" only scheduled once + [InlineData(true, FlowJoinMode.WaitAll, new[] { "A", "B", "C", "E", "F" })] // all Join inbound connections followed, "E" gets scheduled + [InlineData(true, FlowJoinMode.WaitAllActive, new[] { "A", "B", "C", "E", "F" })] // all Join inbound connections followed, "E" gets scheduled + [InlineData(true, FlowJoinMode.WaitAny, new[] { "A", "B", "C", "E", "F" })] // "E" only scheduled once + // Start + // / | \ + // / | \ + // A B C + // / | | + // | | | + // Decision | | + // / \ | | + // | (true) | / + // | \ | / + // (false) \ | / + // | Join + // D | + // \ E + // \ / + // \ / + // \ / + // 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" }; + var a = new WriteLine("A") { Id = "WriteLineA" }; + var b = new WriteLine("B") { Id = "WriteLineB" }; + var c = new WriteLine("C") { Id = "WriteLineC" }; + var decision = new FlowDecision() + { + Condition = new(new Literal(decisionResult)) + }; + var d = new WriteLine("D") { Id = "WriteLineD" }; + var join = new FlowJoin() + { + Mode = new(joinMode) + }; + var e = new WriteLine("E") { Id = "WriteLineE" }; + var f = new WriteLine("F") { Id = "WriteLineF" }; + + workflowBuilder.Root = new Flowchart + { + Activities = + { + start, + a, + b, + c, + decision, + d, + join, + e, + f, + }, + Connections = + { + new(start, a), + new(start, b), + new(start, c), + new(a, decision), + new(new Endpoint(decision, "True"), new Endpoint(join)), + new(new Endpoint(decision, "False"), new Endpoint(d)), + new(b, join), + new(c, join), + new(d, f), + new(join, e), + new(e, f), + } + }; + }); + + await _services.PopulateRegistriesAsync(); + var result = await _workflowRunner.RunAsync(workflow); + 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