From 0eb274e2768eda08e565c9b53798c9bcb12e57d8 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Mon, 12 Feb 2024 15:32:48 +0100 Subject: [PATCH] Added handler for deleting orphaned subscriptions --- .../Elsa.ServerAndStudio.Web.csproj | 3 ++ .../Elsa.MassTransit.AzureServiceBus.csproj | 1 + .../Features/AzureServiceBusFeature.cs | 54 ++++++++++++++++++- .../Handlers/OrphanedSubscriptionRemover.cs | 28 ++++++++++ .../Models/MessageSubscriptionTopology.cs | 6 +++ .../Options/AzureServiceBusOptions.cs | 12 +++++ .../Services/MessageTopologyProvider.cs | 27 ++++++++++ 7 files changed, 130 insertions(+), 1 deletion(-) create mode 100644 src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs create mode 100644 src/modules/Elsa.MassTransit.AzureServiceBus/Models/MessageSubscriptionTopology.cs create mode 100644 src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs create mode 100644 src/modules/Elsa.MassTransit.AzureServiceBus/Services/MessageTopologyProvider.cs diff --git a/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj b/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj index 97e625599..7a7024a6b 100644 --- a/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj +++ b/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj @@ -18,6 +18,9 @@ + + + diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj b/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj index acd31745d..cf92d65df 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj @@ -12,6 +12,7 @@ + diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index a1bb9f706..95dad87d9 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -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. /// public Action? ConfigureServiceBus { get; set; } + + /// + /// A delegate to create a instance. + /// + public Func ServiceBusAdministrationClientFactory { get; set; } = sp => new(GetConnectionString(sp)); /// 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 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]; + 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)); + } + + /// + public override void Apply() + { + 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; + } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs new file mode 100644 index 000000000..cbc2f3c27 --- /dev/null +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs @@ -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; + +/// +/// Class responsible for removing orphaned subscriptions. +/// +public class OrphanedSubscriptionRemover(MessageTopologyProvider topologyProvider, ServiceBusAdministrationClient client) + : INotificationHandler +{ + /// + /// Removes orphaned subscriptions from Azure Service Bus. + /// + 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); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Models/MessageSubscriptionTopology.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Models/MessageSubscriptionTopology.cs new file mode 100644 index 000000000..e0399a48b --- /dev/null +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Models/MessageSubscriptionTopology.cs @@ -0,0 +1,6 @@ +namespace Elsa.MassTransit.AzureServiceBus.Models; + +/// +/// Represents the topology of a message subscription in Azure Service Bus. +/// +public record MessageSubscriptionTopology(string TopicName, string SubscriptionName, bool IsShortLived); \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs new file mode 100644 index 000000000..08f902bef --- /dev/null +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs @@ -0,0 +1,12 @@ +namespace Elsa.MassTransit.AzureServiceBus.Options; + +/// +/// A collection of settings to configure integration with Azure Service Bus. +/// +public class AzureServiceBusOptions +{ + /// + /// Th connection string or connection string name to connect with the service bus. + /// + public string ConnectionStringOrName { get; set; } = default!; +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Services/MessageTopologyProvider.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Services/MessageTopologyProvider.cs new file mode 100644 index 000000000..34276ae20 --- /dev/null +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Services/MessageTopologyProvider.cs @@ -0,0 +1,27 @@ +using Elsa.MassTransit.AzureServiceBus.Models; + +namespace Elsa.MassTransit.AzureServiceBus.Services; + +/// +/// Provides message topology information. +/// +public class MessageTopologyProvider +{ + private readonly IEnumerable _subscriptionTopology; + + /// + /// Provides message topology information. + /// + public MessageTopologyProvider(IEnumerable subscriptionTopology) + { + _subscriptionTopology = subscriptionTopology; + } + + /// + /// Retrieves all the short-lived message subscriptions from the subscription topology. + /// + public IEnumerable GetShortLivedSubscriptions() + { + return _subscriptionTopology.Where(x => x.IsShortLived); + } +} \ No newline at end of file