diff --git a/src/bundles/Elsa.WorkflowServer.Web/Program.cs b/src/bundles/Elsa.WorkflowServer.Web/Program.cs index 478f6d4e2..3d6397ea4 100644 --- a/src/bundles/Elsa.WorkflowServer.Web/Program.cs +++ b/src/bundles/Elsa.WorkflowServer.Web/Program.cs @@ -28,7 +28,7 @@ const bool useProtoActor = true; const bool useHangfire = false; const bool useQuartz = true; const bool useMassTransit = true; -const bool useMassTransitAzureServiceBus = false; +const bool useMassTransitAzureServiceBus = true; const bool useMassTransitRabbitMq = false; var builder = WebApplication.CreateBuilder(args); @@ -202,9 +202,14 @@ services elsa.UseMassTransit(massTransit => { if (useMassTransitAzureServiceBus) - massTransit.UseAzureServiceBus(azureServiceBusConnectionString); + { + massTransit.UseAzureServiceBus(azureServiceBusConnectionString, asb => asb.ConfigureServiceBus = bus => { bus.PrefetchCount = 4; }); + } + if (useMassTransitRabbitMq) - massTransit.UseRabbitMq(rabbitMqConnectionString); + { + massTransit.UseRabbitMq(rabbitMqConnectionString, rabbit => rabbit.ConfigureServiceBus = bus => { bus.PrefetchCount = 4; }); + } }); } diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs index 88413926b..1f801443b 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs @@ -1,8 +1,6 @@ using Elsa.Features.Services; using Elsa.MassTransit.AzureServiceBus.Features; -using Elsa.MassTransit.AzureServiceBus.Options; using Elsa.MassTransit.Features; -using Elsa.MassTransit.Options; using JetBrains.Annotations; // ReSharper disable once CheckNamespace @@ -17,7 +15,7 @@ public static class ModuleExtensions /// /// Enable and configure the Azure Service Bus transport for MassTransit. /// - public static MassTransitFeature UseAzureServiceBus(this MassTransitFeature feature, string? connectionString) + public static MassTransitFeature UseAzureServiceBus(this MassTransitFeature feature, string? connectionString, Action? configure = default) { feature.Module.Configure((Action)Configure); return feature; @@ -25,6 +23,7 @@ public static class ModuleExtensions void Configure(AzureServiceBusFeature bus) { bus.ConnectionString = connectionString; + configure?.Invoke(bus); } } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index b10a3bb54..5050a7c16 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -1,19 +1,13 @@ -using Azure.Messaging.ServiceBus.Administration; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; using Elsa.MassTransit.Features; -using Elsa.MassTransit.Messages; -using Elsa.MassTransit.Options; -using Elsa.Workflows.Runtime.Activities; using MassTransit; -using MassTransit.Configuration; namespace Elsa.MassTransit.AzureServiceBus.Features; -/// /// Configures MassTransit to use the Azure Service Bus transport. -/// +/// See https://masstransit.io/documentation/configuration/transports/azure-service-bus [DependsOn(typeof(MassTransitFeature))] public class AzureServiceBusFeature : FeatureBase { @@ -21,11 +15,14 @@ public class AzureServiceBusFeature : FeatureBase public AzureServiceBusFeature(IModule module) : base(module) { } + + /// An Azure Service Bus connection string. + public string? ConnectionString { get; set; } /// - /// An Azure Service Bus connection string. + /// A delegate that configures the Azure Service Bus transport options. /// - public string? ConnectionString { get; set; } + public Action? ConfigureServiceBus { get; set; } /// public override void Configure() @@ -42,6 +39,7 @@ public class AzureServiceBusFeature : FeatureBase serviceBus.Host(ConnectionString); serviceBus.UseServiceBusMessageScheduler(); + ConfigureServiceBus?.Invoke(serviceBus); serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); }); }; diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs deleted file mode 100644 index 26df271c5..000000000 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs +++ /dev/null @@ -1,8 +0,0 @@ -namespace Elsa.MassTransit.AzureServiceBus.Options; - -/// -/// Provides settings to the RabbitMQ broker for MassTransit. -/// -public class AzureServiceBusOptions -{ -} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs index 432ecce57..f0d5ebee3 100644 --- a/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs @@ -1,8 +1,6 @@ using Elsa.Features.Services; using Elsa.MassTransit.Features; -using Elsa.MassTransit.Options; using Elsa.MassTransit.RabbitMq.Features; -using Elsa.MassTransit.RabbitMq.Options; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; @@ -15,30 +13,36 @@ public static class ModuleExtensions /// /// Enable and configure the RabbitMQ transport for MassTransit. /// - public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) => feature.UseRabbitMq(new Uri(connectionString), null); - - /// - /// Enable and configure the RabbitMQ transport for MassTransit. - /// - public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri connectionString) => feature.UseRabbitMq(connectionString, null); - - /// - /// Enable and configure the RabbitMQ transport for MassTransit. - /// - public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, RabbitMqOptions options) => feature.UseRabbitMq(null, options); - - /// - /// Enable and configure the RabbitMQ transport for MassTransit. - /// - private static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri? connectionString, RabbitMqOptions? options) + public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) { - feature.Module.Configure((Action) Configure); + feature.Module.Configure((Action)Configure); return feature; void Configure(RabbitMqServiceBusFeature bus) { bus.ConnectionString = connectionString; - bus.Options = options; } } + + /// + /// Enable and configure the RabbitMQ transport for MassTransit. + /// + public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Action configure) + { + feature.Module.Configure(configure); + return feature; + } + + /// + /// Enable and configure the RabbitMQ transport for MassTransit. + /// + public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString, Action configure) + { + feature.Module.Configure(rabbitMqFeature => + { + rabbitMqFeature.ConnectionString = connectionString; + configure(rabbitMqFeature); + }); + return feature; + } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs index c488f2994..af1271514 100644 --- a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs @@ -2,9 +2,8 @@ using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; using Elsa.MassTransit.Features; -using Elsa.MassTransit.Options; -using Elsa.MassTransit.RabbitMq.Options; using MassTransit; +using Microsoft.Extensions.DependencyInjection; namespace Elsa.MassTransit.RabbitMq.Features; @@ -19,38 +18,37 @@ public class RabbitMqServiceBusFeature : FeatureBase { } - /// /// A RabbitMQ connection string. - /// - public Uri? ConnectionString { get; set; } + public string? ConnectionString { get; set; } + /// Configures the RabbitMQ transport options. + public Action? TransportOptions { get; set; } + /// - /// RabbitMQ options. + /// Configures the RabbitMQ bus. /// - public RabbitMqOptions? Options { get; set; } + public Action? ConfigureServiceBus { get; set; } /// public override void Configure() { Module.Configure().BusConfigurator = configure => { - configure.UsingRabbitMq((context, configurator) => + configure.UsingRabbitMq((context, serviceBus) => { - if (ConnectionString != null) - { - configurator.Host(ConnectionString); - } - else if (Options != null) - { - configurator.Host(Options.Host, h => - { - h.Username(Options.Username); - h.Password(Options.Password); - }); - } + if (!string.IsNullOrEmpty(ConnectionString)) + serviceBus.Host(ConnectionString); - configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + ConfigureServiceBus?.Invoke(serviceBus); + + serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); }); }; } + + /// + public override void Apply() + { + if (TransportOptions != null) Services.Configure(TransportOptions); + } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Options/RabbitMqOptions.cs b/src/modules/Elsa.MassTransit.RabbitMq/Options/RabbitMqOptions.cs deleted file mode 100644 index 17b26851e..000000000 --- a/src/modules/Elsa.MassTransit.RabbitMq/Options/RabbitMqOptions.cs +++ /dev/null @@ -1,11 +0,0 @@ -namespace Elsa.MassTransit.RabbitMq.Options; - -/// -/// Provides settings to the RabbitMQ broker for MassTransit. -/// -public class RabbitMqOptions -{ - public string Host { get; set; } = default!; - public string Username { get; set; } = default!; - public string Password { get; set; } = default!; -} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchWorkflowRequestConsumerDefinition.cs b/src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchWorkflowRequestConsumerDefinition.cs new file mode 100644 index 000000000..687af34b7 --- /dev/null +++ b/src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchWorkflowRequestConsumerDefinition.cs @@ -0,0 +1,28 @@ +using Elsa.MassTransit.Consumers; +using Elsa.MassTransit.Options; +using MassTransit; +using Microsoft.Extensions.Options; + +namespace Elsa.MassTransit.ConsumerDefinitions; + +/// +/// Configures the endpoint for +/// +public class DispatchWorkflowRequestConsumerDefinition : ConsumerDefinition +{ + private readonly IOptions _options; + + /// + public DispatchWorkflowRequestConsumerDefinition(IOptions options) + { + _options = options; + ConcurrentMessageLimit = _options.Value.ConcurrentMessageLimit; + } + + /// + protected override void ConfigureConsumer(IReceiveEndpointConfigurator endpointConfigurator, IConsumerConfigurator consumerConfigurator, IRegistrationContext context) + { + endpointConfigurator.UseMessageRetry(r => r.Interval(5, 1000)); + endpointConfigurator.UseInMemoryOutbox(context); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs b/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs index d6c681b50..80e03d6b1 100644 --- a/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs +++ b/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs @@ -1,6 +1,7 @@ using Elsa.Features.Services; using Elsa.MassTransit.Features; using Elsa.MassTransit.Implementations; +using Elsa.MassTransit.Models; using MassTransit; // ReSharper disable once CheckNamespace @@ -18,14 +19,26 @@ public static class MassTransitFeatureExtensions /// Registers the specified type for MassTransit service bus consumer discovery. /// public static MassTransitFeature AddConsumer(this MassTransitFeature feature) where T : IConsumer => feature.AddConsumer(typeof(T)); - + /// /// Registers the specified type for MassTransit service bus consumer discovery. /// - public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type) + /// The consumer type. + /// The consumer definition type. + public static MassTransitFeature AddConsumer(this MassTransitFeature feature) + where T : IConsumer + where TDefinition : IConsumerDefinition { - var types = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); - types.Add(type); + return feature.AddConsumer(typeof(T), typeof(TDefinition)); + } + + /// + /// Registers the specified type for MassTransit service bus consumer discovery. + /// + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type, Type? consumerDefinitionType = default) + { + var types = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); + types.Add(new ConsumerTypeDefinition(type, consumerDefinitionType)); return feature; } @@ -33,7 +46,7 @@ public static class MassTransitFeatureExtensions /// Registers a message type which is to be used by the to dynamically provide activities to send and receive these messages. /// public static MassTransitFeature AddMessageType(this MassTransitFeature feature) where T : class => feature.AddMessageType(typeof(T)); - + /// /// Registers a message type which is to be used by the to dynamically provide activities to send and receive these messages. /// @@ -43,12 +56,12 @@ public static class MassTransitFeatureExtensions types.Add(type); return feature; } - + /// /// Returns all collected consumer types. /// - internal static IEnumerable GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); - + internal static IEnumerable GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); + /// /// Returns all collected message types. /// diff --git a/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs index 1b696d661..c0b33e8a0 100644 --- a/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs @@ -24,4 +24,17 @@ public static class ModuleExtensions module.Configure(massTransit => massTransit.AddConsumer()); return module; } + + /// + /// Registers the specified consumer and consumer definition with MassTransit. + /// + public static IModule AddMassTransitConsumer(this IModule module) + where T : IConsumer + where TDefinition : IConsumerDefinition + { + module.Configure(massTransit => massTransit.AddConsumer()); + return module; + } + + } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs index 4c3781cf4..101edf1bc 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs @@ -6,12 +6,14 @@ using Elsa.Features.Abstractions; using Elsa.Features.Services; using Elsa.MassTransit.Consumers; using Elsa.MassTransit.Implementations; +using Elsa.MassTransit.Models; using Elsa.MassTransit.Options; using Elsa.Workflows.Core.Attributes; using Elsa.Workflows.Core.Serialization.Converters; using Elsa.Workflows.Management.Models; using Elsa.Workflows.Management.Options; using MassTransit; +using MassTransit.Configuration; using MassTransit.Serialization; using Microsoft.Extensions.DependencyInjection; @@ -27,35 +29,42 @@ public class MassTransitFeature : FeatureBase { } + /// The number of messages to prefetch. + public int? PrefetchCount { get; set; } + /// - /// A delegate that can be set to configure MassTransit's . + /// A delegate that can be set to configure MassTransit's . Used by transport-level features such as AzureServiceBusFeature and RabbitMqServiceBusFeature. /// public Action? BusConfigurator { get; set; } /// public override void Configure() { - BusConfigurator ??= configure => - { - configure.UsingInMemory((context, configurator) => - { - configurator.ConfigureEndpoints(context); - }); - }; - } /// public override void Apply() { var messageTypes = this.GetMessages(); - + + Services.Configure(x => { }); Services.AddActivityProvider(); - AddMassTransit(BusConfigurator); + + var busConfigurator = BusConfigurator ??= configure => + { + configure.UsingInMemory((context, configurator) => + { + configurator.ConfigureEndpoints(context); + + if (PrefetchCount != null) + configurator.PrefetchCount = PrefetchCount.Value; + }); + }; + AddMassTransit(busConfigurator); // Add collected message types to options. Services.Configure(options => options.MessageTypes = new HashSet(messageTypes)); - + // Add collected message types as available variable types. Services.Configure(options => { @@ -69,31 +78,33 @@ public class MassTransitFeature : FeatureBase options.VariableDescriptors.Add(new VariableDescriptor(messageType, category, description)); } }); - + // Configure message serializer. SystemTextJsonMessageSerializer.Options.Converters.Add(new TypeJsonConverter(WellKnownTypeRegistry.CreateDefault())); } - + /// /// Adds MassTransit to the service container and registers all collected assemblies for discovery of consumers. /// - private void AddMassTransit(Action? config) + private void AddMassTransit(Action busConfigurator) { // For each message type, create a concrete WorkflowMessageConsumer. var workflowMessageConsumerType = typeof(WorkflowMessageConsumer<>); - var workflowMessageConsumers = this.GetMessages().Select(x => workflowMessageConsumerType.MakeGenericType(x)); + var workflowMessageConsumers = this.GetMessages().Select(x => new ConsumerTypeDefinition(workflowMessageConsumerType.MakeGenericType(x))); // Concatenate the manually registered consumers with the workflow message consumers. - var consumerTypes = this.GetConsumers().Concat(workflowMessageConsumers).ToArray(); + var consumerTypeDefinitions = this.GetConsumers().Concat(workflowMessageConsumers).ToArray(); Services.AddMassTransit(bus => { bus.SetKebabCaseEndpointNameFormatter(); - bus.AddConsumers(consumerTypes); - config?.Invoke(bus); + foreach (var definition in consumerTypeDefinitions) + bus.AddConsumer(definition.ConsumerType, definition.ConsumerDefinitionType); + + busConfigurator(bus); }); - + Services.AddOptions().Configure(options => { // Wait until the bus is started before returning from IHostedService.StartAsync. diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs index 4e45ab77a..d09d9d62b 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs @@ -2,9 +2,11 @@ using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.MassTransit.ConsumerDefinitions; using Elsa.MassTransit.Consumers; using Elsa.MassTransit.Implementations; using Elsa.MassTransit.Messages; +using Elsa.MassTransit.Options; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Features; using MassTransit; @@ -27,13 +29,15 @@ public class MassTransitWorkflowDispatcherFeature : FeatureBase /// public override void Configure() { - Module.AddMassTransitConsumer(); + Module.AddMassTransitConsumer(); Module.Configure(f => f.WorkflowDispatcher = sp => sp.GetRequiredService()); } /// public override void Apply() { + Services.AddOptions(); + var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer(); var queueAddress = new Uri($"queue:elsa-{queueName}"); EndpointConvention.Map(queueAddress); diff --git a/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs b/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs new file mode 100644 index 000000000..6d05b7051 --- /dev/null +++ b/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs @@ -0,0 +1,3 @@ +namespace Elsa.MassTransit.Models; + +internal record ConsumerTypeDefinition(Type ConsumerType, Type? ConsumerDefinitionType = default); \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs b/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs new file mode 100644 index 000000000..b848b2382 --- /dev/null +++ b/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs @@ -0,0 +1,10 @@ +using Elsa.MassTransit.ConsumerDefinitions; + +namespace Elsa.MassTransit.Options; + +/// Provides options to the +public class MassTransitWorkflowDispatcherOptions +{ + /// The number of concurrent messages to process. + public int? ConcurrentMessageLimit { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs index 11e3a05cb..63dfd4283 100644 --- a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs @@ -132,7 +132,7 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime var client = _cluster.GetNamedWorkflowGrain(workflowInstanceId); var response = await client.Start(request, options.CancellationTokens.SystemCancellationToken); - return _workflowExecutionResultMapper.Map(response!); + return _workflowExecutionResultMapper.Map(response!); } ///