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.
This commit is contained in:
Sipke Schoorstra 2024-06-10 10:45:12 +02:00 committed by GitHub
parent 5ebf4302d3
commit 29742db513
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 62 additions and 37 deletions

View file

@ -0,0 +1,8 @@
namespace Elsa.Server.Web;
public enum ApplicationRole
{
Hybrid,
Api,
Worker
}

View file

@ -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<ApplicationRole>(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 =>

View file

@ -96,6 +96,7 @@
}
]
},
"AppRole": "Hybrid",
"Runtime": {
"WorkflowInboxCleanup": {
"SweepInterval": "00:00:10:00",

View file

@ -85,19 +85,22 @@ public class AzureServiceBusFeature : FeatureBase
configurator.SetupWorkflowDispatcherEndpoints(context);
ConfigureServiceBus?.Invoke(configurator);
var instanceNameProvider = context.GetRequiredService<IApplicationInstanceNameProvider>();
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));
}
});
};
});

View file

@ -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<MassTransitWorkflowDispatcherFeature>())
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<MassTransitWorkflowDispatcherFeature>())
configurator.SetupWorkflowDispatcherEndpoints(context);
configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
}
});
};
});

View file

@ -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; }
/// <summary>
/// A delegate that can be set to configure MassTransit's <see cref="IBusRegistrationConfigurator"/>. Used by transport-level features such as AzureServiceBusFeature and RabbitMqServiceBusFeature.
@ -130,17 +132,21 @@ public class MassTransitFeature : FeatureBase
{
var options = context.GetRequiredService<IOptions<MassTransitWorkflowDispatcherOptions>>().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<IJsonSerializer>();