From 3611d527d45f1ecd34dcc8bab47acc8828662c72 Mon Sep 17 00:00:00 2001 From: TekH Date: Mon, 27 Jul 2026 13:03:39 +0200 Subject: [PATCH] Refactor RabbitMQ functionality to new project Moved RabbitMQ-related functionality from the `DigitalData.MessagingService.Infrastructure` project to a new dedicated project/namespace `DigitalData.MessagingService.RabbitMQ`. - Updated `DependencyInjection.cs` to use the new namespace. - Added a project reference to `RabbitMQ.csproj` in the `Infrastructure.csproj` file. - Removed `RabbitMqConfiguration.cs` and `RabbitMqConnectionFactory.cs` from the `Infrastructure` project. - Updated namespaces in `OutgoingEmailConsumer.cs`, `OutgoingEmailPublisher.cs`, and `AsyncInitWorker.cs` to use `DigitalData.MessagingService.RabbitMQ`. This refactor improves modularity, maintainability, and separation of concerns by isolating RabbitMQ functionality in its own project. --- .../DependencyInjection.cs | 2 +- ...ata.MessagingService.Infrastructure.csproj | 5 ++ .../Messaging/RabbitMqConfiguration.cs | 54 --------------- .../Queue/OutgoingEmailConsumer.cs | 2 +- .../Queue/OutgoingEmailPublisher.cs | 2 +- .../Queue/RabbitMqConnectionFactory.cs | 67 ------------------- .../Services/Background/AsyncInitWorker.cs | 1 + 7 files changed, 9 insertions(+), 124 deletions(-) delete mode 100644 src/DigitalData.MessagingService.Infrastructure/Messaging/RabbitMqConfiguration.cs delete mode 100644 src/DigitalData.MessagingService.Infrastructure/Queue/RabbitMqConnectionFactory.cs diff --git a/src/DigitalData.MessagingService.Infrastructure/DependencyInjection.cs b/src/DigitalData.MessagingService.Infrastructure/DependencyInjection.cs index 374f69c..a3ee67a 100644 --- a/src/DigitalData.MessagingService.Infrastructure/DependencyInjection.cs +++ b/src/DigitalData.MessagingService.Infrastructure/DependencyInjection.cs @@ -1,8 +1,8 @@ using DigitalData.MessagingService.Application.Common.Interfaces; -using DigitalData.MessagingService.Infrastructure.Messaging; using DigitalData.MessagingService.Infrastructure.Queue; using DigitalData.MessagingService.Infrastructure.Services; using DigitalData.MessagingService.Infrastructure.Services.Background; +using DigitalData.MessagingService.RabbitMQ; using Microsoft.AspNetCore.DataProtection; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; diff --git a/src/DigitalData.MessagingService.Infrastructure/DigitalData.MessagingService.Infrastructure.csproj b/src/DigitalData.MessagingService.Infrastructure/DigitalData.MessagingService.Infrastructure.csproj index f8fa469..a0ed34f 100644 --- a/src/DigitalData.MessagingService.Infrastructure/DigitalData.MessagingService.Infrastructure.csproj +++ b/src/DigitalData.MessagingService.Infrastructure/DigitalData.MessagingService.Infrastructure.csproj @@ -9,6 +9,7 @@ + @@ -33,4 +34,8 @@ + + + + diff --git a/src/DigitalData.MessagingService.Infrastructure/Messaging/RabbitMqConfiguration.cs b/src/DigitalData.MessagingService.Infrastructure/Messaging/RabbitMqConfiguration.cs deleted file mode 100644 index b1e9794..0000000 --- a/src/DigitalData.MessagingService.Infrastructure/Messaging/RabbitMqConfiguration.cs +++ /dev/null @@ -1,54 +0,0 @@ -namespace DigitalData.MessagingService.Infrastructure.Messaging; - -/// -/// Configuration for RabbitMQ connection -/// -public class RabbitMqConfiguration -{ - /// - /// Configuration section name in appsettings.json - /// - public const string SectionName = "RabbitMQ"; - - /// - /// RabbitMQ server hostname - /// - public string HostName { get; set; } = "localhost"; - - /// - /// RabbitMQ AMQP port (default: 5672) - /// - public int Port { get; set; } = 5672; - - /// - /// RabbitMQ username - /// - public string UserName { get; set; } = "guest"; - - /// - /// RabbitMQ password - /// - public string Password { get; set; } = "guest"; - - /// - /// Virtual host (default: /) - /// - public string VirtualHost { get; set; } = "/"; - - /// - /// Enable automatic recovery on connection failure - /// - public bool AutomaticRecoveryEnabled { get; set; } = true; - - /// - /// Network recovery interval in seconds - /// - public int NetworkRecoveryIntervalSeconds { get; set; } = 10; - - public string QueueName { get; set; } = null!; - public string ExchangeName { get; set; } = null!; - public string RoutingKey { get; set; } = null!; - public string DlqQueueName { get; set; } = null!; - public string DlqExchangeName { get; set; } = null!; - public string DlqRoutingKey { get; set; } = null!; -} diff --git a/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailConsumer.cs b/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailConsumer.cs index 7258f63..26fa488 100644 --- a/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailConsumer.cs +++ b/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailConsumer.cs @@ -1,7 +1,7 @@ using System.Text; using System.Text.Json; using DigitalData.MessagingService.Application.Common.Interfaces; -using DigitalData.MessagingService.Infrastructure.Messaging; +using DigitalData.MessagingService.RabbitMQ; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using RabbitMQ.Client; diff --git a/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailPublisher.cs b/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailPublisher.cs index 9578388..545504c 100644 --- a/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailPublisher.cs +++ b/src/DigitalData.MessagingService.Infrastructure/Queue/OutgoingEmailPublisher.cs @@ -1,13 +1,13 @@ using System.Text; using System.Text.Json; using DigitalData.MessagingService.Application.Common.Interfaces; -using DigitalData.MessagingService.Infrastructure.Messaging; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using RabbitMQ.Client; using RabbitMQ.Client.Events; using DigitalData.MessagingService.Application.Common.Events; using DevExpress.CodeParser; +using DigitalData.MessagingService.RabbitMQ; namespace DigitalData.MessagingService.Infrastructure.Queue; diff --git a/src/DigitalData.MessagingService.Infrastructure/Queue/RabbitMqConnectionFactory.cs b/src/DigitalData.MessagingService.Infrastructure/Queue/RabbitMqConnectionFactory.cs deleted file mode 100644 index 1a3cae2..0000000 --- a/src/DigitalData.MessagingService.Infrastructure/Queue/RabbitMqConnectionFactory.cs +++ /dev/null @@ -1,67 +0,0 @@ -using DigitalData.MessagingService.Infrastructure.Messaging; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; -using RabbitMQ.Client; - -namespace DigitalData.MessagingService.Infrastructure.Queue; - -public sealed class RabbitMqConnectionFactory : IAsyncDisposable -{ - private readonly RabbitMqConfiguration _config; - - private readonly ILogger? _logger; - - private readonly Lazy> _lazyConnectionProvider; - - private IConnection? _connection = null; - - private CancellationToken? _cancellationToken; - - public CancellationToken CancellationToken => _cancellationToken - ?? throw new InvalidOperationException("RabbitMqConnectionFactory is not initialized. Call InitAsync() before using this method."); - - public Task GetConnectionAsync() - { - if (_cancellationToken is not null) - return _lazyConnectionProvider.Value; - else - throw new InvalidOperationException("RabbitMqConnectionFactory is not initialized. Call InitAsync() before using this method."); - } - - public RabbitMqConnectionFactory(IOptions config) - { - _config = config.Value; - _lazyConnectionProvider = new (async () => - { - _cancellationToken?.ThrowIfCancellationRequested(); - var factory = new ConnectionFactory - { - HostName = _config.HostName, - Port = _config.Port, - UserName = _config.UserName, - Password = _config.Password, - VirtualHost = _config.VirtualHost, - AutomaticRecoveryEnabled = _config.AutomaticRecoveryEnabled, - NetworkRecoveryInterval = TimeSpan.FromSeconds(_config.NetworkRecoveryIntervalSeconds), - }; - - _connection = await factory.CreateConnectionAsync((CancellationToken)_cancellationToken!); - return _connection; - }); - } - - public async Task InitAsync(CancellationToken cancellationToken = default) - { - _cancellationToken = cancellationToken; - _ = await GetConnectionAsync(); - } - - public async ValueTask DisposeAsync() - { - if (_connection is not null) - { - await _connection.CloseAsync(); - await _connection.DisposeAsync(); - } - } -} diff --git a/src/DigitalData.MessagingService.Infrastructure/Services/Background/AsyncInitWorker.cs b/src/DigitalData.MessagingService.Infrastructure/Services/Background/AsyncInitWorker.cs index 5be68d5..1b1cf00 100644 --- a/src/DigitalData.MessagingService.Infrastructure/Services/Background/AsyncInitWorker.cs +++ b/src/DigitalData.MessagingService.Infrastructure/Services/Background/AsyncInitWorker.cs @@ -1,5 +1,6 @@ using DigitalData.MessagingService.Application.Common.Interfaces; using DigitalData.MessagingService.Infrastructure.Queue; +using DigitalData.MessagingService.RabbitMQ; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging;