Updated distributed workflow definition event handling for activity registry

This commit is contained in:
Raymond den Haan 2024-05-13 12:44:24 +02:00 committed by raymonddenhaan
parent f49bf6e112
commit db079ae41b
10 changed files with 103 additions and 26 deletions

View file

@ -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 ####

View file

@ -14,34 +14,55 @@ public class WorkflowDefinitionEventsConsumer(IActivityRegistryPopulator activit
IConsumer<WorkflowDefinitionCreated>,
IConsumer<WorkflowDefinitionDeleted>,
IConsumer<WorkflowDefinitionPublished>,
IConsumer<WorkflowDefinitionRetracted>,
IConsumer<WorkflowDefinitionsDeleted>,
IConsumer<WorkflowDefinitionVersionDeleted>,
IConsumer<WorkflowDefinitionVersionsDeleted>
{
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionCreated> context) => RefreshAsync();
public Task Consume(ConsumeContext<WorkflowDefinitionCreated> context)
{
return activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionDeleted> context) => RefreshAsync();
/// <inheritdoc />
public async Task Consume(ConsumeContext<WorkflowDefinitionPublished> context)
{
await activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
public Task Consume(ConsumeContext<WorkflowDefinitionDeleted> context)
{
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
return Task.CompletedTask;
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionRetracted> context) => RefreshAsync();
public Task Consume(ConsumeContext<WorkflowDefinitionPublished> context)
{
return activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionsDeleted> context) => RefreshAsync();
public Task Consume(ConsumeContext<WorkflowDefinitionsDeleted> context)
{
foreach (var id in context.Message.Ids)
{
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
}
return Task.CompletedTask;
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionVersionDeleted> context) => RefreshAsync();
public Task Consume(ConsumeContext<WorkflowDefinitionVersionDeleted> context)
{
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
return Task.CompletedTask;
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionVersionsDeleted> context) => RefreshAsync();
private async Task RefreshAsync() => await activityRegistryPopulator.PopulateRegistryAsync(typeof(WorkflowDefinitionActivityProvider));
public Task Consume(ConsumeContext<WorkflowDefinitionVersionsDeleted> context)
{
foreach (var id in context.Message.Ids)
{
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
}
return Task.CompletedTask;
}
}

View file

@ -12,11 +12,6 @@ public interface IDistributedWorkflowDefinitionEventsDispatcher
/// </summary>
Task DispatchAsync(WorkflowDefinitionPublished request, CancellationToken cancellationToken);
/// <summary>
/// Dispatches a workflow definition retracted event.
/// </summary>
Task DispatchAsync(WorkflowDefinitionRetracted request, CancellationToken cancellationToken);
/// <summary>
/// Dispatches a workflow definition deleted event.
/// </summary>

View file

@ -8,7 +8,6 @@ namespace Elsa.MassTransit.Handlers
/// Represents a handler for distributed workflow definition notifications.
public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkflowDefinitionEventsDispatcher distributedEventsDispatcher) :
INotificationHandler<WorkflowDefinitionPublished>,
INotificationHandler<WorkflowDefinitionRetracted>,
INotificationHandler<WorkflowDefinitionDeleted>,
INotificationHandler<WorkflowDefinitionsDeleted>,
INotificationHandler<WorkflowDefinitionCreated>,
@ -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);
/// <inheritdoc />
public async Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) =>
await distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionRetracted(notification.WorkflowDefinition.Id), cancellationToken);
/// <inheritdoc />
public async Task HandleAsync(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken) =>

View file

@ -21,6 +21,13 @@ public interface IActivityRegistry : IActivityProvider
/// <param name="descriptors">The activity descriptors to add.</param>
void AddMany(Type providerType, IEnumerable<ActivityDescriptor> descriptors);
/// <summary>
/// Removes an activity descriptor from the registry.
/// </summary>
/// <param name="providerType">The type of the activity provider.</param>
/// <param name="descriptor">The activity descriptor to remove.</param>
void Remove(Type providerType, ActivityDescriptor descriptor);
/// <summary>
/// Clears all activity descriptors from the registry.
/// </summary>

View file

@ -38,6 +38,13 @@ public class ActivityRegistry : IActivityRegistry
Add(descriptor, target);
}
/// <inheritdoc />
public void Remove(Type providerType, ActivityDescriptor descriptor)
{
_providedActivityDescriptors[providerType].Remove(descriptor);
_activityDescriptors.Remove((descriptor.TypeName, descriptor.Version), out _);
}
/// <inheritdoc />
public void Clear()
{

View file

@ -99,6 +99,7 @@ public class WorkflowDefinitionActivityProvider : IActivityProvider
CustomProperties =
{
["RootType"] = nameof(WorkflowDefinitionActivity),
["WorkflowDefinitionId"] = definition.DefinitionId,
["WorkflowDefinitionVersionId"] = definition.Id
},
ConstructionProperties = new Dictionary<string, object>

View file

@ -27,4 +27,21 @@ public interface IActivityRegistryPopulator
/// <param name="workflowDefinitionId">The ID of the workflow definition.</param>
/// <param name="cancellationToken">The cancellation token.</param>
Task AddToRegistry(Type providerType, string workflowDefinitionId, CancellationToken cancellationToken = default);
/// <summary>
/// Removes workflow definition activities from the <see cref="IActivityRegistry"/>.
/// </summary>
/// <param name="providerType">The type of the Activity Provider.</param>
/// <param name="workflowDefinitionId">The ID of the workflow definition to remove.</param>
/// <param name="cancellationToken">The cancellation token.</param>
void RemoveDefinitionFromRegistry(Type providerType, string workflowDefinitionId, CancellationToken cancellationToken = default);
/// <summary>
/// Removes a workflow definition version activity from the <see cref="IActivityRegistry"/>.
/// </summary>
/// <param name="providerType">The type of the Activity Provider.</param>
/// <param name="workflowDefinitionVersionId">The ID of the workflow definition to remove.</param>
/// <param name="cancellationToken">The cancellation token.</param>
void RemoveDefinitionVersionFromRegistry(Type providerType, string workflowDefinitionVersionId, CancellationToken cancellationToken = default);
}

View file

@ -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);
}
/// <inheritdoc />
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);
}
}
/// <inheritdoc />
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)

View file

@ -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;
}