using Azure.Messaging.ServiceBus.Administration;
using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.Hosting.Management.Contracts;
using Elsa.MassTransit.AzureServiceBus.Handlers;
using Elsa.MassTransit.AzureServiceBus.Models;
using Elsa.MassTransit.AzureServiceBus.Options;
using Elsa.MassTransit.AzureServiceBus.Services;
using Elsa.MassTransit.Features;
using Elsa.MassTransit.Models;
using Elsa.MassTransit.Options;
using MassTransit;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
namespace Elsa.MassTransit.AzureServiceBus.Features;
/// Configures MassTransit to use the Azure Service Bus transport.
/// See https://masstransit.io/documentation/configuration/transports/azure-service-bus
[DependsOn(typeof(MassTransitFeature))]
public class AzureServiceBusFeature : FeatureBase
{
///
public AzureServiceBusFeature(IModule module) : base(module)
{
}
/// An Azure Service Bus connection string.
public string? ConnectionString { get; set; }
///
/// A delegate that configures the Azure Service Bus transport options.
///
public Action? ConfigureServiceBus { get; set; }
///
/// A delegate to configure .
///
public Action AzureServiceBusOptions { get; set; } = _ => { };
///
/// A delegate to create a instance.
///
public Func ServiceBusAdministrationClientFactory { get; set; } = sp => new(GetConnectionString(sp));
///
public override void Configure()
{
Module.Configure(massTransitFeature =>
{
massTransitFeature.BusConfigurator = configure =>
{
var consumers = massTransitFeature.GetConsumers().ToList();
var temporaryConsumers = consumers
.Where(c => c.IsTemporary)
.ToList();
RegisterConsumers(consumers);
configure.AddServiceBusMessageScheduler();
configure.AddConsumers(temporaryConsumers.Select(c => c.ConsumerType).ToArray());
configure.UsingAzureServiceBus((context, serviceBus) =>
{
var options = context.GetRequiredService>().Value;
var instanceNameProvider = context.GetRequiredService();
if (ConnectionString != null)
serviceBus.Host(ConnectionString);
serviceBus.UseServiceBusMessageScheduler();
ConfigureServiceBus?.Invoke(serviceBus);
foreach (var consumer in temporaryConsumers)
{
serviceBus.ReceiveEndpoint($"Elsa-{instanceNameProvider.GetName()}-{consumer.Name}", configurator =>
{
configurator.AutoDeleteOnIdle = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1);
configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit;
configurator.ConfigureConsumer(context, consumer.ConsumerType);
});
}
serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
});
};
});
}
private void RegisterConsumers(List consumers)
{
var subscriptionTopology = new List();
foreach (var consumer in consumers)
{
foreach (var consumerInterface in consumer!.ConsumerType.GetInterfaces())
{
if (!consumerInterface.IsGenericType ||
consumerInterface.GetGenericTypeDefinition() != typeof(IConsumer<>))
{
continue;
}
var genericType = consumerInterface.GetGenericArguments()[0];
//While the name might show up in the Azure portal with a ~, the separator is actually an /
//see https://learn.microsoft.com/en-us/archive/blogs/servicebus/azure-service-bus-azure-resource-manager-and-this-character
var topicName = $"{genericType.Namespace.ToLower()}/{genericType.Name.ToLower()}";
subscriptionTopology.Add(new MessageSubscriptionTopology(topicName,
consumer.Name ?? genericType.Name.ToLower(),
consumer.IsTemporary));
}
}
Services.AddSingleton(new MessageTopologyProvider(subscriptionTopology));
}
///
public override void Apply()
{
Services.Configure(AzureServiceBusOptions);
Services.AddScoped(ServiceBusAdministrationClientFactory);
Services.AddNotificationHandler();
}
private static string GetConnectionString(IServiceProvider serviceProvider)
{
var options = serviceProvider.GetRequiredService>().Value;
var configuration = serviceProvider.GetRequiredService();
return configuration.GetConnectionString(options.ConnectionStringOrName) ?? options.ConnectionStringOrName;
}
}