diff --git a/Directory.Packages.props b/Directory.Packages.props index 370b18e6c..7bb00b30b 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -146,7 +146,6 @@ - @@ -182,8 +181,8 @@ - - - + + + \ No newline at end of file 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 f6a5f3a79..3893a32d1 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 cb73bf98a..0746c037d 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs @@ -87,7 +87,7 @@ internal class Post( PublishWorkflowDefinitionResult? result = null; - if (request.Publish.GetValueOrDefault(false)) + if (request.Publish == true) { result = await workflowDefinitionPublisher.PublishAsync(draft, cancellationToken); @@ -106,7 +106,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 a8e652df5..0a28f964d 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.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs index f6af58e8a..fb81b71b7 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/UpdateReferences/Endpoint.cs @@ -10,7 +10,7 @@ using Microsoft.AspNetCore.Authorization; namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.UpdateReferences; [PublicAPI] -internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinitionPublisher workflowDefinitionPublisher, IAuthorizationService authorizationService) +internal class UpdateReferences(IWorkflowReferenceUpdater workflowReferenceUpdater, IWorkflowDefinitionStore store, IAuthorizationService authorizationService) : ElsaEndpoint { public override void Configure() @@ -43,7 +43,8 @@ internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinit return; } - var affectedWorkflows = await workflowDefinitionPublisher.UpdateReferencesInConsumingWorkflows(definition, cancellationToken); + var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken); + var affectedWorkflows = result.UpdatedWorkflows; var response = new Response(affectedWorkflows.Select(w => w.Name ?? w.DefinitionId)); await SendOkAsync(response, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index b8fb52f3d..e5f69bb47 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -51,7 +51,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/IWorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs index 511ab0003..40c767e3e 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs @@ -67,12 +67,4 @@ public interface IWorkflowDefinitionPublisher /// The cancellation token. /// The saved workflow definition. Task SaveDraftAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); - - /// - /// Updates all referencing workflow definitions to use the version of the specified workflow definition. - /// - /// The workflow definition to update references for. - /// The cancellation token. - /// The updated workflow definitions. - Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default); } \ No newline at end of file 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..e70f7a855 --- /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 that is being referenced. All workflows that reference this definition will be updated to use this newest version. + /// The cancellation token. + /// The result of the operation. + Task UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, 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 d36463f52..62443f526 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -215,6 +215,7 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddSingleton() .AddSingleton() @@ -235,6 +236,7 @@ public class WorkflowManagementFeature : FeatureBase Services .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() ; Services.Configure(options => diff --git a/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs new file mode 100644 index 000000000..5f8e3811c --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Handlers/UpdateConsumingWorkflows.cs @@ -0,0 +1,27 @@ +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; + +namespace Elsa.Workflows.Management.Handlers; + +/// +/// Updates consuming workflows when a workflow definition is published. +/// +public class UpdateConsumingWorkflows(IWorkflowReferenceUpdater workflowReferenceUpdater) : INotificationHandler +{ + /// + public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) + { + var definition = notification.WorkflowDefinition; + 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/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/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/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 f9630d3e7..75a20e2d4 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -1,7 +1,6 @@ using Elsa.Common; using Elsa.Common.Entities; using Elsa.Common.Models; -using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; @@ -21,7 +20,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher private readonly IWorkflowDefinitionStore _workflowDefinitionStore; private readonly INotificationSender _notificationSender; private readonly IIdentityGenerator _identityGenerator; - private readonly IActivityVisitor _activityVisitor; private readonly IActivitySerializer _activitySerializer; private readonly IRequestSender _requestSender; private readonly ISystemClock _systemClock; @@ -34,7 +32,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher IWorkflowDefinitionStore workflowDefinitionStore, INotificationSender notificationSender, IIdentityGenerator identityGenerator, - IActivityVisitor activityVisitor, IActivitySerializer activitySerializer, IRequestSender requestSender, ISystemClock systemClock) @@ -43,7 +40,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher _workflowDefinitionStore = workflowDefinitionStore; _notificationSender = notificationSender; _identityGenerator = identityGenerator; - _activityVisitor = activityVisitor; _activitySerializer = activitySerializer; _requestSender = requestSender; _systemClock = systemClock; @@ -119,16 +115,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); } /// @@ -210,57 +199,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher return draft; } - /// - public async Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default) - { - var updatedWorkflowDefinitions = new List(); - - var workflowDefinitions = (await _workflowDefinitionStore.FindManyAsync(new WorkflowDefinitionFilter - { - VersionOptions = VersionOptions.LatestOrPublished - }, cancellationToken)).ToList(); - - // Remove the dependency from the list of workflow definitions to consider. - workflowDefinitions = workflowDefinitions.Where(x => x.DefinitionId != dependency.DefinitionId).ToList(); - - foreach (var definition in workflowDefinitions) - { - var root = _activitySerializer.Deserialize(definition.StringData!); - var graph = await _activityVisitor.VisitAsync(root, cancellationToken); - var flattenedList = graph.Flatten().ToList(); - var definitionId = dependency.DefinitionId; - var version = dependency.Version; - var nodes = flattenedList - .Where(x => x.Activity is WorkflowDefinitionActivity workflowDefinitionActivity && workflowDefinitionActivity.WorkflowDefinitionId == definitionId) - .ToList(); - - foreach (var node in nodes.Where(activity => activity.Activity.Version < version)) - { - var activity = (WorkflowDefinitionActivity)node.Activity; - activity.Version = version; - activity.WorkflowDefinitionVersionId = dependency.Id; - - if (!updatedWorkflowDefinitions.Contains(definition)) - updatedWorkflowDefinitions.Add(definition); - } - - if (updatedWorkflowDefinitions.Contains(definition)) - { - var serializedData = _activitySerializer.Serialize(root); - definition.StringData = serializedData; - } - } - - if (updatedWorkflowDefinitions.Any()) - { - await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdating(updatedWorkflowDefinitions), cancellationToken); - await _workflowDefinitionStore.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken); - await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdated(updatedWorkflowDefinitions), cancellationToken); - } - - return updatedWorkflowDefinitions; - } - private WorkflowDefinition Initialize(WorkflowDefinition definition) { if (definition.Id == null!) 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..a6fdfab67 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -0,0 +1,131 @@ +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 referencedDefinition, CancellationToken cancellationToken = default) + { + // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. + if (referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) + return new UpdateWorkflowReferencesResult([]); + + // Find all workflow graphs that contain the updated workflow definition. + var consumingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken); + + // Update consuming workflows. + var updatedWorkflows = new List(); + foreach (var workflowGraph in consumingWorkflowGraphs) + { + var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, 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); + + // This is null in case the definition no longer exists in the store. + if (newVersion == null) + return null; + + // Materialize the draft 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; + workflowDefinitionActivity.Version = definition.Version; + } + + // 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, + IsReadonly = false + }; + 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/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs index 6a8b7157d..548b86f5f 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -35,7 +35,7 @@ public class ReloadWorkflowTests : AppComponentTest var client = WorkflowServer.CreateHttpWorkflowClient(); await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); - await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); @@ -92,6 +92,9 @@ public class ReloadWorkflowTests : AppComponentTest // Assert that the activity registry contains a new activity descriptor representing the new workflow version. var activityV2 = _activityRegistry.Find(activityTypeName)!; Assert.Equal(2, activityV2.Version); + + // Cleanup: Delete the workflow definition and its versions. + await _workflowDefinitionManager.DeleteByDefinitionIdAsync(definitionId, CancellationToken.None); } private async Task BuildWorkflowAsync(string definitionId, string definitionVersionId, int version)