From 3b66d7edbbd2952c2b0e26d43a7b02fecad5d64f Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Thu, 25 Jan 2024 14:51:48 +0100 Subject: [PATCH] Made per instance queues temporary --- .../Features/AzureServiceBusFeature.cs | 27 +++++++++++- .../Features/RabbitMqServiceBusFeature.cs | 43 +++++++++++++++---- ...ancelWorkflowsRequestConsumerDefinition.cs | 31 ------------- .../MassTransitFeatureExtensions.cs | 13 +++--- .../Extensions/ModuleExtensions.cs | 8 ++-- .../Features/MassTransitFeature.cs | 7 ++- .../MassTransitWorkflowDispatcherFeature.cs | 2 +- .../Models/ConsumerTypeDefinition.cs | 2 +- .../MassTransitWorkflowDispatcherOptions.cs | 2 + .../Contracts/IInstanceNameRetriever.cs | 12 ++++++ .../Services/RandomInstanceNameRetriever.cs | 19 ++++++++ .../Features/WorkflowRuntimeFeature.cs | 8 ++++ 12 files changed, 120 insertions(+), 54 deletions(-) delete mode 100644 src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchCancelWorkflowsRequestConsumerDefinition.cs create mode 100644 src/modules/Elsa.Workflows.Core/Contracts/IInstanceNameRetriever.cs create mode 100644 src/modules/Elsa.Workflows.Core/Services/RandomInstanceNameRetriever.cs diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index 5050a7c16..d920bf046 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -1,8 +1,14 @@ +using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.MassTransit.Consumers; using Elsa.MassTransit.Features; +using Elsa.MassTransit.Options; +using Elsa.Workflows.Contracts; using MassTransit; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; namespace Elsa.MassTransit.AzureServiceBus.Features; @@ -31,15 +37,34 @@ public class AzureServiceBusFeature : FeatureBase { massTransitFeature.BusConfigurator = configure => { + var tempConsumers = massTransitFeature.GetConsumers() + .Where(c => c.IsTemporary) + .ToList(); + configure.AddServiceBusMessageScheduler(); + configure.AddConsumers(tempConsumers.Select(c => c.ConsumerType).ToArray()); configure.UsingAzureServiceBus((context, serviceBus) => { + var options = context.GetRequiredService>().Value; + var instanceNameRetriever = context.GetRequiredService(); + if (ConnectionString != null) serviceBus.Host(ConnectionString); - serviceBus.UseServiceBusMessageScheduler(); ConfigureServiceBus?.Invoke(serviceBus); + + foreach (var consumer in tempConsumers) + { + //Throw on no known name? + serviceBus.ReceiveEndpoint($"{instanceNameRetriever.GetName()}-{consumer.Name}", configurator => + { + configurator.AutoDeleteOnIdle = options.ShortTermQueueLifetime ?? TimeSpan.FromHours(1); + configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + configurator.ConfigureConsumer(context); + }); + } + serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); }); }; diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs index af1271514..2d3c542e9 100644 --- a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs @@ -1,9 +1,14 @@ +using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.MassTransit.Consumers; using Elsa.MassTransit.Features; +using Elsa.MassTransit.Options; +using Elsa.Workflows.Contracts; using MassTransit; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; namespace Elsa.MassTransit.RabbitMq.Features; @@ -32,18 +37,40 @@ public class RabbitMqServiceBusFeature : FeatureBase /// public override void Configure() { - Module.Configure().BusConfigurator = configure => + Module.Configure(massTransitFeature => { - configure.UsingRabbitMq((context, serviceBus) => + massTransitFeature.BusConfigurator = configure => { - if (!string.IsNullOrEmpty(ConnectionString)) - serviceBus.Host(ConnectionString); + var tempConsumers = massTransitFeature.GetConsumers() + .Where(c => c.IsTemporary) + .ToList(); + + configure.AddConsumers(tempConsumers.Select(c => c.ConsumerType).ToArray()); + + configure.UsingRabbitMq((context, serviceBus) => + { + var options = context.GetRequiredService>().Value; + var instanceNameRetriever = context.GetRequiredService(); + + if (!string.IsNullOrEmpty(ConnectionString)) + serviceBus.Host(ConnectionString); - ConfigureServiceBus?.Invoke(serviceBus); + ConfigureServiceBus?.Invoke(serviceBus); + + foreach (var consumer in tempConsumers) + { + serviceBus.ReceiveEndpoint($"{instanceNameRetriever.GetName()}-{consumer.Name}", configurator => + { + configurator.QueueExpiration = options.ShortTermQueueLifetime ?? TimeSpan.FromHours(1); + configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + configurator.ConfigureConsumer(context); + }); + } - serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); - }); - }; + serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + }); + }; + }); } /// diff --git a/src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchCancelWorkflowsRequestConsumerDefinition.cs b/src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchCancelWorkflowsRequestConsumerDefinition.cs deleted file mode 100644 index 1523d7b55..000000000 --- a/src/modules/Elsa.MassTransit/ConsumerDefinitions/DispatchCancelWorkflowsRequestConsumerDefinition.cs +++ /dev/null @@ -1,31 +0,0 @@ -using Elsa.MassTransit.Consumers; -using Elsa.MassTransit.Options; -using MassTransit; -using MassTransit.Configuration; -using Microsoft.Extensions.Options; - -namespace Elsa.MassTransit.ConsumerDefinitions; - -/// -/// Configures the endpoint for -/// -public class DispatchCancelWorkflowsRequestConsumerDefinition : ConsumerDefinition -{ - private readonly IOptions _options; - - /// - public DispatchCancelWorkflowsRequestConsumerDefinition(IOptions options) - { - _options = options; - ConcurrentMessageLimit = _options.Value.ConcurrentMessageLimit; - EndpointName = "random-queue-name" + DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); - this.EndpointDefinition!.IsTemporary = true; - } - - /// - protected override void ConfigureConsumer(IReceiveEndpointConfigurator endpointConfigurator, IConsumerConfigurator consumerConfigurator, IRegistrationContext context) - { - endpointConfigurator.UseMessageRetry(r => r.Interval(5, 1000)); - endpointConfigurator.tem - } -} \ 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 80e03d6b1..6e97a479a 100644 --- a/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs +++ b/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs @@ -18,27 +18,28 @@ 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)); + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, string? name, bool isTemporary) where T : IConsumer => feature.AddConsumer(typeof(T), name, isTemporary); /// /// Registers the specified type for MassTransit service bus consumer discovery. /// /// The consumer type. /// The consumer definition type. - public static MassTransitFeature AddConsumer(this MassTransitFeature feature) + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, string? name, bool isTemporary) where T : IConsumer where TDefinition : IConsumerDefinition { - return feature.AddConsumer(typeof(T), typeof(TDefinition)); + return feature.AddConsumer(typeof(T), name, isTemporary, typeof(TDefinition)); } /// /// Registers the specified type for MassTransit service bus consumer discovery. /// - public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type, Type? consumerDefinitionType = default) + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type, string? name, bool isTemporary, + Type? consumerDefinitionType = default) { var types = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); - types.Add(new ConsumerTypeDefinition(type, consumerDefinitionType)); + types.Add(new ConsumerTypeDefinition(type, consumerDefinitionType, name, isTemporary)); return feature; } @@ -60,7 +61,7 @@ public static class MassTransitFeatureExtensions /// /// Returns all collected consumer types. /// - internal static IEnumerable GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); + public 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 b1926db21..144315192 100644 --- a/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs @@ -18,20 +18,20 @@ public static class ModuleExtensions /// /// Registers the specified consumer with MassTransit. /// - public static IModule AddMassTransitConsumer(this IModule module) where T : IConsumer + public static IModule AddMassTransitConsumer(this IModule module, string? name = null, bool isTemporary = false) where T : IConsumer { - module.Configure(massTransit => massTransit.AddConsumer()); + module.Configure(massTransit => massTransit.AddConsumer(name, isTemporary)); return module; } /// /// Registers the specified consumer and consumer definition with MassTransit. /// - public static IModule AddMassTransitConsumer(this IModule module) + public static IModule AddMassTransitConsumer(this IModule module, string? name = null, bool isTemporary = false) where T : IConsumer where TDefinition : IConsumerDefinition { - module.Configure(massTransit => massTransit.AddConsumer()); + module.Configure(massTransit => massTransit.AddConsumer(name, isTemporary)); return module; } diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs index d30250b18..dd58c3bd5 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs @@ -48,7 +48,7 @@ public class MassTransitFeature : FeatureBase Services.Configure(x => { }); Services.AddActivityProvider(); - + var busConfigurator = BusConfigurator ??= configure => { configure.UsingInMemory((context, configurator) => @@ -92,7 +92,10 @@ public class MassTransitFeature : FeatureBase var workflowMessageConsumers = this.GetMessages().Select(x => new ConsumerTypeDefinition(workflowMessageConsumerType.MakeGenericType(x))); // Concatenate the manually registered consumers with the workflow message consumers. - var consumerTypeDefinitions = this.GetConsumers().Concat(workflowMessageConsumers).ToArray(); + var consumerTypeDefinitions = this.GetConsumers() + //Temporary queues require implementation specific variables which will be handled in their respective projects + .Where(c => c.IsTemporary == false) + .Concat(workflowMessageConsumers).ToArray(); Services.AddMassTransit(bus => { diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs index c58a652cd..999e1e4ee 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs @@ -35,7 +35,7 @@ public class MassTransitWorkflowDispatcherFeature : FeatureBase public override void Configure() { Module.AddMassTransitConsumer(); - Module.AddMassTransitConsumer(); + Module.AddMassTransitConsumer("dispatch-cancel-workflow", true); Module.Configure(f => f.WorkflowDispatcher = sp => sp.GetRequiredService()); } diff --git a/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs b/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs index 6d05b7051..80a0fc81c 100644 --- a/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs +++ b/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs @@ -1,3 +1,3 @@ namespace Elsa.MassTransit.Models; -internal record ConsumerTypeDefinition(Type ConsumerType, Type? ConsumerDefinitionType = default); \ No newline at end of file +public record ConsumerTypeDefinition(Type ConsumerType, Type? ConsumerDefinitionType = default, string? Name = null, bool IsTemporary = false); \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs b/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs index b848b2382..6f01c967e 100644 --- a/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs +++ b/src/modules/Elsa.MassTransit/Options/MassTransitWorkflowDispatcherOptions.cs @@ -5,6 +5,8 @@ namespace Elsa.MassTransit.Options; /// Provides options to the public class MassTransitWorkflowDispatcherOptions { + /// The TTL of queues that are seen as short lived (typically queues that are created per running instance). + public TimeSpan? ShortTermQueueLifetime { get; set; } /// The number of concurrent messages to process. public int? ConcurrentMessageLimit { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IInstanceNameRetriever.cs b/src/modules/Elsa.Workflows.Core/Contracts/IInstanceNameRetriever.cs new file mode 100644 index 000000000..a75411068 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Contracts/IInstanceNameRetriever.cs @@ -0,0 +1,12 @@ +namespace Elsa.Workflows.Contracts; + +/// +/// Retrieves a name of the current instance. +/// +public interface IInstanceNameRetriever +{ + /// + /// Returns a name for the instance. + /// + public string GetName(); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/RandomInstanceNameRetriever.cs b/src/modules/Elsa.Workflows.Core/Services/RandomInstanceNameRetriever.cs new file mode 100644 index 000000000..493cf5b52 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Services/RandomInstanceNameRetriever.cs @@ -0,0 +1,19 @@ +using Elsa.Workflows.Contracts; + +namespace Elsa.Workflows.Services; + +/// +/// Returns a randomly generated instance name. +/// +public class RandomInstanceNameRetriever : IInstanceNameRetriever +{ + private readonly string _instanceName; + + public RandomInstanceNameRetriever(RandomLongIdentityGenerator identityGenerator) + { + _instanceName = identityGenerator.GenerateId(); + } + + /// + public string GetName() => _instanceName; +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 7dd355726..0d7488fb4 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -16,6 +16,7 @@ using Elsa.Workflows.Runtime.Options; using Elsa.Workflows.Runtime.Providers; using Elsa.Workflows.Runtime.Services; using Elsa.Workflows.Runtime.Stores; +using Elsa.Workflows.Services; using Medallion.Threading; using Medallion.Threading.FileSystem; using Microsoft.Extensions.DependencyInjection; @@ -98,6 +99,11 @@ public class WorkflowRuntimeFeature : FeatureBase /// public Func BackgroundActivityScheduler { get; set; } = sp => ActivatorUtilities.CreateInstance(sp); + /// + /// A factory that instantiates an . + /// + public Func InstanceNameRetriever { get; set; } = sp => ActivatorUtilities.CreateInstance(sp); + /// /// A delegate to configure the . /// @@ -171,6 +177,8 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped(WorkflowExecutionContextStore) .AddSingleton(RunTaskDispatcher) .AddSingleton(BackgroundActivityScheduler) + .AddSingleton(InstanceNameRetriever) + .AddSingleton() .AddScoped() .AddScoped() .AddScoped()