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.