using Elsa.MassTransit.Messages;
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) :
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)
{
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)
{
workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(context.Message.Id);
return Task.CompletedTask;
}
///
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 (KeyValuePair definitionAsActivity in context.Message.DefinitionsAsActivity)
{
await UpdateDefinition(definitionAsActivity.Key, definitionAsActivity.Value);
}
}
private Task UpdateDefinition(string id, bool usableAsActivity)
{
if (usableAsActivity)
return workflowDefinitionActivityRegistryUpdater.AddToRegistry(id);
workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(id);
return Task.CompletedTask;
}
}