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