diff --git a/src/activities/Elsa.Activities.Rebus/StartupTasks/CreateSubscriptions.cs b/src/activities/Elsa.Activities.Rebus/StartupTasks/CreateSubscriptions.cs index b35bde002..49164a890 100644 --- a/src/activities/Elsa.Activities.Rebus/StartupTasks/CreateSubscriptions.cs +++ b/src/activities/Elsa.Activities.Rebus/StartupTasks/CreateSubscriptions.cs @@ -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); } } diff --git a/src/core/Elsa.Abstractions/Services/Messaging/ICommandSender.cs b/src/core/Elsa.Abstractions/Services/Messaging/ICommandSender.cs index 9017d9051..f12cfaedb 100644 --- a/src/core/Elsa.Abstractions/Services/Messaging/ICommandSender.cs +++ b/src/core/Elsa.Abstractions/Services/Messaging/ICommandSender.cs @@ -10,12 +10,12 @@ namespace Elsa.Services /// public interface ICommandSender { - Task SendAsync(object message, string? queue = default, IDictionary? headers = default, CancellationToken cancellationToken = default); - Task DeferAsync(object message, Duration delay, string? queue = default, IDictionary? headers = default, CancellationToken cancellationToken = default); + Task SendAsync(object message, string? queueName = default, IDictionary? headers = default, CancellationToken cancellationToken = default); + Task DeferAsync(object message, Duration delay, string? queueName = default, IDictionary? 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); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Messaging/IServiceBusFactory.cs b/src/core/Elsa.Abstractions/Services/Messaging/IServiceBusFactory.cs index d6030ce65..a76776a67 100644 --- a/src/core/Elsa.Abstractions/Services/Messaging/IServiceBusFactory.cs +++ b/src/core/Elsa.Abstractions/Services/Messaging/IServiceBusFactory.cs @@ -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 GetServiceBusAsync(Type messageType, string? queueName = default, CancellationToken cancellationToken = default); + IBus ConfigureServiceBus(IEnumerable messageTypes, string queueName); + IBus GetServiceBus(Type messageType, string? queueName = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Options/ElsaOptions.cs b/src/core/Elsa.Core/Options/ElsaOptions.cs index c0eddb543..16f51233a 100644 --- a/src/core/Elsa.Core/Options/ElsaOptions.cs +++ b/src/core/Elsa.Core/Options/ElsaOptions.cs @@ -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(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(sp); diff --git a/src/core/Elsa.Core/Options/ServiceBusOptions.cs b/src/core/Elsa.Core/Options/ServiceBusOptions.cs index 2bd7dc7ea..8fef8e6e8 100644 --- a/src/core/Elsa.Core/Options/ServiceBusOptions.cs +++ b/src/core/Elsa.Core/Options/ServiceBusOptions.cs @@ -1,19 +1,32 @@ -using Rebus.Config; +using System; +using Humanizer; +using Rebus.Config; namespace Elsa.Options { public class ServiceBusOptions { + public static string FormatChannelQueueName(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); } } diff --git a/src/core/Elsa.Core/Services/Bookmarks/BookmarkFinder.cs b/src/core/Elsa.Core/Services/Bookmarks/BookmarkFinder.cs index 885817b35..385ee271f 100644 --- a/src/core/Elsa.Core/Services/Bookmarks/BookmarkFinder.cs +++ b/src/core/Elsa.Core/Services/Bookmarks/BookmarkFinder.cs @@ -41,7 +41,7 @@ namespace Elsa.Services.Bookmarks private ISpecification BuildSpecification(string activityType, IEnumerable bookmarks, string? correlationId, string? tenantId) { var specification = bookmarks - .Select(trigger => _hasher.Hash(trigger)) + .Select(bookmark => _hasher.Hash(bookmark)) .Aggregate(Specification.None, (current, hash) => current.Or(new BookmarkHashSpecification(hash, activityType, tenantId))); if (correlationId != null) diff --git a/src/core/Elsa.Core/Services/Dispatch/QueuingWorkflowDispatcher.cs b/src/core/Elsa.Core/Services/Dispatch/QueuingWorkflowDispatcher.cs index e1c2f6bb3..dcc4494e3 100644 --- a/src/core/Elsa.Core/Services/Dispatch/QueuingWorkflowDispatcher.cs +++ b/src/core/Elsa.Core/Services/Dispatch/QueuingWorkflowDispatcher.cs @@ -50,11 +50,11 @@ namespace Elsa.Services.Dispatch } var channel = _workflowChannelOptions.GetChannelOrDefault(workflowBlueprint.Channel); - var queue = ElsaOptions.FormatChannelQueueName(channel); + var queue = ServiceBusOptions.FormatChannelQueueName(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(channel); + var queue = ServiceBusOptions.FormatChannelQueueName(channel); await _commandSender.SendAsync(request, queue, default, cancellationToken); } } diff --git a/src/core/Elsa.Core/Services/Messaging/CommandSender.cs b/src/core/Elsa.Core/Services/Messaging/CommandSender.cs index 7ac9664c1..bf4757c63 100644 --- a/src/core/Elsa.Core/Services/Messaging/CommandSender.cs +++ b/src/core/Elsa.Core/Services/Messaging/CommandSender.cs @@ -15,18 +15,18 @@ namespace Elsa.Services.Messaging _serviceBusFactory = serviceBusFactory; } - public async Task SendAsync(object message, string? queue = default, IDictionary? headers = default, CancellationToken cancellationToken = default) + public async Task SendAsync(object message, string? queueName = default, IDictionary? 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? headers = default, CancellationToken cancellationToken = default) + public async Task DeferAsync(object message, Duration delay, string? queueName = default, IDictionary? 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 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); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/Messaging/EventPublisher.cs b/src/core/Elsa.Core/Services/Messaging/EventPublisher.cs index f829b2280..a90226b01 100644 --- a/src/core/Elsa.Core/Services/Messaging/EventPublisher.cs +++ b/src/core/Elsa.Core/Services/Messaging/EventPublisher.cs @@ -14,7 +14,7 @@ namespace Elsa.Services.Messaging public async Task PublishAsync(object message, IDictionary? headers = default) { - var bus = await _serviceBusFactory.GetServiceBusAsync(message.GetType(), default); + var bus = _serviceBusFactory.GetServiceBus(message.GetType()); await bus.Publish(message, headers); } } diff --git a/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs b/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs index b50af2de9..83e6681d2 100644 --- a/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs +++ b/src/core/Elsa.Core/Services/Messaging/ServiceBusFactory.cs @@ -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 _serviceBuses = new(); + private readonly IDictionary _serviceBuses = new Dictionary(); + private readonly IDictionary _messageTypeQueueDictionary = new Dictionary(); 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 GetServiceBusAsync(Type messageType, string? queueName = default, CancellationToken cancellationToken = default) + public IBus ConfigureServiceBus(IEnumerable 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 { [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}"; diff --git a/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs b/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs index a372bd87c..631d126cc 100644 --- a/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs +++ b/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs @@ -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); } } } diff --git a/src/designer/elsa-workflows-studio/src/components/screens/workflow-definition-editor/elsa-workflow-definition-editor-screen/elsa-workflow-definition-editor-screen.tsx b/src/designer/elsa-workflows-studio/src/components/screens/workflow-definition-editor/elsa-workflow-definition-editor-screen/elsa-workflow-definition-editor-screen.tsx index b77676ba9..be0442a43 100644 --- a/src/designer/elsa-workflows-studio/src/components/screens/workflow-definition-editor/elsa-workflow-definition-editor-screen/elsa-workflow-definition-editor-screen.tsx +++ b/src/designer/elsa-workflows-studio/src/components/screens/workflow-definition-editor/elsa-workflow-definition-editor-screen/elsa-workflow-definition-editor-screen.tsx @@ -143,7 +143,7 @@ export class ElsaWorkflowDefinitionEditorScreen { console.warn(`The specified workflow definition does not exist. Creating a new one.`) } } - + this.updateWorkflowDefinition(workflowDefinition); } diff --git a/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj b/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj index 28717ee27..c4a8ad4ff 100644 --- a/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj +++ b/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj @@ -12,10 +12,12 @@ + + diff --git a/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/Startup.cs b/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/Startup.cs index 9c9884209..bc8b34347 100644 --- a/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/Startup.cs +++ b/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/Startup.cs @@ -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()) diff --git a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs index ee56c7732..fe4868a68 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs @@ -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; diff --git a/src/server/Elsa.Server.Hangfire/Dispatch/HangfireWorkflowDispatcher.cs b/src/server/Elsa.Server.Hangfire/Dispatch/HangfireWorkflowDispatcher.cs index 544f5915a..6ea2fa9a8 100644 --- a/src/server/Elsa.Server.Hangfire/Dispatch/HangfireWorkflowDispatcher.cs +++ b/src/server/Elsa.Server.Hangfire/Dispatch/HangfireWorkflowDispatcher.cs @@ -41,7 +41,7 @@ namespace Elsa.Server.Hangfire.Dispatch } var channel = _workflowChannelOptions.GetChannelOrDefault(workflowBlueprint.Channel); - var queue = ElsaOptions.FormatChannelQueueName(channel); + var queue = ServiceBusOptions.FormatChannelQueueName(channel); EnqueueJob(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(channel); + var queue = ServiceBusOptions.FormatChannelQueueName(channel); EnqueueJob(x => x.ExecuteAsync(request, CancellationToken.None), queue); } diff --git a/src/server/Elsa.Server.Hangfire/Extensions/BackgroundJobServerOptionsExtensions.cs b/src/server/Elsa.Server.Hangfire/Extensions/BackgroundJobServerOptionsExtensions.cs index 4af15d68a..c9b5ef11e 100644 --- a/src/server/Elsa.Server.Hangfire/Extensions/BackgroundJobServerOptionsExtensions.cs +++ b/src/server/Elsa.Server.Hangfire/Extensions/BackgroundJobServerOptionsExtensions.cs @@ -23,8 +23,8 @@ namespace Elsa.Server.Hangfire.Extensions foreach (var channel in channels) { - queues.Add(ElsaOptions.FormatChannelQueueName(channel)); - queues.Add(ElsaOptions.FormatChannelQueueName(channel)); + queues.Add(ServiceBusOptions.FormatChannelQueueName(channel)); + queues.Add(ServiceBusOptions.FormatChannelQueueName(channel)); } options.Queues = queues.ToArray();