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