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.
This commit is contained in:
Sipke Schoorstra 2024-11-08 11:45:25 +01:00
parent 311f59894e
commit 8c7a823374
13 changed files with 120 additions and 157 deletions

View file

@ -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);

View file

@ -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);
}
}

View file

@ -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);
}
}

View file

@ -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)

View file

@ -1,17 +0,0 @@
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Finds child workflows for a given workflow graph.
/// </summary>
public interface IChildWorkflowFinder
{
/// <summary>
/// Finds child workflows for a given workflow graph.
/// </summary>
/// <param name="workflowGraph">The workflow graph.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>A collection of child workflows.</returns>
Task<IEnumerable<WorkflowGraph>> FindChildWorkflowsAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default);
}

View file

@ -1,16 +0,0 @@
using Elsa.Workflows.Management.Models;
namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Defines a visitor that can traverse and process workflow definitions.
/// </summary>
public interface IWorkflowGraphNetworkBuilder
{
/// <summary>
/// Builds a network of workflow graphs and their consumers.
/// </summary>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>A network of workflow graphs and their consumers.</returns>
Task<WorkflowGraphNetwork> BuildAsync(CancellationToken cancellationToken = default);
}

View file

@ -221,7 +221,6 @@ public class WorkflowManagementFeature : FeatureBase
.AddScoped<IWorkflowMaterializer, ClrWorkflowMaterializer>()
.AddScoped<IWorkflowMaterializer, JsonWorkflowMaterializer>()
.AddScoped<IActivityResolver, WorkflowDefinitionActivityResolver>()
.AddScoped<IWorkflowGraphNetworkBuilder, WorkflowGraphNetworkBuilder>()
.AddScoped<WorkflowDefinitionMapper>()
.AddSingleton<VariableDefinitionMapper>()
.AddSingleton<WorkflowStateMapper>()

View file

@ -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<WorkflowDefinitionPublished>
/// <summary>
/// Updates consuming workflows when a workflow definition is published.
/// </summary>
public class UpdateConsumingWorkflows(
IWorkflowDefinitionPublisher publisher,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowDefinitionStore workflowDefinitionStore,
IApiSerializer serializer) : INotificationHandler<WorkflowDefinitionPublished>
{
/// <inheritdoc />
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<WorkflowDefinitionActivity> 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<IEnumerable<WorkflowGraph>> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken)
{
var workflowDefinitionFilter = new WorkflowDefinitionFilter
{
VersionOptions = VersionOptions.LatestOrPublished
};
var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken);
var workflowGraphs = new List<WorkflowGraph>();
foreach (var workflowDefinitionSummary in workflowDefinitionSummaries)
{
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken);
if (workflowGraph != null)
workflowGraphs.Add(workflowGraph);
}
return workflowGraphs;
}
}

View file

@ -0,0 +1,5 @@
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Models;
public record AffectedWorkflows(ICollection<WorkflowDefinition> WorkflowDefinitions);

View file

@ -1,8 +1,6 @@
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Models;
/// <summary>
/// Represents the result of publishing a workflow definition.
/// </summary>
public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection<WorkflowValidationError> ValidationErrors, IEnumerable<WorkflowDefinition>? ConsumingWorkflows);
public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection<WorkflowValidationError> ValidationErrors, AffectedWorkflows AffectedWorkflows);

View file

@ -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;
/// </summary>
/// <param name="WorkflowDefinition">The workflow definition.</param>
[PublicAPI]
public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition) : INotification;
public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition, AffectedWorkflows AffectedWorkflows) : INotification;

View file

@ -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<WorkflowDefinition>();
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<WorkflowDefinition>());
await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition, affectedWorkflows), cancellationToken);
return new PublishWorkflowDefinitionResult(true, validationErrors, affectedWorkflows);
}
/// <inheritdoc />

View file

@ -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;
/// <inheritdoc />
public class WorkflowGraphNetworkBuilder(IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowDefinitionService workflowDefinitionService) : IWorkflowGraphNetworkBuilder
{
/// <inheritdoc />
public async Task<WorkflowGraphNetwork> BuildAsync(CancellationToken cancellationToken = default)
{
var workflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList();
var nodes = new HashSet<WorkflowGraphNode>();
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<IEnumerable<WorkflowGraph>> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken)
{
var workflowDefinitionFilter = new WorkflowDefinitionFilter
{
VersionOptions = VersionOptions.LatestOrPublished
};
var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken);
var workflowGraphs = new List<WorkflowGraph>();
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<WorkflowGraph> WorkflowGraphs { get; set; } = new List<WorkflowGraph>();
public HashSet<WorkflowGraphNode> VisitedNodes { get; set; } = new();
}
}