elsa-core/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs

81 lines
2.9 KiB
C#
Raw Normal View History

2024-01-25 13:51:48 +00:00
using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
2024-01-25 13:51:48 +00:00
using Elsa.MassTransit.Consumers;
using Elsa.MassTransit.Features;
2024-01-25 13:51:48 +00:00
using Elsa.MassTransit.Options;
using Elsa.Workflows.Contracts;
using MassTransit;
using Microsoft.Extensions.DependencyInjection;
2024-01-25 13:51:48 +00:00
using Microsoft.Extensions.Options;
namespace Elsa.MassTransit.RabbitMq.Features;
/// <summary>
/// Configures MassTransit to use the RabbitMQ transport.
/// </summary>
[DependsOn(typeof(MassTransitFeature))]
public class RabbitMqServiceBusFeature : FeatureBase
{
/// <inheritdoc />
public RabbitMqServiceBusFeature(IModule module) : base(module)
{
}
/// A RabbitMQ connection string.
public string? ConnectionString { get; set; }
/// Configures the RabbitMQ transport options.
public Action<RabbitMqTransportOptions>? TransportOptions { get; set; }
/// <summary>
/// Configures the RabbitMQ bus.
/// </summary>
public Action<IRabbitMqBusFactoryConfigurator>? ConfigureServiceBus { get; set; }
/// <inheritdoc />
public override void Configure()
{
2024-01-25 13:51:48 +00:00
Module.Configure<MassTransitFeature>(massTransitFeature =>
{
2024-01-25 13:51:48 +00:00
massTransitFeature.BusConfigurator = configure =>
{
2024-01-25 13:51:48 +00:00
var tempConsumers = massTransitFeature.GetConsumers()
.Where(c => c.IsTemporary)
.ToList();
configure.AddConsumers(tempConsumers.Select(c => c.ConsumerType).ToArray());
configure.UsingRabbitMq((context, serviceBus) =>
{
var options = context.GetRequiredService<IOptions<MassTransitWorkflowDispatcherOptions>>().Value;
var instanceNameRetriever = context.GetRequiredService<IInstanceNameRetriever>();
if (!string.IsNullOrEmpty(ConnectionString))
serviceBus.Host(ConnectionString);
2024-01-25 13:51:48 +00:00
ConfigureServiceBus?.Invoke(serviceBus);
foreach (var consumer in tempConsumers)
{
serviceBus.ReceiveEndpoint($"{instanceNameRetriever.GetName()}-{consumer.Name}", configurator =>
{
configurator.QueueExpiration = options.ShortTermQueueLifetime ?? TimeSpan.FromHours(1);
configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit;
configurator.ConfigureConsumer<DispatchCancelWorkflowsRequestConsumer>(context);
});
}
2024-01-25 13:51:48 +00:00
serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
});
};
});
}
/// <inheritdoc />
public override void Apply()
{
if (TransportOptions != null) Services.Configure(TransportOptions);
}
}