From 6ec61215a17347fa67efd669359b6693b413ccef Mon Sep 17 00:00:00 2001 From: MilePrivate Date: Sun, 9 Aug 2026 11:24:33 +0200 Subject: [PATCH] Add RabbitMQ email worker --- .../BackgroundJobs/OutboxBackgroundService.cs | 3 +- .../Services/Messaging/RabbitMqTopology.cs | 18 +- .../appsettings.Development.json | 6 +- .../Orders/OrderEventRoutingKeys.cs | 10 + .../OrderEventEmailHandlerTests.cs | 100 ++++++ .../OrderProcessing.EmailWorker.Tests.csproj | 26 ++ .../Configuration/RabbitMqOptions.cs | 26 ++ .../Emailing/EmailMessage.cs | 3 + .../Emailing/IEmailSender.cs | 6 + .../Emailing/LoggingEmailSender.cs | 24 ++ .../Messaging/OrderEventEmailHandler.cs | 130 ++++++++ .../UnsupportedIntegrationEventException.cs | 9 + .../OrderProcessing.EmailWorker.csproj | 18 ++ OrderProcessing.EmailWorker/Program.cs | 53 +++ .../Properties/launchSettings.json | 12 + .../RabbitMqEmailConsumer.cs | 301 ++++++++++++++++++ .../appsettings.Development.json | 8 + OrderProcessing.EmailWorker/appsettings.json | 20 ++ OrderProcessingPlatform.slnx | 2 + 19 files changed, 758 insertions(+), 17 deletions(-) create mode 100644 OrderProcessing.Contracts/Orders/OrderEventRoutingKeys.cs create mode 100644 OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs create mode 100644 OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj create mode 100644 OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs create mode 100644 OrderProcessing.EmailWorker/Emailing/EmailMessage.cs create mode 100644 OrderProcessing.EmailWorker/Emailing/IEmailSender.cs create mode 100644 OrderProcessing.EmailWorker/Emailing/LoggingEmailSender.cs create mode 100644 OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs create mode 100644 OrderProcessing.EmailWorker/Messaging/UnsupportedIntegrationEventException.cs create mode 100644 OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj create mode 100644 OrderProcessing.EmailWorker/Program.cs create mode 100644 OrderProcessing.EmailWorker/Properties/launchSettings.json create mode 100644 OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs create mode 100644 OrderProcessing.EmailWorker/appsettings.Development.json create mode 100644 OrderProcessing.EmailWorker/appsettings.json diff --git a/OrderProcessing.Api/BackgroundJobs/OutboxBackgroundService.cs b/OrderProcessing.Api/BackgroundJobs/OutboxBackgroundService.cs index bf36cb0..b6fca05 100644 --- a/OrderProcessing.Api/BackgroundJobs/OutboxBackgroundService.cs +++ b/OrderProcessing.Api/BackgroundJobs/OutboxBackgroundService.cs @@ -3,8 +3,7 @@ namespace OrderProcessing.Api.BackgroundJobs; -public sealed class OutboxBackgroundService - : BackgroundService +public sealed class OutboxBackgroundService : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly OutboxOptions _options; diff --git a/OrderProcessing.Api/Services/Messaging/RabbitMqTopology.cs b/OrderProcessing.Api/Services/Messaging/RabbitMqTopology.cs index c282059..c19f681 100644 --- a/OrderProcessing.Api/Services/Messaging/RabbitMqTopology.cs +++ b/OrderProcessing.Api/Services/Messaging/RabbitMqTopology.cs @@ -4,27 +4,21 @@ namespace OrderProcessing.Api.Services.Messaging; public static class RabbitMqTopology { - public const string OrderCreatedRoutingKey = "order.created"; + public const string OrderCreatedRoutingKey = OrderEventRoutingKeys.Created; - public const string OrderCompletedRoutingKey = "order.completed"; + public const string OrderCompletedRoutingKey = OrderEventRoutingKeys.Completed; - public const string OrderCancelledRoutingKey = "order.cancelled"; + public const string OrderCancelledRoutingKey = OrderEventRoutingKeys.Cancelled; public static string GetRoutingKey(string eventType) { return eventType switch { - var type when type == - typeof(OrderCreatedIntegrationEvent).FullName => - OrderCreatedRoutingKey, + var type when type == typeof(OrderCreatedIntegrationEvent).FullName => OrderCreatedRoutingKey, - var type when type == - typeof(OrderCompletedIntegrationEvent).FullName => - OrderCompletedRoutingKey, + var type when type == typeof(OrderCompletedIntegrationEvent).FullName => OrderCompletedRoutingKey, - var type when type == - typeof(OrderCancelledIntegrationEvent).FullName => - OrderCancelledRoutingKey, + var type when type == typeof(OrderCancelledIntegrationEvent).FullName => OrderCancelledRoutingKey, _ => throw new InvalidOperationException( $"No RabbitMQ routing key is configured " + diff --git a/OrderProcessing.Api/appsettings.Development.json b/OrderProcessing.Api/appsettings.Development.json index b00282e..be76cea 100644 --- a/OrderProcessing.Api/appsettings.Development.json +++ b/OrderProcessing.Api/appsettings.Development.json @@ -3,9 +3,9 @@ "LogLevel": { "Default": "Information", "Microsoft.AspNetCore": "Warning" - }, - "RabbitMq": { - "Enabled": true } + }, + "RabbitMq": { + "Enabled": true } } diff --git a/OrderProcessing.Contracts/Orders/OrderEventRoutingKeys.cs b/OrderProcessing.Contracts/Orders/OrderEventRoutingKeys.cs new file mode 100644 index 0000000..246a0b4 --- /dev/null +++ b/OrderProcessing.Contracts/Orders/OrderEventRoutingKeys.cs @@ -0,0 +1,10 @@ +namespace OrderProcessing.Contracts.Orders; + +public static class OrderEventRoutingKeys +{ + public const string Created = "order.created"; + + public const string Completed = "order.completed"; + + public const string Cancelled = "order.cancelled"; +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs b/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs new file mode 100644 index 0000000..1e767c6 --- /dev/null +++ b/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs @@ -0,0 +1,100 @@ +using System.Text.Json; +using Microsoft.Extensions.Logging.Abstractions; +using OrderProcessing.Contracts.Orders; +using OrderProcessing.EmailWorker.Emailing; +using OrderProcessing.EmailWorker.Messaging; + +namespace OrderProcessing.EmailWorker.Tests; + +public sealed class OrderEventEmailHandlerTests +{ + [Fact] + public async Task HandleAsync_ForOrderCreatedEvent_SendsCreatedEmail() + { + // Arrange + var sender = new TestEmailSender(); + + var handler = new OrderEventEmailHandler( + sender, + NullLogger.Instance); + + var integrationEvent = + new OrderCreatedIntegrationEvent( + MessageId: Guid.NewGuid(), + OccurredAtUtc: DateTime.UtcNow, + OrderId: 1050, + CustomerId: 1001, + CustomerName: "John Smith", + CustomerEmail: "john.smith@example.com", + TotalAmount: 99.99m, + CreatedAtUtc: DateTime.UtcNow, + Items: + [ + new OrderItemIntegrationModel( + ProductId: 2001, + ProductName: "Keyboard", + Quantity: 1, + UnitPrice: 99.99m, + LineTotal: 99.99m) + ]); + + var body = JsonSerializer.SerializeToUtf8Bytes( + integrationEvent, + new JsonSerializerOptions( + JsonSerializerDefaults.Web)); + + // Act + await handler.HandleAsync( + typeof(OrderCreatedIntegrationEvent).FullName!, + body, + CancellationToken.None); + + // Assert + var email = Assert.Single(sender.Messages); + + Assert.Equal( + "john.smith@example.com", + email.Recipient); + + Assert.Contains( + "1050", + email.Subject); + + Assert.Contains( + "created", + email.Subject, + StringComparison.OrdinalIgnoreCase); + } + + [Fact] + public async Task HandleAsync_ForUnknownEvent_ThrowsException() + { + // Arrange + var handler = new OrderEventEmailHandler( + new TestEmailSender(), + NullLogger.Instance); + + // Act + var action = () => handler.HandleAsync( + "UnknownIntegrationEvent", + [], + CancellationToken.None); + + // Assert + await Assert.ThrowsAsync< + UnsupportedIntegrationEventException>( + action); + } + + private sealed class TestEmailSender : IEmailSender + { + public List Messages { get; } = []; + + public Task SendAsync(EmailMessage message, CancellationToken cancellationToken) + { + Messages.Add(message); + + return Task.CompletedTask; + } + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj b/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj new file mode 100644 index 0000000..cc30b6d --- /dev/null +++ b/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj @@ -0,0 +1,26 @@ + + + + net10.0 + enable + enable + false + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs b/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs new file mode 100644 index 0000000..7cfad54 --- /dev/null +++ b/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs @@ -0,0 +1,26 @@ +namespace OrderProcessing.EmailWorker.Configuration; + +public sealed class RabbitMqOptions +{ + public const string SectionName = "RabbitMq"; + + public string HostName { get; set; } = "localhost"; + + public int Port { get; set; } = 5672; + + public string UserName { get; set; } = "guest"; + + public string Password { get; set; } = "guest"; + + public string VirtualHost { get; set; } = "/"; + + public string ExchangeName { get; set; } = "order-processing.events"; + + public string EmailQueueName { get; set; } = "order-processing.email"; + + public string ClientProvidedName { get; set; } = "order-processing-email-worker"; + + public int NetworkRecoveryIntervalSeconds { get; set; } = 5; + + public ushort PrefetchCount { get; set; } = 1; +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Emailing/EmailMessage.cs b/OrderProcessing.EmailWorker/Emailing/EmailMessage.cs new file mode 100644 index 0000000..578ffe6 --- /dev/null +++ b/OrderProcessing.EmailWorker/Emailing/EmailMessage.cs @@ -0,0 +1,3 @@ +namespace OrderProcessing.EmailWorker.Emailing; + +public sealed record EmailMessage(string Recipient, string Subject, string Body); \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Emailing/IEmailSender.cs b/OrderProcessing.EmailWorker/Emailing/IEmailSender.cs new file mode 100644 index 0000000..3093fa7 --- /dev/null +++ b/OrderProcessing.EmailWorker/Emailing/IEmailSender.cs @@ -0,0 +1,6 @@ +namespace OrderProcessing.EmailWorker.Emailing; + +public interface IEmailSender +{ + Task SendAsync(EmailMessage message, CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Emailing/LoggingEmailSender.cs b/OrderProcessing.EmailWorker/Emailing/LoggingEmailSender.cs new file mode 100644 index 0000000..0f0aa14 --- /dev/null +++ b/OrderProcessing.EmailWorker/Emailing/LoggingEmailSender.cs @@ -0,0 +1,24 @@ +namespace OrderProcessing.EmailWorker.Emailing; + +public sealed class LoggingEmailSender : IEmailSender +{ + private readonly ILogger _logger; + + public LoggingEmailSender(ILogger logger) + { + _logger = logger; + } + + public Task SendAsync(EmailMessage message, CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + _logger.LogInformation( + "Simulated sending email to {Recipient} " + + "with subject {Subject}", + message.Recipient, + message.Subject); + + return Task.CompletedTask; + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs b/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs new file mode 100644 index 0000000..7eecce5 --- /dev/null +++ b/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs @@ -0,0 +1,130 @@ +using System.Text.Json; +using OrderProcessing.Contracts.Orders; +using OrderProcessing.EmailWorker.Emailing; + +namespace OrderProcessing.EmailWorker.Messaging; + +public sealed class OrderEventEmailHandler +{ + private static readonly JsonSerializerOptions SerializerOptions = new(JsonSerializerDefaults.Web); + + private readonly IEmailSender _emailSender; + private readonly ILogger _logger; + + public OrderEventEmailHandler(IEmailSender emailSender, ILogger logger) + { + _emailSender = emailSender; + _logger = logger; + } + + public async Task HandleAsync(string eventType, byte[] body, CancellationToken cancellationToken) + { + switch (eventType) + { + case var type when type == typeof(OrderCreatedIntegrationEvent).FullName: + + var createdEvent = Deserialize(body); + + await HandleCreatedAsync(createdEvent, cancellationToken); + + break; + + case var type when type == typeof(OrderCompletedIntegrationEvent).FullName: + + var completedEvent = Deserialize(body); + + await HandleCompletedAsync(completedEvent, cancellationToken); + + break; + + case var type when type == typeof(OrderCancelledIntegrationEvent).FullName: + + var cancelledEvent = Deserialize(body); + + await HandleCancelledAsync(cancelledEvent, cancellationToken); + + break; + + default: + throw new UnsupportedIntegrationEventException(eventType); + } + } + + private async Task HandleCreatedAsync(OrderCreatedIntegrationEvent integrationEvent, CancellationToken cancellationToken) + { + var email = new EmailMessage( + Recipient: integrationEvent.CustomerEmail, + Subject: + $"Order #{integrationEvent.OrderId} created", + Body: + $"Hello {integrationEvent.CustomerName}, " + + $"your order #{integrationEvent.OrderId} " + + $"was created with a total of " + + $"{integrationEvent.TotalAmount:C}."); + + await _emailSender.SendAsync( + email, + cancellationToken); + + _logger.LogInformation( + "Handled order-created event {MessageId} " + + "for order {OrderId}", + integrationEvent.MessageId, + integrationEvent.OrderId); + } + + private async Task HandleCompletedAsync( + OrderCompletedIntegrationEvent integrationEvent, + CancellationToken cancellationToken) + { + var email = new EmailMessage( + Recipient: integrationEvent.CustomerEmail, + Subject: + $"Order #{integrationEvent.OrderId} completed", + Body: + $"Hello {integrationEvent.CustomerName}, " + + $"your order #{integrationEvent.OrderId} " + + "has been completed."); + + await _emailSender.SendAsync( + email, + cancellationToken); + + _logger.LogInformation( + "Handled order-completed event {MessageId} " + + "for order {OrderId}", + integrationEvent.MessageId, + integrationEvent.OrderId); + } + + private async Task HandleCancelledAsync( + OrderCancelledIntegrationEvent integrationEvent, + CancellationToken cancellationToken) + { + var email = new EmailMessage( + Recipient: integrationEvent.CustomerEmail, + Subject: + $"Order #{integrationEvent.OrderId} cancelled", + Body: + $"Hello {integrationEvent.CustomerName}, " + + $"your order #{integrationEvent.OrderId} " + + "has been cancelled."); + + await _emailSender.SendAsync( + email, + cancellationToken); + + _logger.LogInformation( + "Handled order-cancelled event {MessageId} " + + "for order {OrderId}", + integrationEvent.MessageId, + integrationEvent.OrderId); + } + + private static TEvent Deserialize(byte[] body) + { + return JsonSerializer.Deserialize(body, SerializerOptions) ?? throw new JsonException( + $"Could not deserialize " + + $"{typeof(TEvent).Name}."); + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Messaging/UnsupportedIntegrationEventException.cs b/OrderProcessing.EmailWorker/Messaging/UnsupportedIntegrationEventException.cs new file mode 100644 index 0000000..07dcfcb --- /dev/null +++ b/OrderProcessing.EmailWorker/Messaging/UnsupportedIntegrationEventException.cs @@ -0,0 +1,9 @@ +namespace OrderProcessing.EmailWorker.Messaging; + +public sealed class UnsupportedIntegrationEventException : Exception +{ + public UnsupportedIntegrationEventException(string eventType) : + base($"Integration event type '{eventType}' is not supported by the email worker.") + { + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj b/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj new file mode 100644 index 0000000..7ebe150 --- /dev/null +++ b/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj @@ -0,0 +1,18 @@ + + + + net10.0 + enable + enable + dotnet-OrderProcessing.EmailWorker-36a8a289-8560-4b14-8e2e-1d0c3a026ba4 + + + + + + + + + + + diff --git a/OrderProcessing.EmailWorker/Program.cs b/OrderProcessing.EmailWorker/Program.cs new file mode 100644 index 0000000..1f4f08d --- /dev/null +++ b/OrderProcessing.EmailWorker/Program.cs @@ -0,0 +1,53 @@ +using OrderProcessing.EmailWorker; +using OrderProcessing.EmailWorker.Configuration; +using OrderProcessing.EmailWorker.Emailing; +using OrderProcessing.EmailWorker.Messaging; + +var builder = Host.CreateApplicationBuilder(args); + +builder.Services + .AddOptions() + .Bind( + builder.Configuration.GetSection( + RabbitMqOptions.SectionName)) + .Validate( + options => + !string.IsNullOrWhiteSpace(options.HostName), + "RabbitMQ host name is required.") + .Validate( + options => options.Port > 0, + "RabbitMQ port must be greater than zero.") + .Validate( + options => + !string.IsNullOrWhiteSpace(options.UserName), + "RabbitMQ user name is required.") + .Validate( + options => + !string.IsNullOrWhiteSpace(options.Password), + "RabbitMQ password is required.") + .Validate( + options => + !string.IsNullOrWhiteSpace(options.ExchangeName), + "RabbitMQ exchange name is required.") + .Validate( + options => + !string.IsNullOrWhiteSpace(options.EmailQueueName), + "RabbitMQ email queue name is required.") + .Validate( + options => options.PrefetchCount > 0, + "RabbitMQ prefetch count must be greater than zero.") + .ValidateOnStart(); + +builder.Services.AddSingleton< + IEmailSender, + LoggingEmailSender>(); + +builder.Services.AddScoped< + OrderEventEmailHandler>(); + +builder.Services.AddHostedService< + RabbitMqEmailConsumer>(); + +var host = builder.Build(); + +await host.RunAsync(); \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Properties/launchSettings.json b/OrderProcessing.EmailWorker/Properties/launchSettings.json new file mode 100644 index 0000000..7de974d --- /dev/null +++ b/OrderProcessing.EmailWorker/Properties/launchSettings.json @@ -0,0 +1,12 @@ +{ + "$schema": "https://json.schemastore.org/launchsettings.json", + "profiles": { + "OrderProcessing.EmailWorker": { + "commandName": "Project", + "dotnetRunMessages": true, + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + } + } + } +} diff --git a/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs b/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs new file mode 100644 index 0000000..e37d4f9 --- /dev/null +++ b/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs @@ -0,0 +1,301 @@ +using Microsoft.Extensions.Options; +using OrderProcessing.Contracts.Orders; +using OrderProcessing.EmailWorker.Configuration; +using OrderProcessing.EmailWorker.Messaging; +using RabbitMQ.Client; +using RabbitMQ.Client.Events; +using System.Text.Json; +using System.Threading.Channels; + +namespace OrderProcessing.EmailWorker; + +public sealed class RabbitMqEmailConsumer : BackgroundService +{ + private readonly IServiceScopeFactory _scopeFactory; + private readonly RabbitMqOptions _options; + private readonly ILogger _logger; + + private IConnection? _connection; + private IChannel? _channel; + private string? _consumerTag; + + public RabbitMqEmailConsumer( + IServiceScopeFactory scopeFactory, + IOptions options, + ILogger logger) + { + _scopeFactory = scopeFactory; + _options = options.Value; + _logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + while (!stoppingToken.IsCancellationRequested) + { + try + { + await RunConsumerAsync(stoppingToken); + return; + } + catch (OperationCanceledException) + when (stoppingToken.IsCancellationRequested) + { + return; + } + catch (Exception exception) + { + _logger.LogError( + exception, + "RabbitMQ email consumer failed to start. " + + "Retrying in {RetryIntervalSeconds} seconds", + _options.NetworkRecoveryIntervalSeconds); + + await CloseRabbitMqAsync(); + + await Task.Delay( + TimeSpan.FromSeconds( + _options.NetworkRecoveryIntervalSeconds), + stoppingToken); + } + } + } + + private async Task RunConsumerAsync(CancellationToken stoppingToken) + { + var connectionFactory = new ConnectionFactory + { + HostName = _options.HostName, + Port = _options.Port, + UserName = _options.UserName, + Password = _options.Password, + VirtualHost = _options.VirtualHost, + + AutomaticRecoveryEnabled = true, + TopologyRecoveryEnabled = true, + + NetworkRecoveryInterval = TimeSpan.FromSeconds(_options.NetworkRecoveryIntervalSeconds), + + ConsumerDispatchConcurrency = 1 + }; + + _connection = await connectionFactory.CreateConnectionAsync(_options.ClientProvidedName, stoppingToken); + + _channel = await _connection.CreateChannelAsync(cancellationToken: stoppingToken); + + await DeclareTopologyAsync(_channel, stoppingToken); + + await _channel.BasicQosAsync( + prefetchSize: 0, + prefetchCount: _options.PrefetchCount, + global: false, + cancellationToken: stoppingToken); + + var consumer = new AsyncEventingBasicConsumer(_channel); + + consumer.ReceivedAsync += async (_, eventArgs) => + await HandleDeliveryAsync(eventArgs); + + _consumerTag = + await _channel.BasicConsumeAsync( + queue: _options.EmailQueueName, + autoAck: false, + consumer: consumer, + cancellationToken: stoppingToken); + + _logger.LogInformation( + "Email worker is consuming queue {QueueName} " + + "on RabbitMQ {HostName}:{Port} " + + "with consumer tag {ConsumerTag}", + _options.EmailQueueName, + _options.HostName, + _options.Port, + _consumerTag); + + try + { + await Task.Delay( + Timeout.InfiniteTimeSpan, + stoppingToken); + } + catch (OperationCanceledException) + when (stoppingToken.IsCancellationRequested) + { + _logger.LogInformation( + "RabbitMQ email consumer is stopping"); + } + finally + { + await CloseRabbitMqAsync(); + } + } + + private async Task HandleDeliveryAsync( + BasicDeliverEventArgs eventArgs) + { + var channel = _channel + ?? throw new InvalidOperationException( + "RabbitMQ channel is not available."); + + // RabbitMQ.Client 7 uses ReadOnlyMemory. + // Copy it before this callback returns. + var body = eventArgs.Body.ToArray(); + + var eventType = eventArgs.BasicProperties.Type; + + try + { + if (string.IsNullOrWhiteSpace(eventType)) + { + throw new UnsupportedIntegrationEventException( ""); + } + + await using var scope = _scopeFactory.CreateAsyncScope(); + + var handler = scope.ServiceProvider + .GetRequiredService(); + + await handler.HandleAsync( + eventType, + body, + CancellationToken.None); + + await channel.BasicAckAsync( + deliveryTag: eventArgs.DeliveryTag, + multiple: false); + + _logger.LogInformation( + "Acknowledged RabbitMQ message {MessageId} " + + "with routing key {RoutingKey}", + eventArgs.BasicProperties.MessageId, + eventArgs.RoutingKey); + } + catch (JsonException exception) + { + await RejectPermanentFailureAsync( + channel, + eventArgs, + exception); + } + catch (UnsupportedIntegrationEventException exception) + { + await RejectPermanentFailureAsync( + channel, + eventArgs, + exception); + } + catch (Exception exception) + { + _logger.LogError( + exception, + "Temporary failure processing RabbitMQ " + + "message {MessageId}. The message will be requeued", + eventArgs.BasicProperties.MessageId); + + await channel.BasicNackAsync( + deliveryTag: eventArgs.DeliveryTag, + multiple: false, + requeue: true); + } + } + + private async Task RejectPermanentFailureAsync( + IChannel channel, + BasicDeliverEventArgs eventArgs, + Exception exception) + { + _logger.LogError( + exception, + "Permanent failure processing RabbitMQ " + + "message {MessageId}. The message will not be requeued", + eventArgs.BasicProperties.MessageId); + + await channel.BasicNackAsync( + deliveryTag: eventArgs.DeliveryTag, + multiple: false, + requeue: false); + } + + private async Task DeclareTopologyAsync(IChannel channel, CancellationToken cancellationToken) + { + await channel.ExchangeDeclareAsync( + exchange: _options.ExchangeName, + type: ExchangeType.Topic, + durable: true, + autoDelete: false, + cancellationToken: cancellationToken); + + await channel.QueueDeclareAsync( + queue: _options.EmailQueueName, + durable: true, + exclusive: false, + autoDelete: false, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.EmailQueueName, + exchange: _options.ExchangeName, + routingKey: OrderEventRoutingKeys.Created, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.EmailQueueName, + exchange: _options.ExchangeName, + routingKey: OrderEventRoutingKeys.Completed, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.EmailQueueName, + exchange: _options.ExchangeName, + routingKey: OrderEventRoutingKeys.Cancelled, + cancellationToken: cancellationToken); + } + + private async Task CloseRabbitMqAsync() + { + var channel = _channel; + _channel = null; + _consumerTag = null; + + if (channel is not null) + { + try + { + if (channel.IsOpen) + { + await channel.CloseAsync(CancellationToken.None); + } + } + catch (Exception exception) + { + _logger.LogDebug( + exception, + "Error closing RabbitMQ consumer channel"); + } + + await channel.DisposeAsync(); + } + + var connection = _connection; + _connection = null; + + if (connection is not null) + { + try + { + if (connection.IsOpen) + { + await connection.CloseAsync(CancellationToken.None); + } + } + catch (Exception exception) + { + _logger.LogDebug( + exception, + "Error closing RabbitMQ consumer connection"); + } + + await connection.DisposeAsync(); + } + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/appsettings.Development.json b/OrderProcessing.EmailWorker/appsettings.Development.json new file mode 100644 index 0000000..b2dcdb6 --- /dev/null +++ b/OrderProcessing.EmailWorker/appsettings.Development.json @@ -0,0 +1,8 @@ +{ + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft.Hosting.Lifetime": "Information" + } + } +} diff --git a/OrderProcessing.EmailWorker/appsettings.json b/OrderProcessing.EmailWorker/appsettings.json new file mode 100644 index 0000000..6e3ee7c --- /dev/null +++ b/OrderProcessing.EmailWorker/appsettings.json @@ -0,0 +1,20 @@ +{ + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft.Hosting.Lifetime": "Information" + } + }, + "RabbitMq": { + "HostName": "localhost", + "Port": 5672, + "UserName": "guest", + "Password": "guest", + "VirtualHost": "/", + "ExchangeName": "order-processing.events", + "EmailQueueName": "order-processing.email", + "ClientProvidedName": "order-processing-email-worker", + "NetworkRecoveryIntervalSeconds": 5, + "PrefetchCount": 1 + } +} diff --git a/OrderProcessingPlatform.slnx b/OrderProcessingPlatform.slnx index d5b542c..c1d1b00 100644 --- a/OrderProcessingPlatform.slnx +++ b/OrderProcessingPlatform.slnx @@ -1,7 +1,9 @@ + +