From 29742db513d20761e630fb17e656bae2c1bcc840 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 10 Jun 2024 10:45:12 +0200 Subject: [PATCH] Add application roles and configure them in MassTransit (#5561) A new enum ApplicationRole has been added for distinguishing among different roles (Hybrid, Api, Worker) an application can take. In the configuration of MassTransit, it is now possible to disable the consumers based on application role, which can help optimize the usage of resources and increase application efficiency. --- .../Elsa.Server.Web/Enums/ApplicationRole.cs | 8 +++++ src/bundles/Elsa.Server.Web/Program.cs | 6 +++- src/bundles/Elsa.Server.Web/appsettings.json | 1 + .../Features/AzureServiceBusFeature.cs | 27 +++++++------- .../Features/RabbitMqServiceBusFeature.cs | 35 ++++++++++--------- .../Features/MassTransitFeature.cs | 22 +++++++----- 6 files changed, 62 insertions(+), 37 deletions(-) create mode 100644 src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs diff --git a/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs new file mode 100644 index 000000000..60cfbe097 --- /dev/null +++ b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs @@ -0,0 +1,8 @@ +namespace Elsa.Server.Web; + +public enum ApplicationRole +{ + Hybrid, + Api, + Worker +} \ No newline at end of file diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index f55311e24..87df5a0f3 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -49,7 +49,7 @@ const bool useMemoryStores = false; const bool useCaching = true; const bool useReadOnlyMode = false; const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit; -const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory; +const MassTransitBroker useMassTransitBroker = MassTransitBroker.RabbitMq; var builder = WebApplication.CreateBuilder(args); var services = builder.Services; @@ -65,6 +65,7 @@ var azureServiceBusConnectionString = configuration.GetConnectionString("AzureSe var rabbitMqConnectionString = configuration.GetConnectionString("RabbitMq")!; var redisConnectionString = configuration.GetConnectionString("Redis")!; var distributedLockProviderName = configuration.GetSection("Runtime")["DistributedLockProvider"]; +var appRole = Enum.Parse(configuration["AppRole"]); // Add Elsa services. services @@ -325,6 +326,9 @@ services { elsa.UseMassTransit(massTransit => { + if (appRole == ApplicationRole.Api) + massTransit.DisableConsumers = true; + if (useMassTransitBroker == MassTransitBroker.AzureServiceBus) { massTransit.UseAzureServiceBus(azureServiceBusConnectionString, serviceBusFeature => serviceBusFeature.ConfigureServiceBus = bus => diff --git a/src/bundles/Elsa.Server.Web/appsettings.json b/src/bundles/Elsa.Server.Web/appsettings.json index 595ffd131..247a862e5 100644 --- a/src/bundles/Elsa.Server.Web/appsettings.json +++ b/src/bundles/Elsa.Server.Web/appsettings.json @@ -96,6 +96,7 @@ } ] }, + "AppRole": "Hybrid", "Runtime": { "WorkflowInboxCleanup": { "SweepInterval": "00:00:10:00", diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index 1c3f4ab3e..5e70a9713 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -85,19 +85,22 @@ public class AzureServiceBusFeature : FeatureBase configurator.SetupWorkflowDispatcherEndpoints(context); ConfigureServiceBus?.Invoke(configurator); var instanceNameProvider = context.GetRequiredService(); - - foreach (var consumer in temporaryConsumers) - { - var queueName = $"{consumer.Name}-{instanceNameProvider.GetName()}"; - configurator.ReceiveEndpoint(queueName, endpointConfigurator => - { - endpointConfigurator.AutoDeleteOnIdle = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); - endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); - }); - } - configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + if (!massTransitFeature.DisableConsumers) + { + foreach (var consumer in temporaryConsumers) + { + var queueName = $"{consumer.Name}-{instanceNameProvider.GetName()}"; + configurator.ReceiveEndpoint(queueName, endpointConfigurator => + { + endpointConfigurator.AutoDeleteOnIdle = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); + endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); + }); + } + + configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + } }); }; }); diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs index 8fb594585..626c06372 100644 --- a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs @@ -68,24 +68,27 @@ public class RabbitMqServiceBusFeature : FeatureBase ConfigureServiceBus?.Invoke(configurator); - foreach (var consumer in temporaryConsumers) + if (!massTransitFeature.DisableConsumers) { - configurator.ReceiveEndpoint($"{instanceNameProvider.GetName()}-{consumer.Name}", - endpointConfigurator => - { - endpointConfigurator.QueueExpiration = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); - endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - endpointConfigurator.Durable = false; - endpointConfigurator.AutoDelete = true; - endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); - }); - } + foreach (var consumer in temporaryConsumers) + { + configurator.ReceiveEndpoint($"{instanceNameProvider.GetName()}-{consumer.Name}", + endpointConfigurator => + { + endpointConfigurator.QueueExpiration = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); + endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + endpointConfigurator.Durable = false; + endpointConfigurator.AutoDelete = true; + endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); + }); + } - // Only configure the dispatcher endpoints if the Masstransit Workflow Dispatcher feature is enabled. - if (Module.HasFeature()) - configurator.SetupWorkflowDispatcherEndpoints(context); - - configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + // Only configure the dispatcher endpoints if the Masstransit Workflow Dispatcher feature is enabled. + if (Module.HasFeature()) + configurator.SetupWorkflowDispatcherEndpoints(context); + + configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + } }); }; }); diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs index d2e33bda0..20412485a 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs @@ -35,6 +35,8 @@ public class MassTransitFeature : FeatureBase /// The number of messages to prefetch. [Obsolete("PrefetchCount has been moved to be included in MassTransitOptions")] public int? PrefetchCount { get; set; } + + public bool DisableConsumers { get; set; } /// /// A delegate that can be set to configure MassTransit's . Used by transport-level features such as AzureServiceBusFeature and RabbitMqServiceBusFeature. @@ -130,17 +132,21 @@ public class MassTransitFeature : FeatureBase { var options = context.GetRequiredService>().Value; - foreach (var consumer in temporaryConsumers) + if(!DisableConsumers) { - busFactoryConfigurator.ReceiveEndpoint(consumer.Name!, endpoint => + foreach (var consumer in temporaryConsumers) { - endpoint.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - endpoint.ConfigureConsumer(context, consumer.ConsumerType); - }); + busFactoryConfigurator.ReceiveEndpoint(consumer.Name!, endpoint => + { + endpoint.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + endpoint.ConfigureConsumer(context, consumer.ConsumerType); + }); + } + + busFactoryConfigurator.SetupWorkflowDispatcherEndpoints(context); + busFactoryConfigurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); } - - busFactoryConfigurator.SetupWorkflowDispatcherEndpoints(context); - busFactoryConfigurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + busFactoryConfigurator.ConfigureJsonSerializerOptions(serializerOptions => { var serializer = context.GetRequiredService();