Re-register Azure Service bus workers when workflows get published or retracted

This commit is contained in:
Sipke Schoorstra 2021-04-13 15:55:28 +02:00
parent c120dbd975
commit 3a175936a9
20 changed files with 341 additions and 178 deletions

View file

@ -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<IActivityExecutionResult> 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))

View file

@ -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<MessageBusFactory>()
.AddSingleton<IMessageSenderFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddSingleton<IMessageReceiverFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddSingleton<BusClientFactory>()
.AddSingleton<IQueueMessageSenderFactory>(sp => sp.GetRequiredService<BusClientFactory>())
.AddSingleton<IQueueMessageReceiverClientFactory>(sp => sp.GetRequiredService<BusClientFactory>())
.AddSingleton<ITopicMessageSenderFactory>(sp => sp.GetRequiredService<BusClientFactory>())
.AddSingleton<ITopicMessageReceiverFactory>(sp => sp.GetRequiredService<BusClientFactory>())
.AddSingleton<IServiceBusQueuesStarter, ServiceBusQueuesStarter>()
.AddSingleton<IServiceBusTopicsStarter, ServiceBusTopicsStarter>()
.AddStartupTask<StartServiceBusQueues>()
.AddStartupTask<StartServiceBusTopics>()
.AddNotificationHandlers(typeof(RestartServiceBusQueues))
.AddBookmarkProvider<QueueMessageReceivedBookmarkProvider>()
.AddSingleton<ITopicMessageSenderFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddSingleton<ITopicMessageReceiverFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddStartupTask<StartServiceBusSubscription>()
.AddBookmarkProvider<TopicMessageReceivedBookmarkProvider>()
;
@ -41,7 +44,7 @@ namespace Elsa.Activities.AzureServiceBus.Extensions
.AddActivity<SendAzureServiceBusTopicMessage>()
.AddActivity<AzureServiceBusTopicMessageReceived>()
;
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<IOptions<AzureServiceBusOptions>>().Value;

View file

@ -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<WorkflowDefinitionPublished>, INotificationHandler<WorkflowDefinitionRetracted>
{
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);
}
}

View file

@ -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<WorkflowDefinitionPublished>, INotificationHandler<WorkflowDefinitionRetracted>
{
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);
}
}

View file

@ -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<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
private readonly IDictionary<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
private readonly IDictionary<(string topicName, string queueName), IReceiverClient> _topicReceivers = new Dictionary<(string topicName, string queueName), IReceiverClient>();
private readonly IDictionary<string, ISenderClient> _senders = new Dictionary<string, ISenderClient>();
private readonly IDictionary<string, IReceiverClient> _receivers = new Dictionary<string, IReceiverClient>();
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<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken)
public async Task<ISenderClient> GetSenderAsync(string queueName, CancellationToken cancellationToken)
{
await _semaphore.WaitAsync(cancellationToken);
@ -43,7 +42,7 @@ namespace Elsa.Activities.AzureServiceBus.Services
}
}
public async Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken)
public async Task<IReceiverClient> 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<IMessageSender> GetTopicSenderAsync(string topicName, CancellationToken cancellationToken)
public async Task<ISenderClient> 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)
};
}
}

View file

@ -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<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken = default);
}
}

View file

@ -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<IReceiverClient> GetReceiverAsync(string queueName, CancellationToken cancellationToken = default);
Task DisposeReceiverAsync(IReceiverClient receiverClient, CancellationToken cancellationToken = default);
}
}

View file

@ -4,8 +4,8 @@ using Microsoft.Azure.ServiceBus.Core;
namespace Elsa.Activities.AzureServiceBus.Services
{
public interface IMessageSenderFactory
public interface IQueueMessageSenderFactory
{
Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken = default);
Task<ISenderClient> GetSenderAsync(string queueName, CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,10 @@
using System.Threading;
using System.Threading.Tasks;
namespace Elsa.Activities.AzureServiceBus.Services
{
public interface IServiceBusQueuesStarter
{
Task CreateWorkersAsync(CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,10 @@
using System.Threading;
using System.Threading.Tasks;
namespace Elsa.Activities.AzureServiceBus.Services
{
public interface IServiceBusTopicsStarter
{
Task CreateWorkersAsync(CancellationToken stoppingToken);
}
}

View file

@ -7,5 +7,6 @@ namespace Elsa.Activities.AzureServiceBus.Services
public interface ITopicMessageReceiverFactory
{
Task<IReceiverClient> GetTopicReceiverAsync(string topicName, string subscriptionName, CancellationToken cancellationToken = default);
Task DisposeReceiverAsync(IReceiverClient receiverClient, CancellationToken cancellationToken = default);
}
}

View file

@ -6,6 +6,6 @@ namespace Elsa.Activities.AzureServiceBus.Services
{
public interface ITopicMessageSenderFactory
{
Task<IMessageSender> GetTopicSenderAsync(string topicName, CancellationToken cancellationToken = default);
Task<ISenderClient> GetTopicSenderAsync(string topicName, CancellationToken cancellationToken = default);
}
}

View file

@ -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<AzureServiceBusOptions> options,
ILogger<QueueWorker> logger) : base(messageReceiver, workflowDispatcher, options, logger)
Func<IReceiverClient, Task> disposeReceiverAction,
ILogger<QueueWorker> logger) : base(messageReceiver, workflowDispatcher, options, disposeReceiverAction, logger)
{
}

View file

@ -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<QueueWorker> _workers;
public ServiceBusQueuesStarter(
IQueueMessageReceiverClientFactory messageReceiverClientFactory,
IServiceScopeFactory scopeFactory,
IServiceProvider serviceProvider,
ILogger<ServiceBusQueuesStarter> logger)
{
_messageReceiverClientFactory = messageReceiverClientFactory;
_scopeFactory = scopeFactory;
_serviceProvider = serviceProvider;
_logger = logger;
_workers = new List<QueueWorker>();
}
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<QueueWorker>(_serviceProvider, receiver, (Func<IReceiverClient, Task>) 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<string> GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken)
{
using var scope = _scopeFactory.CreateScope();
var workflowRegistry = scope.ServiceProvider.GetRequiredService<IWorkflowRegistry>();
var workflowBlueprintReflector = scope.ServiceProvider.GetRequiredService<IWorkflowBlueprintReflector>();
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<AzureServiceBusQueueMessageReceived>())
{
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!;
}
}
}
}
}

View file

@ -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<TopicWorker> _workers;
public ServiceBusTopicsStarter(
ITopicMessageReceiverFactory receiverFactory,
IServiceScopeFactory scopeFactory,
IServiceProvider serviceProvider)
{
_receiverFactory = receiverFactory;
_scopeFactory = scopeFactory;
_serviceProvider = serviceProvider;
_workers = new List<TopicWorker>();
}
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<TopicWorker>(_serviceProvider, receiver, (Func<IReceiverClient, Task>) 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<IWorkflowRegistry>();
var workflowBlueprintReflector = scope.ServiceProvider.GetRequiredService<IWorkflowBlueprintReflector>();
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<AzureServiceBusTopicMessageReceived>())
{
var topicName = await activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken);
var subscriptionName = await activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken);
yield return (topicName, subscriptionName)!;
}
}
}
}
}

View file

@ -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<AzureServiceBusOptions> options,
ILogger<TopicWorker> logger) : base(receiverClient, workflowDispatcher, options, logger)
Func<IReceiverClient, Task> disposeReceiverAction,
ILogger<TopicWorker> logger) : base(receiverClient, workflowDispatcher, options, disposeReceiverAction, logger)
{
}

View file

@ -18,18 +18,21 @@ namespace Elsa.Activities.AzureServiceBus.Services
private const string TenantId = default;
private readonly IWorkflowDispatcher _workflowDispatcher;
private readonly Func<IReceiverClient, Task> _disposeReceiverAction;
private readonly ILogger _logger;
protected WorkerBase(
IReceiverClient receiverClient,
IWorkflowDispatcher workflowDispatcher,
IOptions<AzureServiceBusOptions> options,
Func<IReceiverClient, Task> 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);

View file

@ -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<StartServiceBusQueues> _logger;
public StartServiceBusQueues(IWorkflowBlueprintReflector workflowBlueprintReflector, IMessageReceiverFactory messageReceiverFactory, IServiceProvider serviceProvider, ILogger<StartServiceBusQueues> 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<QueueWorker>(_serviceProvider, receiver);
}
}
private async IAsyncEnumerable<string> GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken)
{
var workflowRegistry = _serviceProvider.GetRequiredService<IWorkflowRegistry>();
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<AzureServiceBusQueueMessageReceived>())
{
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);
}
}

View file

@ -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<TopicWorker>(_serviceProvider, receiver);
}
}
private async IAsyncEnumerable<(string topicName, string subscriptionName)> GetTopicSubscriptionNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken)
{
var workflowRegistry = _serviceProvider.GetRequiredService<IWorkflowRegistry>();
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<AzureServiceBusTopicMessageReceived>())
{
var topicName = await activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken);
var subscriptionName = await activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken);
yield return (topicName, subscriptionName)!;
}
}
}
}
}

View file

@ -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);
}
}