From cc86ee9f583c2c0d84d3b4bad1906b92fa9357fc Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 17:16:21 +0100 Subject: [PATCH] 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