diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceived.cs similarity index 75% rename from src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceived.cs index 696426a3f..dd83385c1 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceived.cs @@ -7,7 +7,7 @@ using Elsa.Services; using Elsa.Services.Models; using Microsoft.Azure.ServiceBus; -namespace Elsa.Activities.AzureServiceBus.Activities +namespace Elsa.Activities.AzureServiceBus { [Trigger(Category = "Azure Service Bus", DisplayName = "Service Bus Message Received", Description = "Triggered when a message is received on the specified queue", Outcomes = new[] { OutcomeNames.Done })] public class AzureServiceBusMessageReceived : Activity @@ -22,9 +22,10 @@ namespace Elsa.Activities.AzureServiceBus.Activities [ActivityProperty] public string QueueName { get; set; } = default!; [ActivityProperty] public Type MessageType { get; set; } = default!; - protected override IActivityExecutionResult OnExecute() => Suspend(); + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => context.WorkflowExecutionContext.IsFirstPass ? ExecuteInternal(context) : Suspend(); + protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) => ExecuteInternal(context); - protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) + protected IActivityExecutionResult ExecuteInternal(ActivityExecutionContext context) { var message = (Message) context.Input!; var bytes = message.Body; diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs new file mode 100644 index 000000000..179865f9e --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs @@ -0,0 +1,11 @@ +using System; +using Elsa.Builders; + +namespace Elsa.Activities.AzureServiceBus +{ + public static class AzureServiceBusMessageReceivedBuilderExtensions + { + public static IActivityBuilder MessageReceived(this IBuilder builder, Action> setup) => builder.Then(setup); + public static IActivityBuilder MessageReceived(this IBuilder builder, string queueName) => builder.MessageReceived(setup => setup.WithQueueName(queueName).WithMessageType()); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs new file mode 100644 index 000000000..f091195ce --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs @@ -0,0 +1,23 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.AzureServiceBus +{ + public static class AzureServiceBusMessageReceivedExtensions + { + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.QueueName, value!); + + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived) => messageReceived.WithMessageType(typeof(T)); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs new file mode 100644 index 000000000..123ffdd76 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs @@ -0,0 +1,16 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +namespace Elsa.Activities.AzureServiceBus +{ + public static class SendAzureServiceBusMessageBuilderExtensions + { + public static IActivityBuilder SendMessage(this IBuilder builder, Action> setup) => builder.Then(setup); + public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func> message) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message)); + public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func message) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message)); + public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func message) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message)); + public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, object message) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message)); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs new file mode 100644 index 000000000..49e336f76 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs @@ -0,0 +1,22 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.AzureServiceBus +{ + public static class SendAzureServiceBusMessageExtensions + { + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func> value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func> value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, string value) => activity.Set(x => x.QueueName, value!); + + public static ISetupActivity WithMessage(this ISetupActivity activity, Func> value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, object value) => activity.Set(x => x.Message, value!); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusMessage.cs similarity index 74% rename from src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusMessage.cs index 8fe9f0a7d..a1f5e8e22 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusMessage.cs @@ -1,25 +1,25 @@ using System.Text; using System.Threading; using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Services; using Elsa.ActivityResults; using Elsa.Attributes; using Elsa.Serialization; using Elsa.Services; using Elsa.Services.Models; using Microsoft.Azure.ServiceBus; -using IServiceBusFactory = Elsa.Activities.AzureServiceBus.Services.IServiceBusFactory; -namespace Elsa.Activities.AzureServiceBus.Activities +namespace Elsa.Activities.AzureServiceBus { [Trigger(Category = "Azure Service Bus", DisplayName = "Send Service Bus Message", Description = "Sends a message to the specified queue", Outcomes = new[] { OutcomeNames.Done })] public class SendAzureServiceBusMessage : Activity { - private readonly IServiceBusFactory _serviceBusFactory; + private readonly IMessageSenderFactory _messageSenderFactory; private readonly IContentSerializer _serializer; - public SendAzureServiceBusMessage(IServiceBusFactory serviceBusFactory, IContentSerializer serializer) + public SendAzureServiceBusMessage(IMessageSenderFactory messageSenderFactory, IContentSerializer serializer) { - _serviceBusFactory = serviceBusFactory; + _messageSenderFactory = messageSenderFactory; _serializer = serializer; } @@ -28,7 +28,7 @@ namespace Elsa.Activities.AzureServiceBus.Activities protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) { - var sender = await _serviceBusFactory.GetSenderAsync(QueueName, cancellationToken); + var sender = await _messageSenderFactory.GetSenderAsync(QueueName, cancellationToken); var json = _serializer.Serialize(Message); var bytes = Encoding.UTF8.GetBytes(json); var message = new Message(bytes); diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj b/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj index d064f2ff9..1e03d9a0d 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj +++ b/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj @@ -2,17 +2,34 @@ netstandard2.0 + 1.0.0 latest + Elsa Contributors + + Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application. + This package provides activities to send and receive messages using Azure Service Bus. + + 2020 + https://github.com/elsa-workflows/elsa-core + https://github.com/elsa-workflows/elsa-core + GitHub + elsa, workflows + icon.png enable - - + - + + + + True + + + diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj.DotSettings b/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj.DotSettings new file mode 100644 index 000000000..5433deafa --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj.DotSettings @@ -0,0 +1,4 @@ + + True + True + True \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs new file mode 100644 index 000000000..02f43aae2 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ManagementClientExtensions.cs @@ -0,0 +1,17 @@ +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.ServiceBus.Management; + +namespace Elsa.Activities.AzureServiceBus.Extensions +{ + public static class ManagementClientExtensions + { + public static async Task EnsureQueueExistsAsync(this ManagementClient managementClient, string queueName, CancellationToken cancellationToken) + { + if (await managementClient.QueueExistsAsync(queueName, cancellationToken)) + return; + + await managementClient.CreateQueueAsync(queueName, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs index c46d48006..86604cd56 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -1,9 +1,8 @@ using System; -using Elsa.Activities.AzureServiceBus.Activities; using Elsa.Activities.AzureServiceBus.Options; using Elsa.Activities.AzureServiceBus.Services; using Elsa.Activities.AzureServiceBus.StartupTasks; -using Elsa.Runtime; +using Elsa.Activities.AzureServiceBus.Triggers; using Microsoft.Azure.ServiceBus; using Microsoft.Azure.ServiceBus.Management; using Microsoft.Extensions.DependencyInjection; @@ -23,8 +22,10 @@ namespace Elsa.Activities.AzureServiceBus.Extensions return services .AddSingleton(CreateServiceBusConnection) .AddSingleton(CreateServiceBusManagementClient) - .AddSingleton() - .AddStartupTask() + .AddSingleton() + .AddSingleton() + .AddHostedService() + .AddTriggerProvider() .AddActivity() .AddActivity(); } diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs similarity index 78% rename from src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusFactory.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs index 4d0e0dd55..c56792c64 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IMessageSenderFactory.cs @@ -4,9 +4,13 @@ using Microsoft.Azure.ServiceBus.Core; namespace Elsa.Activities.AzureServiceBus.Services { - public interface IServiceBusFactory + public interface IMessageSenderFactory { Task GetSenderAsync(string queueName, CancellationToken cancellationToken = default); + } + + 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/MessageReceiverFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageReceiverFactory.cs new file mode 100644 index 000000000..4ef78a7c2 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageReceiverFactory.cs @@ -0,0 +1,45 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Extensions; +using Microsoft.Azure.ServiceBus; +using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Azure.ServiceBus.Management; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class MessageReceiverFactory : IMessageReceiverFactory + { + private readonly ServiceBusConnection _connection; + private readonly ManagementClient _managementClient; + private readonly IDictionary _receivers = new Dictionary(); + private readonly SemaphoreSlim _semaphore = new(1); + + public MessageReceiverFactory(ServiceBusConnection connection, ManagementClient managementClient) + { + _connection = connection; + _managementClient = managementClient; + } + + public async Task GetReceiverAsync(string queueName, CancellationToken cancellationToken) + { + if (_receivers.TryGetValue(queueName, out var messageReceiver)) + return messageReceiver; + + await _semaphore.WaitAsync(cancellationToken); + + try + { + await _managementClient.EnsureQueueExistsAsync(queueName, cancellationToken); + var newMessageReceiver = new MessageReceiver(_connection, queueName); + _receivers.Add(queueName, newMessageReceiver); + return newMessageReceiver; + } + finally + { + _semaphore.Release(); + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs new file mode 100644 index 000000000..ebfd37c49 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageSenderFactory.cs @@ -0,0 +1,44 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Extensions; +using Microsoft.Azure.ServiceBus; +using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Azure.ServiceBus.Management; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class MessageSenderFactory : IMessageSenderFactory + { + private readonly ServiceBusConnection _connection; + private readonly ManagementClient _managementClient; + private readonly IDictionary _senders = new Dictionary(); + private readonly SemaphoreSlim _semaphore = new(1); + + public MessageSenderFactory(ServiceBusConnection connection, ManagementClient managementClient) + { + _connection = connection; + _managementClient = managementClient; + } + + public async Task GetSenderAsync(string queueName, CancellationToken cancellationToken) + { + if (_senders.TryGetValue(queueName, out var messageSender)) + return messageSender; + + await _semaphore.WaitAsync(cancellationToken); + + try + { + await _managementClient.EnsureQueueExistsAsync(queueName, cancellationToken); + var newMessageSender = new MessageSender(_connection, queueName); + _senders.Add(queueName, newMessageSender); + return newMessageSender; + } + finally + { + _semaphore.Release(); + } + } + } +} \ 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 14089b2f9..5dec1fea1 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -1,40 +1,58 @@ -using System.Threading; +using System; +using System.Threading; using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Triggers; using Elsa.Services; +using Microsoft.Azure.ServiceBus; using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; namespace Elsa.Activities.AzureServiceBus.Services { public class QueueWorker { private readonly IMessageReceiver _messageReceiver; - private readonly IWorkflowScheduler _workflowScheduler; + private readonly IServiceProvider _serviceProvider; + private readonly ILogger _logger; - public QueueWorker(IMessageReceiver messageReceiver, IWorkflowScheduler workflowScheduler) + public QueueWorker(IMessageReceiver messageReceiver, IServiceProvider serviceProvider, ILogger logger) { _messageReceiver = messageReceiver; - _workflowScheduler = workflowScheduler; - } - - public async Task StartAsync(CancellationToken cancellationToken) - { - var cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - await Task.Factory.StartNew(() => ReadQueueAsync(cancellationTokenSource.Token), cancellationToken); - } - - private async Task ReadQueueAsync(CancellationToken cancellationToken) - { - while (!cancellationToken.IsCancellationRequested) + _serviceProvider = serviceProvider; + _logger = logger; + + _messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler) { - var message = await _messageReceiver.ReceiveAsync(); + AutoComplete = false, + MaxConcurrentCalls = 10 + }); + } - if(message == null) - continue; + private async Task OnMessageReceived(Message message, CancellationToken cancellationToken) + { + using (var scope = _serviceProvider.CreateScope()) + { + var workflowScheduler = scope.ServiceProvider.GetRequiredService(); - await _workflowScheduler.TriggerWorkflowsAsync(x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), message, + await workflowScheduler.TriggerWorkflowsAsync(x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), message, message.CorrelationId, cancellationToken: cancellationToken); } + + await _messageReceiver.CompleteAsync(message.SystemProperties.LockToken); + } + + private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e) + { + var context = e.ExceptionReceivedContext; + + _logger.LogError("Message handler encountered an exception {Exception}.", e.Exception); + _logger.LogError("Exception context for troubleshooting:"); + _logger.LogError("- Endpoint: {Endpoint}", context.Endpoint); + _logger.LogError("- Entity Path: {EntityPath}", context.EntityPath); + _logger.LogError("- Executing Action: {Action}", context.Action); + + return Task.CompletedTask; } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusFactory.cs deleted file mode 100644 index bfec991c2..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusFactory.cs +++ /dev/null @@ -1,54 +0,0 @@ -using System.Collections.Generic; -using System.Threading; -using System.Threading.Tasks; -using Microsoft.Azure.ServiceBus; -using Microsoft.Azure.ServiceBus.Core; -using Microsoft.Azure.ServiceBus.Management; - -namespace Elsa.Activities.AzureServiceBus.Services -{ - public class ServiceBusFactory : IServiceBusFactory - { - private readonly ServiceBusConnection _connection; - private readonly ManagementClient _managementClient; - private readonly IDictionary _senders = new Dictionary(); - private readonly IDictionary _receivers = new Dictionary(); - - public ServiceBusFactory(ServiceBusConnection connection, ManagementClient managementClient) - { - _connection = connection; - _managementClient = managementClient; - } - - public async Task GetSenderAsync(string queueName, CancellationToken cancellationToken) - { - if (_senders.TryGetValue(queueName, out var messageSender)) - return messageSender; - - await EnsureQueueExistsAsync(queueName, cancellationToken); - - var newMessageSender = new MessageSender(_connection, queueName); - _senders.Add(queueName, newMessageSender); - return newMessageSender; - } - - public async Task GetReceiverAsync(string queueName, CancellationToken cancellationToken) - { - if (_receivers.TryGetValue(queueName, out var messageReceiver)) - return messageReceiver; - - await EnsureQueueExistsAsync(queueName, cancellationToken); - var newMessageReceiver = new MessageReceiver(_connection, queueName); - _receivers.Add(queueName, newMessageReceiver); - return newMessageReceiver; - } - - private async Task EnsureQueueExistsAsync(string queueName, CancellationToken cancellationToken) - { - if (await _managementClient.QueueExistsAsync(queueName, cancellationToken)) - return; - - await _managementClient.CreateQueueAsync(queueName, cancellationToken); - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs index 7f51ae4b0..a4a50a3a4 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs @@ -4,39 +4,37 @@ using System.Linq; using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; -using Elsa.Activities.AzureServiceBus.Activities; using Elsa.Activities.AzureServiceBus.Services; using Elsa.Services; using Microsoft.Extensions.DependencyInjection; -using IServiceBusFactory = Elsa.Activities.AzureServiceBus.Services.IServiceBusFactory; +using Microsoft.Extensions.Hosting; namespace Elsa.Activities.AzureServiceBus.StartupTasks { - public class StartServiceBusQueues : IStartupTask + public class StartServiceBusQueues : BackgroundService { private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector; - private readonly IServiceBusFactory _serviceBusFactory; + private readonly IMessageReceiverFactory _messageReceiverFactory; private readonly IServiceProvider _serviceProvider; - public StartServiceBusQueues(IWorkflowRegistry workflowRegistry, IWorkflowBlueprintReflector workflowBlueprintReflector, IServiceBusFactory serviceBusFactory, IServiceProvider serviceProvider) + public StartServiceBusQueues(IWorkflowRegistry workflowRegistry, IWorkflowBlueprintReflector workflowBlueprintReflector, IMessageReceiverFactory messageReceiverFactory, IServiceProvider serviceProvider) { _workflowRegistry = workflowRegistry; _workflowBlueprintReflector = workflowBlueprintReflector; - _serviceBusFactory = serviceBusFactory; + _messageReceiverFactory = messageReceiverFactory; _serviceProvider = serviceProvider; } - public async Task ExecuteAsync(CancellationToken cancellationToken = default) + protected override async Task ExecuteAsync(CancellationToken stoppingToken) { + var cancellationToken = stoppingToken; var queueNames = await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken); foreach (var queueName in queueNames) { - var receiver = await _serviceBusFactory.GetReceiverAsync(queueName, cancellationToken); - var worker = ActivatorUtilities.CreateInstance(_serviceProvider, receiver); - - await worker.StartAsync(cancellationToken); + var receiver = await _messageReceiverFactory.GetReceiverAsync(queueName, cancellationToken); + ActivatorUtilities.CreateInstance(_serviceProvider, receiver); } } diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs b/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs index 602cfe0c0..46d352288 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs @@ -1,6 +1,5 @@ using System.Threading; using System.Threading.Tasks; -using Elsa.Activities.AzureServiceBus.Activities; using Elsa.Triggers; namespace Elsa.Activities.AzureServiceBus.Triggers diff --git a/src/activities/Elsa.Activities.AzureServiceBus/icon.png b/src/activities/Elsa.Activities.AzureServiceBus/icon.png new file mode 100644 index 000000000..28dbaafbf Binary files /dev/null and b/src/activities/Elsa.Activities.AzureServiceBus/icon.png differ