From 171f1478018faa6e2448bec66bd4a0b8a942cc13 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Mon, 24 Jun 2024 10:45:53 +0200 Subject: [PATCH] Update workflow definition registry and retraction handling This update enhances the handling of workflow definition retraction across multiple files and includes new notifications for workflow definition version retraction. The logic for updating the workflow definition registry has been adapted to keep published workflows in the registry unless they are no longer marked as an activity. --- Directory.Packages.props | 6 +++--- .../WorkflowDefinitionEventsConsumer.cs | 19 +++++++++++++------ ...butedWorkflowDefinitionEventsDispatcher.cs | 5 +++++ ...dWorkflowDefinitionNotificationsHandler.cs | 11 +++++++++-- .../Messages/WorkflowDefinitionRetracted.cs | 3 ++- .../WorkflowDefinitionVersionRetracted.cs | 10 ++++++++++ .../MassTransitDistributedEventsDispatcher.cs | 6 ++++++ .../Handlers/RefreshActivityRegistry.cs | 19 +++++++++++++------ .../WorkflowDefinitionVersionRetracted.cs | 12 ++++++++++++ .../Services/WorkflowDefinitionPublisher.cs | 4 ++++ 10 files changed, 77 insertions(+), 18 deletions(-) create mode 100644 src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionVersionRetracted.cs create mode 100644 src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionVersionRetracted.cs diff --git a/Directory.Packages.props b/Directory.Packages.props index a45a005f4..581449a8e 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -23,9 +23,9 @@ - - - + + + diff --git a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs index 0722687c2..ce3757c72 100644 --- a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs +++ b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs @@ -13,6 +13,7 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr IConsumer, IConsumer, IConsumer, + IConsumer, IConsumer, IConsumer, IConsumer, @@ -28,14 +29,19 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr /// public Task Consume(ConsumeContext context) { - return UpdateDefinition(context.Message.Id, true, context.Message.UsableAsActivity); + return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity); } /// public Task Consume(ConsumeContext context) { - workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(context.Message.Id); - return Task.CompletedTask; + return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity); + } + + /// + public Task Consume(ConsumeContext context) + { + return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity); } /// @@ -72,13 +78,14 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr { foreach (WorkflowDefinitionVersionUpdate definitionUpdate in context.Message.WorkflowDefinitionVersionUpdates) { - await UpdateDefinition(definitionUpdate.Id, definitionUpdate.IsPublished, definitionUpdate.UsableAsActivity); + await UpdateDefinition(definitionUpdate.Id, definitionUpdate.UsableAsActivity); } } - private Task UpdateDefinition(string id, bool isPublished, bool usableAsActivity) + private Task UpdateDefinition(string id, bool usableAsActivity) { - if (isPublished && usableAsActivity) + // Once a workflow has been published it should remain in the activity registry unless no longer being marked as an activity. + if (usableAsActivity) return workflowDefinitionActivityRegistryUpdater.AddToRegistry(id); workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(id); diff --git a/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs b/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs index 27746cf15..46a670eec 100644 --- a/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs +++ b/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs @@ -17,6 +17,11 @@ public interface IDistributedWorkflowDefinitionEventsDispatcher /// Task DispatchAsync(WorkflowDefinitionRetracted request, CancellationToken cancellationToken); + /// + /// Dispatches a workflow definition version retracted event. + /// + Task DispatchAsync(WorkflowDefinitionVersionRetracted request, CancellationToken cancellationToken); + /// /// Dispatches a workflow definition deleted event. /// diff --git a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs index 657ebca18..1c2f38774 100644 --- a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs +++ b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs @@ -9,6 +9,7 @@ namespace Elsa.MassTransit.Handlers; public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkflowDefinitionEventsDispatcher distributedEventsDispatcher) : INotificationHandler, INotificationHandler, + INotificationHandler, INotificationHandler, INotificationHandler, INotificationHandler, @@ -21,8 +22,14 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkf notification.WorkflowDefinition.Options.UsableAsActivity.GetValueOrDefault()), cancellationToken); /// - public Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => - distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionRetracted(notification.WorkflowDefinition.Id), cancellationToken); + public Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => + distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionRetracted(notification.WorkflowDefinition.Id, + notification.WorkflowDefinition.Options.UsableAsActivity.GetValueOrDefault()), cancellationToken); + + /// + public Task HandleAsync(WorkflowDefinitionVersionRetracted notification, CancellationToken cancellationToken) => + distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionVersionRetracted(notification.WorkflowDefinition.Id, + notification.WorkflowDefinition.Options.UsableAsActivity.GetValueOrDefault()), cancellationToken); /// public Task HandleAsync(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken) => diff --git a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionRetracted.cs b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionRetracted.cs index f6106c707..aea258fdc 100644 --- a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionRetracted.cs +++ b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionRetracted.cs @@ -3,7 +3,8 @@ namespace Elsa.MassTransit.Messages; /// /// Represents a distributed message that a workflow definition has been retracted. /// -public class WorkflowDefinitionRetracted(string id) +public class WorkflowDefinitionRetracted(string id, bool usableAsActivity) { public string Id { get; set; } = id; + public bool UsableAsActivity { get; set; } = usableAsActivity; } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionVersionRetracted.cs b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionVersionRetracted.cs new file mode 100644 index 000000000..9d315a6ab --- /dev/null +++ b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionVersionRetracted.cs @@ -0,0 +1,10 @@ +namespace Elsa.MassTransit.Messages; + +/// +/// Represents a distributed message that a workflow definition version has been retracted. +/// +public class WorkflowDefinitionVersionRetracted(string id, bool usableAsActivity) +{ + public string Id { get; set; } = id; + public bool UsableAsActivity { get; set; } = usableAsActivity; +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Services/MassTransitDistributedEventsDispatcher.cs b/src/modules/Elsa.MassTransit/Services/MassTransitDistributedEventsDispatcher.cs index a976aa2b0..9ea97cf70 100644 --- a/src/modules/Elsa.MassTransit/Services/MassTransitDistributedEventsDispatcher.cs +++ b/src/modules/Elsa.MassTransit/Services/MassTransitDistributedEventsDispatcher.cs @@ -21,6 +21,12 @@ public class MassTransitDistributedEventsDispatcher(IBus bus) : IDistributedWork return bus.Publish(request, cancellationToken); } + /// + public Task DispatchAsync(WorkflowDefinitionVersionRetracted request, CancellationToken cancellationToken) + { + return bus.Publish(request, cancellationToken); + } + /// public Task DispatchAsync(WorkflowDefinitionDeleted request, CancellationToken cancellationToken) { diff --git a/src/modules/Elsa.Workflows.Management/Handlers/RefreshActivityRegistry.cs b/src/modules/Elsa.Workflows.Management/Handlers/RefreshActivityRegistry.cs index 95713d2b4..399c84620 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/RefreshActivityRegistry.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/RefreshActivityRegistry.cs @@ -15,6 +15,7 @@ namespace Elsa.Workflows.Management.Handlers; public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater) : INotificationHandler, INotificationHandler, + INotificationHandler, INotificationHandler, INotificationHandler, INotificationHandler, @@ -24,14 +25,19 @@ public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater /// public Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) { - return UpdateDefinition(notification.WorkflowDefinition.Id, true, notification.WorkflowDefinition.Options.UsableAsActivity); + return UpdateDefinition(notification.WorkflowDefinition.Id, notification.WorkflowDefinition.Options.UsableAsActivity); } /// public Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) { - workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(notification.WorkflowDefinition.Id); - return Task.CompletedTask; + return UpdateDefinition(notification.WorkflowDefinition.Id, notification.WorkflowDefinition.Options.UsableAsActivity); + } + + /// + public Task HandleAsync(WorkflowDefinitionVersionRetracted notification, CancellationToken cancellationToken) + { + return UpdateDefinition(notification.WorkflowDefinition.Id, notification.WorkflowDefinition.Options.UsableAsActivity); } /// @@ -75,13 +81,14 @@ public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater { foreach (var definition in notification.WorkflowDefinitions) { - await UpdateDefinition(definition.Id, definition.IsPublished, definition.Options.UsableAsActivity); + await UpdateDefinition(definition.Id, definition.Options.UsableAsActivity); } } - private Task UpdateDefinition(string id, bool isPublished, bool? usableAsActivity) + private Task UpdateDefinition(string id, bool? usableAsActivity) { - if (isPublished && usableAsActivity.GetValueOrDefault()) + // Once a workflow has been published it should remain in the activity registry unless no longer being marked as an activity. + if (usableAsActivity.GetValueOrDefault()) return workflowDefinitionActivityRegistryUpdater.AddToRegistry(id); workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(id); diff --git a/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionVersionRetracted.cs b/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionVersionRetracted.cs new file mode 100644 index 000000000..c6eb03f0f --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Notifications/WorkflowDefinitionVersionRetracted.cs @@ -0,0 +1,12 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Management.Entities; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Management.Notifications; + +/// +/// A notification that is sent when a workflow definition version is retracted. +/// +/// The workflow definition. +[PublicAPI] +public record WorkflowDefinitionVersionRetracted(WorkflowDefinition WorkflowDefinition) : 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 275bc88a3..173c3b29d 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -107,9 +107,13 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher foreach (var publishedAndOrLatestWorkflow in publishedWorkflows) { + var isPublished = publishedAndOrLatestWorkflow.IsPublished; publishedAndOrLatestWorkflow.IsPublished = false; publishedAndOrLatestWorkflow.IsLatest = false; await _workflowDefinitionStore.SaveAsync(publishedAndOrLatestWorkflow, cancellationToken); + + if (isPublished) + await _notificationSender.SendAsync(new WorkflowDefinitionVersionRetracted(publishedAndOrLatestWorkflow), cancellationToken); } // Save the new published definition.