Refactor MassTransit & AzureServiceBus configuration and delete unnecessary options

AzureServiceBusOption.cs and RabbitMqOptions.cs files were removed as they're no longer needed. The system now defers the configuration of MassTransit and AzureServiceBus to the application layer (Program.cs). The configuration for PrefetchCount is also allowed at the transport-level. A new abstraction, ConsumerTypeDefinition, collects all the consumer types with their respective definitions if provided. Adjustments have been made in related extension methods and classes to reflect these changes.
This commit is contained in:
Sipke Schoorstra 2023-12-11 19:59:13 +01:00
parent e9f1aade15
commit 065765307b
15 changed files with 172 additions and 105 deletions

View file

@ -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; });
}
});
}

View file

@ -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
/// <summary>
/// Enable and configure the Azure Service Bus transport for MassTransit.
/// </summary>
public static MassTransitFeature UseAzureServiceBus(this MassTransitFeature feature, string? connectionString)
public static MassTransitFeature UseAzureServiceBus(this MassTransitFeature feature, string? connectionString, Action<AzureServiceBusFeature>? configure = default)
{
feature.Module.Configure((Action<AzureServiceBusFeature>)Configure);
return feature;
@ -25,6 +23,7 @@ public static class ModuleExtensions
void Configure(AzureServiceBusFeature bus)
{
bus.ConnectionString = connectionString;
configure?.Invoke(bus);
}
}
}

View file

@ -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;
/// <summary>
/// Configures MassTransit to use the Azure Service Bus transport.
/// </summary>
/// 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; }
/// <summary>
/// An Azure Service Bus connection string.
/// A delegate that configures the Azure Service Bus transport options.
/// </summary>
public string? ConnectionString { get; set; }
public Action<IServiceBusBusFactoryConfigurator>? ConfigureServiceBus { get; set; }
/// <inheritdoc />
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));
});
};

View file

@ -1,8 +0,0 @@
namespace Elsa.MassTransit.AzureServiceBus.Options;
/// <summary>
/// Provides settings to the RabbitMQ broker for MassTransit.
/// </summary>
public class AzureServiceBusOptions
{
}

View file

@ -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
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) => feature.UseRabbitMq(new Uri(connectionString), null);
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri connectionString) => feature.UseRabbitMq(connectionString, null);
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, RabbitMqOptions options) => feature.UseRabbitMq(null, options);
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
private static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri? connectionString, RabbitMqOptions? options)
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString)
{
feature.Module.Configure((Action<RabbitMqServiceBusFeature>) Configure);
feature.Module.Configure((Action<RabbitMqServiceBusFeature>)Configure);
return feature;
void Configure(RabbitMqServiceBusFeature bus)
{
bus.ConnectionString = connectionString;
bus.Options = options;
}
}
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Action<RabbitMqServiceBusFeature> configure)
{
feature.Module.Configure(configure);
return feature;
}
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString, Action<RabbitMqServiceBusFeature> configure)
{
feature.Module.Configure<RabbitMqServiceBusFeature>(rabbitMqFeature =>
{
rabbitMqFeature.ConnectionString = connectionString;
configure(rabbitMqFeature);
});
return feature;
}
}

View file

@ -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
{
}
/// <summary>
/// A RabbitMQ connection string.
/// </summary>
public Uri? ConnectionString { get; set; }
public string? ConnectionString { get; set; }
/// Configures the RabbitMQ transport options.
public Action<RabbitMqTransportOptions>? TransportOptions { get; set; }
/// <summary>
/// RabbitMQ options.
/// Configures the RabbitMQ bus.
/// </summary>
public RabbitMqOptions? Options { get; set; }
public Action<IRabbitMqBusFactoryConfigurator>? ConfigureServiceBus { get; set; }
/// <inheritdoc />
public override void Configure()
{
Module.Configure<MassTransitFeature>().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));
});
};
}
/// <inheritdoc />
public override void Apply()
{
if (TransportOptions != null) Services.Configure(TransportOptions);
}
}

View file

@ -1,11 +0,0 @@
namespace Elsa.MassTransit.RabbitMq.Options;
/// <summary>
/// Provides settings to the RabbitMQ broker for MassTransit.
/// </summary>
public class RabbitMqOptions
{
public string Host { get; set; } = default!;
public string Username { get; set; } = default!;
public string Password { get; set; } = default!;
}

View file

@ -0,0 +1,28 @@
using Elsa.MassTransit.Consumers;
using Elsa.MassTransit.Options;
using MassTransit;
using Microsoft.Extensions.Options;
namespace Elsa.MassTransit.ConsumerDefinitions;
/// <summary>
/// Configures the endpoint for <see cref="DispatchWorkflowRequestConsumer"/>
/// </summary>
public class DispatchWorkflowRequestConsumerDefinition : ConsumerDefinition<DispatchWorkflowRequestConsumer>
{
private readonly IOptions<MassTransitWorkflowDispatcherOptions> _options;
/// <inheritdoc />
public DispatchWorkflowRequestConsumerDefinition(IOptions<MassTransitWorkflowDispatcherOptions> options)
{
_options = options;
ConcurrentMessageLimit = _options.Value.ConcurrentMessageLimit;
}
/// <inheritdoc />
protected override void ConfigureConsumer(IReceiveEndpointConfigurator endpointConfigurator, IConsumerConfigurator<DispatchWorkflowRequestConsumer> consumerConfigurator, IRegistrationContext context)
{
endpointConfigurator.UseMessageRetry(r => r.Interval(5, 1000));
endpointConfigurator.UseInMemoryOutbox(context);
}
}

View file

@ -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.
/// </summary>
public static MassTransitFeature AddConsumer<T>(this MassTransitFeature feature) where T : IConsumer => feature.AddConsumer(typeof(T));
/// <summary>
/// Registers the specified type for MassTransit service bus consumer discovery.
/// </summary>
public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type)
/// <typeparam name="T">The consumer type.</typeparam>
/// <typeparam name="TDefinition">The consumer definition type.</typeparam>
public static MassTransitFeature AddConsumer<T, TDefinition>(this MassTransitFeature feature)
where T : IConsumer
where TDefinition : IConsumerDefinition
{
var types = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet<Type>());
types.Add(type);
return feature.AddConsumer(typeof(T), typeof(TDefinition));
}
/// <summary>
/// Registers the specified type for MassTransit service bus consumer discovery.
/// </summary>
public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type, Type? consumerDefinitionType = default)
{
var types = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet<ConsumerTypeDefinition>());
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 <see cref="MassTransitActivityTypeProvider"/> to dynamically provide activities to send and receive these messages.
/// </summary>
public static MassTransitFeature AddMessageType<T>(this MassTransitFeature feature) where T : class => feature.AddMessageType(typeof(T));
/// <summary>
/// Registers a message type which is to be used by the <see cref="MassTransitActivityTypeProvider"/> to dynamically provide activities to send and receive these messages.
/// </summary>
@ -43,12 +56,12 @@ public static class MassTransitFeatureExtensions
types.Add(type);
return feature;
}
/// <summary>
/// Returns all collected consumer types.
/// </summary>
internal static IEnumerable<Type> GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet<Type>());
internal static IEnumerable<ConsumerTypeDefinition> GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet<ConsumerTypeDefinition>());
/// <summary>
/// Returns all collected message types.
/// </summary>

View file

@ -24,4 +24,17 @@ public static class ModuleExtensions
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T>());
return module;
}
/// <summary>
/// Registers the specified consumer and consumer definition with MassTransit.
/// </summary>
public static IModule AddMassTransitConsumer<T, TDefinition>(this IModule module)
where T : IConsumer
where TDefinition : IConsumerDefinition
{
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T, TDefinition>());
return module;
}
}

View file

@ -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; }
/// <summary>
/// A delegate that can be set to configure MassTransit's <see cref="IBusRegistrationConfigurator"/>.
/// A delegate that can be set to configure MassTransit's <see cref="IBusRegistrationConfigurator"/>. Used by transport-level features such as AzureServiceBusFeature and RabbitMqServiceBusFeature.
/// </summary>
public Action<IBusRegistrationConfigurator>? BusConfigurator { get; set; }
/// <inheritdoc />
public override void Configure()
{
BusConfigurator ??= configure =>
{
configure.UsingInMemory((context, configurator) =>
{
configurator.ConfigureEndpoints(context);
});
};
}
/// <inheritdoc />
public override void Apply()
{
var messageTypes = this.GetMessages();
Services.Configure<MassTransitWorkflowDispatcherOptions>(x => { });
Services.AddActivityProvider<MassTransitActivityTypeProvider>();
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<MassTransitActivityOptions>(options => options.MessageTypes = new HashSet<Type>(messageTypes));
// Add collected message types as available variable types.
Services.Configure<ManagementOptions>(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()));
}
/// <summary>
/// Adds MassTransit to the service container and registers all collected assemblies for discovery of consumers.
/// </summary>
private void AddMassTransit(Action<IBusRegistrationConfigurator>? config)
private void AddMassTransit(Action<IBusRegistrationConfigurator> busConfigurator)
{
// For each message type, create a concrete WorkflowMessageConsumer<T>.
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<MassTransitHostOptions>().Configure(options =>
{
// Wait until the bus is started before returning from IHostedService.StartAsync.

View file

@ -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
/// <inheritdoc />
public override void Configure()
{
Module.AddMassTransitConsumer<DispatchWorkflowRequestConsumer>();
Module.AddMassTransitConsumer<DispatchWorkflowRequestConsumer, DispatchWorkflowRequestConsumerDefinition>();
Module.Configure<WorkflowRuntimeFeature>(f => f.WorkflowDispatcher = sp => sp.GetRequiredService<MassTransitWorkflowDispatcher>());
}
/// <inheritdoc />
public override void Apply()
{
Services.AddOptions<MassTransitWorkflowDispatcherOptions>();
var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer<DispatchWorkflowRequestConsumer>();
var queueAddress = new Uri($"queue:elsa-{queueName}");
EndpointConvention.Map<DispatchWorkflowDefinition>(queueAddress);

View file

@ -0,0 +1,3 @@
namespace Elsa.MassTransit.Models;
internal record ConsumerTypeDefinition(Type ConsumerType, Type? ConsumerDefinitionType = default);

View file

@ -0,0 +1,10 @@
using Elsa.MassTransit.ConsumerDefinitions;
namespace Elsa.MassTransit.Options;
/// Provides options to the <see cref="DispatchWorkflowRequestConsumerDefinition"/>
public class MassTransitWorkflowDispatcherOptions
{
/// The number of concurrent messages to process.
public int? ConcurrentMessageLimit { get; set; }
}

View file

@ -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!);
}
/// <inheritdoc />