diff --git a/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj b/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj index d3103fc5a..f92ae959c 100644 --- a/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj +++ b/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj @@ -52,14 +52,10 @@ - + + - - - - - - + diff --git a/src/modules/Elsa.Caching.Distributed.MassTransit/Features/MassTransitDistributedCacheFeature.cs b/src/modules/Elsa.Caching.Distributed.MassTransit/Features/MassTransitDistributedCacheFeature.cs index ea0c09900..f325911b6 100644 --- a/src/modules/Elsa.Caching.Distributed.MassTransit/Features/MassTransitDistributedCacheFeature.cs +++ b/src/modules/Elsa.Caching.Distributed.MassTransit/Features/MassTransitDistributedCacheFeature.cs @@ -20,7 +20,7 @@ public class MassTransitDistributedCacheFeature(IModule module) : FeatureBase(mo /// public override void Configure() { - Module.AddMassTransitConsumer("elsa-trigger-change-token-signal", true); + Module.AddMassTransitConsumer("elsa-trigger-change-token-signal", true, true); Module.Use(feature => feature.WithChangeTokenSignalPublisher(sp => sp.GetRequiredService())); } diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs index 5e70a9713..7b2073d47 100644 --- a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs @@ -63,16 +63,16 @@ public class AzureServiceBusFeature : FeatureBase RegisterConsumers(consumers); configure.AddServiceBusMessageScheduler(); - + // Consumers need to be added before the UsingAzureServiceBus statement to prevent exceptions. foreach (var consumer in temporaryConsumers) configure.AddConsumer(consumer.ConsumerType).ExcludeFromConfigureEndpoints(); configure.UsingAzureServiceBus((context, configurator) => { - if (ConnectionString != null) + if (ConnectionString != null) configurator.Host(ConnectionString); - + var options = context.GetRequiredService>().Value; if (options.PrefetchCount is not null) @@ -80,32 +80,34 @@ public class AzureServiceBusFeature : FeatureBase if (options.MaxAutoRenewDuration is not null) configurator.MaxAutoRenewDuration = options.MaxAutoRenewDuration.Value; configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - + configurator.UseServiceBusMessageScheduler(); - configurator.SetupWorkflowDispatcherEndpoints(context); ConfigureServiceBus?.Invoke(configurator); var instanceNameProvider = context.GetRequiredService(); + foreach (var consumer in temporaryConsumers) + { + var queueName = $"{consumer.Name}-{instanceNameProvider.GetName()}"; + configurator.ReceiveEndpoint(queueName, endpointConfigurator => + { + endpointConfigurator.AutoDeleteOnIdle = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); + endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); + }); + } + if (!massTransitFeature.DisableConsumers) { - foreach (var consumer in temporaryConsumers) - { - var queueName = $"{consumer.Name}-{instanceNameProvider.GetName()}"; - configurator.ReceiveEndpoint(queueName, endpointConfigurator => - { - endpointConfigurator.AutoDeleteOnIdle = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); - endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); - }); - } - - configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + if (Module.HasFeature()) + configurator.SetupWorkflowDispatcherEndpoints(context); } + + configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); }); }; }); } - + /// public override void Apply() { @@ -119,7 +121,7 @@ public class AzureServiceBusFeature : FeatureBase var configuration = serviceProvider.GetRequiredService(); return configuration.GetConnectionString(options.ConnectionStringOrName) ?? options.ConnectionStringOrName; } - + private void RegisterConsumers(List consumers) { var subscriptionTopology = ( @@ -134,5 +136,4 @@ public class AzureServiceBusFeature : FeatureBase Services.AddSingleton(new MessageTopologyProvider(subscriptionTopology)); Services.AddNotificationHandler(); } - } \ 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 6ebd9abb7..064d0a948 100644 --- a/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs +++ b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs @@ -48,11 +48,11 @@ public class RabbitMqServiceBusFeature : FeatureBase var temporaryConsumers = massTransitFeature.GetConsumers() .Where(c => c.IsTemporary) .ToList(); - + // Consumers need to be added before the UsingRabbitMq statement to prevent exceptions. foreach (var consumer in temporaryConsumers) configure.AddConsumer(consumer.ConsumerType).ExcludeFromConfigureEndpoints(); - + configure.UsingRabbitMq((context, configurator) => { var options = context.GetRequiredService>().Value; @@ -64,30 +64,29 @@ public class RabbitMqServiceBusFeature : FeatureBase if (options.PrefetchCount is not null) configurator.PrefetchCount = options.PrefetchCount.Value; configurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - + ConfigureServiceBus?.Invoke(configurator); + foreach (var consumer in temporaryConsumers) + { + configurator.ReceiveEndpoint($"{instanceNameProvider.GetName()}-{consumer.Name}", + endpointConfigurator => + { + endpointConfigurator.QueueExpiration = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); + endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + endpointConfigurator.Durable = false; + endpointConfigurator.AutoDelete = true; + endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); + }); + } + if (!massTransitFeature.DisableConsumers) { - foreach (var consumer in temporaryConsumers) - { - configurator.ReceiveEndpoint($"{instanceNameProvider.GetName()}-{consumer.Name}", - endpointConfigurator => - { - endpointConfigurator.QueueExpiration = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1); - endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - endpointConfigurator.Durable = false; - endpointConfigurator.AutoDelete = true; - endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType); - }); - } - - // Only configure the dispatcher endpoints if the Masstransit Workflow Dispatcher feature is enabled. if (Module.HasFeature()) configurator.SetupWorkflowDispatcherEndpoints(context); - - configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); } + + configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); }); }; }); diff --git a/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs b/src/modules/Elsa.MassTransit/Extensions/MassTransitFeatureExtensions.cs index d71775196..5f1347364 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, string? name, bool isTemporary) where T : IConsumer => feature.AddConsumer(typeof(T), name, isTemporary); + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, string? name, bool isTemporary, bool ignoreConsumersDisabled = false) where T : IConsumer => + feature.AddConsumer(typeof(T), name, isTemporary, ignoreConsumersDisabled: ignoreConsumersDisabled); /// /// Registers the specified type for MassTransit service bus consumer discovery. /// /// The consumer type. /// The consumer definition type. - public static MassTransitFeature AddConsumer(this MassTransitFeature feature, string? name, bool isTemporary) - where T : IConsumer + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, string? name, bool isTemporary, bool ignoreConsumersDisabled = false) + where T : IConsumer where TDefinition : IConsumerDefinition { - return feature.AddConsumer(typeof(T), name, isTemporary, typeof(TDefinition)); + return feature.AddConsumer(typeof(T), name, isTemporary, typeof(TDefinition), ignoreConsumersDisabled); } /// /// Registers the specified type for MassTransit service bus consumer discovery. /// - public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type, string? name, bool isTemporary, Type? consumerDefinitionType = default) + public static MassTransitFeature AddConsumer(this MassTransitFeature feature, Type type, string? name, bool isTemporary, Type? consumerDefinitionType = default, bool ignoreConsumersDisabled = false) { var types = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); - types.Add(new ConsumerTypeDefinition(type, consumerDefinitionType, name, isTemporary)); + types.Add(new ConsumerTypeDefinition(type, consumerDefinitionType, name, isTemporary, ignoreConsumersDisabled)); return feature; } @@ -60,7 +61,12 @@ public static class MassTransitFeatureExtensions /// /// Returns all collected consumer types. /// - public static IEnumerable GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); + public static IEnumerable GetConsumers(this MassTransitFeature feature) + { + var definitions = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet()); + var disableConsumers = feature.DisableConsumers; + return !disableConsumers ? definitions : definitions.Where(x => x.IgnoreConsumersDisabled).ToList(); + } /// /// Returns all collected message types. diff --git a/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs index 97481d740..6f889f772 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, string? name = null, bool isTemporary = false) where T : IConsumer + public static IModule AddMassTransitConsumer(this IModule module, string? name = null, bool isTemporary = false, bool ignoreConsumersDisabled = false) where T : IConsumer { - module.Configure(massTransit => massTransit.AddConsumer(name, isTemporary)); + module.Configure(massTransit => massTransit.AddConsumer(name, isTemporary, ignoreConsumersDisabled)); return module; } /// /// Registers the specified consumer and consumer definition with MassTransit. /// - public static IModule AddMassTransitConsumer(this IModule module, string? name = null, bool isTemporary = false) + public static IModule AddMassTransitConsumer(this IModule module, string? name = null, bool isTemporary = false, bool ignoreConsumersDisabled = false) where T : IConsumer where TDefinition : IConsumerDefinition { - module.Configure(massTransit => massTransit.AddConsumer(name, isTemporary)); + module.Configure(massTransit => massTransit.AddConsumer(name, isTemporary, ignoreConsumersDisabled)); 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 20412485a..efcf866cd 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitFeature.cs @@ -35,7 +35,7 @@ public class MassTransitFeature : FeatureBase /// The number of messages to prefetch. [Obsolete("PrefetchCount has been moved to be included in MassTransitOptions")] public int? PrefetchCount { get; set; } - + public bool DisableConsumers { get; set; } /// @@ -132,21 +132,23 @@ public class MassTransitFeature : FeatureBase { var options = context.GetRequiredService>().Value; - if(!DisableConsumers) + foreach (var consumer in temporaryConsumers) { - foreach (var consumer in temporaryConsumers) + busFactoryConfigurator.ReceiveEndpoint(consumer.Name!, endpoint => { - busFactoryConfigurator.ReceiveEndpoint(consumer.Name!, endpoint => - { - endpoint.ConcurrentMessageLimit = options.ConcurrentMessageLimit; - endpoint.ConfigureConsumer(context, consumer.ConsumerType); - }); - } - - busFactoryConfigurator.SetupWorkflowDispatcherEndpoints(context); - busFactoryConfigurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + endpoint.ConcurrentMessageLimit = options.ConcurrentMessageLimit; + endpoint.ConfigureConsumer(context, consumer.ConsumerType); + }); } - + + if (!DisableConsumers) + { + if (Module.HasFeature()) + busFactoryConfigurator.SetupWorkflowDispatcherEndpoints(context); + } + + busFactoryConfigurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false)); + busFactoryConfigurator.ConfigureJsonSerializerOptions(serializerOptions => { var serializer = context.GetRequiredService(); diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowManagementFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowManagementFeature.cs index df0aa50e1..837c53ecd 100644 --- a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowManagementFeature.cs +++ b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowManagementFeature.cs @@ -22,7 +22,7 @@ public class MassTransitWorkflowManagementFeature(IModule module) : FeatureBase( [RequiresUnreferencedCode("The assembly containing the specified marker type will be scanned for activity types.")] public override void Configure() { - Module.AddMassTransitConsumer("elsa-workflow-definition-updates", true); + Module.AddMassTransitConsumer("elsa-workflow-definition-updates", true, true); } /// diff --git a/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs b/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs index 28ba0b701..391ac9dbc 100644 --- a/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs +++ b/src/modules/Elsa.MassTransit/Models/ConsumerTypeDefinition.cs @@ -7,4 +7,5 @@ public record ConsumerTypeDefinition( Type ConsumerType, Type? ConsumerDefinitionType = default, string? Name = null, - bool IsTemporary = false); \ No newline at end of file + bool IsTemporary = false, + bool IgnoreConsumersDisabled = false); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionActivityRegistryUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionActivityRegistryUpdater.cs index dc9eb96db..e3b8c4a23 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionActivityRegistryUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionActivityRegistryUpdater.cs @@ -17,7 +17,7 @@ public class WorkflowDefinitionActivityRegistryUpdater(WorkflowDefinitionActivit { var descriptors = await provider.GetDescriptorsAsync(cancellationToken); var descriptorToAdd = descriptors - .SingleOrDefault(d => + .FirstOrDefault(d => d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && val.ToString() == workflowDefinitionVersionId); @@ -47,7 +47,7 @@ public class WorkflowDefinitionActivityRegistryUpdater(WorkflowDefinitionActivit var providerDescriptors = registry.ListByProvider(_providerType); var descriptorToRemove = providerDescriptors - .SingleOrDefault(d => + .FirstOrDefault(d => d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && val.ToString() == workflowDefinitionVersionId);