From fb3b5339f447fbb6838d79574ea9d9cc446d1a94 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Wed, 29 May 2024 15:39:15 +0200 Subject: [PATCH] Made Updating workflow references part of Publish logic --- .../Contracts/IWorkflowDefinitionsApi.cs | 4 +- .../BulkPublishWorkflowDefinitionsResponse.cs | 3 +- .../SaveWorkflowDefinitionResponse.cs | 8 +++ .../BulkPublish/Endpoint.cs | 10 ++- .../WorkflowDefinitions/BulkPublish/Models.cs | 3 +- .../WorkflowDefinitions/Post/Endpoint.cs | 9 ++- .../WorkflowDefinitions/Post/Models.cs | 5 ++ .../WorkflowDefinitions/Publish/Endpoint.cs | 7 +- .../WorkflowDefinitions/Publish/Models.cs | 6 +- .../UpdateReferences/Endpoint.cs | 4 +- .../Contracts/IWorkflowDefinitionManager.cs | 8 --- .../Contracts/IWorkflowDefinitionPublisher.cs | 8 +++ .../Models/PublishWorkflowDefinitionResult.cs | 4 +- .../Services/WorkflowDefinitionManager.cs | 57 +--------------- .../Services/WorkflowDefinitionPublisher.cs | 66 ++++++++++++++++++- .../DependencyWorkflowsPublishing/Tests.cs | 8 +-- 16 files changed, 123 insertions(+), 87 deletions(-) create mode 100644 src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/SaveWorkflowDefinitionResponse.cs create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Models.cs diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Contracts/IWorkflowDefinitionsApi.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Contracts/IWorkflowDefinitionsApi.cs index b165cec8f..f9e8d9e20 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Contracts/IWorkflowDefinitionsApi.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Contracts/IWorkflowDefinitionsApi.cs @@ -73,7 +73,7 @@ public interface IWorkflowDefinitionsApi /// The request containing the workflow definition to save. /// The cancellation token. [Post("/workflow-definitions")] - Task SaveAsync(SaveWorkflowDefinitionRequest request, CancellationToken cancellationToken = default); + Task SaveAsync(SaveWorkflowDefinitionRequest request, CancellationToken cancellationToken = default); /// /// Deletes a workflow definition. @@ -99,7 +99,7 @@ public interface IWorkflowDefinitionsApi /// The cancellation token. [Post("/workflow-definitions/{definitionId}/publish")] [Headers(MediaTypeNames.Application.Json)] - Task PublishAsync(string definitionId, PublishWorkflowDefinitionRequest request, CancellationToken cancellationToken = default); + Task PublishAsync(string definitionId, PublishWorkflowDefinitionRequest request, CancellationToken cancellationToken = default); /// /// Retracts a workflow definition. diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/BulkPublishWorkflowDefinitionsResponse.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/BulkPublishWorkflowDefinitionsResponse.cs index ed77df362..732e11682 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/BulkPublishWorkflowDefinitionsResponse.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/BulkPublishWorkflowDefinitionsResponse.cs @@ -6,4 +6,5 @@ namespace Elsa.Api.Client.Resources.WorkflowDefinitions.Responses; public record BulkPublishWorkflowDefinitionsResponse( ICollection Published, ICollection AlreadyPublished, - ICollection NotFound); \ No newline at end of file + ICollection NotFound, + ICollection ConsumingUpdated); \ No newline at end of file diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/SaveWorkflowDefinitionResponse.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/SaveWorkflowDefinitionResponse.cs new file mode 100644 index 000000000..52fef142c --- /dev/null +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Responses/SaveWorkflowDefinitionResponse.cs @@ -0,0 +1,8 @@ +using Elsa.Api.Client.Resources.WorkflowDefinitions.Models; + +namespace Elsa.Api.Client.Resources.WorkflowDefinitions.Responses; + +/// +/// Represents the response for saving a workflow definition. +/// +public record SaveWorkflowDefinitionResponse(WorkflowDefinition WorkflowDefinition, int ConsumingWorkflowCount); \ 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 9e12132a6..99115adad 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Endpoint.cs @@ -33,6 +33,7 @@ internal class BulkPublish(IWorkflowDefinitionStore store, IWorkflowDefinitionPu var notFound = new List(); var alreadyPublished = new List(); var skipped = new List(); + var consumingUpdated = new List(); var definitions = (await store.FindManyAsync(new WorkflowDefinitionFilter { @@ -62,10 +63,15 @@ internal class BulkPublish(IWorkflowDefinitionStore store, IWorkflowDefinitionPu continue; } - await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken); + var result = await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken); published.Add(definitionId); + + if (result.ConsumingWorkflows?.Any() == true) + { + consumingUpdated.AddRange(result.ConsumingWorkflows.Select(x => x.DefinitionId)); + } } - return new Response(published, alreadyPublished, notFound, skipped); + return new Response(published, alreadyPublished, notFound, skipped, consumingUpdated); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Models.cs index 0cb2e8f01..b2ac516f6 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Models.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkPublish/Models.cs @@ -5,10 +5,11 @@ internal class Request public ICollection DefinitionIds { get; set; } = default!; } -internal class Response(ICollection published, ICollection alreadyPublished, ICollection notFound, ICollection skipped) +internal class Response(ICollection published, ICollection alreadyPublished, ICollection notFound, ICollection skipped, ICollection consumingUpdated) { public ICollection Published { get; } = published; public ICollection AlreadyPublished { get; } = alreadyPublished; public ICollection NotFound { get; } = notFound; public ICollection Skipped { get; } = skipped; + public ICollection ConsumingUpdated { get; } = consumingUpdated; } \ No newline at end of file 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 5c26bba3c..75d4ad93b 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs @@ -92,9 +92,11 @@ internal class Post( draft.Outcomes = outcomes; draft.Options = model.Options ?? new WorkflowOptions(); + PublishWorkflowDefinitionResult? result = null; + if (request.Publish.GetValueOrDefault(false)) { - var result = await workflowDefinitionPublisher.PublishAsync(draft, cancellationToken); + result = await workflowDefinitionPublisher.PublishAsync(draft, cancellationToken); if (!result.Succeeded) { @@ -110,10 +112,11 @@ internal class Post( await workflowDefinitionPublisher.SaveDraftAsync(draft, cancellationToken); } - var response = await linker.MapAsync(draft, cancellationToken); + var mappedDefinition = await linker.MapAsync(draft, cancellationToken); + var response = new Response(mappedDefinition, result?.ConsumingWorkflows?.Count() ?? 0); if (isNew) - await SendCreatedAtAsync(new { definitionId }, response, cancellation: cancellationToken); + await SendCreatedAtAsync(new { definitionId }, mappedDefinition, cancellation: cancellationToken); else { await HttpContext.Response.WriteAsJsonAsync(response, serializerOptions, cancellationToken); diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Models.cs new file mode 100644 index 000000000..b41310a68 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Models.cs @@ -0,0 +1,5 @@ +using Elsa.Workflows.Api.Models; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Post; + +internal record Response(LinkedWorkflowDefinitionModel WorkflowDefinition, int ConsumingWorkflowCount); \ 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 87b0bb0d3..833c49d96 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Endpoint.cs @@ -55,10 +55,11 @@ internal class Publish(IWorkflowDefinitionStore store, IWorkflowDefinitionPublis return; } - await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken); - - var response = await linker.MapAsync(definition, cancellationToken); + var result = await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken); + var mappedDefinition = await linker.MapAsync(definition, cancellationToken); + var response = new Response(mappedDefinition, result.ConsumingWorkflows?.Count() ?? 0); + // We do not want to include composite root activities in the response. var serializerOptions = serializer.GetOptions().Clone(); serializerOptions.Converters.Add(new JsonIgnoreCompositeRootConverterFactory()); diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Models.cs index fe670b1b6..d16834475 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Models.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Publish/Models.cs @@ -1,6 +1,10 @@ +using Elsa.Workflows.Api.Models; + namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Publish; internal class Request { public string DefinitionId { get; set; } = default!; -} \ No newline at end of file +} + +internal record Response(LinkedWorkflowDefinitionModel WorkflowDefinition, int ConsumingWorkflowCount); \ 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 ad26e3639..4efbeb133 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, IWorkflowDefinitionManager workflowDefinitionManager, IAuthorizationService authorizationService) +internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinitionPublisher workflowDefinitionPublisher, IAuthorizationService authorizationService) : ElsaEndpoint { public override void Configure() @@ -43,7 +43,7 @@ internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinit return; } - var affectedWorkflows = await workflowDefinitionManager.UpdateReferencesInConsumingWorkflows(definition, cancellationToken); + var affectedWorkflows = await workflowDefinitionPublisher.UpdateReferencesInConsumingWorkflows(definition, cancellationToken); var response = new Response(affectedWorkflows.Select(w => w.Name ?? w.DefinitionId)); await SendOkAsync(response, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionManager.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionManager.cs index 7fddbc39d..3faac3ac1 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionManager.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionManager.cs @@ -64,12 +64,4 @@ public interface IWorkflowDefinitionManager /// The cancellation token. /// The new workflow definition. Task RevertVersionAsync(string definitionId, int version, 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/IWorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs index 402243f73..487074525 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionPublisher.cs @@ -68,4 +68,12 @@ 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/Models/PublishWorkflowDefinitionResult.cs b/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs index 2c646ba4b..87848df5a 100644 --- a/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs +++ b/src/modules/Elsa.Workflows.Management/Models/PublishWorkflowDefinitionResult.cs @@ -1,6 +1,8 @@ +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); \ No newline at end of file +public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection ValidationErrors, IEnumerable? ConsumingWorkflows); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs index d62374635..2c24dccda 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs @@ -1,8 +1,6 @@ using Elsa.Common.Models; -using Elsa.Extensions; using Elsa.Mediator.Contracts; 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; @@ -17,8 +15,6 @@ public class WorkflowDefinitionManager : IWorkflowDefinitionManager private readonly INotificationSender _notificationSender; private readonly IWorkflowDefinitionPublisher _workflowPublisher; private readonly IIdentityGenerator _identityGenerator; - private readonly IActivitySerializer _activitySerializer; - private readonly IActivityVisitor _activityVisitor; /// /// Constructor. @@ -27,16 +23,12 @@ public class WorkflowDefinitionManager : IWorkflowDefinitionManager IWorkflowDefinitionStore store, INotificationSender notificationSender, IWorkflowDefinitionPublisher workflowPublisher, - IIdentityGenerator identityGenerator, - IActivitySerializer activitySerializer, - IActivityVisitor activityVisitor) + IIdentityGenerator identityGenerator) { _store = store; _notificationSender = notificationSender; _workflowPublisher = workflowPublisher; _identityGenerator = identityGenerator; - _activitySerializer = activitySerializer; - _activityVisitor = activityVisitor; } /// @@ -143,53 +135,6 @@ public class WorkflowDefinitionManager : IWorkflowDefinitionManager return draft; } - /// - public async Task> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default) - { - var updatedWorkflowDefinitions = new List(); - - var workflowDefinitions = (await _store.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 _store.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken); - - return updatedWorkflowDefinitions; - } - private async Task EnsureLastVersionIsLatestAsync(IEnumerable definitionIds, CancellationToken cancellationToken) { foreach (var definitionId in definitionIds) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs index bdba1ab26..a0008149e 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -1,9 +1,11 @@ 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; @@ -21,6 +23,7 @@ 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; @@ -33,6 +36,7 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher IWorkflowDefinitionStore workflowDefinitionStore, INotificationSender notificationSender, IIdentityGenerator identityGenerator, + IActivityVisitor activityVisitor, IActivitySerializer activitySerializer, IRequestSender requestSender, ISystemClock systemClock) @@ -41,6 +45,7 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher _workflowDefinitionStore = workflowDefinitionStore; _notificationSender = notificationSender; _identityGenerator = identityGenerator; + _activityVisitor = activityVisitor; _activitySerializer = activitySerializer; _requestSender = requestSender; _systemClock = systemClock; @@ -77,7 +82,7 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher return new PublishWorkflowDefinitionResult(false, new List { new("Workflow definition not found.") - }); + }, null); return await PublishAsync(definition, cancellationToken); } @@ -90,7 +95,7 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher var validationErrors = responses.SelectMany(r => r.ValidationErrors).ToList(); if (validationErrors.Any()) - return new PublishWorkflowDefinitionResult(false, validationErrors); + return new PublishWorkflowDefinitionResult(false, validationErrors, null); await _notificationSender.SendAsync(new WorkflowDefinitionPublishing(definition), cancellationToken); @@ -113,7 +118,15 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher await _workflowDefinitionStore.SaveAsync(definition, cancellationToken); await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition), cancellationToken); - return new PublishWorkflowDefinitionResult(true, validationErrors); + + var consumingWorkflows = new List(); + + if (definition.Options.UsableAsActivity == true) + { + consumingWorkflows.AddRange(await UpdateReferencesInConsumingWorkflows(definition, cancellationToken)); + } + + return new PublishWorkflowDefinitionResult(true, validationErrors, consumingWorkflows); } /// @@ -196,6 +209,53 @@ 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 _workflowDefinitionStore.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken); + + return updatedWorkflowDefinitions; + } + private WorkflowDefinition Initialize(WorkflowDefinition definition) { if (definition.Id == null!) diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DependencyWorkflowsPublishing/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DependencyWorkflowsPublishing/Tests.cs index a312e7cdd..7b0f28b12 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DependencyWorkflowsPublishing/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DependencyWorkflowsPublishing/Tests.cs @@ -21,7 +21,7 @@ public class Tests private readonly IWorkflowDefinitionPublisher _workflowDefinitionPublisher; private readonly IActivitySerializer _activitySerializer; private readonly IActivityVisitor _activityVisitor; - private readonly IWorkflowDefinitionManager _workflowManager; + // private readonly IWorkflowDefinitionManager _workflowManager; /// /// Initializes a new instance of the class. @@ -34,7 +34,7 @@ public class Tests .Build(); _workflowDefinitionPublisher = _services.GetRequiredService(); - _workflowManager = _services.GetRequiredService(); + // _workflowManager = _services.GetRequiredService(); _activitySerializer = _services.GetRequiredService(); _activityVisitor = _services.GetRequiredService(); } @@ -59,8 +59,8 @@ public class Tests var childDefinitionV2 = (await _workflowDefinitionPublisher.GetDraftAsync(childDefinitionV1.DefinitionId, VersionOptions.Published))!; await _workflowDefinitionPublisher.PublishAsync(childDefinitionV2); - // Update consuming workflows to point to the new version of the child workflow. - await _workflowManager.UpdateReferencesInConsumingWorkflows(childDefinitionV2); + // // Update consuming workflows to point to the new version of the child workflow. + // await _workflowDefinitionPublisher.UpdateReferencesInConsumingWorkflows(childDefinitionV2); // Assert that the parent workflow now points to the new version of the child workflow. parentDefinition = await _services.GetWorkflowDefinitionAsync("parent", VersionOptions.Latest);