From dc429700abc7686dbee32896b26d96d7c2d47023 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 9 Oct 2025 21:33:23 +0200 Subject: [PATCH] Fixes deadlock issue with implicit joins (#6955) * Fix outbound connection processing logic Corrected `completedActivityExcecutedByBackwardConnection` to `completedActivityExecutedByBackwardConnection`. Improved flowgraph outbound connection handling by separating visitation and processing logic, ensuring skipped connections are propagated consistently. * Add tests for decision implicit join workflows Introduce new integration tests to verify workflows with implicit joins on both decision outcomes. Added corresponding workflow definitions and updated the test project to ensure compatibility. Refactored connection visit logic for better readability and maintainability. --- .../Elsa.Server.Web/Workflows/triggers.json | 105 ++++++++++++++++ .../Flowchart/Activities/Flowchart.cs | 26 ++-- .../Activities/Flowchart/Models/FlowScope.cs | 12 +- .../Elsa.Workflows.IntegrationTests.csproj | 3 + .../DecisionBothOutcomesJoinedTests.cs | 31 +++++ ...cision-implicit-join-on-both-outcomes.json | 115 ++++++++++++++++++ 6 files changed, 276 insertions(+), 16 deletions(-) create mode 100644 src/apps/Elsa.Server.Web/Workflows/triggers.json create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/DecisionBothOutcomesJoinedTests.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/Workflows/decision-implicit-join-on-both-outcomes.json diff --git a/src/apps/Elsa.Server.Web/Workflows/triggers.json b/src/apps/Elsa.Server.Web/Workflows/triggers.json new file mode 100644 index 000000000..5150e9cbb --- /dev/null +++ b/src/apps/Elsa.Server.Web/Workflows/triggers.json @@ -0,0 +1,105 @@ +{ + "$schema": "https://elsaworkflows.io/schemas/workflow-definition/v3.0.0/schema.json", + "id": "3b42d3276206e00f", + "definitionId": "fb27085e78433f79", + "name": "Triggers", + "createdAt": "2025-10-09T07:37:31.521686\u002B00:00", + "version": 3, + "toolVersion": "3.5.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": {}, + "isReadonly": false, + "isSystem": false, + "isLatest": true, + "isPublished": true, + "options": { + "usableAsActivity": false, + "autoUpdateConsumingWorkflows": false + }, + "root": { + "id": "83f95a4ca749c950", + "nodeId": "Workflow1:83f95a4ca749c950", + "name": "Flowchart1", + "type": "Elsa.Flowchart", + "version": 1, + "customProperties": { + "notFoundConnections": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": {}, + "activities": [ + { + "path": { + "typeName": "String", + "expression": { + "type": "Literal", + "value": "triggers" + } + }, + "supportedMethods": { + "typeName": "List\u003CString\u003E", + "expression": { + "type": "Object", + "value": "[\u0022GET\u0022]" + } + }, + "authorize": { + "typeName": "Boolean", + "expression": { + "type": "Literal", + "value": false + } + }, + "policy": { + "typeName": "String", + "expression": { + "type": "Literal" + } + }, + "requestTimeout": null, + "requestSizeLimit": null, + "fileSizeLimit": null, + "allowedFileExtensions": null, + "blockedFileExtensions": null, + "allowedMimeTypes": null, + "exposeRequestTooLargeOutcome": false, + "exposeFileTooLargeOutcome": false, + "exposeInvalidFileExtensionOutcome": false, + "exposeInvalidFileMimeTypeOutcome": false, + "parsedContent": null, + "files": null, + "routeData": null, + "queryStringData": null, + "headers": null, + "result": null, + "id": "607c2042a187fe5a", + "nodeId": "Workflow1:83f95a4ca749c950:607c2042a187fe5a", + "name": "HttpEndpoint1", + "type": "Elsa.HttpEndpoint", + "version": 1, + "customProperties": { + "canStartWorkflow": true, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -316.431640625, + "y": -470.017578125 + }, + "size": { + "width": 195.1328125, + "height": 67.9765625 + } + } + } + } + ], + "variables": [], + "connections": [] + } +} \ 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 37903d923..1d24195ee 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs @@ -173,8 +173,8 @@ public class Flowchart : Container // Schedule the outbound activities var flowGraph = GetFlowGraph(flowchartContext); var flowScope = GetFlowScope(flowchartContext); - var completedActivityExcecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault(BackwardConnectionActivityInput); - var hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExcecutedByBackwardConnection); + var completedActivityExecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault(BackwardConnectionActivityInput); + var hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExecutedByBackwardConnection); // If there are not any outbound connections, complete the flowchart activity if there is no other pending work. if (!hasScheduledActivity) @@ -211,15 +211,23 @@ public class Flowchart : Container { flowScope.RegisterActivityVisit(activity); } + + var outboundConnections = flowGraph.GetOutboundConnections(activity); - // Process each outbound connection from the current activity - foreach (var outboundConnection in flowGraph.GetOutboundConnections(activity)) + // Register the outbound connections as visited. + foreach (var outboundConnection in outboundConnections) { var connectionFollowed = outcomes.Names.Contains(outboundConnection.Source.Port); flowScope.RegisterConnectionVisit(outboundConnection, connectionFollowed); + } + + // Process each outbound connection from the current activity + foreach (var outboundConnection in outboundConnections) + { + var connectionFollowed = flowScope.GetConnectionLastVisitFollowed(outboundConnection); if(!connectionFollowed) - continue; // Skip if connection was not followed. + continue; // Skip if the connection was not followed. var outboundActivity = outboundConnection.Target.Activity; @@ -280,11 +288,9 @@ public class Flowchart : Container await flowchartContext.ScheduleActivityAsync(outboundActivity, OnChildCompletedAsync); return true; } - else - { - // Propagate skipped connections by scheduling with Outcomes.Empty - return await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty); - } + + // Propagate skipped connections by scheduling with Outcomes.Empty + return await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, outboundActivity, Outcomes.Empty); } /// 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 0fa2790a3..742fd3222 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/FlowScope.cs @@ -37,7 +37,7 @@ public class FlowScope /// /// The activity to check. /// The visit count of the activity. - private long GetActivityVisitCount(IActivity activity) => ActivitiesVisitCount.TryGetValue(activity.Id, out var count) ? count : 0; + private long GetActivityVisitCount(IActivity activity) => ActivitiesVisitCount.GetValueOrDefault(activity.Id, 0); /// /// Registers a visit to the specified connection and records whether it was followed. @@ -46,7 +46,7 @@ public class FlowScope /// Indicates whether the connection was followed. public void RegisterConnectionVisit(Connection connection, bool followed) { - string connectionId = connection.ToString(); + var connectionId = connection.ToString(); ConnectionVisitCount.TryAdd(connectionId, 0); ConnectionVisitCount[connectionId]++; ConnectionLastVisitFollowed[connectionId] = followed; @@ -57,14 +57,14 @@ public class FlowScope /// /// The connection to check. /// The visit count of the connection. - private long GetConnectionVisitCount(Connection connection) => ConnectionVisitCount.TryGetValue(connection.ToString(), out var count) ? count : 0; + private long GetConnectionVisitCount(Connection connection) => ConnectionVisitCount.GetValueOrDefault(connection.ToString(), 0); /// /// Determines whether the last visit to the specified connection was followed. /// /// The connection to check. /// True if the connection was followed on the last visit, otherwise false. - private bool GetConnectionLastVisitFollowed(Connection connection) => ConnectionLastVisitFollowed.TryGetValue(connection.ToString(), out var followed) ? followed : false; + public bool GetConnectionLastVisitFollowed(Connection connection) => ConnectionLastVisitFollowed.GetValueOrDefault(connection.ToString(), false); /// /// Determines whether all inbound connections to the specified activity have been visited. @@ -76,7 +76,7 @@ public class FlowScope { var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity); var outboundActivityVisitCount = GetActivityVisitCount(activity); - var minConnectionVisitCount = forwardInboundConnections.Min(c => GetConnectionVisitCount(c)); + var minConnectionVisitCount = forwardInboundConnections.Min(GetConnectionVisitCount); return minConnectionVisitCount > outboundActivityVisitCount; } @@ -90,7 +90,7 @@ public class FlowScope { var forwardInboundConnections = flowGraph.GetForwardInboundConnections(activity); var outboundActivityVisitCount = GetActivityVisitCount(activity); - var maxConnectionVisitCount = forwardInboundConnections.Max(c => GetConnectionVisitCount(c)); + var maxConnectionVisitCount = forwardInboundConnections.Max(GetConnectionVisitCount); return maxConnectionVisitCount > outboundActivityVisitCount && forwardInboundConnections.Any(c => GetConnectionVisitCount(c) == maxConnectionVisitCount && GetConnectionLastVisitFollowed(c)); } diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj b/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj index 74bf453d5..97fcec8c4 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj +++ b/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj @@ -28,5 +28,8 @@ Always + + Always + diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/DecisionBothOutcomesJoinedTests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/DecisionBothOutcomesJoinedTests.cs new file mode 100644 index 000000000..178d03bbf --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/DecisionBothOutcomesJoinedTests.cs @@ -0,0 +1,31 @@ +using Elsa.Testing.Shared; +using Xunit.Abstractions; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.ImplicitJoins; + +public class DecisionBothOutcomesJoinedTests +{ + private readonly CapturingTextWriter _capturingTextWriter = new(); + private readonly IServiceProvider _services; + + public DecisionBothOutcomesJoinedTests(ITestOutputHelper testOutputHelper) + { + _services = new TestApplicationBuilder(testOutputHelper).WithCapturingTextWriter(_capturingTextWriter).Build(); + } + + [Fact(DisplayName = "Connecting both outcomes of a decision to an implicit join should result in the workflow completing.")] + public async Task Test1() + { + // Populate registries. + await _services.PopulateRegistriesAsync(); + + // Import workflow. + var workflowDefinition = await _services.ImportWorkflowDefinitionAsync("Scenarios/ImplicitJoins/Workflows/decision-implicit-join-on-both-outcomes.json"); + + // Execute. + var state = await _services.RunWorkflowUntilEndAsync(workflowDefinition.DefinitionId); + + // Assert. + Assert.Equal(WorkflowStatus.Finished, state.Status); + } +} \ No newline at end of file diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/Workflows/decision-implicit-join-on-both-outcomes.json b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/Workflows/decision-implicit-join-on-both-outcomes.json new file mode 100644 index 000000000..78a39b091 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/ImplicitJoins/Workflows/decision-implicit-join-on-both-outcomes.json @@ -0,0 +1,115 @@ +{ + "$schema": "https://elsaworkflows.io/schemas/workflow-definition/v3.0.0/schema.json", + "id": "e8e53204f45ca599", + "definitionId": "b8c92c080faf734f", + "name": "Decision Implicit Join On Both Outcomes", + "createdAt": "2025-10-09T12:06:51.388915\u002B00:00", + "version": 1, + "toolVersion": "3.6.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": {}, + "isReadonly": false, + "isSystem": false, + "isLatest": true, + "isPublished": false, + "options": { + "autoUpdateConsumingWorkflows": false + }, + "root": { + "id": "990d07168c5223de", + "nodeId": "Workflow1:990d07168c5223de", + "name": "Flowchart1", + "type": "Elsa.Flowchart", + "version": 1, + "customProperties": { + "notFoundConnections": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": {}, + "activities": [ + { + "text": null, + "id": "36874900ac3c0b3e", + "nodeId": "Workflow1:990d07168c5223de:36874900ac3c0b3e", + "name": "WriteLine1", + "type": "Elsa.WriteLine", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": 1.7409133911132812, + "y": -481.0423583984375 + }, + "size": { + "width": 159.6171875, + "height": 67.9765625 + } + } + } + }, + { + "condition": { + "typeName": "Boolean", + "expression": { + "type": "JavaScript", + "value": "true" + } + }, + "id": "da8f7fd55383b02a", + "nodeId": "Workflow1:990d07168c5223de:da8f7fd55383b02a", + "name": "FlowDecision1", + "type": "Elsa.FlowDecision", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -300, + "y": -480 + }, + "size": { + "width": 149.84375, + "height": 67.9765625 + } + } + } + } + ], + "variables": [], + "connections": [ + { + "source": { + "activity": "da8f7fd55383b02a", + "port": "False" + }, + "target": { + "activity": "36874900ac3c0b3e", + "port": "In" + }, + "vertices": [] + }, + { + "source": { + "activity": "da8f7fd55383b02a", + "port": "True" + }, + "target": { + "activity": "36874900ac3c0b3e", + "port": "In" + }, + "vertices": [] + } + ] + } +} \ No newline at end of file