Service Bus refactorings (#1613)

Fixes an issue where message types were registered with their own service bus, instead of these message types becoming part of the same bus configuration.

Fixes #1464
This commit is contained in:
Sipke Schoorstra 2021-10-08 16:45:23 +02:00 committed by GitHub
parent 0be4d3e156
commit 4934de6ac8
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
17 changed files with 98 additions and 93 deletions

View file

@ -23,7 +23,8 @@ namespace Elsa.Activities.Rebus.StartupTasks
{
foreach (var messageType in _messageTypes)
{
var bus = await _serviceBusFactory.GetServiceBusAsync(messageType, cancellationToken: cancellationToken);
var queueName = messageType.Name;
var bus = _serviceBusFactory.ConfigureServiceBus(new[] { messageType }, queueName);
await bus.Subscribe(messageType);
}
}

View file

@ -10,12 +10,12 @@ namespace Elsa.Services
/// </summary>
public interface ICommandSender
{
Task SendAsync(object message, string? queue = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default);
Task DeferAsync(object message, Duration delay, string? queue = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default);
Task SendAsync(object message, string? queueName = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default);
Task DeferAsync(object message, Duration delay, string? queueName = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default);
}
public static class CommandSenderExtensions
{
public static Task SendAsync(this ICommandSender commandSender, object message, CancellationToken cancellationToken = default) => commandSender.SendAsync(message, default, default, cancellationToken);
public static Task SendAsync(this ICommandSender commandSender, object message, string? queueName = default, CancellationToken cancellationToken = default) => commandSender.SendAsync(message, queueName, default, cancellationToken);
}
}

View file

@ -1,4 +1,5 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Rebus.Bus;
@ -7,6 +8,7 @@ namespace Elsa.Services
{
public interface IServiceBusFactory
{
Task<IBus> GetServiceBusAsync(Type messageType, string? queueName = default, CancellationToken cancellationToken = default);
IBus ConfigureServiceBus(IEnumerable<Type> messageTypes, string queueName);
IBus GetServiceBus(Type messageType, string? queueName = default);
}
}

View file

@ -19,21 +19,10 @@ using Rebus.Transport.InMem;
namespace Elsa.Options
{
public record MessageTypeConfig(Type MessageType, string? QueueName = default);
public record MessageTypeConfig(Type MessageType, string QueueName);
public class ElsaOptions
{
public static string FormatChannelQueueName<TMessage>(string channel) => FormatChannelQueueName(typeof(TMessage), channel);
public static string FormatChannelQueueName(Type messageType, string channel) => FormatChannelQueueName(messageType.Name, channel);
public static string FormatChannelQueueName(string queueName, string channel)
{
var queue = !string.IsNullOrWhiteSpace(channel) ? $"{queueName}{channel}" : queueName;
return FormatQueueName(queue);
}
public static string FormatQueueName(string queue) => queue.Humanize().Dehumanize().Underscore().Dasherize();
internal ElsaOptions()
{
WorkflowDefinitionStoreFactory = sp => ActivatorUtilities.CreateInstance<InMemoryWorkflowDefinitionStore>(sp);

View file

@ -1,19 +1,32 @@
using Rebus.Config;
using System;
using Humanizer;
using Rebus.Config;
namespace Elsa.Options
{
public class ServiceBusOptions
{
public static string FormatChannelQueueName<TMessage>(string channel) => FormatChannelQueueName(typeof(TMessage), channel);
public static string FormatChannelQueueName(Type messageType, string channel) => FormatChannelQueueName(messageType.Name, channel);
public static string FormatChannelQueueName(string queueName, string channel)
{
var queue = !string.IsNullOrWhiteSpace(channel) ? $"{queueName}{channel}" : queueName;
return FormatQueueName(queue);
}
public static string FormatQueueName(string queue) => queue.Humanize().Dehumanize().Underscore().Dasherize();
public int? NumberOfWorkers { get; set; }
public int? MaxParallelism { get; set; }
public string? QueuePrefix { get; set; }
public void Apply(OptionsConfigurer configurer)
{
if(NumberOfWorkers != null)
if (NumberOfWorkers != null)
configurer.SetNumberOfWorkers(NumberOfWorkers.Value);
if(MaxParallelism != null)
if (MaxParallelism != null)
configurer.SetMaxParallelism(MaxParallelism.Value);
}
}

View file

@ -41,7 +41,7 @@ namespace Elsa.Services.Bookmarks
private ISpecification<Bookmark> BuildSpecification(string activityType, IEnumerable<IBookmark> bookmarks, string? correlationId, string? tenantId)
{
var specification = bookmarks
.Select(trigger => _hasher.Hash(trigger))
.Select(bookmark => _hasher.Hash(bookmark))
.Aggregate(Specification<Bookmark>.None, (current, hash) => current.Or(new BookmarkHashSpecification(hash, activityType, tenantId)));
if (correlationId != null)

View file

@ -50,11 +50,11 @@ namespace Elsa.Services.Dispatch
}
var channel = _workflowChannelOptions.GetChannelOrDefault(workflowBlueprint.Channel);
var queue = ElsaOptions.FormatChannelQueueName<ExecuteWorkflowInstanceRequest>(channel);
var queue = ServiceBusOptions.FormatChannelQueueName<ExecuteWorkflowInstanceRequest>(channel);
await _commandSender.SendAsync(request, queue, cancellationToken: cancellationToken);
}
public async Task DispatchAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken = default) => await _commandSender.SendAsync(request, cancellationToken);
public async Task DispatchAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken = default) => await _commandSender.SendAsync(request, cancellationToken: cancellationToken);
public async Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken cancellationToken = default)
{
@ -67,7 +67,7 @@ namespace Elsa.Services.Dispatch
}
var channel = _workflowChannelOptions.GetChannelOrDefault(workflowBlueprint.Channel);
var queue = ElsaOptions.FormatChannelQueueName<ExecuteWorkflowDefinitionRequest>(channel);
var queue = ServiceBusOptions.FormatChannelQueueName<ExecuteWorkflowDefinitionRequest>(channel);
await _commandSender.SendAsync(request, queue, default, cancellationToken);
}
}

View file

@ -15,18 +15,18 @@ namespace Elsa.Services.Messaging
_serviceBusFactory = serviceBusFactory;
}
public async Task SendAsync(object message, string? queue = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default)
public async Task SendAsync(object message, string? queueName = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default)
{
var bus = await GetBusAsync(message, queue, cancellationToken);
var bus = GetBus(message, queueName);
await bus.Send(message, headers);
}
public async Task DeferAsync(object message, Duration delay, string? queue = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default)
public async Task DeferAsync(object message, Duration delay, string? queueName = default, IDictionary<string, string>? headers = default, CancellationToken cancellationToken = default)
{
var bus = await GetBusAsync(message, queue, cancellationToken);
var bus = GetBus(message, queueName);
await bus.Defer(delay.ToTimeSpan(), message, headers);
}
private async Task<IBus> GetBusAsync(object message, string? queue, CancellationToken cancellationToken) => await _serviceBusFactory.GetServiceBusAsync(message.GetType(), queue, cancellationToken);
private IBus GetBus(object message, string? queueName) => _serviceBusFactory.GetServiceBus(message.GetType(), queueName);
}
}

View file

@ -14,7 +14,7 @@ namespace Elsa.Services.Messaging
public async Task PublishAsync(object message, IDictionary<string, string>? headers = default)
{
var bus = await _serviceBusFactory.GetServiceBusAsync(message.GetType(), default);
var bus = _serviceBusFactory.GetServiceBus(message.GetType());
await bus.Publish(message, headers);
}
}

View file

@ -1,8 +1,6 @@
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using System.Linq;
using Elsa.Options;
using Elsa.Serialization;
using Microsoft.Extensions.Logging;
@ -19,9 +17,9 @@ namespace Elsa.Services.Messaging
private readonly ElsaOptions _elsaOptions;
private readonly ILoggerFactory _loggerFactory;
private readonly IServiceProvider _serviceProvider;
private readonly ConcurrentDictionary<string, IBus> _serviceBuses = new();
private readonly IDictionary<string, IBus> _serviceBuses = new Dictionary<string, IBus>();
private readonly IDictionary<Type, string> _messageTypeQueueDictionary = new Dictionary<Type, string>();
private readonly DependencyInjectionHandlerActivator _handlerActivator;
private readonly SemaphoreSlim _semaphore = new(1);
public ServiceBusFactory(ElsaOptions elsaOptions, ILoggerFactory loggerFactory, IServiceProvider serviceProvider)
{
@ -33,53 +31,42 @@ namespace Elsa.Services.Messaging
public void Dispose()
{
foreach (var key in _serviceBuses.Keys)
{
if (_serviceBuses.Remove(key, out var bus))
{
bus.Dispose();
}
}
_serviceBuses.Clear();
foreach (var bus in _serviceBuses.Values)
bus.Dispose();
}
public async Task<IBus> GetServiceBusAsync(Type messageType, string? queueName = default, CancellationToken cancellationToken = default)
public IBus ConfigureServiceBus(IEnumerable<Type> messageTypes, string queueName)
{
queueName = string.IsNullOrWhiteSpace(queueName)
? ElsaOptions.FormatChannelQueueName(messageType, _elsaOptions.WorkflowChannelOptions.Default)
: ElsaOptions.FormatQueueName(queueName);
queueName = ServiceBusOptions.FormatQueueName(queueName);
var prefixedQueueName = PrefixQueueName(queueName);
await _semaphore.WaitAsync(cancellationToken);
try
{
if (_serviceBuses.TryGetValue(prefixedQueueName, out var bus))
return bus;
var messageTypeList = messageTypes.ToList();
var configurer = Configure.With(_handlerActivator);
var map = messageTypeList.ToDictionary(x => x, _ => prefixedQueueName);
var configureContext = new ServiceBusEndpointConfigurationContext(configurer, prefixedQueueName, map, _serviceProvider);
var configurer = Configure.With(_handlerActivator);
var map = new Dictionary<Type, string> { [messageType] = prefixedQueueName };
var configureContext = new ServiceBusEndpointConfigurationContext(configurer, prefixedQueueName, map, _serviceProvider);
// Default options.
configurer
.Serialization(serializer => serializer.UseNewtonsoftJson(DefaultContentSerializer.CreateDefaultJsonSerializationSettings()))
.Logging(l => l.MicrosoftExtensionsLogging(_loggerFactory))
.Routing(r => r.TypeBased().Map(map))
.Options(options => options.Apply(_elsaOptions.ServiceBusOptions));
// Default options.
configurer
.Serialization(serializer => serializer.UseNewtonsoftJson(DefaultContentSerializer.CreateDefaultJsonSerializationSettings()))
.Logging(l => l.MicrosoftExtensionsLogging(_loggerFactory))
.Routing(r => r.TypeBased().Map(map))
.Options(options => options.Apply(_elsaOptions.ServiceBusOptions));
// Configure transport.
_elsaOptions.ConfigureServiceBusEndpoint(configureContext);
var newBus = configurer.Start();
_serviceBuses.TryAdd(prefixedQueueName, newBus);
// Configure transport.
_elsaOptions.ConfigureServiceBusEndpoint(configureContext);
return newBus;
}
finally
{
_semaphore.Release();
}
var newBus = configurer.Start();
_serviceBuses.Add(prefixedQueueName, newBus);
foreach (var messageType in messageTypeList)
_messageTypeQueueDictionary[messageType] = prefixedQueueName;
return newBus;
}
public IBus GetServiceBus(Type messageType, string? queueName = default)
{
queueName ??= _messageTypeQueueDictionary[messageType];
return _serviceBuses[queueName];
}
private string PrefixQueueName(string name) => $"{_elsaOptions.ServiceBusOptions.QueuePrefix}{name}";

View file

@ -32,8 +32,8 @@ namespace Elsa.StartupTasks
// For each workflow channel, register a competing message type for workflow definition and workflow instance consumers.
foreach (var workflowChannel in workflowChannels)
{
_competingMessageTypes.Add(new MessageTypeConfig(typeof(ExecuteWorkflowDefinitionRequest), ElsaOptions.FormatChannelQueueName("ExecuteWorkflow", workflowChannel)));
_competingMessageTypes.Add(new MessageTypeConfig(typeof(ExecuteWorkflowInstanceRequest), ElsaOptions.FormatChannelQueueName("ExecuteWorkflow", workflowChannel)));
_competingMessageTypes.Add(new MessageTypeConfig(typeof(ExecuteWorkflowDefinitionRequest), ServiceBusOptions.FormatChannelQueueName("ExecuteWorkflow", workflowChannel)));
_competingMessageTypes.Add(new MessageTypeConfig(typeof(ExecuteWorkflowInstanceRequest), ServiceBusOptions.FormatChannelQueueName("ExecuteWorkflow", workflowChannel)));
}
}
@ -45,19 +45,30 @@ namespace Elsa.StartupTasks
if (handle == null)
throw new Exception("Could not acquire a lock within the maximum amount of time configured");
foreach (var messageType in _competingMessageTypes)
var competingMessageTypeGroups = _competingMessageTypes.GroupBy(x => x.QueueName);
foreach (var messageTypeGroup in competingMessageTypeGroups)
{
var bus = await _serviceBusFactory.GetServiceBusAsync(messageType.MessageType, messageType.QueueName, cancellationToken);
await bus.Subscribe(messageType.MessageType);
var queueName = messageTypeGroup.Key;
var messageTypes = messageTypeGroup.Select(x => x.MessageType).ToList();
var bus = _serviceBusFactory.ConfigureServiceBus(messageTypes, queueName);
foreach (var messageType in messageTypes)
await bus.Subscribe(messageType);
}
var containerName = _containerNameAccessor.GetContainerName();
foreach (var messageType in _pubSubMessageTypes)
var pubSubMessageTypeGroups = _pubSubMessageTypes.GroupBy(x => x.QueueName);
foreach (var messageTypeGroup in pubSubMessageTypeGroups)
{
var queueName = $"{containerName}:{messageType.QueueName ?? messageType.MessageType.Name}";
var bus = await _serviceBusFactory.GetServiceBusAsync(messageType.MessageType, queueName, cancellationToken);
await bus.Subscribe(messageType.MessageType);
var queueName = $"{containerName}:{messageTypeGroup.Key}";
var messageTypes = messageTypeGroup.Select(x => x.MessageType).ToList();
var bus = _serviceBusFactory.ConfigureServiceBus(messageTypes, queueName);
foreach (var messageType in messageTypes)
await bus.Subscribe(messageType);
}
}
}

View file

@ -143,7 +143,7 @@ export class ElsaWorkflowDefinitionEditorScreen {
console.warn(`The specified workflow definition does not exist. Creating a new one.`)
}
}
this.updateWorkflowDefinition(workflowDefinition);
}

View file

@ -12,10 +12,12 @@
<ProjectReference Include="..\..\..\..\activities\Elsa.Activities.Email\Elsa.Activities.Email.csproj" />
<ProjectReference Include="..\..\..\..\activities\Elsa.Activities.Http\Elsa.Activities.Http.csproj" />
<ProjectReference Include="..\..\..\..\activities\Elsa.Activities.Temporal.Quartz\Elsa.Activities.Temporal.Quartz.csproj" />
<ProjectReference Include="..\..\..\..\caching\Elsa.Caching.Rebus\Elsa.Caching.Rebus.csproj" />
<ProjectReference Include="..\..\..\..\core\Elsa\Elsa.csproj" />
<ProjectReference Include="..\..\..\..\designer\bindings\aspnet\Elsa.Designer.Components.Web\Elsa.Designer.Components.Web.csproj" />
<ProjectReference Include="..\..\..\..\persistence\Elsa.Persistence.EntityFramework\Elsa.Persistence.EntityFramework.Sqlite\Elsa.Persistence.EntityFramework.Sqlite.csproj" />
<ProjectReference Include="..\..\..\..\server\Elsa.Server.Api\Elsa.Server.Api.csproj" />
<ProjectReference Include="..\..\..\..\servicebus\Elsa.Rebus.AzureServiceBus\Elsa.Rebus.AzureServiceBus.csproj" />
<ProjectReference Include="..\..\..\persistence\Elsa.Samples.Persistence.EntityFramework\Elsa.Samples.Persistence.EntityFramework.csproj" />
</ItemGroup>

View file

@ -26,8 +26,6 @@ namespace ElsaDashboard.Samples.AspNetCore.Monolith
// Elsa Server.
var elsaSection = Configuration.GetSection("Elsa");
var serviceBusConnectionString = Configuration.GetConnectionString("ASB");
services
.AddElsa(options => options
.UseEntityFrameworkPersistence(ef => ef.UseSqlite())

View file

@ -1,4 +1,6 @@
using System.Collections.Generic;
using Elsa.Caching.Rebus.Extensions;
using Elsa.Rebus.AzureServiceBus;
using Elsa.Retention.Extensions;
using Elsa.Server.Api.Extensions;
using Elsa.Server.Api.Hubs;

View file

@ -41,7 +41,7 @@ namespace Elsa.Server.Hangfire.Dispatch
}
var channel = _workflowChannelOptions.GetChannelOrDefault(workflowBlueprint.Channel);
var queue = ElsaOptions.FormatChannelQueueName<ExecuteWorkflowDefinitionRequest>(channel);
var queue = ServiceBusOptions.FormatChannelQueueName<ExecuteWorkflowDefinitionRequest>(channel);
EnqueueJob<WorkflowDefinitionJob>(x => x.ExecuteAsync(request, CancellationToken.None), queue);
}
@ -68,7 +68,7 @@ namespace Elsa.Server.Hangfire.Dispatch
}
var channel = _workflowChannelOptions.GetChannelOrDefault(workflowBlueprint.Channel);
var queue = ElsaOptions.FormatChannelQueueName<ExecuteWorkflowInstanceRequest>(channel);
var queue = ServiceBusOptions.FormatChannelQueueName<ExecuteWorkflowInstanceRequest>(channel);
EnqueueJob<WorkflowInstanceJob>(x => x.ExecuteAsync(request, CancellationToken.None), queue);
}

View file

@ -23,8 +23,8 @@ namespace Elsa.Server.Hangfire.Extensions
foreach (var channel in channels)
{
queues.Add(ElsaOptions.FormatChannelQueueName<ExecuteWorkflowDefinitionRequest>(channel));
queues.Add(ElsaOptions.FormatChannelQueueName<ExecuteWorkflowInstanceRequest>(channel));
queues.Add(ServiceBusOptions.FormatChannelQueueName<ExecuteWorkflowDefinitionRequest>(channel));
queues.Add(ServiceBusOptions.FormatChannelQueueName<ExecuteWorkflowInstanceRequest>(channel));
}
options.Queues = queues.ToArray();