From c3b6decf669cb1bfbc6740b71e7aafe0f5dcabf1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 25 Oct 2024 22:49:49 +0200 Subject: [PATCH 01/17] Update default version to 3.2.2 in workflows Changed the fallback version from 3.2.1-preview to 3.2.2-preview in the GitHub Actions workflow configuration. This ensures that new preview versions are correctly labeled with the latest version number. --- .github/workflows/packages.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index c4461693c..39eb047a6 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -62,7 +62,7 @@ jobs: TAG_NAME=${TAG_NAME#refs/tags/} # remove the refs/tags/ prefix echo "VERSION=${TAG_NAME}" >> $GITHUB_ENV else - echo "VERSION=3.2.1-preview.${{github.run_number}}" >> $GITHUB_ENV + echo "VERSION=3.2.2-preview.${{github.run_number}}" >> $GITHUB_ENV fi - name: Set up JDK 17 uses: actions/setup-java@v2 From 91b7e0c800c06afbdb3324453d29ebb7acf5306f Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 31 Oct 2024 19:53:26 +0100 Subject: [PATCH 02/17] Add authorization to WorkflowInstanceHub This commit decorates the WorkflowInstanceHub class with the [Authorize] attribute to ensure that only authorized users can connect to the SignalR hub. This change improves the security of the workflow event notification system. --- .../Elsa.Workflows.Api/RealTime/Hubs/WorkflowInstanceHub.cs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/modules/Elsa.Workflows.Api/RealTime/Hubs/WorkflowInstanceHub.cs b/src/modules/Elsa.Workflows.Api/RealTime/Hubs/WorkflowInstanceHub.cs index 835a955f7..3426c07ef 100644 --- a/src/modules/Elsa.Workflows.Api/RealTime/Hubs/WorkflowInstanceHub.cs +++ b/src/modules/Elsa.Workflows.Api/RealTime/Hubs/WorkflowInstanceHub.cs @@ -1,6 +1,7 @@ using Elsa.Workflows.Api.RealTime.Contracts; using Elsa.Workflows.Runtime.Contracts; using JetBrains.Annotations; +using Microsoft.AspNetCore.Authorization; using Microsoft.AspNetCore.SignalR; namespace Elsa.Workflows.Api.RealTime.Hubs; @@ -9,6 +10,7 @@ namespace Elsa.Workflows.Api.RealTime.Hubs; /// Represents a SignalR hub for receiving workflow events on the client. /// [PublicAPI] +[Authorize] public class WorkflowInstanceHub : Hub { private readonly IWorkflowRuntime _workflowRuntime; From b3c73e99eab12587336e6d40b8bd87168d6749b8 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 31 Oct 2024 19:59:42 +0100 Subject: [PATCH 03/17] Disable SignalR in Elsa Server Updated `useSignalR` flag to false in `Program.cs` due to Elsa Studio's current inability to send authenticated requests to the SignalR hub. This change ensures better security and stability until the necessary update is implemented. --- src/bundles/Elsa.Server.Web/Program.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index 33328dace..209451473 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -48,7 +48,7 @@ const bool runEFCoreMigrations = true; const bool useMemoryStores = false; const bool useCaching = true; const bool useReadOnlyMode = false; -const bool useSignalR = true; +const bool useSignalR = false; // Disable until Elsa Studio is updated to send authenticated requests to the SignalR hub. const bool useAzureServiceBus = false; const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit; const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory; From 311f59894e2b9a75c47a13ab90a434eddfc6422e Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 6 Nov 2024 11:58:13 +0100 Subject: [PATCH 04/17] Add workflow graph network builder and child workflow finder Introduce IWorkflowGraphNetworkBuilder and IChildWorkflowFinder interfaces with their implementations to build and manage workflow graph networks. Additionally, register new services and handlers to support these functionalities in the Workflow Management feature. --- .../Contracts/IChildWorkflowFinder.cs | 17 ++++ .../Contracts/IWorkflowGraphNetworkBuilder.cs | 16 +++ .../Features/WorkflowManagementFeature.cs | 2 + .../Handlers/UpdateConsumingWorkflows.cs | 13 +++ .../Models/WorkflowGraphNetwork.cs | 3 + .../Models/WorkflowGraphNode.cs | 24 +++++ .../Services/WorkflowGraphNetworkBuilder.cs | 99 +++++++++++++++++++ 7 files changed, 174 insertions(+) create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs create mode 100644 src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs b/src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs new file mode 100644 index 000000000..351493fc5 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs @@ -0,0 +1,17 @@ +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Finds child workflows for a given workflow graph. +/// +public interface IChildWorkflowFinder +{ + /// + /// Finds child workflows for a given workflow graph. + /// + /// The workflow graph. + /// The cancellation token. + /// A collection of child workflows. + Task> FindChildWorkflowsAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs new file mode 100644 index 000000000..a22fd8aac --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs @@ -0,0 +1,16 @@ +using Elsa.Workflows.Management.Models; + +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Defines a visitor that can traverse and process workflow definitions. +/// +public interface IWorkflowGraphNetworkBuilder +{ + /// + /// Builds a network of workflow graphs and their consumers. + /// + /// The cancellation token. + /// A network of workflow graphs and their consumers. + Task BuildAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 3fdc5adf8..75b36b111 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -221,6 +221,7 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddSingleton() .AddSingleton() @@ -233,6 +234,7 @@ public class WorkflowManagementFeature : FeatureBase Services .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() ; Services.Configure(options => diff --git a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs new file mode 100644 index 000000000..c031fe183 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs @@ -0,0 +1,13 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Notifications; + +namespace Elsa.Workflows.Management.Handlers; + +public class UpdateConsumingWorkflows(IWorkflowGraphNetworkBuilder workflowGraphNetworkBuilder) : INotificationHandler +{ + public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) + { + var network = await workflowGraphNetworkBuilder.BuildAsync(cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs new file mode 100644 index 000000000..bb6ddc551 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs @@ -0,0 +1,3 @@ +namespace Elsa.Workflows.Management.Models; + +public record WorkflowGraphNetwork(HashSet Nodes); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs new file mode 100644 index 000000000..765f1c92b --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs @@ -0,0 +1,24 @@ +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Models; + +/// +/// Represents an activity in the context of an hierarchical tree structure, providing access to its siblings, parents and children. +/// +public class WorkflowGraphNode(WorkflowGraph workflowGraph) +{ + /// + /// Gets the workflow graph associated with this node. + /// + public WorkflowGraph WorkflowGraph { get; } = workflowGraph; + + /// + /// Gets the parents of this node. + /// + public HashSet Predecessors { get; set; } = new(); + + /// + /// Gets the children of this node. + /// + public HashSet Successors { get; set; } = new(); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs new file mode 100644 index 000000000..c725fbe3e --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs @@ -0,0 +1,99 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Services; + +/// +public class WorkflowGraphNetworkBuilder(IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowDefinitionService workflowDefinitionService) : IWorkflowGraphNetworkBuilder +{ + /// + public async Task BuildAsync(CancellationToken cancellationToken = default) + { + var workflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList(); + var nodes = new HashSet(); + + var context = new WorkflowGraphNetworkBuilderContext + { + WorkflowGraphs = workflowGraphs, + VisitedNodes = nodes + }; + + foreach (var workflowGraph in context.WorkflowGraphs) + VisitWorkflowGraph(workflowGraph, context); + + return new WorkflowGraphNetwork(nodes); + } + + private void VisitWorkflowGraph(WorkflowGraph workflowGraph, WorkflowGraphNetworkBuilderContext context) + { + var activityNodes = workflowGraph.Nodes; + var node = context.VisitedNodes.FirstOrDefault(x => x.WorkflowGraph == workflowGraph); + + if (node == null) + { + node = new WorkflowGraphNode(workflowGraph); + context.VisitedNodes.Add(node); + } + + var workflowDefinitionNodes = activityNodes + .Where(x => x.Activity is WorkflowDefinitionActivity) + .Select(x => (WorkflowDefinitionActivity)x.Activity) + .ToList(); + + foreach (var workflowDefinitionNode in workflowDefinitionNodes) + { + VisitConsumingActivityNode(node, workflowDefinitionNode, context); + } + } + + private void VisitConsumingActivityNode(WorkflowGraphNode node, WorkflowDefinitionActivity workflowDefinitionNode, WorkflowGraphNetworkBuilderContext context) + { + var consumedWorkflowDefinitionId = workflowDefinitionNode.WorkflowDefinitionId; + var consumedWorkflowGraph = context.WorkflowGraphs.FirstOrDefault(x => x.Workflow.Identity.DefinitionId == consumedWorkflowDefinitionId); + + if (consumedWorkflowGraph == null) + return; + + var childNode = context.VisitedNodes.FirstOrDefault(x => x.WorkflowGraph == consumedWorkflowGraph); + + if (childNode == null) + { + childNode = new WorkflowGraphNode(consumedWorkflowGraph); + context.VisitedNodes.Add(childNode); + VisitWorkflowGraph(childNode.WorkflowGraph, context); + } + + node.Successors.Add(childNode); + childNode.Predecessors.Add(node); + } + + private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) + { + var workflowDefinitionFilter = new WorkflowDefinitionFilter + { + VersionOptions = VersionOptions.LatestOrPublished + }; + var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); + var workflowGraphs = new List(); + + foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) + { + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + if (workflowGraph != null) + workflowGraphs.Add(workflowGraph); + } + + return workflowGraphs; + } + + private class WorkflowGraphNetworkBuilderContext + { + public ICollection WorkflowGraphs { get; set; } = new List(); + public HashSet VisitedNodes { get; set; } = new(); + } +} \ No newline at end of file From 9c7ee16e9615da38e037f4accbf48ded71541181 Mon Sep 17 00:00:00 2001 From: MariusVuscanNx <96233009+MariusVuscanNx@users.noreply.github.com> Date: Wed, 6 Nov 2024 16:04:28 +0200 Subject: [PATCH 05/17] Updated refit (#6098) --- Directory.Packages.props | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/Directory.Packages.props b/Directory.Packages.props index d61250a20..ccb99ec62 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -84,8 +84,6 @@ - - @@ -139,6 +137,8 @@ + + @@ -170,5 +170,7 @@ + + \ No newline at end of file From 8c7a8233745d5d0890fa759a71bb2b3f957d81eb Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 11:45:25 +0100 Subject: [PATCH 06/17] Remove unused workflow graph network builder Deleted the unused `IWorkflowGraphNetworkBuilder` interface, its implementation, and related files. Updated affected classes to replace dependency on the removed builder with direct methods. Enhanced affected workflows tracking during publication. --- .../BulkPublish/Endpoint.cs | 6 +- .../WorkflowDefinitions/Post/Endpoint.cs | 6 +- .../WorkflowDefinitions/Publish/Endpoint.cs | 2 +- .../WorkflowDefinitionActivity.cs | 2 +- .../Contracts/IChildWorkflowFinder.cs | 17 --- .../Contracts/IWorkflowGraphNetworkBuilder.cs | 16 --- .../Features/WorkflowManagementFeature.cs | 1 - .../Handlers/UpdateConsumingWorkflows.cs | 103 +++++++++++++++++- .../Models/AffectedWorkflows.cs | 5 + .../Models/PublishWorkflowDefinitionResult.cs | 4 +- .../WorkflowDefinitionPublished.cs | 3 +- .../Services/WorkflowDefinitionPublisher.cs | 13 +-- .../Services/WorkflowGraphNetworkBuilder.cs | 99 ----------------- 13 files changed, 120 insertions(+), 157 deletions(-) delete mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs delete mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs delete mode 100644 src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs index 675c80e79..2b11c9125 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs @@ -66,10 +66,8 @@ internal class BulkPublish(IWorkflowDefinitionStore store, IWorkflowDefinitionPu var result = await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken); published.Add(definitionId); - if (result.ConsumingWorkflows?.Any() == true) - { - updatedConsumers.AddRange(result.ConsumingWorkflows.Select(x => x.DefinitionId)); - } + if (result.AffectedWorkflows.WorkflowDefinitions.Count > 0) + updatedConsumers.AddRange(result.AffectedWorkflows.WorkflowDefinitions.Select(x => x.DefinitionId)); } return new Response(published, alreadyPublished, notFound, skipped, updatedConsumers); diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs index df0e87a15..a3e84dc0b 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs @@ -7,6 +7,7 @@ using Elsa.Workflows.Api.Models; using Elsa.Workflows.Api.Requirements; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Mappers; using Elsa.Workflows.Management.Materializers; using Elsa.Workflows.Management.Models; @@ -88,7 +89,7 @@ internal class Post( PublishWorkflowDefinitionResult? result = null; - if (request.Publish.GetValueOrDefault(false)) + if (request.Publish == true) { result = await workflowDefinitionPublisher.PublishAsync(draft, cancellationToken); @@ -107,7 +108,8 @@ internal class Post( } var mappedDefinition = await linker.MapAsync(draft, cancellationToken); - var response = new Response(mappedDefinition, false, result?.ConsumingWorkflows?.Count() ?? 0); + var affectedWorkflows = result?.AffectedWorkflows?.WorkflowDefinitions ?? []; + var response = new Response(mappedDefinition, false, affectedWorkflows.Count); await HttpContext.Response.WriteAsJsonAsync(response, serializerOptions, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs index 2f9cd2ca6..26e0ad9b3 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs @@ -46,7 +46,7 @@ internal class Publish(IWorkflowDefinitionStore store, IWorkflowDefinitionPublis var isPublished = definition.IsPublished; var result = !isPublished ? await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken) : null; var mappedDefinition = await linker.MapAsync(definition, cancellationToken); - var response = new Response(mappedDefinition, isPublished, result?.ConsumingWorkflows?.Count() ?? 0); + var response = new Response(mappedDefinition, isPublished, result?.AffectedWorkflows.WorkflowDefinitions.Count ?? 0); await SendOkAsync(response, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index abddaf098..cc1934b36 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -52,7 +52,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable var serviceProvider = context.ServiceProvider; var cancellationToken = context.CancellationToken; - // Find the workflow definition and not the graph; the graph must be computed at runtime, since one NodeIds will vary across graphs. + // Find the workflow definition and not the graph; the graph must be computed at runtime, since NodeIds will vary across graphs. var workflowDefinition = await GetWorkflowDefinitionAsync(serviceProvider, cancellationToken); if (workflowDefinition == null) diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs b/src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs deleted file mode 100644 index 351493fc5..000000000 --- a/src/modules/Elsa.Workflows.Management/Contracts/IChildWorkflowFinder.cs +++ /dev/null @@ -1,17 +0,0 @@ -using Elsa.Workflows.Models; - -namespace Elsa.Workflows.Management.Contracts; - -/// -/// Finds child workflows for a given workflow graph. -/// -public interface IChildWorkflowFinder -{ - /// - /// Finds child workflows for a given workflow graph. - /// - /// The workflow graph. - /// The cancellation token. - /// A collection of child workflows. - Task> FindChildWorkflowsAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default); -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs deleted file mode 100644 index a22fd8aac..000000000 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowGraphNetworkBuilder.cs +++ /dev/null @@ -1,16 +0,0 @@ -using Elsa.Workflows.Management.Models; - -namespace Elsa.Workflows.Management.Contracts; - -/// -/// Defines a visitor that can traverse and process workflow definitions. -/// -public interface IWorkflowGraphNetworkBuilder -{ - /// - /// Builds a network of workflow graphs and their consumers. - /// - /// The cancellation token. - /// A network of workflow graphs and their consumers. - Task BuildAsync(CancellationToken cancellationToken = default); -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 75b36b111..e8be69a28 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -221,7 +221,6 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() - .AddScoped() .AddScoped() .AddSingleton() .AddSingleton() diff --git a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs index c031fe183..c43173352 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs @@ -1,13 +1,112 @@ +using Elsa.Common.Models; using Elsa.Mediator.Contracts; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Notifications; +using Elsa.Workflows.Models; namespace Elsa.Workflows.Management.Handlers; -public class UpdateConsumingWorkflows(IWorkflowGraphNetworkBuilder workflowGraphNetworkBuilder) : INotificationHandler +/// +/// Updates consuming workflows when a workflow definition is published. +/// +public class UpdateConsumingWorkflows( + IWorkflowDefinitionPublisher publisher, + IWorkflowDefinitionService workflowDefinitionService, + IWorkflowDefinitionStore workflowDefinitionStore, + IApiSerializer serializer) : INotificationHandler { + /// public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) { - var network = await workflowGraphNetworkBuilder.BuildAsync(cancellationToken); + var definition = notification.WorkflowDefinition; + + // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. + if (definition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) + return; + + // Build a network of workflow graphs. + var allWorkflowDefinitionGraphs = await GetAllWorkflowGraphsAsync(cancellationToken); + var definitionId = definition.DefinitionId; + + // Find all root workflow definition activity nodes. + foreach (var workflowGraph in allWorkflowDefinitionGraphs) + { + var consumerDefinitionId = workflowGraph.Workflow.Identity.DefinitionId; + + var workflowDefinitionActivities = + FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root) + .Where(x => x.WorkflowDefinitionId == definitionId) + .ToList(); + + // Skip if the published workflow definition is not used in any workflow graph. + if (workflowDefinitionActivities.Count == 0) + continue; + + // Create a new version of the published workflow definition. + var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished; + var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + if(newVersion == null) + continue; + + var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); + var newWorkflowDefinitionActivities = + FindWorkflowActivityDefinitionActivityNodes(newWorkflowGraph.Root) + .Where(x => x.WorkflowDefinitionId == definitionId) + .ToList(); + + foreach (var workflowDefinitionActivity in newWorkflowDefinitionActivities) + { + // Update the consuming workflow graph to use the new version of the published workflow definition. + workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id; + workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id; + } + + // Update the new version of the published workflow definition. + if(newWorkflowGraph.Root.Activity is Workflow newWorkflow) + newVersion.StringData = serializer.Serialize(newWorkflow.Root); + + // If the draft is new, publish it. + if (originalVersionIsPublished) + await publisher.PublishAsync(newVersion, cancellationToken); + else + await publisher.SaveDraftAsync(newVersion, cancellationToken); + } + } + + private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) + { + foreach (var child in parent.Children) + { + if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity) + yield return workflowDefinitionActivity; + + foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child)) + yield return grandChild; + } + } + + private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) + { + var workflowDefinitionFilter = new WorkflowDefinitionFilter + { + VersionOptions = VersionOptions.LatestOrPublished + }; + var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); + var workflowGraphs = new List(); + + foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) + { + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + if (workflowGraph != null) + workflowGraphs.Add(workflowGraph); + } + + return workflowGraphs; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs b/src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs new file mode 100644 index 000000000..da5a0d71b --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs @@ -0,0 +1,5 @@ +using Elsa.Workflows.Management.Entities; + +namespace Elsa.Workflows.Management.Models; + +public record AffectedWorkflows(ICollection WorkflowDefinitions); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs b/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs index 87848df5a..31a9f37f3 100644 --- a/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs +++ b/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs @@ -1,8 +1,6 @@ -using Elsa.Workflows.Management.Entities; - namespace Elsa.Workflows.Management.Models; /// /// Represents the result of publishing a workflow definition. /// -public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection ValidationErrors, IEnumerable? ConsumingWorkflows); \ No newline at end of file +public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection ValidationErrors, AffectedWorkflows AffectedWorkflows); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs b/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs index 81db3b296..62c84c96c 100644 --- a/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs +++ b/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs @@ -1,5 +1,6 @@ using Elsa.Mediator.Contracts; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; using JetBrains.Annotations; namespace Elsa.Workflows.Management.Notifications; @@ -9,4 +10,4 @@ namespace Elsa.Workflows.Management.Notifications; /// /// The workflow definition. [PublicAPI] -public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition) : INotification; \ No newline at end of file +public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition, AffectedWorkflows AffectedWorkflows) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs index 173c3b29d..62bfc8e94 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -121,16 +121,9 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher definition = Initialize(definition); await _workflowDefinitionStore.SaveAsync(definition, cancellationToken); - await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition), cancellationToken); - - var consumingWorkflows = new List(); - - if (definition.Options is { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) - { - consumingWorkflows.AddRange(await UpdateReferencesInConsumingWorkflows(definition, cancellationToken)); - } - - return new PublishWorkflowDefinitionResult(true, validationErrors, consumingWorkflows); + var affectedWorkflows = new AffectedWorkflows(new List()); + await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition, affectedWorkflows), cancellationToken); + return new PublishWorkflowDefinitionResult(true, validationErrors, affectedWorkflows); } /// diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs deleted file mode 100644 index c725fbe3e..000000000 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowGraphNetworkBuilder.cs +++ /dev/null @@ -1,99 +0,0 @@ -using Elsa.Common.Models; -using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; -using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Management.Filters; -using Elsa.Workflows.Management.Models; -using Elsa.Workflows.Models; - -namespace Elsa.Workflows.Management.Services; - -/// -public class WorkflowGraphNetworkBuilder(IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowDefinitionService workflowDefinitionService) : IWorkflowGraphNetworkBuilder -{ - /// - public async Task BuildAsync(CancellationToken cancellationToken = default) - { - var workflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList(); - var nodes = new HashSet(); - - var context = new WorkflowGraphNetworkBuilderContext - { - WorkflowGraphs = workflowGraphs, - VisitedNodes = nodes - }; - - foreach (var workflowGraph in context.WorkflowGraphs) - VisitWorkflowGraph(workflowGraph, context); - - return new WorkflowGraphNetwork(nodes); - } - - private void VisitWorkflowGraph(WorkflowGraph workflowGraph, WorkflowGraphNetworkBuilderContext context) - { - var activityNodes = workflowGraph.Nodes; - var node = context.VisitedNodes.FirstOrDefault(x => x.WorkflowGraph == workflowGraph); - - if (node == null) - { - node = new WorkflowGraphNode(workflowGraph); - context.VisitedNodes.Add(node); - } - - var workflowDefinitionNodes = activityNodes - .Where(x => x.Activity is WorkflowDefinitionActivity) - .Select(x => (WorkflowDefinitionActivity)x.Activity) - .ToList(); - - foreach (var workflowDefinitionNode in workflowDefinitionNodes) - { - VisitConsumingActivityNode(node, workflowDefinitionNode, context); - } - } - - private void VisitConsumingActivityNode(WorkflowGraphNode node, WorkflowDefinitionActivity workflowDefinitionNode, WorkflowGraphNetworkBuilderContext context) - { - var consumedWorkflowDefinitionId = workflowDefinitionNode.WorkflowDefinitionId; - var consumedWorkflowGraph = context.WorkflowGraphs.FirstOrDefault(x => x.Workflow.Identity.DefinitionId == consumedWorkflowDefinitionId); - - if (consumedWorkflowGraph == null) - return; - - var childNode = context.VisitedNodes.FirstOrDefault(x => x.WorkflowGraph == consumedWorkflowGraph); - - if (childNode == null) - { - childNode = new WorkflowGraphNode(consumedWorkflowGraph); - context.VisitedNodes.Add(childNode); - VisitWorkflowGraph(childNode.WorkflowGraph, context); - } - - node.Successors.Add(childNode); - childNode.Predecessors.Add(node); - } - - private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) - { - var workflowDefinitionFilter = new WorkflowDefinitionFilter - { - VersionOptions = VersionOptions.LatestOrPublished - }; - var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); - var workflowGraphs = new List(); - - foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) - { - var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken); - - if (workflowGraph != null) - workflowGraphs.Add(workflowGraph); - } - - return workflowGraphs; - } - - private class WorkflowGraphNetworkBuilderContext - { - public ICollection WorkflowGraphs { get; set; } = new List(); - public HashSet VisitedNodes { get; set; } = new(); - } -} \ No newline at end of file From c7fa4468be1d156d8d02d2f70fb931ba9f8571e1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 11:57:15 +0100 Subject: [PATCH 07/17] Optimize workflow update logic to skip redundant checks Added condition to ensure only relevant activities are processed by verifying their version IDs. This reduces unnecessary computations and ensures the system skips already up-to-date workflow activities. --- .../Handlers/UpdateConsumingWorkflows.cs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs index c43173352..a596382a8 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs @@ -56,9 +56,13 @@ public class UpdateConsumingWorkflows( var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); var newWorkflowDefinitionActivities = FindWorkflowActivityDefinitionActivityNodes(newWorkflowGraph.Root) - .Where(x => x.WorkflowDefinitionId == definitionId) + .Where(x => x.WorkflowDefinitionId == definitionId && x.WorkflowDefinitionVersionId != definition.Id) .ToList(); + // Skip if the new version of the published workflow definition is not used in any workflow graph, or if the activity is already up-to-date. + if(newWorkflowDefinitionActivities.Count == 0) + continue; + foreach (var workflowDefinitionActivity in newWorkflowDefinitionActivities) { // Update the consuming workflow graph to use the new version of the published workflow definition. From cc86ee9f583c2c0d84d3b4bad1906b92fa9357fc Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 17:16:21 +0100 Subject: [PATCH 08/17] Add workflow reference updater implementation Introduced IWorkflowReferenceUpdater interface to manage workflow references updates. Implemented the corresponding service and integrated it into existing event handling logic. This improves maintainability and reduces code duplication by isolating the reference update logic. --- .../Contracts/IWorkflowReferenceUpdater.cs | 18 +++ .../Features/WorkflowManagementFeature.cs | 1 + .../Handlers/UpdateConsumingWorkflows.cs | 99 +------------- .../Models/UpdateWorkflowReferencesResult.cs | 18 +++ .../Services/WorkflowReferenceUpdater.cs | 128 ++++++++++++++++++ 5 files changed, 170 insertions(+), 94 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs new file mode 100644 index 000000000..81bed0900 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs @@ -0,0 +1,18 @@ +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; + +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Updates references to the specified workflow of all workflows that reference it. +/// +public interface IWorkflowReferenceUpdater +{ + /// + /// Updates references to the specified workflow of all workflows that reference it. + /// + /// The workflow definition to update references for. + /// The cancellation token. + /// The result of the operation. + Task UpdateWorkflowReferencesAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index e8be69a28..47de8144c 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -213,6 +213,7 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddSingleton() .AddSingleton() diff --git a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs index a596382a8..5f8e3811c 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs @@ -1,9 +1,11 @@ using Elsa.Common.Models; +using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Notifications; using Elsa.Workflows.Models; @@ -13,104 +15,13 @@ namespace Elsa.Workflows.Management.Handlers; /// /// Updates consuming workflows when a workflow definition is published. /// -public class UpdateConsumingWorkflows( - IWorkflowDefinitionPublisher publisher, - IWorkflowDefinitionService workflowDefinitionService, - IWorkflowDefinitionStore workflowDefinitionStore, - IApiSerializer serializer) : INotificationHandler +public class UpdateConsumingWorkflows(IWorkflowReferenceUpdater workflowReferenceUpdater) : INotificationHandler { /// public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) { var definition = notification.WorkflowDefinition; - - // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. - if (definition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) - return; - - // Build a network of workflow graphs. - var allWorkflowDefinitionGraphs = await GetAllWorkflowGraphsAsync(cancellationToken); - var definitionId = definition.DefinitionId; - - // Find all root workflow definition activity nodes. - foreach (var workflowGraph in allWorkflowDefinitionGraphs) - { - var consumerDefinitionId = workflowGraph.Workflow.Identity.DefinitionId; - - var workflowDefinitionActivities = - FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root) - .Where(x => x.WorkflowDefinitionId == definitionId) - .ToList(); - - // Skip if the published workflow definition is not used in any workflow graph. - if (workflowDefinitionActivities.Count == 0) - continue; - - // Create a new version of the published workflow definition. - var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished; - var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken); - - if(newVersion == null) - continue; - - var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); - var newWorkflowDefinitionActivities = - FindWorkflowActivityDefinitionActivityNodes(newWorkflowGraph.Root) - .Where(x => x.WorkflowDefinitionId == definitionId && x.WorkflowDefinitionVersionId != definition.Id) - .ToList(); - - // Skip if the new version of the published workflow definition is not used in any workflow graph, or if the activity is already up-to-date. - if(newWorkflowDefinitionActivities.Count == 0) - continue; - - foreach (var workflowDefinitionActivity in newWorkflowDefinitionActivities) - { - // Update the consuming workflow graph to use the new version of the published workflow definition. - workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id; - workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id; - } - - // Update the new version of the published workflow definition. - if(newWorkflowGraph.Root.Activity is Workflow newWorkflow) - newVersion.StringData = serializer.Serialize(newWorkflow.Root); - - // If the draft is new, publish it. - if (originalVersionIsPublished) - await publisher.PublishAsync(newVersion, cancellationToken); - else - await publisher.SaveDraftAsync(newVersion, cancellationToken); - } - } - - private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) - { - foreach (var child in parent.Children) - { - if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity) - yield return workflowDefinitionActivity; - - foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child)) - yield return grandChild; - } - } - - private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) - { - var workflowDefinitionFilter = new WorkflowDefinitionFilter - { - VersionOptions = VersionOptions.LatestOrPublished - }; - var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); - var workflowGraphs = new List(); - - foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) - { - var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken); - - if (workflowGraph != null) - workflowGraphs.Add(workflowGraph); - } - - return workflowGraphs; + var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken); + notification.AffectedWorkflows.WorkflowDefinitions.AddRange(result.UpdatedWorkflows); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs b/src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs new file mode 100644 index 000000000..eebddc7aa --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs @@ -0,0 +1,18 @@ +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Models; + +/// +/// Represents the result of updating workflow references. +/// +public class UpdateWorkflowReferencesResult(IEnumerable updatedWorkflows) +{ + /// + /// Gets a collection of workflow graphs that have been updated. + /// + /// + /// A read-only collection of instances representing the updated workflows. + /// + public IReadOnlyCollection UpdatedWorkflows { get; } = updatedWorkflows.ToList(); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs new file mode 100644 index 000000000..db7cbb3dc --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -0,0 +1,128 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Services; + +/// +public class WorkflowReferenceUpdater( + IWorkflowDefinitionPublisher publisher, + IWorkflowDefinitionService workflowDefinitionService, + IWorkflowDefinitionStore workflowDefinitionStore, + IApiSerializer serializer) : IWorkflowReferenceUpdater +{ + /// + public async Task UpdateWorkflowReferencesAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + { + // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. + if (definition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) + return new UpdateWorkflowReferencesResult([]); + + // Find all workflow graphs that contain the updated workflow definition. + var matchingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(definition.DefinitionId, cancellationToken); + + // Find all root workflow definition activity nodes. + var updatedWorkflows = new List(); + foreach (var workflowGraph in matchingWorkflowGraphs) + { + var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, definition, cancellationToken); + + if (newDefinition != null) + updatedWorkflows.Add(newDefinition); + } + + return new UpdateWorkflowReferencesResult(updatedWorkflows); + } + + private async Task UpdateConsumingWorkflowAsync(WorkflowGraph workflowGraph, WorkflowDefinition definition, CancellationToken cancellationToken) + { + // Create a new version of the published workflow definition or get the existing draft. + var consumerDefinitionId = workflowGraph.Workflow.Identity.DefinitionId; + var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished; + var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + if (newVersion == null) + return null; + + // Materialize the new/draft version to find all workflow definition activities that use the updated workflow definition. + var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); + var outdatedWorkflowDefinitionActivities = FindOutdatedWorkflowDefinitionActivities(newWorkflowGraph, definition).ToList(); + + // Skip if the new version of the published workflow definition is not used in the workflow or if the activity is already up to date. + if (outdatedWorkflowDefinitionActivities.Count == 0) + return null; + + // Update the consuming workflow graph to use the new version of the published workflow definition. + foreach (var workflowDefinitionActivity in outdatedWorkflowDefinitionActivities) + { + workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id; + workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id; + } + + // Update the new version of the published workflow definition. + if (newWorkflowGraph.Root.Activity is Workflow newWorkflow) + newVersion.StringData = serializer.Serialize(newWorkflow.Root); + + // If the draft is new, publish it. + if (originalVersionIsPublished) + await publisher.PublishAsync(newVersion, cancellationToken); + else + await publisher.SaveDraftAsync(newVersion, cancellationToken); + + return newVersion; + } + + private IEnumerable FindOutdatedWorkflowDefinitionActivities(WorkflowGraph workflowGraph, WorkflowDefinition updatedDefinition) + { + return FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root) + .Where(x => x.WorkflowDefinitionId == updatedDefinition.DefinitionId && x.WorkflowDefinitionVersionId != updatedDefinition.Id); + } + + private async Task> FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(string updatedWorkflowDefinitionId, CancellationToken cancellationToken) + { + var allWorkflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList(); + var filteredWorkflowGraphs = allWorkflowGraphs + .Where(x => FindWorkflowActivityDefinitionActivityNodes(x.Root).Any(y => y.WorkflowDefinitionId == updatedWorkflowDefinitionId)) + .ToList(); + + return filteredWorkflowGraphs; + } + + private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) + { + foreach (var child in parent.Children) + { + if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity) + yield return workflowDefinitionActivity; + + foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child)) + yield return grandChild; + } + } + + private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) + { + var workflowDefinitionFilter = new WorkflowDefinitionFilter + { + VersionOptions = VersionOptions.LatestOrPublished + }; + var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); + var workflowGraphs = new List(); + + foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) + { + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + if (workflowGraph != null) + workflowGraphs.Add(workflowGraph); + } + + return workflowGraphs; + } +} \ No newline at end of file From a7b0d5630de58979db5837ed8d7c3b6294178ff9 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 17:20:40 +0100 Subject: [PATCH 09/17] Set workflow activity version during update Added a missing assignment of the workflow activity version in the `UpdateWorkflowDefinition` method. This ensures that the `Version` field is properly set when workflow definitions are updated. --- .../Services/WorkflowReferenceUpdater.cs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs index db7cbb3dc..b244b13b2 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -63,6 +63,7 @@ public class WorkflowReferenceUpdater( { workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id; workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id; + workflowDefinitionActivity.Version = definition.Version; } // Update the new version of the published workflow definition. From 8fe0147abd7ea92b2e851337a25dfa59d442f826 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:00:21 +0100 Subject: [PATCH 10/17] Remove UpdateReferencesInConsumingWorkflows method Eliminated the UpdateReferencesInConsumingWorkflows method from IWorkflowDefinitionPublisher and its implementation from WorkflowDefinitionPublisher. The responsibility of updating workflow references is now transitioned to IWorkflowReferenceUpdater used in the `UpdateReferences` endpoint. --- .../UpdateReferences/Endpoint.cs | 5 +- .../Contracts/IWorkflowDefinitionPublisher.cs | 8 --- .../Services/WorkflowDefinitionPublisher.cs | 56 ------------------- 3 files changed, 3 insertions(+), 66 deletions(-) diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs index 4efbeb133..c10119f22 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs @@ -10,7 +10,7 @@ using Microsoft.AspNetCore.Authorization; namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.UpdateReferences; [PublicAPI] -internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinitionPublisher workflowDefinitionPublisher, IAuthorizationService authorizationService) +internal class UpdateReferences(IWorkflowReferenceUpdater workflowReferenceUpdater, IWorkflowDefinitionStore store, IAuthorizationService authorizationService) : ElsaEndpoint { public override void Configure() @@ -43,7 +43,8 @@ internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinit return; } - var affectedWorkflows = await workflowDefinitionPublisher.UpdateReferencesInConsumingWorkflows(definition, cancellationToken); + var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken); + var affectedWorkflows = result.UpdatedWorkflows; var response = new Response(affectedWorkflows.Select(w => w.Name ?? w.DefinitionId)); await SendOkAsync(response, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs index 487074525..402243f73 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs @@ -68,12 +68,4 @@ public interface IWorkflowDefinitionPublisher /// The cancellation token. /// The saved workflow definition. Task SaveDraftAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); - - /// - /// Updates all referencing workflow definitions to use the version of the specified workflow definition. - /// - /// The workflow definition to update references for. - /// The cancellation token. - /// The updated workflow definitions. - Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs index 62bfc8e94..7f5bf7bfd 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -1,11 +1,9 @@ using Elsa.Common.Contracts; using Elsa.Common.Entities; using Elsa.Common.Models; -using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; using Elsa.Workflows.Contracts; -using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; @@ -23,7 +21,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher private readonly IWorkflowDefinitionStore _workflowDefinitionStore; private readonly INotificationSender _notificationSender; private readonly IIdentityGenerator _identityGenerator; - private readonly IActivityVisitor _activityVisitor; private readonly IActivitySerializer _activitySerializer; private readonly IRequestSender _requestSender; private readonly ISystemClock _systemClock; @@ -36,7 +33,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher IWorkflowDefinitionStore workflowDefinitionStore, INotificationSender notificationSender, IIdentityGenerator identityGenerator, - IActivityVisitor activityVisitor, IActivitySerializer activitySerializer, IRequestSender requestSender, ISystemClock systemClock) @@ -45,7 +41,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher _workflowDefinitionStore = workflowDefinitionStore; _notificationSender = notificationSender; _identityGenerator = identityGenerator; - _activityVisitor = activityVisitor; _activitySerializer = activitySerializer; _requestSender = requestSender; _systemClock = systemClock; @@ -205,57 +200,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher return draft; } - /// - public async Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default) - { - var updatedWorkflowDefinitions = new List(); - - var workflowDefinitions = (await _workflowDefinitionStore.FindManyAsync(new WorkflowDefinitionFilter - { - VersionOptions = VersionOptions.LatestOrPublished - }, cancellationToken)).ToList(); - - // Remove the dependency from the list of workflow definitions to consider. - workflowDefinitions = workflowDefinitions.Where(x => x.DefinitionId != dependency.DefinitionId).ToList(); - - foreach (var definition in workflowDefinitions) - { - var root = _activitySerializer.Deserialize(definition.StringData!); - var graph = await _activityVisitor.VisitAsync(root, cancellationToken); - var flattenedList = graph.Flatten().ToList(); - var definitionId = dependency.DefinitionId; - var version = dependency.Version; - var nodes = flattenedList - .Where(x => x.Activity is WorkflowDefinitionActivity workflowDefinitionActivity && workflowDefinitionActivity.WorkflowDefinitionId == definitionId) - .ToList(); - - foreach (var node in nodes.Where(activity => activity.Activity.Version < version)) - { - var activity = (WorkflowDefinitionActivity)node.Activity; - activity.Version = version; - activity.WorkflowDefinitionVersionId = dependency.Id; - - if (!updatedWorkflowDefinitions.Contains(definition)) - updatedWorkflowDefinitions.Add(definition); - } - - if (updatedWorkflowDefinitions.Contains(definition)) - { - var serializedData = _activitySerializer.Serialize(root); - definition.StringData = serializedData; - } - } - - if (updatedWorkflowDefinitions.Any()) - { - await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdating(updatedWorkflowDefinitions), cancellationToken); - await _workflowDefinitionStore.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken); - await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdated(updatedWorkflowDefinitions), cancellationToken); - } - - return updatedWorkflowDefinitions; - } - private WorkflowDefinition Initialize(WorkflowDefinition definition) { if (definition.Id == null!) From 57cafd69a5ddb89cc4a16c5f4100360aedb3d83a Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:07:59 +0100 Subject: [PATCH 11/17] Remove WorkflowGraphNetwork and WorkflowGraphNode models These models were deleted as they are no longer needed in the system. The removal helps in reducing code clutter and improves maintainability by eliminating unused components. --- .../Models/WorkflowGraphNetwork.cs | 3 --- .../Models/WorkflowGraphNode.cs | 24 ------------------- 2 files changed, 27 deletions(-) delete mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs delete mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs deleted file mode 100644 index bb6ddc551..000000000 --- a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNetwork.cs +++ /dev/null @@ -1,3 +0,0 @@ -namespace Elsa.Workflows.Management.Models; - -public record WorkflowGraphNetwork(HashSet Nodes); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs deleted file mode 100644 index 765f1c92b..000000000 --- a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphNode.cs +++ /dev/null @@ -1,24 +0,0 @@ -using Elsa.Workflows.Models; - -namespace Elsa.Workflows.Management.Models; - -/// -/// Represents an activity in the context of an hierarchical tree structure, providing access to its siblings, parents and children. -/// -public class WorkflowGraphNode(WorkflowGraph workflowGraph) -{ - /// - /// Gets the workflow graph associated with this node. - /// - public WorkflowGraph WorkflowGraph { get; } = workflowGraph; - - /// - /// Gets the parents of this node. - /// - public HashSet Predecessors { get; set; } = new(); - - /// - /// Gets the children of this node. - /// - public HashSet Successors { get; set; } = new(); -} \ No newline at end of file From 6503612f7f9f916301b5684c85a580a9db427e8b Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:10:30 +0100 Subject: [PATCH 12/17] Refactor parameter naming in IWorkflowReferenceUpdater Renamed 'definition' to 'referencedDefinition' for clarity in the IWorkflowReferenceUpdater contract and its implementation. This change makes it explicit that the parameter refers to the workflow definition being referenced, ensuring the code is more understandable and maintainable. --- .../Contracts/IWorkflowReferenceUpdater.cs | 4 ++-- .../Services/WorkflowReferenceUpdater.cs | 8 ++++---- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs index 81bed0900..e70f7a855 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs @@ -11,8 +11,8 @@ public interface IWorkflowReferenceUpdater /// /// Updates references to the specified workflow of all workflows that reference it. /// - /// The workflow definition to update references for. + /// The workflow definition that is being referenced. All workflows that reference this definition will be updated to use this newest version. /// The cancellation token. /// The result of the operation. - Task UpdateWorkflowReferencesAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); + Task UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs index b244b13b2..f6b6e57d1 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -18,20 +18,20 @@ public class WorkflowReferenceUpdater( IApiSerializer serializer) : IWorkflowReferenceUpdater { /// - public async Task UpdateWorkflowReferencesAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + public async Task UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default) { // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. - if (definition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) + if (referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) return new UpdateWorkflowReferencesResult([]); // Find all workflow graphs that contain the updated workflow definition. - var matchingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(definition.DefinitionId, cancellationToken); + var matchingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken); // Find all root workflow definition activity nodes. var updatedWorkflows = new List(); foreach (var workflowGraph in matchingWorkflowGraphs) { - var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, definition, cancellationToken); + var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, cancellationToken); if (newDefinition != null) updatedWorkflows.Add(newDefinition); From 6aa40143f342c9fe6f87bfe5c077561230693c24 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:13:59 +0100 Subject: [PATCH 13/17] Update WorkflowReferenceUpdater for clarity and accuracy Renamed variables and updated comments for better clarity. The term "matchingWorkflowGraphs" was changed to "consumingWorkflowGraphs" to more accurately describe its purpose. Adjusted comment to clarify the materialization process of the draft version. --- .../Services/WorkflowReferenceUpdater.cs | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs index f6b6e57d1..ff43414ed 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -25,11 +25,11 @@ public class WorkflowReferenceUpdater( return new UpdateWorkflowReferencesResult([]); // Find all workflow graphs that contain the updated workflow definition. - var matchingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken); - - // Find all root workflow definition activity nodes. + var consumingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken); + + // Update consuming workflows. var updatedWorkflows = new List(); - foreach (var workflowGraph in matchingWorkflowGraphs) + foreach (var workflowGraph in consumingWorkflowGraphs) { var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, cancellationToken); @@ -47,10 +47,11 @@ public class WorkflowReferenceUpdater( var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished; var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + // This is null in case the definition no longer exists in the store. if (newVersion == null) return null; - // Materialize the new/draft version to find all workflow definition activities that use the updated workflow definition. + // Materialize the draft to find all workflow definition activities that use the updated workflow definition. var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); var outdatedWorkflowDefinitionActivities = FindOutdatedWorkflowDefinitionActivities(newWorkflowGraph, definition).ToList(); From e46fa846470c4c4d9ae90695badbb1f37bcaf4ed Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:18:47 +0100 Subject: [PATCH 14/17] Add IsReadonly filter to WorkflowDefinitionFilter This ensures that only non-readonly workflow definitions are fetched. It enhances the accuracy of workflow updates by filtering out readonly versions. --- .../Services/WorkflowReferenceUpdater.cs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs index ff43414ed..a6fdfab67 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -112,7 +112,8 @@ public class WorkflowReferenceUpdater( { var workflowDefinitionFilter = new WorkflowDefinitionFilter { - VersionOptions = VersionOptions.LatestOrPublished + VersionOptions = VersionOptions.LatestOrPublished, + IsReadonly = false }; var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); var workflowGraphs = new List(); From a0aa420a12b41f013226e10ed716fa6ad975d871 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:19:42 +0100 Subject: [PATCH 15/17] Keep Published Workflows Immutable by Updating Drafts Instead (#6106) * Add workflow graph network builder and child workflow finder Introduce IWorkflowGraphNetworkBuilder and IChildWorkflowFinder interfaces with their implementations to build and manage workflow graph networks. Additionally, register new services and handlers to support these functionalities in the Workflow Management feature. * Remove unused workflow graph network builder Deleted the unused `IWorkflowGraphNetworkBuilder` interface, its implementation, and related files. Updated affected classes to replace dependency on the removed builder with direct methods. Enhanced affected workflows tracking during publication. * Optimize workflow update logic to skip redundant checks Added condition to ensure only relevant activities are processed by verifying their version IDs. This reduces unnecessary computations and ensures the system skips already up-to-date workflow activities. * Add workflow reference updater implementation Introduced IWorkflowReferenceUpdater interface to manage workflow references updates. Implemented the corresponding service and integrated it into existing event handling logic. This improves maintainability and reduces code duplication by isolating the reference update logic. * Set workflow activity version during update Added a missing assignment of the workflow activity version in the `UpdateWorkflowDefinition` method. This ensures that the `Version` field is properly set when workflow definitions are updated. * Remove UpdateReferencesInConsumingWorkflows method Eliminated the UpdateReferencesInConsumingWorkflows method from IWorkflowDefinitionPublisher and its implementation from WorkflowDefinitionPublisher. The responsibility of updating workflow references is now transitioned to IWorkflowReferenceUpdater used in the `UpdateReferences` endpoint. * Remove WorkflowGraphNetwork and WorkflowGraphNode models These models were deleted as they are no longer needed in the system. The removal helps in reducing code clutter and improves maintainability by eliminating unused components. * Refactor parameter naming in IWorkflowReferenceUpdater Renamed 'definition' to 'referencedDefinition' for clarity in the IWorkflowReferenceUpdater contract and its implementation. This change makes it explicit that the parameter refers to the workflow definition being referenced, ensuring the code is more understandable and maintainable. * Update WorkflowReferenceUpdater for clarity and accuracy Renamed variables and updated comments for better clarity. The term "matchingWorkflowGraphs" was changed to "consumingWorkflowGraphs" to more accurately describe its purpose. Adjusted comment to clarify the materialization process of the draft version. * Add IsReadonly filter to WorkflowDefinitionFilter This ensures that only non-readonly workflow definitions are fetched. It enhances the accuracy of workflow updates by filtering out readonly versions. --- .../BulkPublish/Endpoint.cs | 6 +- .../WorkflowDefinitions/Post/Endpoint.cs | 6 +- .../WorkflowDefinitions/Publish/Endpoint.cs | 2 +- .../UpdateReferences/Endpoint.cs | 5 +- .../WorkflowDefinitionActivity.cs | 2 +- .../Contracts/IWorkflowDefinitionPublisher.cs | 8 -- .../Contracts/IWorkflowReferenceUpdater.cs | 18 +++ .../Features/WorkflowManagementFeature.cs | 2 + .../Handlers/UpdateConsumingWorkflows.cs | 27 ++++ .../Models/AffectedWorkflows.cs | 5 + .../Models/PublishWorkflowDefinitionResult.cs | 4 +- .../Models/UpdateWorkflowReferencesResult.cs | 18 +++ .../WorkflowDefinitionPublished.cs | 3 +- .../Services/WorkflowDefinitionPublisher.cs | 69 +-------- .../Services/WorkflowReferenceUpdater.cs | 131 ++++++++++++++++++ 15 files changed, 218 insertions(+), 88 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs create mode 100644 src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs index 675c80e79..2b11c9125 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs @@ -66,10 +66,8 @@ internal class BulkPublish(IWorkflowDefinitionStore store, IWorkflowDefinitionPu var result = await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken); published.Add(definitionId); - if (result.ConsumingWorkflows?.Any() == true) - { - updatedConsumers.AddRange(result.ConsumingWorkflows.Select(x => x.DefinitionId)); - } + if (result.AffectedWorkflows.WorkflowDefinitions.Count > 0) + updatedConsumers.AddRange(result.AffectedWorkflows.WorkflowDefinitions.Select(x => x.DefinitionId)); } return new Response(published, alreadyPublished, notFound, skipped, updatedConsumers); diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs index df0e87a15..a3e84dc0b 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs @@ -7,6 +7,7 @@ using Elsa.Workflows.Api.Models; using Elsa.Workflows.Api.Requirements; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Mappers; using Elsa.Workflows.Management.Materializers; using Elsa.Workflows.Management.Models; @@ -88,7 +89,7 @@ internal class Post( PublishWorkflowDefinitionResult? result = null; - if (request.Publish.GetValueOrDefault(false)) + if (request.Publish == true) { result = await workflowDefinitionPublisher.PublishAsync(draft, cancellationToken); @@ -107,7 +108,8 @@ internal class Post( } var mappedDefinition = await linker.MapAsync(draft, cancellationToken); - var response = new Response(mappedDefinition, false, result?.ConsumingWorkflows?.Count() ?? 0); + var affectedWorkflows = result?.AffectedWorkflows?.WorkflowDefinitions ?? []; + var response = new Response(mappedDefinition, false, affectedWorkflows.Count); await HttpContext.Response.WriteAsJsonAsync(response, serializerOptions, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs index 2f9cd2ca6..26e0ad9b3 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs @@ -46,7 +46,7 @@ internal class Publish(IWorkflowDefinitionStore store, IWorkflowDefinitionPublis var isPublished = definition.IsPublished; var result = !isPublished ? await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken) : null; var mappedDefinition = await linker.MapAsync(definition, cancellationToken); - var response = new Response(mappedDefinition, isPublished, result?.ConsumingWorkflows?.Count() ?? 0); + var response = new Response(mappedDefinition, isPublished, result?.AffectedWorkflows.WorkflowDefinitions.Count ?? 0); await SendOkAsync(response, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs index 4efbeb133..c10119f22 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs @@ -10,7 +10,7 @@ using Microsoft.AspNetCore.Authorization; namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.UpdateReferences; [PublicAPI] -internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinitionPublisher workflowDefinitionPublisher, IAuthorizationService authorizationService) +internal class UpdateReferences(IWorkflowReferenceUpdater workflowReferenceUpdater, IWorkflowDefinitionStore store, IAuthorizationService authorizationService) : ElsaEndpoint { public override void Configure() @@ -43,7 +43,8 @@ internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinit return; } - var affectedWorkflows = await workflowDefinitionPublisher.UpdateReferencesInConsumingWorkflows(definition, cancellationToken); + var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken); + var affectedWorkflows = result.UpdatedWorkflows; var response = new Response(affectedWorkflows.Select(w => w.Name ?? w.DefinitionId)); await SendOkAsync(response, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index abddaf098..cc1934b36 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -52,7 +52,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable var serviceProvider = context.ServiceProvider; var cancellationToken = context.CancellationToken; - // Find the workflow definition and not the graph; the graph must be computed at runtime, since one NodeIds will vary across graphs. + // Find the workflow definition and not the graph; the graph must be computed at runtime, since NodeIds will vary across graphs. var workflowDefinition = await GetWorkflowDefinitionAsync(serviceProvider, cancellationToken); if (workflowDefinition == null) diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs index 487074525..402243f73 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs @@ -68,12 +68,4 @@ public interface IWorkflowDefinitionPublisher /// The cancellation token. /// The saved workflow definition. Task SaveDraftAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); - - /// - /// Updates all referencing workflow definitions to use the version of the specified workflow definition. - /// - /// The workflow definition to update references for. - /// The cancellation token. - /// The updated workflow definitions. - Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs new file mode 100644 index 000000000..e70f7a855 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceUpdater.cs @@ -0,0 +1,18 @@ +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; + +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Updates references to the specified workflow of all workflows that reference it. +/// +public interface IWorkflowReferenceUpdater +{ + /// + /// Updates references to the specified workflow of all workflows that reference it. + /// + /// The workflow definition that is being referenced. All workflows that reference this definition will be updated to use this newest version. + /// The cancellation token. + /// The result of the operation. + Task UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 3fdc5adf8..47de8144c 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -213,6 +213,7 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddSingleton() .AddSingleton() @@ -233,6 +234,7 @@ public class WorkflowManagementFeature : FeatureBase Services .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() ; Services.Configure(options => diff --git a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs new file mode 100644 index 000000000..5f8e3811c --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs @@ -0,0 +1,27 @@ +using Elsa.Common.Models; +using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Notifications; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Handlers; + +/// +/// Updates consuming workflows when a workflow definition is published. +/// +public class UpdateConsumingWorkflows(IWorkflowReferenceUpdater workflowReferenceUpdater) : INotificationHandler +{ + /// + public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) + { + var definition = notification.WorkflowDefinition; + var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken); + notification.AffectedWorkflows.WorkflowDefinitions.AddRange(result.UpdatedWorkflows); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs b/src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs new file mode 100644 index 000000000..da5a0d71b --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/AffectedWorkflows.cs @@ -0,0 +1,5 @@ +using Elsa.Workflows.Management.Entities; + +namespace Elsa.Workflows.Management.Models; + +public record AffectedWorkflows(ICollection WorkflowDefinitions); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs b/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs index 87848df5a..31a9f37f3 100644 --- a/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs +++ b/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs @@ -1,8 +1,6 @@ -using Elsa.Workflows.Management.Entities; - namespace Elsa.Workflows.Management.Models; /// /// Represents the result of publishing a workflow definition. /// -public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection ValidationErrors, IEnumerable? ConsumingWorkflows); \ No newline at end of file +public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection ValidationErrors, AffectedWorkflows AffectedWorkflows); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs b/src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs new file mode 100644 index 000000000..eebddc7aa --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/UpdateWorkflowReferencesResult.cs @@ -0,0 +1,18 @@ +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Models; + +/// +/// Represents the result of updating workflow references. +/// +public class UpdateWorkflowReferencesResult(IEnumerable updatedWorkflows) +{ + /// + /// Gets a collection of workflow graphs that have been updated. + /// + /// + /// A read-only collection of instances representing the updated workflows. + /// + public IReadOnlyCollection UpdatedWorkflows { get; } = updatedWorkflows.ToList(); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs b/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs index 81db3b296..62c84c96c 100644 --- a/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs +++ b/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionPublished.cs @@ -1,5 +1,6 @@ using Elsa.Mediator.Contracts; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; using JetBrains.Annotations; namespace Elsa.Workflows.Management.Notifications; @@ -9,4 +10,4 @@ namespace Elsa.Workflows.Management.Notifications; /// /// The workflow definition. [PublicAPI] -public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition) : INotification; \ No newline at end of file +public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition, AffectedWorkflows AffectedWorkflows) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs index 173c3b29d..7f5bf7bfd 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -1,11 +1,9 @@ using Elsa.Common.Contracts; using Elsa.Common.Entities; using Elsa.Common.Models; -using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; using Elsa.Workflows.Contracts; -using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; @@ -23,7 +21,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher private readonly IWorkflowDefinitionStore _workflowDefinitionStore; private readonly INotificationSender _notificationSender; private readonly IIdentityGenerator _identityGenerator; - private readonly IActivityVisitor _activityVisitor; private readonly IActivitySerializer _activitySerializer; private readonly IRequestSender _requestSender; private readonly ISystemClock _systemClock; @@ -36,7 +33,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher IWorkflowDefinitionStore workflowDefinitionStore, INotificationSender notificationSender, IIdentityGenerator identityGenerator, - IActivityVisitor activityVisitor, IActivitySerializer activitySerializer, IRequestSender requestSender, ISystemClock systemClock) @@ -45,7 +41,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher _workflowDefinitionStore = workflowDefinitionStore; _notificationSender = notificationSender; _identityGenerator = identityGenerator; - _activityVisitor = activityVisitor; _activitySerializer = activitySerializer; _requestSender = requestSender; _systemClock = systemClock; @@ -121,16 +116,9 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher definition = Initialize(definition); await _workflowDefinitionStore.SaveAsync(definition, cancellationToken); - await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition), cancellationToken); - - var consumingWorkflows = new List(); - - if (definition.Options is { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) - { - consumingWorkflows.AddRange(await UpdateReferencesInConsumingWorkflows(definition, cancellationToken)); - } - - return new PublishWorkflowDefinitionResult(true, validationErrors, consumingWorkflows); + var affectedWorkflows = new AffectedWorkflows(new List()); + await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition, affectedWorkflows), cancellationToken); + return new PublishWorkflowDefinitionResult(true, validationErrors, affectedWorkflows); } /// @@ -212,57 +200,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher return draft; } - /// - public async Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default) - { - var updatedWorkflowDefinitions = new List(); - - var workflowDefinitions = (await _workflowDefinitionStore.FindManyAsync(new WorkflowDefinitionFilter - { - VersionOptions = VersionOptions.LatestOrPublished - }, cancellationToken)).ToList(); - - // Remove the dependency from the list of workflow definitions to consider. - workflowDefinitions = workflowDefinitions.Where(x => x.DefinitionId != dependency.DefinitionId).ToList(); - - foreach (var definition in workflowDefinitions) - { - var root = _activitySerializer.Deserialize(definition.StringData!); - var graph = await _activityVisitor.VisitAsync(root, cancellationToken); - var flattenedList = graph.Flatten().ToList(); - var definitionId = dependency.DefinitionId; - var version = dependency.Version; - var nodes = flattenedList - .Where(x => x.Activity is WorkflowDefinitionActivity workflowDefinitionActivity && workflowDefinitionActivity.WorkflowDefinitionId == definitionId) - .ToList(); - - foreach (var node in nodes.Where(activity => activity.Activity.Version < version)) - { - var activity = (WorkflowDefinitionActivity)node.Activity; - activity.Version = version; - activity.WorkflowDefinitionVersionId = dependency.Id; - - if (!updatedWorkflowDefinitions.Contains(definition)) - updatedWorkflowDefinitions.Add(definition); - } - - if (updatedWorkflowDefinitions.Contains(definition)) - { - var serializedData = _activitySerializer.Serialize(root); - definition.StringData = serializedData; - } - } - - if (updatedWorkflowDefinitions.Any()) - { - await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdating(updatedWorkflowDefinitions), cancellationToken); - await _workflowDefinitionStore.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken); - await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdated(updatedWorkflowDefinitions), cancellationToken); - } - - return updatedWorkflowDefinitions; - } - private WorkflowDefinition Initialize(WorkflowDefinition definition) { if (definition.Id == null!) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs new file mode 100644 index 000000000..a6fdfab67 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -0,0 +1,131 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Services; + +/// +public class WorkflowReferenceUpdater( + IWorkflowDefinitionPublisher publisher, + IWorkflowDefinitionService workflowDefinitionService, + IWorkflowDefinitionStore workflowDefinitionStore, + IApiSerializer serializer) : IWorkflowReferenceUpdater +{ + /// + public async Task UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default) + { + // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. + if (referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) + return new UpdateWorkflowReferencesResult([]); + + // Find all workflow graphs that contain the updated workflow definition. + var consumingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken); + + // Update consuming workflows. + var updatedWorkflows = new List(); + foreach (var workflowGraph in consumingWorkflowGraphs) + { + var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, cancellationToken); + + if (newDefinition != null) + updatedWorkflows.Add(newDefinition); + } + + return new UpdateWorkflowReferencesResult(updatedWorkflows); + } + + private async Task UpdateConsumingWorkflowAsync(WorkflowGraph workflowGraph, WorkflowDefinition definition, CancellationToken cancellationToken) + { + // Create a new version of the published workflow definition or get the existing draft. + var consumerDefinitionId = workflowGraph.Workflow.Identity.DefinitionId; + var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished; + var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + // This is null in case the definition no longer exists in the store. + if (newVersion == null) + return null; + + // Materialize the draft to find all workflow definition activities that use the updated workflow definition. + var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); + var outdatedWorkflowDefinitionActivities = FindOutdatedWorkflowDefinitionActivities(newWorkflowGraph, definition).ToList(); + + // Skip if the new version of the published workflow definition is not used in the workflow or if the activity is already up to date. + if (outdatedWorkflowDefinitionActivities.Count == 0) + return null; + + // Update the consuming workflow graph to use the new version of the published workflow definition. + foreach (var workflowDefinitionActivity in outdatedWorkflowDefinitionActivities) + { + workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id; + workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id; + workflowDefinitionActivity.Version = definition.Version; + } + + // Update the new version of the published workflow definition. + if (newWorkflowGraph.Root.Activity is Workflow newWorkflow) + newVersion.StringData = serializer.Serialize(newWorkflow.Root); + + // If the draft is new, publish it. + if (originalVersionIsPublished) + await publisher.PublishAsync(newVersion, cancellationToken); + else + await publisher.SaveDraftAsync(newVersion, cancellationToken); + + return newVersion; + } + + private IEnumerable FindOutdatedWorkflowDefinitionActivities(WorkflowGraph workflowGraph, WorkflowDefinition updatedDefinition) + { + return FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root) + .Where(x => x.WorkflowDefinitionId == updatedDefinition.DefinitionId && x.WorkflowDefinitionVersionId != updatedDefinition.Id); + } + + private async Task> FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(string updatedWorkflowDefinitionId, CancellationToken cancellationToken) + { + var allWorkflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList(); + var filteredWorkflowGraphs = allWorkflowGraphs + .Where(x => FindWorkflowActivityDefinitionActivityNodes(x.Root).Any(y => y.WorkflowDefinitionId == updatedWorkflowDefinitionId)) + .ToList(); + + return filteredWorkflowGraphs; + } + + private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) + { + foreach (var child in parent.Children) + { + if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity) + yield return workflowDefinitionActivity; + + foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child)) + yield return grandChild; + } + } + + private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) + { + var workflowDefinitionFilter = new WorkflowDefinitionFilter + { + VersionOptions = VersionOptions.LatestOrPublished, + IsReadonly = false + }; + var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); + var workflowGraphs = new List(); + + foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) + { + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken); + + if (workflowGraph != null) + workflowGraphs.Add(workflowGraph); + } + + return workflowGraphs; + } +} \ No newline at end of file From e5e7211504bfbea2fd5cb083787b6143c6ea5498 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 21:39:01 +0100 Subject: [PATCH 16/17] Add workflow definition cache manager to test setup The IWorkflowDefinitionCacheManager has been added to the test class. This ensures proper cache management during workflow definition reload tests. Additionally, a cleanup step to delete the workflow definition and its versions has been included. --- .../WorkflowDefinitionReload/ReloadWorkflowTests.cs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs index 9a0fea9ca..98db0c4eb 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -22,6 +22,7 @@ public class ReloadWorkflowTests : AppComponentTest private readonly TestWorkflowProvider _testWorkflowProvider; private readonly IWorkflowDefinitionService _workflowDefinitionService; private readonly IActivityRegistry _activityRegistry; + private readonly IWorkflowDefinitionCacheManager _workflowDefinitionCacheManager; public ReloadWorkflowTests(App app) : base(app) { @@ -32,6 +33,7 @@ public class ReloadWorkflowTests : AppComponentTest _activityRegistry = Scope.ServiceProvider.GetRequiredService(); var workflowProviders = Scope.ServiceProvider.GetRequiredService>(); _testWorkflowProvider = (TestWorkflowProvider)workflowProviders.First(x => x is TestWorkflowProvider); + _workflowDefinitionCacheManager = Scope.ServiceProvider.GetRequiredService(); } [Fact] @@ -40,7 +42,7 @@ 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")); - await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); @@ -97,6 +99,9 @@ public class ReloadWorkflowTests : AppComponentTest // Assert that the activity registry contains a new activity descriptor representing the new workflow version. var activityV2 = _activityRegistry.Find(activityTypeName)!; Assert.Equal(2, activityV2.Version); + + // Cleanup: Delete the workflow definition and its versions. + await _workflowDefinitionManager.DeleteByDefinitionIdAsync(definitionId, CancellationToken.None); } private async Task BuildWorkflowAsync(string definitionId, string definitionVersionId, int version) From 111fa10aab5415dc84f91acdcd006e9aa1704296 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 21:39:27 +0100 Subject: [PATCH 17/17] Remove unused dependency from ReloadWorkflowTests Eliminated the IWorkflowDefinitionCacheManager dependency from the ReloadWorkflowTests constructor and fields. This cleanup helps streamline the code and maintainability by removing an unnecessary service. --- .../Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs | 2 -- 1 file changed, 2 deletions(-) diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs index 98db0c4eb..71804fa7c 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -22,7 +22,6 @@ public class ReloadWorkflowTests : AppComponentTest private readonly TestWorkflowProvider _testWorkflowProvider; private readonly IWorkflowDefinitionService _workflowDefinitionService; private readonly IActivityRegistry _activityRegistry; - private readonly IWorkflowDefinitionCacheManager _workflowDefinitionCacheManager; public ReloadWorkflowTests(App app) : base(app) { @@ -33,7 +32,6 @@ public class ReloadWorkflowTests : AppComponentTest _activityRegistry = Scope.ServiceProvider.GetRequiredService(); var workflowProviders = Scope.ServiceProvider.GetRequiredService>(); _testWorkflowProvider = (TestWorkflowProvider)workflowProviders.First(x => x is TestWorkflowProvider); - _workflowDefinitionCacheManager = Scope.ServiceProvider.GetRequiredService(); } [Fact]