diff --git a/src/bundles/Elsa.ServerAndStudio.Web/Program.cs b/src/bundles/Elsa.ServerAndStudio.Web/Program.cs
index 68d57ce17..a6ca02528 100644
--- a/src/bundles/Elsa.ServerAndStudio.Web/Program.cs
+++ b/src/bundles/Elsa.ServerAndStudio.Web/Program.cs
@@ -120,7 +120,11 @@ services
switch (useMassTransitBroker)
{
case MassTransitBroker.AzureServiceBus:
- massTransit.UseAzureServiceBus(azureServiceBusConnectionString);
+ massTransit.UseAzureServiceBus(azureServiceBusConnectionString, asb =>
+ {
+ asb.SubscriptionCleanupOptions = options => options.Interval = TimeSpan.FromMinutes(5);
+ asb.EnableAutomatedSubscriptionCleanup = true;
+ });
break;
case MassTransitBroker.RabbitMq:
massTransit.UseRabbitMq(rabbitMqConnectionString);
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Commands/CleanupSubscriptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Commands/CleanupSubscriptions.cs
new file mode 100644
index 000000000..398ebcd49
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Commands/CleanupSubscriptions.cs
@@ -0,0 +1,8 @@
+using Elsa.Mediator.Contracts;
+
+namespace Elsa.MassTransit.AzureServiceBus.Commands;
+
+///
+/// Notification to clean up the Azure Service Bus subscriptions.
+///
+public record CleanupSubscriptions : ICommand;
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs
index 7b2073d47..0fc93e992 100644
--- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs
@@ -6,6 +6,7 @@ using Elsa.Features.Services;
using Elsa.Hosting.Management.Contracts;
using Elsa.Hosting.Management.Features;
using Elsa.MassTransit.AzureServiceBus.Handlers;
+using Elsa.MassTransit.AzureServiceBus.HostedServices;
using Elsa.MassTransit.AzureServiceBus.Models;
using Elsa.MassTransit.AzureServiceBus.Options;
using Elsa.MassTransit.AzureServiceBus.Services;
@@ -34,6 +35,16 @@ public class AzureServiceBusFeature : FeatureBase
/// An Azure Service Bus connection string.
public string? ConnectionString { get; set; }
+ ///
+ /// Delete subscriptions where their connected queues could not be found.
+ /// Topics without any subscriptions will be deleted as well.
+ ///
+ ///
+ /// - All subscriptions will be cleaned up, including those that were not created by Elsa.
+ /// - Queues in other namespaces will not be found and the subscription will therefore be removed.
+ ///
+ public bool EnableAutomatedSubscriptionCleanup { get; set; }
+
///
/// A delegate that configures the Azure Service Bus transport options.
///
@@ -44,6 +55,11 @@ public class AzureServiceBusFeature : FeatureBase
///
public Action AzureServiceBusOptions { get; set; } = _ => { };
+ ///
+ /// A delegate to configure .
+ ///
+ public Action SubscriptionCleanupOptions { get; set; } = _ => { };
+
///
/// A delegate to create a instance.
///
@@ -87,7 +103,9 @@ public class AzureServiceBusFeature : FeatureBase
foreach (var consumer in temporaryConsumers)
{
- var queueName = $"{consumer.Name}-{instanceNameProvider.GetName()}";
+ // Start with the instance name since we will be using that to delete queues / subscriptions that are no longer needed.
+ // This is the only way to guarantee we can match subscriptions to an application instance, since the Azure Service Bus transport for MassTransit trims names that are too large.
+ var queueName = $"{instanceNameProvider.GetName()}-{consumer.Name}";
configurator.ReceiveEndpoint(queueName, endpointConfigurator =>
{
endpointConfigurator.AutoDeleteOnIdle = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1);
@@ -112,9 +130,19 @@ public class AzureServiceBusFeature : FeatureBase
public override void Apply()
{
Services.Configure(AzureServiceBusOptions);
+ Services.Configure(SubscriptionCleanupOptions);
Services.AddSingleton(ServiceBusAdministrationClientFactory);
}
+ ///
+ public override void ConfigureHostedServices()
+ {
+ if (EnableAutomatedSubscriptionCleanup)
+ {
+ Module.ConfigureHostedService();
+ }
+ }
+
private static string GetConnectionString(IServiceProvider serviceProvider)
{
var options = serviceProvider.GetRequiredService>().Value;
@@ -135,5 +163,6 @@ public class AzureServiceBusFeature : FeatureBase
Services.AddSingleton(new MessageTopologyProvider(subscriptionTopology));
Services.AddNotificationHandler();
+ Services.AddCommandHandler();
}
-}
\ No newline at end of file
+}
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/CleanupSubscriptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/CleanupSubscriptions.cs
new file mode 100644
index 000000000..144d963fe
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/CleanupSubscriptions.cs
@@ -0,0 +1,59 @@
+using Azure.Messaging.ServiceBus.Administration;
+using Elsa.MassTransit.AzureServiceBus.Commands;
+using Elsa.Mediator.Contracts;
+using Elsa.Mediator.Models;
+using JetBrains.Annotations;
+using Microsoft.Extensions.Logging;
+
+namespace Elsa.MassTransit.AzureServiceBus.Handlers;
+
+///
+/// Cleans up Azure Service Bus resources for which no children can be found.
+///
+///
+/// This will clean up topics without subscriptions and subscription of which the queue can not be found.
+/// - This handler is unable to differentiate between topologies created by Elsa and others. It will clean up all topology in the namespace.
+/// - Queues in other namespaces will not be found by this handler and the subscription will therefore be removed.
+///
+[UsedImplicitly]
+public class CleanupOrphanedTopology(ServiceBusAdministrationClient client, ILogger logger) : ICommandHandler
+{
+ ///
+ public async Task HandleAsync(CleanupSubscriptions notification, CancellationToken cancellationToken)
+ {
+ var queues = await client.GetQueuesAsync().ToListAsync(cancellationToken);
+ var queueNames = queues.Select(q => q.Name).ToList();
+
+ await foreach (var topic in client.GetTopicsAsync(cancellationToken))
+ {
+ // Only clean up any topics that have the elsa prefix.
+ if (!topic.Name.StartsWith("elsa"))
+ continue;
+
+ var subscriptions = await client.GetSubscriptionsAsync(topic.Name, cancellationToken).ToListAsync(cancellationToken);
+
+ // Delete topics which have no active subscriptions.
+ if (subscriptions.Count == 0)
+ {
+ logger.LogWarning("Deleting topic {name}", topic.Name);
+ await client.DeleteTopicAsync(topic.Name, cancellationToken);
+ continue;
+ }
+
+ var subscriptionsToDelete = subscriptions.Where(sub =>
+ {
+ var queueName = sub.ForwardTo[(sub.ForwardTo.LastIndexOf('/') + 1)..];
+ return !queueNames.Contains(queueName);
+ }).ToList();
+
+ // Remove subscriptions for which the queue to be forwarded to can not be found in the same namespace.
+ foreach (SubscriptionProperties? asbSubscription in subscriptionsToDelete)
+ {
+ logger.LogWarning("Deleting subscription {subscriptionName} on topic {topicName}", asbSubscription.SubscriptionName, topic.Name);
+ await client.DeleteSubscriptionAsync(topic.Name, asbSubscription.SubscriptionName, cancellationToken);
+ }
+ }
+
+ return Unit.Instance;
+ }
+}
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/RemoveOrphanedSubscriptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/RemoveOrphanedSubscriptions.cs
index e7276a8d6..a95cb6af9 100644
--- a/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/RemoveOrphanedSubscriptions.cs
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/RemoveOrphanedSubscriptions.cs
@@ -9,7 +9,8 @@ using Microsoft.Extensions.Logging;
namespace Elsa.MassTransit.AzureServiceBus.Handlers;
///
-/// A handler for the notification that removes orphaned subscriptions from Azure Service Bus.
+/// A handler for the notification that removes subscriptions from Azure Service Bus
+/// when there has not been any heartbeat received lately for this instance.
///
[UsedImplicitly]
public class RemoveOrphanedSubscriptions(MessageTopologyProvider topologyProvider,
@@ -28,7 +29,7 @@ public class RemoveOrphanedSubscriptions(MessageTopologyProvider topologyProvide
{
try
{
- // Get subscriptions based on topics instead of the topology since when names are longer than 50 characters.
+ // Get subscriptions based on topics instead of the topology since when names are longer than 50 characters
// MassTransit automatically truncates them.
await foreach (var asbSubscription in client.GetSubscriptionsAsync(subscription.TopicName, cancellationToken))
{
@@ -44,7 +45,7 @@ public class RemoveOrphanedSubscriptions(MessageTopologyProvider topologyProvide
}
catch (ServiceBusException ex) when(ex.Reason == ServiceBusFailureReason.MessagingEntityNotFound)
{
- logger.LogWarning(ex, $"Service bus entity {ex.EntityPath} was not found");
+ logger.LogWarning(ex, "Service bus entity {entityPath} was not found", ex.EntityPath);
}
}
}
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/HostedServices/CleanSubscriptionsWithoutQueues.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/HostedServices/CleanSubscriptionsWithoutQueues.cs
new file mode 100644
index 000000000..b83fff786
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/HostedServices/CleanSubscriptionsWithoutQueues.cs
@@ -0,0 +1,57 @@
+using Elsa.MassTransit.AzureServiceBus.Commands;
+using Elsa.MassTransit.AzureServiceBus.Options;
+using Elsa.Mediator.Contracts;
+using Medallion.Threading;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Hosting;
+using Microsoft.Extensions.Options;
+
+namespace Elsa.MassTransit.AzureServiceBus.HostedServices;
+
+///
+/// Represents a hosted service that cleans up subscriptions without queues in Azure Service Bus.
+///
+public class CleanSubscriptionsWithoutQueues(IServiceProvider serviceProvider, IOptions subscriptionCleanupOptions) : IHostedService, IDisposable
+{
+ private Timer? _timer;
+
+ ///
+ public Task StartAsync(CancellationToken cancellationToken)
+ {
+ // Start the service a bit later to make sure it is not running simultaneously with the Heartbeat on startup.
+ _timer = new Timer(CleanUpSubscriptions, null, TimeSpan.FromMinutes(2), subscriptionCleanupOptions.Value.Interval);
+ return Task.CompletedTask;
+ }
+
+ ///
+ public Task StopAsync(CancellationToken cancellationToken)
+ {
+ _timer?.Change(Timeout.Infinite, Timeout.Infinite);
+ return Task.CompletedTask;
+ }
+
+ ///
+ public void Dispose()
+ {
+ _timer?.Dispose();
+ }
+
+ private void CleanUpSubscriptions(object? state)
+ {
+ _ = Task.Run(async () => await CleanUpSubscriptionsAsync());
+ }
+
+ private async Task CleanUpSubscriptionsAsync()
+ {
+ using var scope = serviceProvider.CreateScope();
+ var lockProvider = scope.ServiceProvider.GetRequiredService();
+ var commandSender = scope.ServiceProvider.GetRequiredService();
+
+ const string lockKey = "SubscriptionCleanupService";
+ await using var monitorLock = await lockProvider.TryAcquireLockAsync(lockKey, TimeSpan.Zero);
+ if (monitorLock == null)
+ return;
+
+ await commandSender.SendAsync(new CleanupSubscriptions());
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Options/SubscriptionCleanupOptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/SubscriptionCleanupOptions.cs
new file mode 100644
index 000000000..2d20a6f84
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/SubscriptionCleanupOptions.cs
@@ -0,0 +1,6 @@
+namespace Elsa.MassTransit.AzureServiceBus.Options;
+
+public class SubscriptionCleanupOptions
+{
+ public TimeSpan Interval { get; set; } = TimeSpan.FromDays(7);
+}
\ No newline at end of file