From 42753359d2f10e1f9f9dec2f00c40f3c5bdfc79a Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 6 Oct 2025 15:39:16 +0200 Subject: [PATCH 1/5] Adjust recurring task schedules and bookmark queue TTL for optimized performance. --- src/apps/Elsa.Server.Web/Program.cs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index 64fbe93c4..3e20794bb 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -806,13 +806,13 @@ services.AddActivityStateFilter(); services.Configure(options => { options.Schedule.ConfigureTask(TimeSpan.FromSeconds(300)); - options.Schedule.ConfigureTask(TimeSpan.FromSeconds(300)); + options.Schedule.ConfigureTask(TimeSpan.FromSeconds(3600)); options.Schedule.ConfigureTask(TimeSpan.FromHours(4)); options.Schedule.ConfigureTask(TimeSpan.FromSeconds(15)); }); services.Configure(options => { options.InactivityThreshold = TimeSpan.FromSeconds(15); }); -services.Configure(options => options.Ttl = TimeSpan.FromSeconds(10)); +services.Configure(options => options.Ttl = TimeSpan.FromSeconds(3600)); services.Configure(options => options.CacheDuration = TimeSpan.FromDays(1)); services.Configure(options => options.DefaultIncidentStrategy = typeof(ContinueWithIncidentsStrategy)); services.AddHealthChecks(); From dc429700abc7686dbee32896b26d96d7c2d47023 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 9 Oct 2025 21:33:23 +0200 Subject: [PATCH 2/5] 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 From a74502d3ad7977485d3d851df8156b78857d16b4 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 10 Oct 2025 09:21:11 +0200 Subject: [PATCH 3/5] Fix trigger indexing logic and add tests for trigger persistence after workflow reload. (#6958) * 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. * Refactor trigger indexing logic and add tests for trigger persistence after workflow reload. Streamlined conditional trigger indexing in `DefaultWorkflowDefinitionStorePopulator`. Introduced a test to validate trigger persistence across reloads after publishing a new workflow version. --- ...DefaultWorkflowDefinitionStorePopulator.cs | 6 +-- .../Services/TriggerIndexer.cs | 2 +- .../ReloadWorkflowTests.cs | 41 +++++++++++++++++-- 3 files changed, 41 insertions(+), 8 deletions(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs index d21d298fe..9fb839321 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -83,8 +83,8 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP await AssignIdentities(materializedWorkflow.Workflow, cancellationToken); var workflowDefinition = await AddOrUpdateAsync(materializedWorkflow, cancellationToken); - if (indexTriggers) - await IndexTriggersAsync(workflowDefinition, cancellationToken); + if (indexTriggers && workflowDefinition.IsPublished) + await _triggerIndexer.IndexTriggersAsync(workflowDefinition, cancellationToken); return workflowDefinition; } @@ -260,8 +260,6 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } } - private async Task IndexTriggersAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken) => await _triggerIndexer.IndexTriggersAsync(workflowDefinition, cancellationToken); - /// /// Syncs the items in the primary list with existing items in the secondary list, even when the object instances are not the same (but their IDs are). /// diff --git a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs index 4a7f6f48f..474dff551 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs @@ -89,7 +89,7 @@ public class TriggerIndexer : ITriggerIndexer // Collect new triggers **if the workflow is published**. var newTriggers = workflow.Publication.IsPublished ? await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken) - : new List(0); + : new(0); // Diff triggers. var diff = Diff.For(currentTriggers, newTriggers, new WorkflowTriggerEqualityComparer()); diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs index 3f96d9dbb..d8399a075 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -7,6 +7,7 @@ using Elsa.Workflows.ComponentTests.Materializers; using Elsa.Workflows.ComponentTests.WorkflowProviders; using Elsa.Workflows.Management; using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Filters; using Humanizer; using Microsoft.Extensions.DependencyInjection; @@ -20,6 +21,8 @@ public class ReloadWorkflowTests : AppComponentTest private readonly IWorkflowDefinitionManager _workflowDefinitionManager; private readonly IWorkflowDefinitionService _workflowDefinitionService; private readonly IWorkflowDefinitionsReloader _workflowDefinitionsReloader; + private readonly IWorkflowDefinitionPublisher _workflowDefinitionPublisher; + private readonly ITriggerStore _triggerStore; public ReloadWorkflowTests(App app) : base(app) { @@ -27,7 +30,9 @@ public class ReloadWorkflowTests : AppComponentTest _workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService(); _workflowBuilderFactory = Scope.ServiceProvider.GetRequiredService(); _workflowDefinitionService = Scope.ServiceProvider.GetRequiredService(); + _workflowDefinitionPublisher = Scope.ServiceProvider.GetRequiredService(); _activityRegistry = Scope.ServiceProvider.GetRequiredService(); + _triggerStore = Scope.ServiceProvider.GetRequiredService(); var workflowsProviders = Scope.ServiceProvider.GetRequiredService>(); _testWorkflowProvider = (TestWorkflowProvider)workflowsProviders.First(x => x is TestWorkflowProvider); } @@ -37,9 +42,9 @@ public class ReloadWorkflowTests : AppComponentTest { var client = WorkflowServer.CreateHttpWorkflowClient(); await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); - var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + var firstResponse = await client.SendAsync(new(HttpMethod.Get, "reload-test")); await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); - var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + var secondResponse = await client.SendAsync(new(HttpMethod.Get, "reload-test")); Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); } @@ -103,6 +108,36 @@ public class ReloadWorkflowTests : AppComponentTest await _workflowDefinitionManager.DeleteByDefinitionIdAsync(definitionId, CancellationToken.None); } + [Fact] + public async Task Reloading_AfterPublishingNewVersion_ShouldPersistTriggers() + { + // Get the initial workflow definition. + const string definitionId = "f68b09bc-2013-4617-b82f-d76b6819a624"; + var initialDefinition = await _workflowDefinitionService.FindWorkflowDefinitionAsync(definitionId, VersionOptions.Published, CancellationToken.None); + Assert.NotNull(initialDefinition); + + // Assert that triggers exist initially. + var initialTrigger = await _triggerStore.FindAsync(new(){ WorkflowDefinitionId = definitionId}, CancellationToken.None); + Assert.NotNull(initialTrigger); + + // Publish a new version of the workflow. + var draftDefinition = await _workflowDefinitionPublisher.GetDraftAsync(definitionId, VersionOptions.Latest); + Assert.NotNull(draftDefinition); + await _workflowDefinitionPublisher.PublishAsync(draftDefinition, CancellationToken.None); + + // Assert we are at version 2. + var v2Definition = await _workflowDefinitionService.FindWorkflowDefinitionAsync(definitionId, VersionOptions.Published, CancellationToken.None); + Assert.NotNull(v2Definition); + Assert.Equal(2, v2Definition.Version); + + // Reload the workflow definitions. + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + + // Assert that triggers still exist after reload. + var reloadedTrigger = await _triggerStore.FindAsync(new(){ WorkflowDefinitionId = definitionId}, CancellationToken.None); + Assert.NotNull(reloadedTrigger); + } + private async Task BuildWorkflowAsync(string definitionId, string definitionVersionId, int version) { var builder = _workflowBuilderFactory.CreateBuilder(); @@ -114,6 +149,6 @@ public class ReloadWorkflowTests : AppComponentTest builder.WorkflowOptions.UsableAsActivity = true; var workflow = await builder.BuildWorkflowAsync(); workflow.Name = definitionId; - return new MaterializedWorkflow(workflow, _testWorkflowProvider.Name, TestWorkflowMaterializer.MaterializerName); + return new(workflow, _testWorkflowProvider.Name, TestWorkflowMaterializer.MaterializerName); } } \ No newline at end of file From 59bd21f899ee6b93efb194c9e3c80d1d3f861dbb Mon Sep 17 00:00:00 2001 From: bobhauser Date: Thu, 16 Oct 2025 07:51:37 -0400 Subject: [PATCH 4/5] Fix ObjectConverter.ConvertTo to avoid using sourceTypeConverter.IsValid (#6972) * Fix ObjectConverter.ConvertTo to avoid using sourceTypeConverter.IsValid * Update test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/PersonTypeConverter.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Bob Hauser Co-authored-by: Sipke Schoorstra Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../Helpers/ObjectConverter.cs | 15 ++- .../ObjectConversion/Person.cs | 13 ++- .../ObjectConversion/PersonTypeConverter.cs | 50 ++++++++++ .../ObjectConversion/Tests.cs | 96 +++++++++++++++++++ 4 files changed, 162 insertions(+), 12 deletions(-) create mode 100644 test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/PersonTypeConverter.cs diff --git a/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs b/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs index 40c80e57c..5bc699e2f 100644 --- a/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs +++ b/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs @@ -217,11 +217,16 @@ public static class ObjectConverter var sourceTypeConverter = TypeDescriptor.GetConverter(underlyingSourceType); if (sourceTypeConverter.CanConvertTo(underlyingTargetType)) - { - var isValid = targetTypeConverter.IsValid(value); - - if (isValid) - return sourceTypeConverter.ConvertTo(value, underlyingTargetType); + { + // TypeConverter.IsValid is not supported for ConvertTo, so we have to try and ignore any exceptions. + try + { + return sourceTypeConverter.ConvertTo(value, underlyingTargetType); + } + catch + { + // Ignore and try other conversion strategies. + } } if (underlyingTargetType.IsEnum) diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Person.cs b/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Person.cs index b5583db8b..6d9bbcdcd 100644 --- a/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Person.cs +++ b/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Person.cs @@ -1,9 +1,8 @@ -namespace Elsa.Workflows.Core.UnitTests.ObjectConversion -{ - public class Person - { - public double? Age { get; set; } +namespace Elsa.Workflows.Core.UnitTests.ObjectConversion; - public string? Name { get; set; } - } +public class Person +{ + public double? Age { get; set; } + + public string? Name { get; set; } } diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/PersonTypeConverter.cs b/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/PersonTypeConverter.cs new file mode 100644 index 000000000..7cd107d5a --- /dev/null +++ b/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/PersonTypeConverter.cs @@ -0,0 +1,50 @@ +using System.ComponentModel; +using System.Globalization; +using System.Text.Json; + +namespace Elsa.Workflows.Core.UnitTests.ObjectConversion; + +public class PersonTypeConverter : TypeConverter +{ + public override bool CanConvertFrom(ITypeDescriptorContext? context, Type sourceType) + { + // Allow conversion from string to Person + return sourceType == typeof(string) || base.CanConvertFrom(context, sourceType); + } + + public override bool CanConvertTo(ITypeDescriptorContext? context, Type? destinationType) + { + // Allow conversion from Person to string + return destinationType == typeof(string) || base.CanConvertTo(context, destinationType); + } + + public override object? ConvertFrom(ITypeDescriptorContext? context, CultureInfo? culture, object value) + { + if (value == null) + { + return null; + } + + if (value is string json) + { + if (string.IsNullOrEmpty(json)) + { + return null; + } + + return JsonSerializer.Deserialize(json); + } + + return base.ConvertFrom(context, culture, value); + } + + public override object? ConvertTo(ITypeDescriptorContext? context, CultureInfo? culture, object? value, Type destinationType) + { + if (destinationType == typeof(string) && value is Person person) + { + return JsonSerializer.Serialize(person); + } + + return base.ConvertTo(context, culture, value, destinationType); + } +} diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Tests.cs b/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Tests.cs index 629988fd4..23890b203 100644 --- a/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Tests.cs +++ b/test/unit/Elsa.Workflows.Core.UnitTests/ObjectConversion/Tests.cs @@ -1,3 +1,4 @@ +using System.ComponentModel; using System.Dynamic; using System.Text.Json; using System.Text.Json.Nodes; @@ -352,4 +353,99 @@ public class Tests if(shouldThrow) Assert.Throws(() => input.ConvertTo(targetType, options)); } + + [Fact] + public void ConvertTo_WithRegisteredPersonTypeConverter_ConvertsFromJsonToPerson() + { + try + { + // Arrange + TypeDescriptor.AddAttributes(typeof(Person), new TypeConverterAttribute(typeof(PersonTypeConverter))); + + var json = "{\"Name\":\"Alice\",\"Age\":30}"; + + // Act + var result = json.ConvertTo(_objectConverterOptions); + + // Assert + Assert.NotNull(result); + Assert.IsType(result); + Assert.Equal("Alice", result.Name); + Assert.Equal(30, result.Age); + } + finally + { + // Clean up type descriptor cache so other tests aren't affected + TypeDescriptor.Refresh(typeof(Person)); + } + } + + [Fact] + public void ConvertTo_WithRegisteredPersonTypeConverter_NullInput_ReturnsNull() + { + try + { + // Arrange + TypeDescriptor.AddAttributes(typeof(Person), new TypeConverterAttribute(typeof(PersonTypeConverter))); + + string? json = null; + + // Act + var result = json.ConvertTo(_objectConverterOptions); + + // Assert + Assert.Null(result); + } + finally + { + TypeDescriptor.Refresh(typeof(Person)); + } + } + + [Fact] + public void ConvertTo_WithRegisteredPersonTypeConverter_ConvertsFromPersonToJson() + { + try + { + // Arrange + TypeDescriptor.AddAttributes(typeof(Person), new TypeConverterAttribute(typeof(PersonTypeConverter))); + + var person = new Person { Name = "Bob", Age = 42 }; + + // Act + var result = person.ConvertTo(_objectConverterOptions); + + // Assert + Assert.NotNull(result); + var json = Assert.IsType(result); + Assert.Contains("\"Name\":\"Bob\"", json); + Assert.Contains("\"Age\":42", json); + } + finally + { + // Clean up to avoid polluting TypeDescriptor globally + TypeDescriptor.Refresh(typeof(Person)); + } + } + + [Fact] + public void ConvertTo_WithRegisteredPersonTypeConverter_NullPersonToString_ReturnsNull() + { + try + { + // Arrange + TypeDescriptor.AddAttributes(typeof(Person), new TypeConverterAttribute(typeof(PersonTypeConverter))); + Person? person = null; + + // Act + var result = person.ConvertTo(_objectConverterOptions); + + // Assert + Assert.Null(result); + } + finally + { + TypeDescriptor.Refresh(typeof(Person)); + } + } } \ No newline at end of file From a12a71f36747b7544956cfe0ce140cad18dcc199 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 11 Nov 2025 16:56:38 +0100 Subject: [PATCH 5/5] Renames and makes async cancel method. Renames the `CancelAncestorActivatesAsync` method to `CancelAncestorActivitiesAsync` for clarity and makes it a proper async method, ensuring correct asynchronous execution. This prevents potential issues when cancelling ancestor activities. --- .../Activities/Flowchart/Activities/FlowJoin.cs | 4 ++-- .../Activities/Flowchart/Activities/Flowchart.cs | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) 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 65390c828..067efafc8 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs @@ -43,7 +43,7 @@ public class FlowJoin : Activity, IJoinNode { if (Flowchart.CanWaitAllProceed(context)) { - Flowchart.CancelAncestorActivatesAsync(context); + await Flowchart.CancelAncestorActivitiesAsync(context); await context.CompleteActivityAsync(); } @@ -51,7 +51,7 @@ public class FlowJoin : Activity, IJoinNode } case FlowJoinMode.WaitAny: { - Flowchart.CancelAncestorActivatesAsync(context); + await Flowchart.CancelAncestorActivitiesAsync(context); await context.CompleteActivityAsync(); break; } 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 1d24195ee..217a9315f 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs @@ -351,7 +351,7 @@ public class Flowchart : Container return flowScope.AllInboundConnectionsVisited(flowGraph, activity); } - public static async void CancelAncestorActivatesAsync(ActivityExecutionContext context) + public static async Task CancelAncestorActivitiesAsync(ActivityExecutionContext context) { var flowchartContext = context.ParentActivityExecutionContext!; var flowchart = (Flowchart)flowchartContext.Activity;