Add Azure ServiceBus activities
This commit is contained in:
parent
123e49c19d
commit
e7178ca28e
|
|
@ -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;
|
||||
|
|
@ -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<ISetupActivity<AzureServiceBusMessageReceived>> setup) => builder.Then(setup);
|
||||
public static IActivityBuilder MessageReceived<T>(this IBuilder builder, string queueName) => builder.MessageReceived(setup => setup.WithQueueName(queueName).WithMessageType<T>());
|
||||
}
|
||||
}
|
||||
|
|
@ -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<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<string>> value) => messageReceived.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ValueTask<string>> value) => messageReceived.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<string> value) => messageReceived.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, string> value) => messageReceived.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, string value) => messageReceived.Set(x => x.QueueName, value!);
|
||||
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<Type>> value) => messageReceived.Set(x => x.MessageType, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, Type> value) => messageReceived.Set(x => x.MessageType, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<Type> value) => messageReceived.Set(x => x.MessageType, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!);
|
||||
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType<T>(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived) => messageReceived.WithMessageType(typeof(T));
|
||||
}
|
||||
}
|
||||
|
|
@ -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<ISetupActivity<SendAzureServiceBusMessage>> setup) => builder.Then(setup);
|
||||
public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func<ActivityExecutionContext, ValueTask<object>> message) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message));
|
||||
public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func<ActivityExecutionContext, object> message) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message));
|
||||
public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func<object> 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));
|
||||
}
|
||||
}
|
||||
|
|
@ -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<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, ValueTask<string>> value) => activity.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ValueTask<string>> value) => activity.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<string> value) => activity.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, string> value) => activity.Set(x => x.QueueName, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, string value) => activity.Set(x => x.QueueName, value!);
|
||||
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, ValueTask<object>> value) => activity.Set(x => x.Message, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, object> value) => activity.Set(x => x.Message, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<object> value) => activity.Set(x => x.Message, value!);
|
||||
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, object value) => activity.Set(x => x.Message, value!);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<IActivityExecutionResult> 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);
|
||||
|
|
@ -2,17 +2,34 @@
|
|||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>netstandard2.0</TargetFramework>
|
||||
<PackageVersion>1.0.0</PackageVersion>
|
||||
<LangVersion>latest</LangVersion>
|
||||
<Authors>Elsa Contributors</Authors>
|
||||
<Description>
|
||||
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.
|
||||
</Description>
|
||||
<Copyright>2020</Copyright>
|
||||
<PackageProjectUrl>https://github.com/elsa-workflows/elsa-core</PackageProjectUrl>
|
||||
<RepositoryUrl>https://github.com/elsa-workflows/elsa-core</RepositoryUrl>
|
||||
<RepositoryType>GitHub</RepositoryType>
|
||||
<PackageTags>elsa, workflows</PackageTags>
|
||||
<PackageIcon>icon.png</PackageIcon>
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<!-- <PackageReference Include="Microsoft.Azure.Management.ServiceBus.Fluent" Version="1.35.0" />-->
|
||||
<PackageReference Include="Microsoft.Azure.ServiceBus" Version="5.1.0" />
|
||||
<PackageReference Include="Microsoft.Azure.ServiceBus" Version="5.1.0" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
|
||||
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<None Include="icon.png">
|
||||
<Pack>True</Pack>
|
||||
<PackagePath />
|
||||
</None>
|
||||
</ItemGroup>
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -0,0 +1,4 @@
|
|||
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
|
||||
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Cazureservicebusmessagereceived/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Csendazureservicebusmessage/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>
|
||||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<IServiceBusFactory, ServiceBusFactory>()
|
||||
.AddStartupTask<StartServiceBusQueues>()
|
||||
.AddSingleton<IMessageSenderFactory, MessageSenderFactory>()
|
||||
.AddSingleton<IMessageReceiverFactory, MessageReceiverFactory>()
|
||||
.AddHostedService<StartServiceBusQueues>()
|
||||
.AddTriggerProvider<MessageReceivedTriggerProvider>()
|
||||
.AddActivity<AzureServiceBusMessageReceived>()
|
||||
.AddActivity<SendAzureServiceBusMessage>();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,9 +4,13 @@ using Microsoft.Azure.ServiceBus.Core;
|
|||
|
||||
namespace Elsa.Activities.AzureServiceBus.Services
|
||||
{
|
||||
public interface IServiceBusFactory
|
||||
public interface IMessageSenderFactory
|
||||
{
|
||||
Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
||||
public interface IMessageReceiverFactory
|
||||
{
|
||||
Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
|
||||
private readonly SemaphoreSlim _semaphore = new(1);
|
||||
|
||||
public MessageReceiverFactory(ServiceBusConnection connection, ManagementClient managementClient)
|
||||
{
|
||||
_connection = connection;
|
||||
_managementClient = managementClient;
|
||||
}
|
||||
|
||||
public async Task<IMessageReceiver> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
|
||||
private readonly SemaphoreSlim _semaphore = new(1);
|
||||
|
||||
public MessageSenderFactory(ServiceBusConnection connection, ManagementClient managementClient)
|
||||
{
|
||||
_connection = connection;
|
||||
_managementClient = managementClient;
|
||||
}
|
||||
|
||||
public async Task<IMessageSender> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<QueueWorker> 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<IWorkflowScheduler>();
|
||||
|
||||
await _workflowScheduler.TriggerWorkflowsAsync<MessageReceivedTrigger>(x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), message,
|
||||
await workflowScheduler.TriggerWorkflowsAsync<MessageReceivedTrigger>(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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
|
||||
private readonly IDictionary<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
|
||||
|
||||
public ServiceBusFactory(ServiceBusConnection connection, ManagementClient managementClient)
|
||||
{
|
||||
_connection = connection;
|
||||
_managementClient = managementClient;
|
||||
}
|
||||
|
||||
public async Task<IMessageSender> 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<IMessageReceiver> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<QueueWorker>(_serviceProvider, receiver);
|
||||
|
||||
await worker.StartAsync(cancellationToken);
|
||||
var receiver = await _messageReceiverFactory.GetReceiverAsync(queueName, cancellationToken);
|
||||
ActivatorUtilities.CreateInstance<QueueWorker>(_serviceProvider, receiver);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,6 +1,5 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.AzureServiceBus.Activities;
|
||||
using Elsa.Triggers;
|
||||
|
||||
namespace Elsa.Activities.AzureServiceBus.Triggers
|
||||
|
|
|
|||
BIN
src/activities/Elsa.Activities.AzureServiceBus/icon.png
Normal file
BIN
src/activities/Elsa.Activities.AzureServiceBus/icon.png
Normal file
Binary file not shown.
|
After Width: | Height: | Size: 9.6 KiB |
Loading…
Reference in a new issue