Added handler for deleting orphaned subscriptions

This commit is contained in:
Raymond den Haan 2024-02-12 15:32:48 +01:00
parent 25305a9e37
commit 0eb274e276
7 changed files with 130 additions and 1 deletions

View file

@ -18,6 +18,9 @@
<ProjectReference Include="..\..\modules\Elsa.Identity\Elsa.Identity.csproj"/>
<ProjectReference Include="..\..\modules\Elsa.Liquid\Elsa.Liquid.csproj"/>
<ProjectReference Include="..\..\modules\Elsa.EntityFrameworkCore\Elsa.EntityFrameworkCore.csproj"/>
<ProjectReference Include="..\..\modules\Elsa.MassTransit.AzureServiceBus\Elsa.MassTransit.AzureServiceBus.csproj" />
<ProjectReference Include="..\..\modules\Elsa.MassTransit.RabbitMq\Elsa.MassTransit.RabbitMq.csproj" />
<ProjectReference Include="..\..\modules\Elsa.MassTransit\Elsa.MassTransit.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Python\Elsa.Python.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Quartz\Elsa.Quartz.csproj"/>
<ProjectReference Include="..\..\modules\Elsa.Webhooks\Elsa.Webhooks.csproj"/>

View file

@ -12,6 +12,7 @@
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Hosting.Management\Elsa.Hosting.Management.csproj" />
<ProjectReference Include="..\Elsa.MassTransit\Elsa.MassTransit.csproj" />
</ItemGroup>

View file

@ -1,12 +1,18 @@
using Azure.Messaging.ServiceBus.Administration;
using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.MassTransit.Consumers;
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 Elsa.Workflows.Contracts;
using MassTransit;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
@ -29,6 +35,11 @@ public class AzureServiceBusFeature : FeatureBase
/// A delegate that configures the Azure Service Bus transport options.
/// </summary>
public Action<IServiceBusBusFactoryConfigurator>? ConfigureServiceBus { get; set; }
/// <summary>
/// A delegate to create a <see cref="ServiceBusAdministrationClient"/> instance.
/// </summary>
public Func<IServiceProvider, ServiceBusAdministrationClient> ServiceBusAdministrationClientFactory { get; set; } = sp => new(GetConnectionString(sp));
/// <inheritdoc />
public override void Configure()
@ -41,6 +52,7 @@ public class AzureServiceBusFeature : FeatureBase
var shortLivedConsumers = consumers
.Where(c => c.IsShortLived)
.ToList();
RegisterConsumers(consumers);
configure.AddServiceBusMessageScheduler();
configure.AddConsumers(shortLivedConsumers.Select(c => c.ConsumerType).ToArray());
@ -70,4 +82,44 @@ public class AzureServiceBusFeature : FeatureBase
};
});
}
private void RegisterConsumers(List<ConsumerTypeDefinition> consumers)
{
var subscriptionTopology = new List<MessageSubscriptionTopology>();
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];
var topicName = $"{genericType.Namespace.ToLower()}~{genericType.Name.ToLower()}";
subscriptionTopology.Add(new MessageSubscriptionTopology(topicName,
consumer.Name ?? genericType.Name.ToLower(),
consumer.IsShortLived));
}
}
Services.AddSingleton(new MessageTopologyProvider(subscriptionTopology));
}
/// <inheritdoc />
public override void Apply()
{
Services.AddScoped(ServiceBusAdministrationClientFactory);
Services.AddNotificationHandler<OrphanedSubscriptionRemover>();
}
private static string GetConnectionString(IServiceProvider serviceProvider)
{
var options = serviceProvider.GetRequiredService<IOptions<AzureServiceBusOptions>>().Value;
var configuration = serviceProvider.GetRequiredService<IConfiguration>();
return configuration.GetConnectionString(options.ConnectionStringOrName) ?? options.ConnectionStringOrName;
}
}

View file

@ -0,0 +1,28 @@
using Azure.Messaging.ServiceBus.Administration;
using Elsa.Hosting.Management.Notifications;
using Elsa.MassTransit.AzureServiceBus.Services;
using Elsa.Mediator.Contracts;
namespace Elsa.MassTransit.AzureServiceBus.Handlers;
/// <summary>
/// Class responsible for removing orphaned subscriptions.
/// </summary>
public class OrphanedSubscriptionRemover(MessageTopologyProvider topologyProvider, ServiceBusAdministrationClient client)
: INotificationHandler<InstanceDeactivated>
{
/// <summary>
/// Removes orphaned subscriptions from Azure Service Bus.
/// </summary>
public async Task HandleAsync(InstanceDeactivated notification, CancellationToken cancellationToken)
{
var subscriptions = topologyProvider.GetShortLivedSubscriptions().ToList();
foreach (var subscription in subscriptions)
{
await client.DeleteSubscriptionAsync(subscription.TopicName,
$"Elsa-{notification.InstanceName}-{subscription.SubscriptionName}",
cancellationToken);
}
}
}

View file

@ -0,0 +1,6 @@
namespace Elsa.MassTransit.AzureServiceBus.Models;
/// <summary>
/// Represents the topology of a message subscription in Azure Service Bus.
/// </summary>
public record MessageSubscriptionTopology(string TopicName, string SubscriptionName, bool IsShortLived);

View file

@ -0,0 +1,12 @@
namespace Elsa.MassTransit.AzureServiceBus.Options;
/// <summary>
/// A collection of settings to configure integration with Azure Service Bus.
/// </summary>
public class AzureServiceBusOptions
{
/// <summary>
/// Th connection string or connection string name to connect with the service bus.
/// </summary>
public string ConnectionStringOrName { get; set; } = default!;
}

View file

@ -0,0 +1,27 @@
using Elsa.MassTransit.AzureServiceBus.Models;
namespace Elsa.MassTransit.AzureServiceBus.Services;
/// <summary>
/// Provides message topology information.
/// </summary>
public class MessageTopologyProvider
{
private readonly IEnumerable<MessageSubscriptionTopology> _subscriptionTopology;
/// <summary>
/// Provides message topology information.
/// </summary>
public MessageTopologyProvider(IEnumerable<MessageSubscriptionTopology> subscriptionTopology)
{
_subscriptionTopology = subscriptionTopology;
}
/// <summary>
/// Retrieves all the short-lived message subscriptions from the subscription topology.
/// </summary>
public IEnumerable<MessageSubscriptionTopology> GetShortLivedSubscriptions()
{
return _subscriptionTopology.Where(x => x.IsShortLived);
}
}