using Elsa.MassTransit.Messages;
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
using JetBrains.Annotations;
using MassTransit;
namespace Elsa.MassTransit.Consumers;
///
/// Consumes messages related to workflow definition changes.
///
[PublicAPI]
public class WorkflowDefinitionEventsConsumer(IActivityRegistryUpdateService activityRegistryUpdateService) :
IConsumer,
IConsumer,
IConsumer,
IConsumer,
IConsumer,
IConsumer,
IConsumer,
IConsumer
{
///
public Task Consume(ConsumeContext context)
{
return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity);
}
///
public Task Consume(ConsumeContext context)
{
activityRegistryUpdateService.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
return Task.CompletedTask;
}
///
public Task Consume(ConsumeContext context)
{
return UpdateDefinition(context.Message.Id, context.Message.UsableAsActivity);
}
///
public Task Consume(ConsumeContext context)
{
activityRegistryUpdateService.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
return Task.CompletedTask;
}
///
public Task Consume(ConsumeContext context)
{
foreach (var id in context.Message.Ids)
{
activityRegistryUpdateService.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
}
return Task.CompletedTask;
}
///
public Task Consume(ConsumeContext context)
{
activityRegistryUpdateService.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
return Task.CompletedTask;
}
///
public Task Consume(ConsumeContext context)
{
foreach (var id in context.Message.Ids)
{
activityRegistryUpdateService.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
}
return Task.CompletedTask;
}
///
public async Task Consume(ConsumeContext context)
{
foreach (KeyValuePair definitionAsActivity in context.Message.DefinitionsAsActivity)
{
await UpdateDefinition(definitionAsActivity.Key, definitionAsActivity.Value);
}
}
private Task UpdateDefinition(string id, bool usableAsActivity)
{
if (usableAsActivity)
return activityRegistryUpdateService.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), id);
activityRegistryUpdateService.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
return Task.CompletedTask;
}
}