From b1380616618ebddefc835b919ad722fdfd03fb88 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 27 Jul 2021 21:47:30 +0200 Subject: [PATCH] Reuse certain queues for certain event types (#1317) --- .../Extensions/ServiceCollectionExtensions.cs | 8 +++--- .../ElsaOptionsBuilderExtensions.cs | 4 +-- src/core/Elsa.Core/ElsaOptions.cs | 17 ++++++------ src/core/Elsa.Core/ElsaOptionsBuilder.cs | 12 ++++----- .../ElsaServiceCollectionExtensions.cs | 24 +++++++---------- .../Services/Messaging/ServiceBusFactory.cs | 5 ++-- .../StartupTasks/CreateSubscriptions.cs | 27 ++++++++++++------- 7 files changed, 52 insertions(+), 45 deletions(-) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs index 8dd0b7deb..b3fa3a26d 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -40,10 +40,10 @@ namespace Elsa.Activities.AzureServiceBus.Extensions .AddBookmarkProvider() ; - options.AddCompetingConsumer(); - options.AddCompetingConsumer(); - options.AddCompetingConsumer(); - options.AddCompetingConsumer(); + options.AddPubSubConsumer("WorkflowDefinitionEvents"); + options.AddPubSubConsumer("WorkflowDefinitionEvents"); + options.AddPubSubConsumer("WorkflowDefinitionEvents"); + options.AddPubSubConsumer("WorkflowDefinitionEvents"); options .AddActivity() diff --git a/src/activities/Elsa.Activities.Conductor/Extensions/ElsaOptionsBuilderExtensions.cs b/src/activities/Elsa.Activities.Conductor/Extensions/ElsaOptionsBuilderExtensions.cs index 745993ea7..d07759803 100644 --- a/src/activities/Elsa.Activities.Conductor/Extensions/ElsaOptionsBuilderExtensions.cs +++ b/src/activities/Elsa.Activities.Conductor/Extensions/ElsaOptionsBuilderExtensions.cs @@ -43,8 +43,8 @@ namespace Elsa.Activities.Conductor.Extensions elsa .AddActivitiesFrom() - .AddCompetingConsumer() - .AddCompetingConsumer(); + .AddCompetingConsumer("ConductorCommand") + .AddCompetingConsumer("ConductorCommand"); return elsa; } diff --git a/src/core/Elsa.Core/ElsaOptions.cs b/src/core/Elsa.Core/ElsaOptions.cs index e0bd103ac..72a566f62 100644 --- a/src/core/Elsa.Core/ElsaOptions.cs +++ b/src/core/Elsa.Core/ElsaOptions.cs @@ -22,19 +22,20 @@ using Storage.Net.Blobs; namespace Elsa { - public record CompetingMessageType(Type MessageType, string? Queue = default); - + public record MessageTypeConfig(Type MessageType, string? QueueName = default); + public class ElsaOptions { public static string FormatChannelQueueName(string channel) => FormatChannelQueueName(typeof(TMessage), channel); - - public static string FormatChannelQueueName(Type messageType, string channel) + public static string FormatChannelQueueName(Type messageType, string channel) => FormatChannelQueueName(messageType.Name, channel); + + public static string FormatChannelQueueName(string queueName, string channel) { - var queue = !string.IsNullOrWhiteSpace(channel) ? $"{messageType.Name}{channel}" : messageType.Name; + var queue = !string.IsNullOrWhiteSpace(channel) ? $"{queueName}{channel}" : queueName; return FormatQueueName(queue); } - public static string FormatQueueName(string queue) => queue.Dehumanize().Underscore().Dasherize(); + public static string FormatQueueName(string queue) => queue.Humanize().Dehumanize().Underscore().Dasherize(); internal ElsaOptions() { @@ -65,8 +66,8 @@ namespace Elsa public IEnumerable ActivityTypes => ActivityFactory.Types; public IList WorkflowTypes { get; } = new List(); - public IList CompetingMessageTypes { get; } = new List(); - public IList PubSubMessageTypes { get; } = new List(); + public IList CompetingMessageTypes { get; } = new List(); + public IList PubSubMessageTypes { get; } = new List(); public ServiceBusOptions ServiceBusOptions { get; } = new(); public DistributedLockingOptions DistributedLockingOptions { get; set; } diff --git a/src/core/Elsa.Core/ElsaOptionsBuilder.cs b/src/core/Elsa.Core/ElsaOptionsBuilder.cs index 3712b6176..2cff35032 100644 --- a/src/core/Elsa.Core/ElsaOptionsBuilder.cs +++ b/src/core/Elsa.Core/ElsaOptionsBuilder.cs @@ -150,21 +150,21 @@ namespace Elsa return this; } - public ElsaOptionsBuilder AddCompetingMessageType(Type messageType, string? queue = default) + public ElsaOptionsBuilder AddCompetingMessageType(Type messageType, string? queueName = default) { - ElsaOptions.CompetingMessageTypes.Add(new CompetingMessageType(messageType, queue)); + ElsaOptions.CompetingMessageTypes.Add(new MessageTypeConfig(messageType, queueName)); return this; } - public ElsaOptionsBuilder AddCompetingMessageType(string? queue = default) => AddCompetingMessageType(typeof(T), queue); + public ElsaOptionsBuilder AddCompetingMessageType(string? queueName = default) => AddCompetingMessageType(typeof(T), queueName); - public ElsaOptionsBuilder AddPubSubMessageType(Type messageType) + public ElsaOptionsBuilder AddPubSubMessageType(Type messageType, string? queueName = default) { - ElsaOptions.PubSubMessageTypes.Add(messageType); + ElsaOptions.PubSubMessageTypes.Add(new MessageTypeConfig(messageType, queueName)); return this; } - public ElsaOptionsBuilder AddPubSubMessageType() => AddPubSubMessageType(typeof(T)); + public ElsaOptionsBuilder AddPubSubMessageType(string? queueName = default) => AddPubSubMessageType(typeof(T), queueName); public ElsaOptionsBuilder ConfigureDistributedLockProvider(Action configureOptions) { diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 7a99a9ac8..7ecec0eb5 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -109,14 +109,10 @@ namespace Microsoft.Extensions.DependencyInjection /// Registers a consumer and associated message type using the competing consumer pattern. With the competing consumer pattern, only the first consumer on a given node to obtain a message will handle that message. /// This is in contrast to the Publisher-Subscriber pattern, where a message will be delivered to the consumer on all nodes in a cluster. To register a Publisher-Subscriber consumer, use /// - /// - /// - /// - /// - public static ElsaOptionsBuilder AddCompetingConsumer(this ElsaOptionsBuilder elsaOptions) where TConsumer : class, IHandleMessages + public static ElsaOptionsBuilder AddCompetingConsumer(this ElsaOptionsBuilder elsaOptions, string? queueName = default) where TConsumer : class, IHandleMessages { elsaOptions.AddCompetingConsumerService(); - elsaOptions.AddCompetingMessageType(); + elsaOptions.AddCompetingMessageType(queueName); return elsaOptions; } @@ -126,10 +122,10 @@ namespace Microsoft.Extensions.DependencyInjection return elsaOptions; } - public static ElsaOptionsBuilder AddPubSubConsumer(this ElsaOptionsBuilder elsaOptions) where TConsumer : class, IHandleMessages + public static ElsaOptionsBuilder AddPubSubConsumer(this ElsaOptionsBuilder elsaOptions, string? queueName = default) where TConsumer : class, IHandleMessages { elsaOptions.Services.AddTransient, TConsumer>(); - elsaOptions.AddPubSubMessageType(); + elsaOptions.AddPubSubMessageType(queueName); return elsaOptions; } @@ -249,12 +245,12 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton(); options - .AddCompetingConsumer() - .AddCompetingConsumerService() - .AddCompetingConsumerService() - .AddPubSubConsumer() - .AddPubSubConsumer() - .AddPubSubConsumer(); + .AddCompetingConsumer("ExecuteWorkflow") + .AddCompetingConsumer("ExecuteWorkflow") + .AddCompetingConsumer("ExecuteWorkflow") + .AddPubSubConsumer("WorkflowDefinitionEvents") + .AddPubSubConsumer("WorkflowDefinitionEvents") + .AddPubSubConsumer("WorkflowDefinitionEvents"); // AutoMapper. services diff --git a/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs b/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs index bd9aa9093..58d77c007 100644 --- a/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs +++ b/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs @@ -32,8 +32,9 @@ namespace Elsa.Services.Messaging public async Task GetServiceBusAsync(Type messageType, string? queueName = default, CancellationToken cancellationToken = default) { - if (string.IsNullOrWhiteSpace(queueName)) - queueName = ElsaOptions.FormatChannelQueueName(messageType, _elsaOptions.WorkflowChannelOptions.Default); + queueName = string.IsNullOrWhiteSpace(queueName) + ? ElsaOptions.FormatChannelQueueName(messageType, _elsaOptions.WorkflowChannelOptions.Default) + : ElsaOptions.FormatQueueName(queueName); var prefixedQueueName = PrefixQueueName(queueName); await _semaphore.WaitAsync(cancellationToken); diff --git a/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs b/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs index d95bf4460..0d7957db1 100644 --- a/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs +++ b/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs @@ -10,14 +10,18 @@ namespace Elsa.StartupTasks public class CreateSubscriptions : IStartupTask { private readonly IServiceBusFactory _serviceBusFactory; + private readonly ElsaOptions _elsaOptions; private readonly IContainerNameAccessor _containerNameAccessor; - private readonly IList _competingMessageTypes; - private readonly IEnumerable _pubSubMessageTypes; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly IList _competingMessageTypes; + private readonly IEnumerable _pubSubMessageTypes; - public CreateSubscriptions(IServiceBusFactory serviceBusFactory, ElsaOptions elsaOptions, IContainerNameAccessor containerNameAccessor) + public CreateSubscriptions(IServiceBusFactory serviceBusFactory, ElsaOptions elsaOptions, IContainerNameAccessor containerNameAccessor, IDistributedLockProvider distributedLockProvider) { _serviceBusFactory = serviceBusFactory; + _elsaOptions = elsaOptions; _containerNameAccessor = containerNameAccessor; + _distributedLockProvider = distributedLockProvider; _competingMessageTypes = elsaOptions.CompetingMessageTypes.ToList(); _pubSubMessageTypes = elsaOptions.PubSubMessageTypes; @@ -27,8 +31,8 @@ namespace Elsa.StartupTasks // For each workflow channel, register a competing message type for workflow definition and workflow instance consumers. foreach (var workflowChannel in workflowChannels) { - _competingMessageTypes.Add(new CompetingMessageType(typeof(ExecuteWorkflowDefinitionRequest), ElsaOptions.FormatChannelQueueName(workflowChannel))); - _competingMessageTypes.Add(new CompetingMessageType(typeof(ExecuteWorkflowInstanceRequest), ElsaOptions.FormatChannelQueueName(workflowChannel))); + _competingMessageTypes.Add(new MessageTypeConfig(typeof(ExecuteWorkflowDefinitionRequest), ElsaOptions.FormatChannelQueueName("ExecuteWorkflow", workflowChannel))); + _competingMessageTypes.Add(new MessageTypeConfig(typeof(ExecuteWorkflowInstanceRequest), ElsaOptions.FormatChannelQueueName("ExecuteWorkflow", workflowChannel))); } } @@ -36,18 +40,23 @@ namespace Elsa.StartupTasks public async Task ExecuteAsync(CancellationToken cancellationToken) { + await using var handle = await _distributedLockProvider.AcquireLockAsync(nameof(CreateSubscriptions), _elsaOptions.DistributedLockTimeout, cancellationToken); + + if (handle == null) + throw new Exception("Could not acquire a lock within the maximum amount of time configured"); + foreach (var messageType in _competingMessageTypes) { - var bus = await _serviceBusFactory.GetServiceBusAsync(messageType.MessageType, messageType.Queue, cancellationToken); + var bus = await _serviceBusFactory.GetServiceBusAsync(messageType.MessageType, messageType.QueueName, cancellationToken); await bus.Subscribe(messageType.MessageType); } var containerName = _containerNameAccessor.GetContainerName(); foreach (var messageType in _pubSubMessageTypes) { - var queueName = $"{containerName}:{messageType.Name}"; - var bus = await _serviceBusFactory.GetServiceBusAsync(messageType, queueName, cancellationToken); - await bus.Subscribe(messageType); + var queueName = $"{containerName}:{messageType.QueueName ?? messageType.MessageType.Name}"; + var bus = await _serviceBusFactory.GetServiceBusAsync(messageType.MessageType, queueName, cancellationToken); + await bus.Subscribe(messageType.MessageType); } } }