Merge remote-tracking branch 'origin/patch/3.2.x'

This commit is contained in:
Sipke Schoorstra 2024-06-12 20:23:51 +02:00
commit a7c23d33c3
10 changed files with 80 additions and 75 deletions

View file

@ -52,14 +52,10 @@
<PackageReference Include="Proto.Persistence.Sqlite"/>
<PackageReference Include="Proto.Persistence.SqlServer"/>
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions of Npgsql.-->
<ItemGroup>
<PackageReference Include="Npgsql" VersionOverride="8.0.3" />
</ItemGroup>
<ItemGroup>
<Folder Include="App_Data\"/>
<Folder Include="Workflows\" />
<PackageReference Include="Npgsql" VersionOverride="8.0.3"/>
</ItemGroup>
</Project>

View file

@ -20,7 +20,7 @@ public class MassTransitDistributedCacheFeature(IModule module) : FeatureBase(mo
/// <inheritdoc />
public override void Configure()
{
Module.AddMassTransitConsumer<TriggerChangeTokenSignalConsumer>("elsa-trigger-change-token-signal", true);
Module.AddMassTransitConsumer<TriggerChangeTokenSignalConsumer>("elsa-trigger-change-token-signal", true, true);
Module.Use<DistributedCacheFeature>(feature => feature.WithChangeTokenSignalPublisher(sp => sp.GetRequiredService<MassTransitChangeTokenSignalPublisher>()));
}

View file

@ -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<IOptions<MassTransitOptions>>().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<IApplicationInstanceNameProvider>();
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<MassTransitWorkflowDispatcherFeature>())
configurator.SetupWorkflowDispatcherEndpoints(context);
}
configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
});
};
});
}
/// <inheritdoc />
public override void Apply()
{
@ -119,7 +121,7 @@ public class AzureServiceBusFeature : FeatureBase
var configuration = serviceProvider.GetRequiredService<IConfiguration>();
return configuration.GetConnectionString(options.ConnectionStringOrName) ?? options.ConnectionStringOrName;
}
private void RegisterConsumers(List<ConsumerTypeDefinition> consumers)
{
var subscriptionTopology = (
@ -134,5 +136,4 @@ public class AzureServiceBusFeature : FeatureBase
Services.AddSingleton(new MessageTopologyProvider(subscriptionTopology));
Services.AddNotificationHandler<RemoveOrphanedSubscriptions>();
}
}

View file

@ -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<IOptions<MassTransitOptions>>().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<MassTransitWorkflowDispatcherFeature>())
configurator.SetupWorkflowDispatcherEndpoints(context);
configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
}
configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
});
};
});

View file

@ -18,27 +18,28 @@ public static class MassTransitFeatureExtensions
/// <summary>
/// Registers the specified type for MassTransit service bus consumer discovery.
/// </summary>
public static MassTransitFeature AddConsumer<T>(this MassTransitFeature feature, string? name, bool isTemporary) where T : IConsumer => feature.AddConsumer(typeof(T), name, isTemporary);
public static MassTransitFeature AddConsumer<T>(this MassTransitFeature feature, string? name, bool isTemporary, bool ignoreConsumersDisabled = false) where T : IConsumer =>
feature.AddConsumer(typeof(T), name, isTemporary, ignoreConsumersDisabled: ignoreConsumersDisabled);
/// <summary>
/// Registers the specified type for MassTransit service bus consumer discovery.
/// </summary>
/// <typeparam name="T">The consumer type.</typeparam>
/// <typeparam name="TDefinition">The consumer definition type.</typeparam>
public static MassTransitFeature AddConsumer<T, TDefinition>(this MassTransitFeature feature, string? name, bool isTemporary)
where T : IConsumer
public static MassTransitFeature AddConsumer<T, TDefinition>(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);
}
/// <summary>
/// Registers the specified type for MassTransit service bus consumer discovery.
/// </summary>
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<ConsumerTypeDefinition>());
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
/// <summary>
/// Returns all collected consumer types.
/// </summary>
public static IEnumerable<ConsumerTypeDefinition> GetConsumers(this MassTransitFeature feature) => feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet<ConsumerTypeDefinition>());
public static IEnumerable<ConsumerTypeDefinition> GetConsumers(this MassTransitFeature feature)
{
var definitions = feature.Module.Properties.GetOrAdd(ServiceBusConsumerTypesKey, () => new HashSet<ConsumerTypeDefinition>());
var disableConsumers = feature.DisableConsumers;
return !disableConsumers ? definitions : definitions.Where(x => x.IgnoreConsumersDisabled).ToList();
}
/// <summary>
/// Returns all collected message types.

View file

@ -18,20 +18,20 @@ public static class ModuleExtensions
/// <summary>
/// Registers the specified consumer with MassTransit.
/// </summary>
public static IModule AddMassTransitConsumer<T>(this IModule module, string? name = null, bool isTemporary = false) where T : IConsumer
public static IModule AddMassTransitConsumer<T>(this IModule module, string? name = null, bool isTemporary = false, bool ignoreConsumersDisabled = false) where T : IConsumer
{
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T>(name, isTemporary));
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T>(name, isTemporary, ignoreConsumersDisabled));
return module;
}
/// <summary>
/// Registers the specified consumer and consumer definition with MassTransit.
/// </summary>
public static IModule AddMassTransitConsumer<T, TDefinition>(this IModule module, string? name = null, bool isTemporary = false)
public static IModule AddMassTransitConsumer<T, TDefinition>(this IModule module, string? name = null, bool isTemporary = false, bool ignoreConsumersDisabled = false)
where T : IConsumer
where TDefinition : IConsumerDefinition
{
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T, TDefinition>(name, isTemporary));
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T, TDefinition>(name, isTemporary, ignoreConsumersDisabled));
return module;
}
}

View file

@ -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; }
/// <summary>
@ -132,21 +132,23 @@ public class MassTransitFeature : FeatureBase
{
var options = context.GetRequiredService<IOptions<MassTransitWorkflowDispatcherOptions>>().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<MassTransitWorkflowDispatcherFeature>())
busFactoryConfigurator.SetupWorkflowDispatcherEndpoints(context);
}
busFactoryConfigurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
busFactoryConfigurator.ConfigureJsonSerializerOptions(serializerOptions =>
{
var serializer = context.GetRequiredService<IJsonSerializer>();

View file

@ -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<WorkflowDefinitionEventsConsumer>("elsa-workflow-definition-updates", true);
Module.AddMassTransitConsumer<WorkflowDefinitionEventsConsumer>("elsa-workflow-definition-updates", true, true);
}
/// <inheritdoc />

View file

@ -7,4 +7,5 @@ public record ConsumerTypeDefinition(
Type ConsumerType,
Type? ConsumerDefinitionType = default,
string? Name = null,
bool IsTemporary = false);
bool IsTemporary = false,
bool IgnoreConsumersDisabled = false);

View file

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