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