diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs index 29a06fb1c..8eebd471a 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs @@ -12,12 +12,12 @@ namespace Elsa.Activities.AzureServiceBus [Action(Category = "Azure Service Bus", DisplayName = "Send Service Bus Message", Description = "Sends a message to the specified queue", Outcomes = new[] { OutcomeNames.Done })] public class SendAzureServiceBusQueueMessage : Activity { - private readonly IMessageSenderFactory _messageSenderFactory; + private readonly IQueueMessageSenderFactory _queueMessageSenderFactory; private readonly IContentSerializer _serializer; - public SendAzureServiceBusQueueMessage(IMessageSenderFactory messageSenderFactory, IContentSerializer serializer) + public SendAzureServiceBusQueueMessage(IQueueMessageSenderFactory queueMessageSenderFactory, IContentSerializer serializer) { - _messageSenderFactory = messageSenderFactory; + _queueMessageSenderFactory = queueMessageSenderFactory; _serializer = serializer; } @@ -29,7 +29,7 @@ namespace Elsa.Activities.AzureServiceBus protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) { - var sender = await _messageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken); + var sender = await _queueMessageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken); var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer, Message); if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId)) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs index 58541c07d..0d5d2b05b 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -1,5 +1,6 @@ using System; using Elsa.Activities.AzureServiceBus.Bookmarks; +using Elsa.Activities.AzureServiceBus.Handlers; using Elsa.Activities.AzureServiceBus.Options; using Elsa.Activities.AzureServiceBus.Services; using Elsa.Activities.AzureServiceBus.StartupTasks; @@ -23,15 +24,17 @@ namespace Elsa.Activities.AzureServiceBus.Extensions options.Services .AddSingleton(CreateServiceBusConnection) .AddSingleton(CreateServiceBusManagementClient) - .AddSingleton() - .AddSingleton(sp => sp.GetRequiredService()) - .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton() + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton() + .AddSingleton() .AddStartupTask() + .AddStartupTask() + .AddNotificationHandlers(typeof(RestartServiceBusQueues)) .AddBookmarkProvider() - - .AddSingleton(sp => sp.GetRequiredService()) - .AddSingleton(sp => sp.GetRequiredService()) - .AddStartupTask() .AddBookmarkProvider() ; @@ -41,7 +44,7 @@ namespace Elsa.Activities.AzureServiceBus.Extensions .AddActivity() .AddActivity() ; - + return options; } @@ -51,7 +54,7 @@ namespace Elsa.Activities.AzureServiceBus.Extensions var connectionString = options.ConnectionString; return new ServiceBusConnection(connectionString, RetryPolicy.Default); } - + private static ManagementClient CreateServiceBusManagementClient(IServiceProvider serviceProvider) { var options = serviceProvider.GetRequiredService>().Value; diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Handlers/RestartServiceBusQueues.cs b/src/activities/Elsa.Activities.AzureServiceBus/Handlers/RestartServiceBusQueues.cs new file mode 100644 index 000000000..3fdec6d9e --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Handlers/RestartServiceBusQueues.cs @@ -0,0 +1,16 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.Events; +using MediatR; + +namespace Elsa.Activities.AzureServiceBus.Handlers +{ + public class RestartServiceBusQueues : INotificationHandler, INotificationHandler + { + private readonly IServiceBusQueuesStarter _serviceBusQueuesStarter; + public RestartServiceBusQueues(IServiceBusQueuesStarter serviceBusQueuesStarter) => _serviceBusQueuesStarter = serviceBusQueuesStarter; + public Task Handle(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) => _serviceBusQueuesStarter.CreateWorkersAsync(cancellationToken); + public Task Handle(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => _serviceBusQueuesStarter.CreateWorkersAsync(cancellationToken); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Handlers/RestartServiceBusTopics.cs b/src/activities/Elsa.Activities.AzureServiceBus/Handlers/RestartServiceBusTopics.cs new file mode 100644 index 000000000..550d4ed79 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Handlers/RestartServiceBusTopics.cs @@ -0,0 +1,16 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.Events; +using MediatR; + +namespace Elsa.Activities.AzureServiceBus.Handlers +{ + public class RestartServiceBusTopics : INotificationHandler, INotificationHandler + { + private readonly IServiceBusTopicsStarter _serviceBusTopicsStarter; + public RestartServiceBusTopics(IServiceBusTopicsStarter serviceBusTopicsStarter) => _serviceBusTopicsStarter = serviceBusTopicsStarter; + public Task Handle(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) => _serviceBusTopicsStarter.CreateWorkersAsync(cancellationToken); + public Task Handle(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => _serviceBusTopicsStarter.CreateWorkersAsync(cancellationToken); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/BusClientFactory.cs similarity index 61% rename from src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Services/BusClientFactory.cs index 6af468778..09c45a316 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/BusClientFactory.cs @@ -1,3 +1,4 @@ +using System; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; @@ -7,23 +8,21 @@ using Microsoft.Azure.ServiceBus.Management; namespace Elsa.Activities.AzureServiceBus.Services { - public class MessageBusFactory : IMessageSenderFactory, IMessageReceiverFactory, ITopicMessageReceiverFactory, ITopicMessageSenderFactory + public class BusClientFactory : IQueueMessageSenderFactory, IQueueMessageReceiverClientFactory, ITopicMessageReceiverFactory, ITopicMessageSenderFactory { private readonly ServiceBusConnection _connection; private readonly ManagementClient _managementClient; - private readonly IDictionary _senders = new Dictionary(); - private readonly IDictionary _receivers = new Dictionary(); - - private readonly IDictionary<(string topicName, string queueName), IReceiverClient> _topicReceivers = new Dictionary<(string topicName, string queueName), IReceiverClient>(); + private readonly IDictionary _senders = new Dictionary(); + private readonly IDictionary _receivers = new Dictionary(); private readonly SemaphoreSlim _semaphore = new(1); - public MessageBusFactory(ServiceBusConnection connection, ManagementClient managementClient) + public BusClientFactory(ServiceBusConnection connection, ManagementClient managementClient) { _connection = connection; _managementClient = managementClient; } - public async Task GetSenderAsync(string queueName, CancellationToken cancellationToken) + public async Task GetSenderAsync(string queueName, CancellationToken cancellationToken) { await _semaphore.WaitAsync(cancellationToken); @@ -43,7 +42,7 @@ namespace Elsa.Activities.AzureServiceBus.Services } } - public async Task GetReceiverAsync(string queueName, CancellationToken cancellationToken) + public async Task GetReceiverAsync(string queueName, CancellationToken cancellationToken) { await _semaphore.WaitAsync(cancellationToken); @@ -63,6 +62,24 @@ namespace Elsa.Activities.AzureServiceBus.Services } } + public async Task DisposeReceiverAsync(IReceiverClient receiverClient, CancellationToken cancellationToken = default) + { + await _semaphore.WaitAsync(cancellationToken); + + var key = GetKeyFor(receiverClient); + + try + { + _receivers.Remove(key); + await receiverClient.UnregisterMessageHandlerAsync(TimeSpan.FromSeconds(1)); + await receiverClient.CloseAsync(); + } + finally + { + _semaphore.Release(); + } + } + private async Task EnsureQueueExistsAsync(string queueName, CancellationToken cancellationToken) { if (await _managementClient.QueueExistsAsync(queueName, cancellationToken)) @@ -71,7 +88,7 @@ namespace Elsa.Activities.AzureServiceBus.Services await _managementClient.CreateQueueAsync(queueName, cancellationToken); } - public async Task GetTopicSenderAsync(string topicName, CancellationToken cancellationToken) + public async Task GetTopicSenderAsync(string topicName, CancellationToken cancellationToken) { await _semaphore.WaitAsync(cancellationToken); @@ -95,7 +112,9 @@ namespace Elsa.Activities.AzureServiceBus.Services { await _semaphore.WaitAsync(cancellationToken); - if (_topicReceivers.TryGetValue((topicName,subscriptionName), out var messageReceiver)) + var key = $"{topicName}:{subscriptionName}"; + + if (_receivers.TryGetValue(key, out var messageReceiver)) return messageReceiver; try @@ -104,9 +123,12 @@ namespace Elsa.Activities.AzureServiceBus.Services var newTopicMessageReceiver = new SubscriptionClient( _connection, - topicPath: topicName, subscriptionName,ReceiveMode.PeekLock,RetryPolicy.Default) ; - - _topicReceivers.Add((topicName, subscriptionName), newTopicMessageReceiver); + topicName, + subscriptionName, + ReceiveMode.PeekLock, + RetryPolicy.Default); + + _receivers.Add(key, newTopicMessageReceiver); return newTopicMessageReceiver; } finally @@ -121,12 +143,20 @@ namespace Elsa.Activities.AzureServiceBus.Services await _managementClient.CreateTopicAsync(topicName, cancellationToken); } - private async Task EnsureTopicAndSubscriptionExistsAsync(string topicName, string subscriptionName ,CancellationToken cancellationToken) + private async Task EnsureTopicAndSubscriptionExistsAsync(string topicName, string subscriptionName, CancellationToken cancellationToken) { await EnsureTopicExistsAsync(topicName, cancellationToken); - if(!await _managementClient.SubscriptionExistsAsync(topicName, subscriptionName, cancellationToken)) + if (!await _managementClient.SubscriptionExistsAsync(topicName, subscriptionName, cancellationToken)) await _managementClient.CreateSubscriptionAsync(topicName, subscriptionName, cancellationToken); } + + private static string GetKeyFor(IReceiverClient receiverClient) => + receiverClient switch + { + IMessageReceiver messageReceiver => messageReceiver.Path, + ISubscriptionClient subscriptionClient => $"{subscriptionClient.TopicPath}:{subscriptionClient.SubscriptionName}", + _ => throw new ArgumentOutOfRangeException(nameof(receiverClient), receiverClient, null) + }; } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs deleted file mode 100644 index 716a3bfb8..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageReceiverFactory.cs +++ /dev/null @@ -1,11 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Microsoft.Azure.ServiceBus.Core; - -namespace Elsa.Activities.AzureServiceBus.Services -{ - public interface IMessageReceiverFactory - { - Task GetReceiverAsync(string queueName, CancellationToken cancellationToken = default); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IQueueMessageReceiverClientFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IQueueMessageReceiverClientFactory.cs new file mode 100644 index 000000000..ec2ce9c4d --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IQueueMessageReceiverClientFactory.cs @@ -0,0 +1,12 @@ +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.ServiceBus.Core; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public interface IQueueMessageReceiverClientFactory + { + Task GetReceiverAsync(string queueName, CancellationToken cancellationToken = default); + Task DisposeReceiverAsync(IReceiverClient receiverClient, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IQueueMessageSenderFactory.cs similarity index 50% rename from src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Services/IQueueMessageSenderFactory.cs index 250346053..5d3975d4c 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IQueueMessageSenderFactory.cs @@ -4,8 +4,8 @@ using Microsoft.Azure.ServiceBus.Core; namespace Elsa.Activities.AzureServiceBus.Services { - public interface IMessageSenderFactory + public interface IQueueMessageSenderFactory { - Task GetSenderAsync(string queueName, CancellationToken cancellationToken = default); + Task GetSenderAsync(string queueName, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusQueuesStarter.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusQueuesStarter.cs new file mode 100644 index 000000000..083a1af78 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusQueuesStarter.cs @@ -0,0 +1,10 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public interface IServiceBusQueuesStarter + { + Task CreateWorkersAsync(CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusTopicsStarter.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusTopicsStarter.cs new file mode 100644 index 000000000..99549f9cb --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusTopicsStarter.cs @@ -0,0 +1,10 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public interface IServiceBusTopicsStarter + { + Task CreateWorkersAsync(CancellationToken stoppingToken); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs index e747140db..7424498c5 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs @@ -7,5 +7,6 @@ namespace Elsa.Activities.AzureServiceBus.Services public interface ITopicMessageReceiverFactory { Task GetTopicReceiverAsync(string topicName, string subscriptionName, CancellationToken cancellationToken = default); + Task DisposeReceiverAsync(IReceiverClient receiverClient, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs index a0d16a411..380513552 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs @@ -6,6 +6,6 @@ namespace Elsa.Activities.AzureServiceBus.Services { public interface ITopicMessageSenderFactory { - Task GetTopicSenderAsync(string topicName, CancellationToken cancellationToken = default); + Task GetTopicSenderAsync(string topicName, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index dcabb8c96..686c5fc40 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -1,4 +1,6 @@ -using Elsa.Activities.AzureServiceBus.Bookmarks; +using System; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Bookmarks; using Elsa.Activities.AzureServiceBus.Options; using Elsa.Bookmarks; using Elsa.Dispatch; @@ -15,7 +17,8 @@ namespace Elsa.Activities.AzureServiceBus.Services IReceiverClient messageReceiver, IWorkflowDispatcher workflowDispatcher, IOptions options, - ILogger logger) : base(messageReceiver, workflowDispatcher, options, logger) + Func disposeReceiverAction, + ILogger logger) : base(messageReceiver, workflowDispatcher, options, disposeReceiverAction, logger) { } diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs new file mode 100644 index 000000000..3e7204219 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusQueuesStarter.cs @@ -0,0 +1,97 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Runtime.CompilerServices; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services; +using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class ServiceBusQueuesStarter : IServiceBusQueuesStarter + { + private readonly IQueueMessageReceiverClientFactory _messageReceiverClientFactory; + private readonly IServiceScopeFactory _scopeFactory; + private readonly IServiceProvider _serviceProvider; + private readonly ILogger _logger; + private readonly ICollection _workers; + + public ServiceBusQueuesStarter( + IQueueMessageReceiverClientFactory messageReceiverClientFactory, + IServiceScopeFactory scopeFactory, + IServiceProvider serviceProvider, + ILogger logger) + { + _messageReceiverClientFactory = messageReceiverClientFactory; + _scopeFactory = scopeFactory; + _serviceProvider = serviceProvider; + _logger = logger; + _workers = new List(); + } + + public async Task CreateWorkersAsync(CancellationToken cancellationToken = default) + { + await DisposeExistingWorkersAsync(); + var queueNames = (await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); + + foreach (var queueName in queueNames) + { + var receiver = await _messageReceiverClientFactory.GetReceiverAsync(queueName, cancellationToken); + var worker = ActivatorUtilities.CreateInstance(_serviceProvider, receiver, (Func) DisposeReceiverAsync); + _workers.Add(worker); + } + } + + private async Task DisposeExistingWorkersAsync() + { + foreach (var worker in _workers.ToList()) + { + await worker.DisposeAsync(); + _workers.Remove(worker); + } + } + + private async Task DisposeReceiverAsync(IReceiverClient messageReceiver) => await _messageReceiverClientFactory.DisposeReceiverAsync(messageReceiver); + + private async IAsyncEnumerable GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + using var scope = _scopeFactory.CreateScope(); + var workflowRegistry = scope.ServiceProvider.GetRequiredService(); + var workflowBlueprintReflector = scope.ServiceProvider.GetRequiredService(); + var workflows = await workflowRegistry.ListAsync(cancellationToken); + + var query = + from workflow in workflows + from activity in workflow.Activities + where activity.Type == nameof(AzureServiceBusQueueMessageReceived) + select workflow; + + foreach (var workflow in query) + { + var workflowBlueprintWrapper = await workflowBlueprintReflector.ReflectAsync(scope.ServiceProvider, workflow, cancellationToken); + + foreach (var activity in workflowBlueprintWrapper.Filter()) + { + var queueName = await activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken); + + if (string.IsNullOrWhiteSpace(queueName)) + { + _logger.LogWarning( + "Encountered a queue name that is null or empty in activity {ActivityType}:{ActivityId} in workflow {WorkflowDefinitionId}:v{WorkflowDefinitionVersion}", + activity.ActivityBlueprint.Type, + activity.ActivityBlueprint.Id, + workflow.Id, + workflow.Version); + + continue; + } + + yield return queueName!; + } + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusTopicsStarter.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusTopicsStarter.cs new file mode 100644 index 000000000..4acf4c0af --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusTopicsStarter.cs @@ -0,0 +1,82 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Runtime.CompilerServices; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services; +using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class ServiceBusTopicsStarter : IServiceBusTopicsStarter + { + private readonly ITopicMessageReceiverFactory _receiverFactory; + private readonly IServiceScopeFactory _scopeFactory; + private readonly IServiceProvider _serviceProvider; + private readonly ICollection _workers; + + public ServiceBusTopicsStarter( + ITopicMessageReceiverFactory receiverFactory, + IServiceScopeFactory scopeFactory, + IServiceProvider serviceProvider) + { + _receiverFactory = receiverFactory; + _scopeFactory = scopeFactory; + _serviceProvider = serviceProvider; + _workers = new List(); + } + + public async Task CreateWorkersAsync(CancellationToken stoppingToken) + { + var cancellationToken = stoppingToken; + await DisposeExistingWorkersAsync(); + var entities = (await GetTopicSubscriptionNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); + + foreach (var entity in entities) + { + var receiver = await _receiverFactory.GetTopicReceiverAsync(entity.topicName, entity.subscriptionName, cancellationToken); + var worker = ActivatorUtilities.CreateInstance(_serviceProvider, receiver, (Func) DisposeReceiverAsync); + _workers.Add(worker); + } + } + + private async Task DisposeExistingWorkersAsync() + { + foreach (var worker in _workers.ToList()) + { + await worker.DisposeAsync(); + _workers.Remove(worker); + } + } + + private async Task DisposeReceiverAsync(IReceiverClient messageReceiver) => await _receiverFactory.DisposeReceiverAsync(messageReceiver); + + private async IAsyncEnumerable<(string topicName, string subscriptionName)> GetTopicSubscriptionNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + using var scope = _scopeFactory.CreateScope(); + var workflowRegistry = scope.ServiceProvider.GetRequiredService(); + var workflowBlueprintReflector = scope.ServiceProvider.GetRequiredService(); + var workflows = await workflowRegistry.ListAsync(cancellationToken); + + var query = + from workflow in workflows + from activity in workflow.Activities + where activity.Type == nameof(AzureServiceBusTopicMessageReceived) + select workflow; + + foreach (var workflow in query) + { + var workflowBlueprintWrapper = await workflowBlueprintReflector.ReflectAsync(scope.ServiceProvider, workflow, cancellationToken); + + foreach (var activity in workflowBlueprintWrapper.Filter()) + { + var topicName = await activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken); + var subscriptionName = await activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken); + yield return (topicName, subscriptionName)!; + } + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs index f6178b24f..fdb890583 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs @@ -1,3 +1,5 @@ +using System; +using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Bookmarks; using Elsa.Activities.AzureServiceBus.Options; using Elsa.Bookmarks; @@ -15,7 +17,8 @@ namespace Elsa.Activities.AzureServiceBus.Services IReceiverClient receiverClient, IWorkflowDispatcher workflowDispatcher, IOptions options, - ILogger logger) : base(receiverClient, workflowDispatcher, options, logger) + Func disposeReceiverAction, + ILogger logger) : base(receiverClient, workflowDispatcher, options, disposeReceiverAction, logger) { } diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/WorkerBase.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/WorkerBase.cs index 117975d83..c9e7a1163 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/WorkerBase.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/WorkerBase.cs @@ -18,18 +18,21 @@ namespace Elsa.Activities.AzureServiceBus.Services private const string TenantId = default; private readonly IWorkflowDispatcher _workflowDispatcher; + private readonly Func _disposeReceiverAction; private readonly ILogger _logger; protected WorkerBase( IReceiverClient receiverClient, IWorkflowDispatcher workflowDispatcher, IOptions options, + Func disposeReceiverAction, ILogger logger) { ReceiverClient = receiverClient; _workflowDispatcher = workflowDispatcher; + _disposeReceiverAction = disposeReceiverAction; _logger = logger; - + ReceiverClient.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler) { AutoComplete = false, @@ -39,8 +42,7 @@ namespace Elsa.Activities.AzureServiceBus.Services protected IReceiverClient ReceiverClient { get; } protected abstract string ActivityType { get; } - - public async ValueTask DisposeAsync() => await ReceiverClient.CloseAsync(); + public async ValueTask DisposeAsync() => await _disposeReceiverAction(ReceiverClient); protected abstract IBookmark CreateBookmark(Message message); protected abstract IBookmark CreateTrigger(Message message); diff --git a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs index e1bcbdff8..fdf51cb29 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs @@ -1,77 +1,15 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Runtime.CompilerServices; -using System.Threading; +using System.Threading; using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Services; using Elsa.Services; -using Microsoft.Extensions.DependencyInjection; -using Microsoft.Extensions.Logging; namespace Elsa.Activities.AzureServiceBus.StartupTasks { public class StartServiceBusQueues : IStartupTask { - private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector; - private readonly IMessageReceiverFactory _messageReceiverFactory; - private readonly IServiceProvider _serviceProvider; - private readonly ILogger _logger; - - public StartServiceBusQueues(IWorkflowBlueprintReflector workflowBlueprintReflector, IMessageReceiverFactory messageReceiverFactory, IServiceProvider serviceProvider, ILogger logger) - { - _workflowBlueprintReflector = workflowBlueprintReflector; - _messageReceiverFactory = messageReceiverFactory; - _serviceProvider = serviceProvider; - _logger = logger; - } - + private readonly IServiceBusQueuesStarter _serviceBusQueuesStarter; + public StartServiceBusQueues(IServiceBusQueuesStarter serviceBusQueuesStarter) => _serviceBusQueuesStarter = serviceBusQueuesStarter; public int Order => 2000; - - public async Task ExecuteAsync(CancellationToken stoppingToken) - { - var cancellationToken = stoppingToken; - var queueNames = (await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); - - foreach (var queueName in queueNames) - { - var receiver = await _messageReceiverFactory.GetReceiverAsync(queueName, cancellationToken); - ActivatorUtilities.CreateInstance(_serviceProvider, receiver); - } - } - - private async IAsyncEnumerable GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken) - { - var workflowRegistry = _serviceProvider.GetRequiredService(); - var workflows = await workflowRegistry.ListAsync(cancellationToken); - - var query = - from workflow in workflows - from activity in workflow.Activities - where activity.Type == nameof(AzureServiceBusQueueMessageReceived) - select workflow; - - foreach (var workflow in query) - { - var workflowBlueprintWrapper = await _workflowBlueprintReflector.ReflectAsync(_serviceProvider, workflow, cancellationToken); - - foreach (var activity in workflowBlueprintWrapper.Filter()) - { - var queueName = await activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken); - - if (string.IsNullOrWhiteSpace(queueName)) - { - _logger.LogWarning( - "Encountered a queue name that is null or empty in activity {ActivityType}:{ActivityId} in workflow {WorkflowDefinitionId}", - activity.ActivityBlueprint.Type, - activity.ActivityBlueprint.Id, - workflow.Id); - continue; - } - - yield return queueName!; - } - } - } + public Task ExecuteAsync(CancellationToken stoppingToken) => _serviceBusQueuesStarter.CreateWorkersAsync(stoppingToken); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs deleted file mode 100644 index 8aba8b1b2..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs +++ /dev/null @@ -1,64 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Runtime.CompilerServices; -using System.Threading; -using System.Threading.Tasks; -using Elsa.Activities.AzureServiceBus.Services; -using Elsa.Services; -using Microsoft.Extensions.DependencyInjection; - -namespace Elsa.Activities.AzureServiceBus.StartupTasks -{ - public class StartServiceBusSubscription : IStartupTask - { - private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector; - private readonly ITopicMessageReceiverFactory _messageReceiverFactory; - private readonly IServiceProvider _serviceProvider; - - public StartServiceBusSubscription(IWorkflowBlueprintReflector workflowBlueprintReflector, ITopicMessageReceiverFactory messageReceiverFactory, IServiceProvider serviceProvider) - { - _workflowBlueprintReflector = workflowBlueprintReflector; - _messageReceiverFactory = messageReceiverFactory; - _serviceProvider = serviceProvider; - } - - public int Order => 2000; - - public async Task ExecuteAsync(CancellationToken stoppingToken) - { - var cancellationToken = stoppingToken; - var entities = (await GetTopicSubscriptionNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); - - foreach (var entity in entities) - { - var receiver = await _messageReceiverFactory.GetTopicReceiverAsync(entity.topicName, entity.subscriptionName, cancellationToken); - ActivatorUtilities.CreateInstance(_serviceProvider, receiver); - } - } - - private async IAsyncEnumerable<(string topicName, string subscriptionName)> GetTopicSubscriptionNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken) - { - var workflowRegistry = _serviceProvider.GetRequiredService(); - var workflows = await workflowRegistry.ListAsync(cancellationToken); - - var query = - from workflow in workflows - from activity in workflow.Activities - where activity.Type == nameof(AzureServiceBusTopicMessageReceived) - select workflow; - - foreach (var workflow in query) - { - var workflowBlueprintWrapper = await _workflowBlueprintReflector.ReflectAsync(_serviceProvider, workflow, cancellationToken); - - foreach (var activity in workflowBlueprintWrapper.Filter()) - { - var topicName = await activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken); - var subscriptionName = await activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken); - yield return (topicName, subscriptionName)!; - } - } - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusTopics.cs b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusTopics.cs new file mode 100644 index 000000000..8052d7186 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusTopics.cs @@ -0,0 +1,15 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.Services; + +namespace Elsa.Activities.AzureServiceBus.StartupTasks +{ + public class StartServiceBusTopics : IStartupTask + { + private readonly IServiceBusTopicsStarter _serviceBusTopicsStarter; + public StartServiceBusTopics(IServiceBusTopicsStarter serviceBusTopicsStarter) => _serviceBusTopicsStarter = serviceBusTopicsStarter; + public int Order => 2000; + public Task ExecuteAsync(CancellationToken stoppingToken) => _serviceBusTopicsStarter.CreateWorkersAsync(stoppingToken); + } +} \ No newline at end of file