diff --git a/src/bundles/Elsa.ServerAndStudio.Web/Program.cs b/src/bundles/Elsa.ServerAndStudio.Web/Program.cs index d4ee31093..7f064cf4b 100644 --- a/src/bundles/Elsa.ServerAndStudio.Web/Program.cs +++ b/src/bundles/Elsa.ServerAndStudio.Web/Program.cs @@ -24,10 +24,12 @@ var azureServiceBusConnectionString = configuration.GetConnectionString("AzureSe var identitySection = configuration.GetSection("Identity"); var identityTokenSection = identitySection.GetSection("Tokens"); var massTransitSection = configuration.GetSection("MassTransit"); +var massTransitDispatcherSection = configuration.GetSection("MassTransit.Dispatcher"); var heartbeatSection = configuration.GetSection("Heartbeat"); -const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory; +const MassTransitBroker useMassTransitBroker = MassTransitBroker.AzureServiceBus; -services.Configure(massTransitSection); +services.Configure(massTransitSection); +services.Configure(massTransitDispatcherSection); // Add Elsa services. services @@ -87,28 +89,14 @@ services { elsa.UseMassTransit(massTransit => { - if (useMassTransitBroker == MassTransitBroker.AzureServiceBus) + switch (useMassTransitBroker) { - massTransit.UseAzureServiceBus(azureServiceBusConnectionString, serviceBusFeature => serviceBusFeature.ConfigureServiceBus = bus => - { - bus.PrefetchCount = 4; - bus.LockDuration = TimeSpan.FromMinutes(5); - bus.MaxConcurrentCalls = 32; - bus.MaxDeliveryCount = 8; - // etc. - }); - } - - if (useMassTransitBroker == MassTransitBroker.RabbitMq) - { - massTransit.UseRabbitMq(rabbitMqConnectionString, rabbit => rabbit.ConfigureServiceBus = bus => - { - bus.PrefetchCount = 4; - bus.Durable = true; - bus.AutoDelete = false; - bus.ConcurrentMessageLimit = 32; - // etc. - }); + case MassTransitBroker.AzureServiceBus: + massTransit.UseAzureServiceBus(azureServiceBusConnectionString); + break; + case MassTransitBroker.RabbitMq: + massTransit.UseRabbitMq(rabbitMqConnectionString); + break; } } ); diff --git a/src/bundles/Elsa.ServerAndStudio.Web/appsettings.json b/src/bundles/Elsa.ServerAndStudio.Web/appsettings.json index 7a6609301..648b878a7 100644 --- a/src/bundles/Elsa.ServerAndStudio.Web/appsettings.json +++ b/src/bundles/Elsa.ServerAndStudio.Web/appsettings.json @@ -66,7 +66,13 @@ "Timeout": "00:00:05:00" }, "MassTransit": { - "TemporaryQueueTtl": "00:00:05:00" + "TemporaryQueueTtl": "00:00:05:00", + "ConcurrentMessageLimit": 4, + "PrefetchCount": 4, + "MaxAutoRenewDuration": "00:00:00:30", + "Dispatcher": { + "ConcurrentMessageLimit": 4 + } }, "Smtp": { "Host": "localhost", diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index 341fa9713..ecb9900e3 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -73,10 +73,17 @@ public class AzureServiceBusFeature : FeatureBase if (ConnectionString != null) configurator.Host(ConnectionString); + var options = context.GetRequiredService>().Value; + + if (options.PrefetchCount is not null) + configurator.PrefetchCount = options.PrefetchCount.Value; + if (options.MaxAutoRenewDuration is not null) + configurator.MaxAutoRenewDuration = options.MaxAutoRenewDuration.Value; + configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + configurator.UseServiceBusMessageScheduler(); configurator.SetupWorkflowDispatcherEndpoints(context); ConfigureServiceBus?.Invoke(configurator); - var options = context.GetRequiredService>().Value; var instanceNameProvider = context.GetRequiredService(); foreach (var consumer in temporaryConsumers) diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs index dc8348d8a..9e129519b 100644 --- a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs @@ -56,24 +56,27 @@ public class RabbitMqServiceBusFeature : FeatureBase configure.UsingRabbitMq((context, configurator) => { - var options = context.GetRequiredService>().Value; + var options = context.GetRequiredService>().Value; var instanceNameProvider = context.GetRequiredService(); if (!string.IsNullOrEmpty(ConnectionString)) configurator.Host(ConnectionString); + if (options.PrefetchCount is not null) + configurator.PrefetchCount = options.PrefetchCount.Value; + configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + ConfigureServiceBus?.Invoke(configurator); foreach (var consumer in temporaryConsumers) { - configure.AddConsumer(consumer.ConsumerType).ExcludeFromConfigureEndpoints(); - 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); }); } diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs index 8648f0153..bee489e9a 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs @@ -33,6 +33,7 @@ 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; } /// @@ -56,6 +57,7 @@ public class MassTransitFeature : FeatureBase var messageTypes = this.GetMessages(); Services.AddSingleton(ChannelQueueFormatterFactory); + Services.Configure(x => x.PrefetchCount ??= PrefetchCount); Services.Configure(x => { }); Services.AddActivityProvider(); _runInMemory = BusConfigurator is null; diff --git a/src/modules/Elsa.MassTransit/Options/MassTransitOptions.cs b/src/modules/Elsa.MassTransit/Options/MassTransitOptions.cs new file mode 100644 index 000000000..695858658 --- /dev/null +++ b/src/modules/Elsa.MassTransit/Options/MassTransitOptions.cs @@ -0,0 +1,17 @@ +namespace Elsa.MassTransit.Options; + +/// +/// Represents the options for configuring MassTransit. +/// +public class MassTransitOptions +{ + /// The TTL of queues that are seen as temporary (typically queues that are created per running instance). + public TimeSpan? TemporaryQueueTtl { get; set; } + /// The number of concurrent messages to process. + public int? ConcurrentMessageLimit { get; set; } + /// The number of messages to fetch from the bus with each request. + public int? PrefetchCount { get; set; } + /// The maximum duration for auto-renewal of a resource. + /// Only relevant when using Azure Service Bus. + public TimeSpan? MaxAutoRenewDuration { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs b/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs index a872a853c..b848b2382 100644 --- a/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs +++ b/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs @@ -5,8 +5,6 @@ namespace Elsa.MassTransit.Options; /// Provides options to the public class MassTransitWorkflowDispatcherOptions { - /// The TTL of queues that are seen as temporary (typically queues that are created per running instance). - public TimeSpan? TemporaryQueueTtl { get; set; } /// The number of concurrent messages to process. public int? ConcurrentMessageLimit { get; set; } } \ No newline at end of file