elsa-core/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs
jdevillard 2195e43709
LogRecord storage at differents levels (#4911)
* Adding PersistenceStrategy and filter mapping for ActivityExecutionRecord

* Add default PersistenceStrategy provider/service

update ActivityMapper Implementation

* Fix Forgot to use the  Default Persistence from the Server configuration

* refactor the configuration of persistence in WorkflowManagementFeature

* rename PersistenceStrategy to LogPersistenceMode

* use const to defined LogPersistence Key in json and rename the key to logPersistenceMode

* Refactor PersistenceTab and update project references

Updated various aspects of the PersistenceTab class and its functionality to improve code quality and readability. Simplified the handling of persistence configurations and simplified the use of properties. Transitioned project reference for Elsa.Api.Client from package reference to direct project reference for better development experience in Elsa.Studio.Core.

* - Change how to get the Default Persistence Mode for an Activity.

/**
Because the entire workflow is considered as an activity, the schema must be the same
ie with
"logPersistenceMode": {
               "default": "default",
}
**/

- fix logic to get the default persistence mode working for the whole activity.

* fix default change value

---------

Co-authored-by: Jérémie DEVILLARD <jdevillard@users.noreply.github.com>
Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>
2024-03-30 21:22:35 +01:00

93 lines
3.5 KiB
C#

using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.Hosting.Management.Contracts;
using Elsa.Hosting.Management.Features;
using Elsa.MassTransit.Consumers;
using Elsa.MassTransit.Extensions;
using Elsa.MassTransit.Features;
using Elsa.MassTransit.Options;
using MassTransit;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
namespace Elsa.MassTransit.RabbitMq.Features;
/// <summary>
/// Configures MassTransit to use the RabbitMQ transport.
/// </summary>
[DependsOn(typeof(MassTransitFeature))]
[DependsOn(typeof(InstanceManagementFeature))]
public class RabbitMqServiceBusFeature : FeatureBase
{
/// <inheritdoc />
public RabbitMqServiceBusFeature(IModule module) : base(module)
{
}
/// A RabbitMQ connection string.
public string? ConnectionString { get; set; }
/// <summary>
/// Configures the RabbitMQ transport options.
/// </summary>
public Action<RabbitMqTransportOptions>? TransportOptions { get; set; }
/// <summary>
/// Configures the RabbitMQ bus.
/// </summary>
public Action<IRabbitMqBusFactoryConfigurator>? ConfigureServiceBus { get; set; }
/// <inheritdoc />
public override void Configure()
{
Module.Configure<MassTransitFeature>(massTransitFeature =>
{
massTransitFeature.BusConfigurator = configure =>
{
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<MassTransitWorkflowDispatcherOptions>>().Value;
var instanceNameProvider = context.GetRequiredService<IApplicationInstanceNameProvider>();
if (!string.IsNullOrEmpty(ConnectionString))
configurator.Host(ConnectionString);
ConfigureServiceBus?.Invoke(configurator);
foreach (var consumer in temporaryConsumers)
{
configure.AddConsumer(consumer.ConsumerType).ExcludeFromConfigureEndpoints();
configurator.ReceiveEndpoint($"{instanceNameProvider.GetName()}-{consumer.Name}",
endpointConfigurator =>
{
endpointConfigurator.QueueExpiration = options.TemporaryQueueTtl ?? TimeSpan.FromHours(1);
endpointConfigurator.ConcurrentMessageLimit = options.ConcurrentMessageLimit;
endpointConfigurator.ConfigureConsumer(context, consumer.ConsumerType);
});
}
configurator.SetupWorkflowDispatcherEndpoints(context);
configurator.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
});
};
});
}
/// <inheritdoc />
public override void Apply()
{
if (TransportOptions != null) Services.Configure(TransportOptions);
}
}