Implement automated cleanup for Azure Service Bus subscriptions (#5636)
* Implement automated cleanup for Azure Service Bus subscriptions An automated cleanup process has been added for subscriptions that do not have queues connected to Azure Service Bus. This new feature deletes orphaned topics and cleans up other remnants within the namespace. * Removed string interpolation from log message. * Update Program.cs Use Memory transport to ensure docker image for demo purposes functions correctly. * Update AzureServiceBusFeature.cs * Update CleanupSubscriptions.cs * Refactor notifications to commands in MassTransit module Replaced the usage of notifications with commands in the Elsa.MassTransit.AzureServiceBus module. This included changing notification handler to command handler in methods and updating services to use the new command handlers. --------- Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>
This commit is contained in:
parent
ad494c33d1
commit
0ff791a0b9
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -0,0 +1,8 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
|
||||
namespace Elsa.MassTransit.AzureServiceBus.Commands;
|
||||
|
||||
/// <summary>
|
||||
/// Notification to clean up the Azure Service Bus subscriptions.
|
||||
/// </summary>
|
||||
public record CleanupSubscriptions : ICommand;
|
||||
|
|
@ -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; }
|
||||
|
||||
/// <summary>
|
||||
/// Delete subscriptions where their connected queues could not be found.
|
||||
/// Topics without any subscriptions will be deleted as well.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// - 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.
|
||||
/// </remarks>
|
||||
public bool EnableAutomatedSubscriptionCleanup { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// A delegate that configures the Azure Service Bus transport options.
|
||||
/// </summary>
|
||||
|
|
@ -44,6 +55,11 @@ public class AzureServiceBusFeature : FeatureBase
|
|||
/// </summary>
|
||||
public Action<AzureServiceBusOptions> AzureServiceBusOptions { get; set; } = _ => { };
|
||||
|
||||
/// <summary>
|
||||
/// A delegate to configure <see cref="SubscriptionCleanupOptions"/>.
|
||||
/// </summary>
|
||||
public Action<SubscriptionCleanupOptions> SubscriptionCleanupOptions { get; set; } = _ => { };
|
||||
|
||||
/// <summary>
|
||||
/// A delegate to create a <see cref="ServiceBusAdministrationClient"/> instance.
|
||||
/// </summary>
|
||||
|
|
@ -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);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override void ConfigureHostedServices()
|
||||
{
|
||||
if (EnableAutomatedSubscriptionCleanup)
|
||||
{
|
||||
Module.ConfigureHostedService<CleanSubscriptionsWithoutQueues>();
|
||||
}
|
||||
}
|
||||
|
||||
private static string GetConnectionString(IServiceProvider serviceProvider)
|
||||
{
|
||||
var options = serviceProvider.GetRequiredService<IOptions<AzureServiceBusOptions>>().Value;
|
||||
|
|
@ -135,5 +163,6 @@ public class AzureServiceBusFeature : FeatureBase
|
|||
|
||||
Services.AddSingleton(new MessageTopologyProvider(subscriptionTopology));
|
||||
Services.AddNotificationHandler<RemoveOrphanedSubscriptions>();
|
||||
Services.AddCommandHandler<CleanupOrphanedTopology>();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Cleans up Azure Service Bus resources for which no children can be found.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// 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.
|
||||
/// </remarks>
|
||||
[UsedImplicitly]
|
||||
public class CleanupOrphanedTopology(ServiceBusAdministrationClient client, ILogger<CleanupOrphanedTopology> logger) : ICommandHandler<CleanupSubscriptions>
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<Unit> 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;
|
||||
}
|
||||
}
|
||||
|
|
@ -9,7 +9,8 @@ using Microsoft.Extensions.Logging;
|
|||
namespace Elsa.MassTransit.AzureServiceBus.Handlers;
|
||||
|
||||
/// <summary>
|
||||
/// A handler for the <see cref="HeartbeatTimedOut"/> notification that removes orphaned subscriptions from Azure Service Bus.
|
||||
/// A handler for the <see cref="HeartbeatTimedOut"/> notification that removes subscriptions from Azure Service Bus
|
||||
/// when there has not been any heartbeat received lately for this instance.
|
||||
/// </summary>
|
||||
[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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a hosted service that cleans up subscriptions without queues in Azure Service Bus.
|
||||
/// </summary>
|
||||
public class CleanSubscriptionsWithoutQueues(IServiceProvider serviceProvider, IOptions<SubscriptionCleanupOptions> subscriptionCleanupOptions) : IHostedService, IDisposable
|
||||
{
|
||||
private Timer? _timer;
|
||||
|
||||
/// <inheritdoc />
|
||||
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;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task StopAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
_timer?.Change(Timeout.Infinite, Timeout.Infinite);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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<IDistributedLockProvider>();
|
||||
var commandSender = scope.ServiceProvider.GetRequiredService<ICommandSender>();
|
||||
|
||||
const string lockKey = "SubscriptionCleanupService";
|
||||
await using var monitorLock = await lockProvider.TryAcquireLockAsync(lockKey, TimeSpan.Zero);
|
||||
if (monitorLock == null)
|
||||
return;
|
||||
|
||||
await commandSender.SendAsync(new CleanupSubscriptions());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,6 @@
|
|||
namespace Elsa.MassTransit.AzureServiceBus.Options;
|
||||
|
||||
public class SubscriptionCleanupOptions
|
||||
{
|
||||
public TimeSpan Interval { get; set; } = TimeSpan.FromDays(7);
|
||||
}
|
||||
Loading…
Reference in a new issue