From db079ae41b62ec0a4cc909388da6f9570f5c8db3 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Mon, 13 May 2024 12:44:24 +0200 Subject: [PATCH] Updated distributed workflow definition event handling for activity registry --- .editorconfig | 1 + .../WorkflowDefinitionEventsConsumer.cs | 49 +++++++++++++------ ...butedWorkflowDefinitionEventsDispatcher.cs | 5 -- ...dWorkflowDefinitionNotificationsHandler.cs | 4 -- .../Contracts/IActivityRegistry.cs | 7 +++ .../Services/ActivityRegistry.cs | 7 +++ .../WorkflowDefinitionActivityProvider.cs | 1 + .../Contracts/IActivityRegistryPopulator.cs | 17 +++++++ .../Services/ActivityRegistryPopulator.cs | 37 +++++++++++++- .../Services/WorkflowDefinitionManager.cs | 1 - 10 files changed, 103 insertions(+), 26 deletions(-) diff --git a/.editorconfig b/.editorconfig index e558c49f3..209a9ac3d 100644 --- a/.editorconfig +++ b/.editorconfig @@ -14,6 +14,7 @@ tab_width = 4 # New line preferences end_of_line = crlf insert_final_newline = false +max_line_length = 180 #### .NET Coding Conventions #### diff --git a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs index 75bbb73e1..5dc953707 100644 --- a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs +++ b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs @@ -14,34 +14,55 @@ public class WorkflowDefinitionEventsConsumer(IActivityRegistryPopulator activit IConsumer, IConsumer, IConsumer, - IConsumer, IConsumer, IConsumer, IConsumer { /// - public Task Consume(ConsumeContext context) => RefreshAsync(); + public Task Consume(ConsumeContext context) + { + return activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id); + } /// - public Task Consume(ConsumeContext context) => RefreshAsync(); - - /// - public async Task Consume(ConsumeContext context) - { - await activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id); + public Task Consume(ConsumeContext context) + { + activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id); + return Task.CompletedTask; } /// - public Task Consume(ConsumeContext context) => RefreshAsync(); + public Task Consume(ConsumeContext context) + { + return activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id); + } /// - public Task Consume(ConsumeContext context) => RefreshAsync(); + public Task Consume(ConsumeContext context) + { + foreach (var id in context.Message.Ids) + { + activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id); + } + + return Task.CompletedTask; + } /// - public Task Consume(ConsumeContext context) => RefreshAsync(); + public Task Consume(ConsumeContext context) + { + activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id); + return Task.CompletedTask; + } /// - public Task Consume(ConsumeContext context) => RefreshAsync(); - - private async Task RefreshAsync() => await activityRegistryPopulator.PopulateRegistryAsync(typeof(WorkflowDefinitionActivityProvider)); + public Task Consume(ConsumeContext context) + { + foreach (var id in context.Message.Ids) + { + activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id); + } + + return Task.CompletedTask; + } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs b/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs index bd0231b71..5f3967dcd 100644 --- a/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs +++ b/src/modules/Elsa.MassTransit/Contracts/IDistributedWorkflowDefinitionEventsDispatcher.cs @@ -12,11 +12,6 @@ public interface IDistributedWorkflowDefinitionEventsDispatcher /// Task DispatchAsync(WorkflowDefinitionPublished request, CancellationToken cancellationToken); - /// - /// Dispatches a workflow definition retracted event. - /// - Task DispatchAsync(WorkflowDefinitionRetracted 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 0aeaa5b06..7f3b61dce 100644 --- a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs +++ b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs @@ -8,7 +8,6 @@ namespace Elsa.MassTransit.Handlers /// Represents a handler for distributed workflow definition notifications. public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkflowDefinitionEventsDispatcher distributedEventsDispatcher) : INotificationHandler, - INotificationHandler, INotificationHandler, INotificationHandler, INotificationHandler, @@ -19,9 +18,6 @@ namespace Elsa.MassTransit.Handlers public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) => await distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionPublished(notification.WorkflowDefinition.Id), cancellationToken); - /// - public async Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => - await distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionRetracted(notification.WorkflowDefinition.Id), cancellationToken); /// public async Task HandleAsync(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken) => diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IActivityRegistry.cs b/src/modules/Elsa.Workflows.Core/Contracts/IActivityRegistry.cs index be4771710..0e91e44ec 100644 --- a/src/modules/Elsa.Workflows.Core/Contracts/IActivityRegistry.cs +++ b/src/modules/Elsa.Workflows.Core/Contracts/IActivityRegistry.cs @@ -21,6 +21,13 @@ public interface IActivityRegistry : IActivityProvider /// The activity descriptors to add. void AddMany(Type providerType, IEnumerable descriptors); + /// + /// Removes an activity descriptor from the registry. + /// + /// The type of the activity provider. + /// The activity descriptor to remove. + void Remove(Type providerType, ActivityDescriptor descriptor); + /// /// Clears all activity descriptors from the registry. /// diff --git a/src/modules/Elsa.Workflows.Core/Services/ActivityRegistry.cs b/src/modules/Elsa.Workflows.Core/Services/ActivityRegistry.cs index 88c019b03..ccfa3ec69 100644 --- a/src/modules/Elsa.Workflows.Core/Services/ActivityRegistry.cs +++ b/src/modules/Elsa.Workflows.Core/Services/ActivityRegistry.cs @@ -38,6 +38,13 @@ public class ActivityRegistry : IActivityRegistry Add(descriptor, target); } + /// + public void Remove(Type providerType, ActivityDescriptor descriptor) + { + _providedActivityDescriptors[providerType].Remove(descriptor); + _activityDescriptors.Remove((descriptor.TypeName, descriptor.Version), out _); + } + /// public void Clear() { diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs index 4ec7241e8..1a5962a1a 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs @@ -99,6 +99,7 @@ public class WorkflowDefinitionActivityProvider : IActivityProvider CustomProperties = { ["RootType"] = nameof(WorkflowDefinitionActivity), + ["WorkflowDefinitionId"] = definition.DefinitionId, ["WorkflowDefinitionVersionId"] = definition.Id }, ConstructionProperties = new Dictionary diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IActivityRegistryPopulator.cs b/src/modules/Elsa.Workflows.Management/Contracts/IActivityRegistryPopulator.cs index 5988d72b1..2e957f12c 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IActivityRegistryPopulator.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IActivityRegistryPopulator.cs @@ -27,4 +27,21 @@ public interface IActivityRegistryPopulator /// The ID of the workflow definition. /// The cancellation token. Task AddToRegistry(Type providerType, string workflowDefinitionId, CancellationToken cancellationToken = default); + + /// + /// Removes workflow definition activities from the . + /// + /// The type of the Activity Provider. + /// The ID of the workflow definition to remove. + /// The cancellation token. + void RemoveDefinitionFromRegistry(Type providerType, string workflowDefinitionId, CancellationToken cancellationToken = default); + + + /// + /// Removes a workflow definition version activity from the . + /// + /// The type of the Activity Provider. + /// The ID of the workflow definition to remove. + /// The cancellation token. + void RemoveDefinitionVersionFromRegistry(Type providerType, string workflowDefinitionVersionId, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/ActivityRegistryPopulator.cs b/src/modules/Elsa.Workflows.Management/Services/ActivityRegistryPopulator.cs index acbf4f924..5c47899f9 100644 --- a/src/modules/Elsa.Workflows.Management/Services/ActivityRegistryPopulator.cs +++ b/src/modules/Elsa.Workflows.Management/Services/ActivityRegistryPopulator.cs @@ -1,5 +1,6 @@ using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Models; namespace Elsa.Workflows.Management.Services; @@ -43,10 +44,42 @@ public class ActivityRegistryPopulator : IActivityRegistryPopulator var provider = _providers.First(x => x.GetType() == providerType); var descriptors = await provider.GetDescriptorsAsync(cancellationToken); var descriptorToAdd = descriptors - .SingleOrDefault(d => d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && val.ToString() == workflowDefinitionVersionId); + .SingleOrDefault(d => + d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && + val.ToString() == workflowDefinitionVersionId); if (descriptorToAdd is not null) - _registry.Add(providerType, descriptorToAdd!); + _registry.Add(providerType, descriptorToAdd); + } + + /// + public void RemoveDefinitionFromRegistry(Type providerType, string workflowDefinitionId, CancellationToken cancellationToken = default) + { + var providerDescriptors = _registry.ListByProvider(providerType); + + var descriptorsToRemove = providerDescriptors + .Where(d => + d.CustomProperties.TryGetValue("WorkflowDefinitionId", out var val) && + val.ToString() == workflowDefinitionId); + + foreach (ActivityDescriptor activityDescriptor in descriptorsToRemove) + { + _registry.Remove(providerType, activityDescriptor); + } + } + + /// + public void RemoveDefinitionVersionFromRegistry(Type providerType, string workflowDefinitionVersionId, CancellationToken cancellationToken = default) + { + var providerDescriptors = _registry.ListByProvider(providerType); + + var descriptorToRemove = providerDescriptors + .SingleOrDefault(d => + d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && + val.ToString() == workflowDefinitionVersionId); + + if (descriptorToRemove is not null) + _registry.Remove(providerType, descriptorToRemove); } private async Task PopulateRegistryAsync(IActivityProvider provider, CancellationToken cancellationToken = default) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs index 2c24dccda..657f72528 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionManager.cs @@ -37,7 +37,6 @@ public class WorkflowDefinitionManager : IWorkflowDefinitionManager await _notificationSender.SendAsync(new WorkflowDefinitionDeleting(definitionId), cancellationToken); var filter = new WorkflowDefinitionFilter { DefinitionId = definitionId }; var count = await _store.DeleteAsync(filter, cancellationToken); - await EnsureLastVersionIsLatestAsync(definitionId, cancellationToken); await _notificationSender.SendAsync(new WorkflowDefinitionDeleted(definitionId), cancellationToken); return count; }