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.
This commit is contained in:
Raymond den Haan 2024-06-24 10:45:53 +02:00 committed by raymonddenhaan
parent 11d0a53dce
commit 171f147801
10 changed files with 77 additions and 18 deletions

View file

@ -23,9 +23,9 @@
<PackageVersion Include="DistributedLock.Postgres" Version="1.1.0" />
<PackageVersion Include="DistributedLock.Redis" Version="1.0.3" />
<PackageVersion Include="Elastic.Clients.Elasticsearch" Version="8.14.0" />
<PackageVersion Include="Elsa.Studio" Version="3.2.0-rc2.351" />
<PackageVersion Include="Elsa.Studio.Core.BlazorWasm" Version="3.2.0-rc2.351" />
<PackageVersion Include="Elsa.Studio.Login.BlazorWasm" Version="3.2.0-rc2.351" />
<PackageVersion Include="Elsa.Studio" Version="3.2.0-rc2.369" />
<PackageVersion Include="Elsa.Studio.Core.BlazorWasm" Version="3.2.0-rc2.369" />
<PackageVersion Include="Elsa.Studio.Login.BlazorWasm" Version="3.2.0-rc2.369" />
<PackageVersion Include="FastEndpoints" Version="5.26.0" />
<PackageVersion Include="FastEndpoints.Security" Version="5.26.0" />
<PackageVersion Include="FastEndpoints.Swagger" Version="5.26.0" />

View file

@ -13,6 +13,7 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr
IConsumer<WorkflowDefinitionDeleted>,
IConsumer<WorkflowDefinitionPublished>,
IConsumer<WorkflowDefinitionRetracted>,
IConsumer<WorkflowDefinitionVersionRetracted>,
IConsumer<WorkflowDefinitionsDeleted>,
IConsumer<WorkflowDefinitionVersionDeleted>,
IConsumer<WorkflowDefinitionVersionsDeleted>,
@ -28,14 +29,19 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionPublished> context)
{
return UpdateDefinition(context.Message.Id, true, context.Message.UsableAsActivity);
return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity);
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionRetracted> context)
{
workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(context.Message.Id);
return Task.CompletedTask;
return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity);
}
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionVersionRetracted> context)
{
return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity);
}
/// <inheritdoc />
@ -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);

View file

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

View file

@ -9,6 +9,7 @@ namespace Elsa.MassTransit.Handlers;
public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkflowDefinitionEventsDispatcher distributedEventsDispatcher) :
INotificationHandler<WorkflowDefinitionPublished>,
INotificationHandler<WorkflowDefinitionRetracted>,
INotificationHandler<WorkflowDefinitionVersionRetracted>,
INotificationHandler<WorkflowDefinitionDeleted>,
INotificationHandler<WorkflowDefinitionsDeleted>,
INotificationHandler<WorkflowDefinitionVersionDeleted>,
@ -21,8 +22,14 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkf
notification.WorkflowDefinition.Options.UsableAsActivity.GetValueOrDefault()), cancellationToken);
/// <inheritdoc />
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);
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionVersionRetracted notification, CancellationToken cancellationToken) =>
distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionVersionRetracted(notification.WorkflowDefinition.Id,
notification.WorkflowDefinition.Options.UsableAsActivity.GetValueOrDefault()), cancellationToken);
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken) =>

View file

@ -3,7 +3,8 @@ namespace Elsa.MassTransit.Messages;
/// <summary>
/// Represents a distributed message that a workflow definition has been retracted.
/// </summary>
public class WorkflowDefinitionRetracted(string id)
public class WorkflowDefinitionRetracted(string id, bool usableAsActivity)
{
public string Id { get; set; } = id;
public bool UsableAsActivity { get; set; } = usableAsActivity;
}

View file

@ -0,0 +1,10 @@
namespace Elsa.MassTransit.Messages;
/// <summary>
/// Represents a distributed message that a workflow definition version has been retracted.
/// </summary>
public class WorkflowDefinitionVersionRetracted(string id, bool usableAsActivity)
{
public string Id { get; set; } = id;
public bool UsableAsActivity { get; set; } = usableAsActivity;
}

View file

@ -21,6 +21,12 @@ public class MassTransitDistributedEventsDispatcher(IBus bus) : IDistributedWork
return bus.Publish(request, cancellationToken);
}
/// <inheritdoc />
public Task DispatchAsync(WorkflowDefinitionVersionRetracted request, CancellationToken cancellationToken)
{
return bus.Publish(request, cancellationToken);
}
/// <inheritdoc />
public Task DispatchAsync(WorkflowDefinitionDeleted request, CancellationToken cancellationToken)
{

View file

@ -15,6 +15,7 @@ namespace Elsa.Workflows.Management.Handlers;
public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater) :
INotificationHandler<WorkflowDefinitionPublished>,
INotificationHandler<WorkflowDefinitionRetracted>,
INotificationHandler<WorkflowDefinitionVersionRetracted>,
INotificationHandler<WorkflowDefinitionDeleted>,
INotificationHandler<WorkflowDefinitionsDeleted>,
INotificationHandler<WorkflowDefinitionVersionDeleted>,
@ -24,14 +25,19 @@ public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater
/// <inheritdoc />
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);
}
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken)
{
workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(notification.WorkflowDefinition.Id);
return Task.CompletedTask;
return UpdateDefinition(notification.WorkflowDefinition.Id, notification.WorkflowDefinition.Options.UsableAsActivity);
}
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionVersionRetracted notification, CancellationToken cancellationToken)
{
return UpdateDefinition(notification.WorkflowDefinition.Id, notification.WorkflowDefinition.Options.UsableAsActivity);
}
/// <inheritdoc />
@ -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);

View file

@ -0,0 +1,12 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using JetBrains.Annotations;
namespace Elsa.Workflows.Management.Notifications;
/// <summary>
/// A notification that is sent when a workflow definition version is retracted.
/// </summary>
/// <param name="WorkflowDefinition">The workflow definition.</param>
[PublicAPI]
public record WorkflowDefinitionVersionRetracted(WorkflowDefinition WorkflowDefinition) : INotification;

View file

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