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();