From 8fe0147abd7ea92b2e851337a25dfa59d442f826 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 8 Nov 2024 18:00:21 +0100 Subject: [PATCH] Remove UpdateReferencesInConsumingWorkflows method Eliminated the UpdateReferencesInConsumingWorkflows method from IWorkflowDefinitionPublisher and its implementation from WorkflowDefinitionPublisher. The responsibility of updating workflow references is now transitioned to IWorkflowReferenceUpdater used in the `UpdateReferences` endpoint. --- .../UpdateReferences/Endpoint.cs | 5 +- .../Contracts/IWorkflowDefinitionPublisher.cs | 8 --- .../Services/WorkflowDefinitionPublisher.cs | 56 ------------------- 3 files changed, 3 insertions(+), 66 deletions(-) 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 4efbeb133..c10119f22 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/Contracts/IWorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs index 487074525..402243f73 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs @@ -68,12 +68,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/Services/WorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs index 62bfc8e94..7f5bf7bfd 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -1,11 +1,9 @@ using Elsa.Common.Contracts; using Elsa.Common.Entities; 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; @@ -23,7 +21,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; @@ -36,7 +33,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher IWorkflowDefinitionStore workflowDefinitionStore, INotificationSender notificationSender, IIdentityGenerator identityGenerator, - IActivityVisitor activityVisitor, IActivitySerializer activitySerializer, IRequestSender requestSender, ISystemClock systemClock) @@ -45,7 +41,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher _workflowDefinitionStore = workflowDefinitionStore; _notificationSender = notificationSender; _identityGenerator = identityGenerator; - _activityVisitor = activityVisitor; _activitySerializer = activitySerializer; _requestSender = requestSender; _systemClock = systemClock; @@ -205,57 +200,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!)