Merge remote-tracking branch 'origin/develop/3.6.0' into develop/3.6.0

This commit is contained in:
Sipke Schoorstra 2025-12-13 21:33:05 +01:00
commit 6f64739d83
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
26 changed files with 533 additions and 260 deletions

View file

@ -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)

View file

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

View file

@ -43,38 +43,38 @@
</ItemGroup>
<!-- .NET 10 specific package versions -->
<ItemGroup Condition="'$(TargetFramework)' == 'net10.0'">
<PackageVersion Include="Microsoft.AspNetCore.Authorization" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.Components" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.DevServer" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.Server" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.DataProtection.Abstractions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.Mvc.Testing" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Data.Sqlite" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Data.Sqlite.Core" Version="10.0.0"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore" Version="10.0.0"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.Design" Version="10.0.0"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.Relational" Version="10.0.0"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.Sqlite" Version="10.0.0"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.SqlServer" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Caching.Abstractions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Caching.Memory" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Configuration" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Configuration.Abstractions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Configuration.Json" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.DependencyModel" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Hosting.Abstractions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Http" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Http.Polly" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Http.Resilience" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Options" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Options.ConfigurationExtensions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Resilience" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Logging" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Logging.Abstractions" Version="10.0.0"/>
<PackageVersion Include="Microsoft.Extensions.Logging.Console" Version="10.0.0"/>
<PackageVersion Include="Microsoft.AspNetCore.Authorization" Version="10.0.1"/>
<PackageVersion Include="Microsoft.AspNetCore.Components" Version="10.0.1"/>
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly" Version="10.0.1"/>
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.DevServer" Version="10.0.1"/>
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.Server" Version="10.0.1"/>
<PackageVersion Include="Microsoft.AspNetCore.DataProtection.Abstractions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.AspNetCore.Mvc.Testing" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Data.Sqlite" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Data.Sqlite.Core" Version="10.0.1"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore" Version="10.0.1"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.Design" Version="10.0.1"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.Relational" Version="10.0.1"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.Sqlite" Version="10.0.1"/>
<PackageVersion Include="Microsoft.EntityFrameworkCore.SqlServer" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Caching.Abstractions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Caching.Memory" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Configuration" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Configuration.Abstractions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Configuration.Json" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.DependencyModel" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Hosting.Abstractions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Http" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Http.Polly" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Http.Resilience" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Options" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Options.ConfigurationExtensions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Resilience" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Logging" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Logging.Abstractions" Version="10.0.1"/>
<PackageVersion Include="Microsoft.Extensions.Logging.Console" Version="10.0.1"/>
<PackageVersion Include="Npgsql.EntityFrameworkCore.PostgreSQL" Version="10.0.0"/>
<PackageVersion Include="Oracle.EntityFrameworkCore" Version="10.23.26000"/>
<PackageVersion Include="Scrutor" Version="7.0.0"/>
@ -126,7 +126,7 @@
<PackageVersion Include="Moq" Version="4.20.72"/>
<PackageVersion Include="MySql.Data" Version="9.5.0"/>
<PackageVersion Include="Newtonsoft.Json" Version="13.0.4"/>
<PackageVersion Include="Npgsql" Version="10.0.0"/>
<PackageVersion Include="Npgsql" Version="10.0.1"/>
<PackageVersion Include="NSubstitute" Version="5.3.0"/>
<PackageVersion Include="NuGet.Packaging" Version="7.0.1"/>
<PackageVersion Include="NuGet.Protocol" Version="7.0.1"/>
@ -151,20 +151,21 @@
<PackageVersion Include="Refit" Version="9.0.2"/>
<PackageVersion Include="Refit.HttpClientFactory" Version="9.0.2"/>
<PackageVersion Include="Serilog" Version="4.3.0"/>
<PackageVersion Include="Serilog.Extensions.Logging" Version="10.0.0"/>
<PackageVersion Include="Serilog.Extensions.Logging" Version="10.0.1"/>
<PackageVersion Include="Serilog.Sinks.File" Version="7.0.0"/>
<PackageVersion Include="Serilog.Formatting.Compact" Version="3.0.0"/>
<PackageVersion Include="ShortGuid" Version="2.0.1"/>
<PackageVersion Include="StackExchange.Redis" Version="2.10.1"/>
<PackageVersion Include="System.CommandLine" Version="2.0.0"/>
<PackageVersion Include="System.Data.SqlClient" Version="4.9.0"/>
<PackageVersion Include="System.Formats.Asn1" Version="10.0.0"/>
<PackageVersion Include="System.Formats.Asn1" Version="10.0.1"/>
<PackageVersion Include="System.Linq.Async" Version="7.0.0"/>
<PackageVersion Include="System.Linq.Dynamic.Core" Version="1.7.1"/>
<PackageVersion Include="System.Net.Http" Version="4.3.4"/>
<PackageVersion Include="System.Text.Json" Version="10.0.0"/>
<PackageVersion Include="System.Text.Json" Version="10.0.1"/>
<PackageVersion Include="System.Text.RegularExpressions" Version="4.3.1"/>
<PackageVersion Include="Testcontainers" Version="4.9.0"/>
<PackageVersion Include="Testcontainers.MsSql" Version="4.9.0"/>
<PackageVersion Include="Testcontainers.PostgreSql" Version="4.9.0"/>
<PackageVersion Include="Testcontainers.RabbitMq" Version="4.9.0"/>
<PackageVersion Include="Testcontainers.Redis" Version="4.9.0"/>

View file

@ -4,21 +4,10 @@
<packageSources>
<clear />
<add key="NuGet official package source" value="https://api.nuget.org/v3/index.json" />
<add key="Elsa 3 Preview" value="https://f.feedz.io/elsa-workflows/elsa-3/nuget/index.json" />
<add key="Webhooks Core Preview" value="https://f.feedz.io/sfmskywalker/webhooks-core/nuget/index.json" />
</packageSources>
<packageSourceMapping>
<packageSource key="NuGet official package source">
<package pattern="*" />
</packageSource>
<packageSource key="Elsa 3 Preview">
<package pattern="Elsa" />
<package pattern="Elsa.*" />
<package pattern="Elsa.Studio.*" />
</packageSource>
<packageSource key="Webhooks Core Preview">
<package pattern="WebhooksCore" />
<package pattern="WebhooksCore.*" />
</packageSource>
</packageSourceMapping>
</configuration>

View file

@ -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 ./

View file

@ -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 ./

View file

@ -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 ./

View file

@ -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 ./

View file

@ -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<HttpRequest>
{
var path = Path.Get(context);
var methods = SupportedMethods.GetOrDefault(context) ?? new List<string> { HttpMethods.Get };
await context.WaitForHttpRequestAsync(path, methods, OnResumeAsync);
await context.WaitForHttpRequestAsync(path, methods, OnResumeAsync, Elsa.Http.HttpStimulusNames.HttpEndpoint);
}
private async ValueTask OnResumeAsync(ActivityExecutionContext context)

View file

@ -22,8 +22,8 @@ public abstract class HttpEndpointBase<TResult> : Trigger<TResult>
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<object> GetTriggerPayloads(TriggerIndexingContext context)

View file

@ -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<string> methods, ExecuteActivityDelegate? callback = null)
public static async ValueTask WaitForHttpRequestAsync(this ActivityExecutionContext context, string path, IEnumerable<string> 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;
}

View file

@ -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).
/// </summary>
[Input(
Description = "The join mode determines whether this activity should continue as soon as one inbound path comes in (Wait Any), or once all inbound paths have executed (Wait All).",
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
};
}
}

View file

@ -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();
}
/// <summary>
/// Checks if there is any pending work for the flowchart.
/// </summary>
private bool HasPendingWork(ActivityExecutionContext context)
{
var workflowExecutionContext = context.WorkflowExecutionContext;
// Use HashSet for O(1) lookups
var activityIds = new HashSet<string>(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<bool>(BackwardConnectionActivityInput);
bool hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExecutedByBackwardConnection);
var completedActivityExcecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault<bool>(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());
}
/// <summary>
/// Schedules outbound activities based on the flowchart's structure and execution state.
/// This method determines whether an activity should be scheduled based on visited connections,
/// forward traversal rules, and backward connections.
/// 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.
/// </summary>
/// <param name="flowGraph">The graph representation of the flowchart.</param>
/// <param name="flowScope">Tracks activity and connection visits.</param>
/// <param name="flowchartContext">The execution context of the flowchart.</param>
/// <param name="activity">The current activity being processed.</param>
/// <param name="outcomes">The outcomes that determine which connections were followed.</param>
/// <param name="completedActivityExecutedByBackwardConnection">Indicates if the completed activity was executed due to a backward connection.</param>
/// <param name="completionCallback">The callback to invoke upon activity completion.</param>
/// <param name="completedActivityExecutedByBackwardConnection">Indicates if the completed activity
/// was executed due to a backward connection.</param>
/// <returns>True if at least one activity was scheduled; otherwise, false.</returns>
private async ValueTask<bool> ScheduleOutboundActivitiesAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity activity, Outcomes outcomes, bool completedActivityExecutedByBackwardConnection = false)
private static async ValueTask<bool> 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
/// <summary>
/// Schedules an outbound activity that originates from a backward connection.
/// </summary>
private async ValueTask<bool> ScheduleBackwardConnectionActivityAsync(FlowGraph flowGraph, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity, bool connectionFollowed, bool backwardConnectionIsValid)
private static async ValueTask<bool> 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<string, object>()
{
{
BackwardConnectionActivityInput, true
}
}
CompletionCallback = completionCallback,
Input = new Dictionary<string, object>() { { BackwardConnectionActivityInput, true } }
};
await flowchartContext.ScheduleActivityAsync(outboundActivity, scheduleWorkOptions);
@ -142,91 +227,136 @@ public partial class Flowchart
}
/// <summary>
/// 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.
/// </summary>
private async ValueTask<bool> ScheduleNonJoinActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, IActivity outboundActivity)
private static async ValueTask<FlowJoinMode> 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<FlowJoin, FlowJoinMode>(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;
}
}
/// <summary>
/// Schedules a join activity based on inbound connection statuses.
/// </summary>
private async ValueTask<bool> ScheduleJoinActivityAsync(FlowGraph flowGraph, FlowScope flowScope, ActivityExecutionContext flowchartContext, Connection outboundConnection, IActivity outboundActivity)
private static async ValueTask<bool> 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)
/// <summary>
/// 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.
/// </summary>
private static async ValueTask<bool> 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);
}
/// <summary>
/// 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.
/// </summary>
private static async ValueTask<bool> 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);
}
/// <summary>
/// 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.
/// </summary>
private static async ValueTask<bool> 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;
}
/// <summary>
/// Schedules the outbound activity.
/// </summary>
private static async ValueTask<bool> ScheduleOutboundActivityAsync(ActivityExecutionContext flowchartContext, IActivity outboundActivity, ActivityCompletionCallback completionCallback)
{
await flowchartContext.ScheduleActivityAsync(outboundActivity, completionCallback);
return true;
}
/// <summary>
/// Skips the outbound activity by propagating skipped connections.
/// </summary>
private static async ValueTask<bool> 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);
}
}

View file

@ -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
/// <inheritdoc />
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();
}
}
}
}

View file

@ -3,5 +3,6 @@ namespace Elsa.Workflows.Activities.Flowchart.Models;
public enum FlowJoinMode
{
WaitAll,
WaitAny
WaitAllActive,
WaitAny,
}

View file

@ -86,7 +86,7 @@ public class FlowScope
/// <param name="flowGraph">The flow graph containing connections.</param>
/// <param name="activity">The activity to check.</param>
/// <returns>True if any inbound connection has been followed, otherwise false.</returns>
public bool HasFollowedInboundConnection(FlowGraph flowGraph, IActivity activity)
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));
}
/// <summary>
/// Determines whether all inbound connection to the specified activity has been followed.
/// </summary>
/// <param name="flowGraph">The flow graph containing connections.</param>
/// <param name="activity">The activity to check.</param>
/// <returns>True if all inbound connection has been followed, otherwise false.</returns>
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));
}
/// <summary>
/// Determines whether a connection should be ignored based on visit counts.
/// </summary>

View file

@ -371,13 +371,15 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable
/// <param name="payloads">The payloads to create bookmarks for.</param>
/// <param name="callback">An optional callback that is invoked when the bookmark is resumed.</param>
/// <param name="includeActivityInstanceId">Whether or not the activity instance ID should be included in the bookmark payload.</param>
public void CreateBookmarks(IEnumerable<object> payloads, ExecuteActivityDelegate? callback = null, bool includeActivityInstanceId = true)
/// <param name="bookmarkName">An optional name to use for the bookmark. Defaults to the activity type.</param>
public void CreateBookmarks(IEnumerable<object> 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
});
}

View file

@ -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
/// <summary>
/// Returns the <see cref="WorkflowExecutionContext"/> of the specified <see cref="ExpressionExecutionContext"/>
/// </summary>
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.");
}
/// <summary>
/// Returns the <see cref="ActivityExecutionContext"/> of the specified <see cref="ExpressionExecutionContext"/>
/// </summary>
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.");
}
/// <summary>
/// Returns the <see cref="ActivityExecutionContext"/> of the specified <see cref="ExpressionExecutionContext"/>
/// Returns the <see cref="ActivityExecutionContext"/> of the specified <see cref="ExpressionExecutionContext"/>
/// </summary>
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.
/// </summary>
public IEnumerable<string> 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.
/// </summary>
public IEnumerable<Variable> GetVariablesInScope() =>
EnumerateVariablesInScope(context)
context.EnumerateVariablesInScope()
.Where(x => !string.IsNullOrWhiteSpace(x.Name))
.DistinctBy(x => x.Name);
@ -307,7 +315,7 @@ public static class ExpressionExecutionContextExtensions
/// </summary>
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);
}
/// <summary>
@ -384,7 +392,6 @@ public static class ExpressionExecutionContextExtensions
return serializerOptions;
}
/// <param name="context"></param>
extension(ExpressionExecutionContext context)
{
/// <summary>

View file

@ -118,6 +118,14 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module)
/// </summary>
public Func<IServiceProvider, ICommandHandler> DispatchWorkflowCommandHandler { get; set; } = sp => sp.GetRequiredService<DispatchWorkflowCommandHandler>();
/// <summary>
/// A factory that instantiates an <see cref="IWorkflowResumer"/>.
/// </summary>
public Func<IServiceProvider, IWorkflowResumer> WorkflowResumer { get; set; } = sp => sp.GetRequiredService<WorkflowResumer>();
/// <summary>
/// A factory that instantiates an <see cref="IBookmarkQueueWorker"/>.
/// </summary>
public Func<IServiceProvider, IBookmarkQueueWorker> BookmarkQueueWorker { get; set; } = sp => sp.GetRequiredService<BookmarkQueueWorker>();
/// <summary>
@ -249,7 +257,7 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module)
.AddSingleton(BackgroundActivityScheduler)
.AddSingleton<RandomLongIdentityGenerator>()
.AddSingleton<IBookmarkQueueSignaler, BookmarkQueueSignaler>()
.AddScoped<IBookmarkQueueWorker, BookmarkQueueWorker>()
.AddScoped(BookmarkQueueWorker)
.AddScoped<IBookmarkManager, DefaultBookmarkManager>()
.AddScoped<IActivityExecutionManager, DefaultActivityExecutionManager>()
.AddScoped<IActivityExecutionStatsService, ActivityExecutionStatsService>()
@ -275,7 +283,10 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module)
.AddScoped<IBookmarksPersister, BookmarksPersister>()
.AddScoped<IBookmarkResumer, BookmarkResumer>()
.AddScoped<IBookmarkQueue, StoreBookmarkQueue>()
.AddScoped<IWorkflowResumer, WorkflowResumer>()
.AddScoped(WorkflowResumer)
.AddScoped<WorkflowResumer>()
.AddScoped(BookmarkQueueWorker)
.AddScoped<BookmarkQueueWorker>()
.AddScoped<ITriggerInvoker, TriggerInvoker>()
.AddScoped<IWorkflowCanceler, WorkflowCanceler>()
.AddScoped<IWorkflowCancellationService, WorkflowCancellationService>()

View file

@ -1,14 +1,15 @@
<Project>
<Import Project="$([MSBuild]::GetPathOfFileAbove('Directory.Build.props', '$(MSBuildThisFileDirectory)../'))" />
<Import Project="$([MSBuild]::GetPathOfFileAbove('Directory.Build.props', '$(MSBuildThisFileDirectory)../'))"/>
<PropertyGroup>
<Threshold>40</Threshold>
</PropertyGroup>
<PropertyGroup>
<Threshold>40</Threshold>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Testcontainers.PostgreSql"/>
<PackageReference Include="Testcontainers.RabbitMq"/>
</ItemGroup>
<ItemGroup>
<PackageReference Include="Testcontainers.PostgreSql"/>
<PackageReference Include="Testcontainers.RabbitMq"/>
<PackageReference Include="Testcontainers.MsSql"/>
</ItemGroup>
</Project>

View file

@ -22,6 +22,7 @@
<ProjectReference Include="..\..\..\src\clients\Elsa.Api.Client\Elsa.Api.Client.csproj" />
<ProjectReference Include="..\..\..\src\common\Elsa.Testing.Shared.Component\Elsa.Testing.Shared.Component.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Persistence.EFCore.PostgreSql\Elsa.Persistence.EFCore.PostgreSql.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Persistence.EFCore.SqlServer\Elsa.Persistence.EFCore.SqlServer.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Scheduling\Elsa.Scheduling.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.WorkflowProviders.BlobStorage\Elsa.WorkflowProviders.BlobStorage.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Runtime.Distributed\Elsa.Workflows.Runtime.Distributed.csproj" />

View file

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

View file

@ -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()

View file

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

View file

@ -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<IJavaScriptEvaluator>();
}
[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<Guid>(result);
}
}

View file

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