using Elsa.MassTransit.Messages; using Elsa.MassTransit.Services; using Elsa.Mediator.Contracts; using Elsa.Workflows.Management.Contracts; using JetBrains.Annotations; using MassTransit; namespace Elsa.MassTransit.Consumers; /// Consumes messages related to workflow definition changes. [PublicAPI] public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater, INotificationSender notificationSender) : global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer { /// public Task Consume(ConsumeContext context) { workflowDefinitionActivityRegistryUpdater.RemoveDefinitionFromRegistry(context.Message.Id); return Task.CompletedTask; } /// public Task Consume(ConsumeContext context) { return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity); } /// public Task Consume(ConsumeContext context) { return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity); } /// public Task Consume(ConsumeContext context) { return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity); } /// public Task Consume(ConsumeContext context) { foreach (var id in context.Message.Ids) workflowDefinitionActivityRegistryUpdater.RemoveDefinitionFromRegistry(id); return Task.CompletedTask; } /// public Task Consume(ConsumeContext context) { workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(context.Message.Id); return Task.CompletedTask; } /// public Task Consume(ConsumeContext context) { foreach (var id in context.Message.Ids) workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(id); return Task.CompletedTask; } /// public async Task Consume(ConsumeContext context) { foreach (WorkflowDefinitionVersionUpdate definitionUpdate in context.Message.WorkflowDefinitionVersionUpdates) await UpdateDefinition(definitionUpdate.Id, definitionUpdate.UsableAsActivity); } /// public async Task Consume(ConsumeContext context) { var message = context.Message; var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsRefreshed(message.WorkflowDefinitionIds); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = true; await notificationSender.SendAsync(notification, context.CancellationToken); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = false; } /// public async Task Consume(ConsumeContext context) { var message = context.Message; var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsReloaded(message.WorkflowDefinitionIds); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = true; await notificationSender.SendAsync(notification, context.CancellationToken); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = false; } private Task UpdateDefinition(string id, bool 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); return Task.CompletedTask; } }