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.
This commit is contained in:
Sipke Schoorstra 2024-11-08 17:16:21 +01:00
parent c7fa4468be
commit cc86ee9f58
5 changed files with 170 additions and 94 deletions

View file

@ -0,0 +1,18 @@
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Models;
namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Updates references to the specified workflow of all workflows that reference it.
/// </summary>
public interface IWorkflowReferenceUpdater
{
/// <summary>
/// Updates references to the specified workflow of all workflows that reference it.
/// </summary>
/// <param name="definition">The workflow definition to update references for.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The result of the operation.</returns>
Task<UpdateWorkflowReferencesResult> UpdateWorkflowReferencesAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default);
}

View file

@ -213,6 +213,7 @@ public class WorkflowManagementFeature : FeatureBase
.AddScoped<IWorkflowDefinitionImporter, WorkflowDefinitionImporter>()
.AddScoped<IWorkflowDefinitionManager, WorkflowDefinitionManager>()
.AddScoped<IWorkflowInstanceManager, WorkflowInstanceManager>()
.AddScoped<IWorkflowReferenceUpdater, WorkflowReferenceUpdater>()
.AddScoped<IActivityRegistryPopulator, ActivityRegistryPopulator>()
.AddSingleton<IExpressionDescriptorRegistry, ExpressionDescriptorRegistry>()
.AddSingleton<IExpressionDescriptorProvider, DefaultExpressionDescriptorProvider>()

View file

@ -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;
/// <summary>
/// Updates consuming workflows when a workflow definition is published.
/// </summary>
public class UpdateConsumingWorkflows(
IWorkflowDefinitionPublisher publisher,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowDefinitionStore workflowDefinitionStore,
IApiSerializer serializer) : INotificationHandler<WorkflowDefinitionPublished>
public class UpdateConsumingWorkflows(IWorkflowReferenceUpdater workflowReferenceUpdater) : INotificationHandler<WorkflowDefinitionPublished>
{
/// <inheritdoc />
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<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;
var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken);
notification.AffectedWorkflows.WorkflowDefinitions.AddRange(result.UpdatedWorkflows);
}
}

View file

@ -0,0 +1,18 @@
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Models;
/// <summary>
/// Represents the result of updating workflow references.
/// </summary>
public class UpdateWorkflowReferencesResult(IEnumerable<WorkflowDefinition> updatedWorkflows)
{
/// <summary>
/// Gets a collection of workflow graphs that have been updated.
/// </summary>
/// <value>
/// A read-only collection of <see cref="WorkflowGraph"/> instances representing the updated workflows.
/// </value>
public IReadOnlyCollection<WorkflowDefinition> UpdatedWorkflows { get; } = updatedWorkflows.ToList();
}

View file

@ -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;
/// <inheritdoc />
public class WorkflowReferenceUpdater(
IWorkflowDefinitionPublisher publisher,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowDefinitionStore workflowDefinitionStore,
IApiSerializer serializer) : IWorkflowReferenceUpdater
{
/// <inheritdoc />
public async Task<UpdateWorkflowReferencesResult> 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<WorkflowDefinition>();
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<WorkflowDefinition?> 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<WorkflowDefinitionActivity> FindOutdatedWorkflowDefinitionActivities(WorkflowGraph workflowGraph, WorkflowDefinition updatedDefinition)
{
return FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root)
.Where(x => x.WorkflowDefinitionId == updatedDefinition.DefinitionId && x.WorkflowDefinitionVersionId != updatedDefinition.Id);
}
private async Task<IEnumerable<WorkflowGraph>> 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<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;
}
}