diff --git a/src/bundles/Elsa.ServerAndStudio.Web/Program.cs b/src/bundles/Elsa.ServerAndStudio.Web/Program.cs index a83d2cf6a..27933b2ac 100644 --- a/src/bundles/Elsa.ServerAndStudio.Web/Program.cs +++ b/src/bundles/Elsa.ServerAndStudio.Web/Program.cs @@ -47,7 +47,25 @@ services if (useMassTransit) { - elsa.UseMassTransit(); + elsa.UseMassTransit(massTransit => + massTransit.UseAzureServiceBus(azureServiceBusConnectionString, serviceBusFeature => serviceBusFeature.ConfigureServiceBus = bus => + { + bus.PrefetchCount = 4; + bus.LockDuration = TimeSpan.FromMinutes(5); + bus.MaxConcurrentCalls = 32; + bus.MaxDeliveryCount = 8; + // etc. + }) + // massTransit.UseRabbitMq(rabbitMqConnectionString, rabbit => rabbit.ConfigureServiceBus = bus => + // { + // bus.PrefetchCount = 4; + // bus.Durable = true; + // bus.AutoDelete = false; + // bus.ConcurrentMessageLimit = 32; + // // etc. + // })) + ); + } }); diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs index 1f801443b..e07f57402 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs @@ -22,6 +22,7 @@ public static class ModuleExtensions void Configure(AzureServiceBusFeature bus) { + bus.AzureServiceBusOptions = options => options.ConnectionStringOrName = connectionString; bus.ConnectionString = connectionString; configure?.Invoke(bus); } diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index 95dad87d9..e8a79844d 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -36,6 +36,11 @@ public class AzureServiceBusFeature : FeatureBase /// public Action? ConfigureServiceBus { get; set; } + /// + /// A delegate to configure . + /// + public Action AzureServiceBusOptions { get; set; } = _ => { }; + /// /// A delegate to create a instance. /// @@ -98,7 +103,9 @@ public class AzureServiceBusFeature : FeatureBase } var genericType = consumerInterface.GetGenericArguments()[0]; - var topicName = $"{genericType.Namespace.ToLower()}~{genericType.Name.ToLower()}"; + //While the name might show up in the Azure portal with a ~, the separator is actually an / + //see https://learn.microsoft.com/en-us/archive/blogs/servicebus/azure-service-bus-azure-resource-manager-and-this-character + var topicName = $"{genericType.Namespace.ToLower()}/{genericType.Name.ToLower()}"; subscriptionTopology.Add(new MessageSubscriptionTopology(topicName, consumer.Name ?? genericType.Name.ToLower(), @@ -112,6 +119,7 @@ public class AzureServiceBusFeature : FeatureBase /// public override void Apply() { + Services.Configure(AzureServiceBusOptions); Services.AddScoped(ServiceBusAdministrationClientFactory); Services.AddNotificationHandler(); } diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs index cbc2f3c27..34d80d2a3 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Handlers/OrphanedSubscriptionRemover.cs @@ -8,7 +8,9 @@ namespace Elsa.MassTransit.AzureServiceBus.Handlers; /// /// Class responsible for removing orphaned subscriptions. /// -public class OrphanedSubscriptionRemover(MessageTopologyProvider topologyProvider, ServiceBusAdministrationClient client) +public class OrphanedSubscriptionRemover( + MessageTopologyProvider topologyProvider, + ServiceBusAdministrationClient client) : INotificationHandler { /// @@ -16,13 +18,26 @@ public class OrphanedSubscriptionRemover(MessageTopologyProvider topologyProvide /// public async Task HandleAsync(InstanceDeactivated notification, CancellationToken cancellationToken) { - var subscriptions = topologyProvider.GetShortLivedSubscriptions().ToList(); - - foreach (var subscription in subscriptions) + var subscriptionTopology = topologyProvider.GetShortLivedSubscriptions().ToList(); + + foreach (var subscription in subscriptionTopology) { - await client.DeleteSubscriptionAsync(subscription.TopicName, - $"Elsa-{notification.InstanceName}-{subscription.SubscriptionName}", - cancellationToken); + // 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)) + { + if (asbSubscription.SubscriptionName.StartsWith($"Elsa-{notification.InstanceName}-")) + { + var queueName = asbSubscription.ForwardTo.Substring(asbSubscription.ForwardTo.LastIndexOf('/') + 1); + + await Task.WhenAll( + client.DeleteSubscriptionAsync(subscription.TopicName, + asbSubscription.SubscriptionName, + cancellationToken), + client.DeleteQueueAsync(queueName, cancellationToken)); + } + } } } } \ 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 index 08f902bef..4d678dd61 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs @@ -8,5 +8,5 @@ public class AzureServiceBusOptions /// /// Th connection string or connection string name to connect with the service bus. /// - public string ConnectionStringOrName { get; set; } = default!; + public string? ConnectionStringOrName { get; set; } } \ No newline at end of file