diff --git a/OrderProcessing.ReadModelWorker.Tests/OrderProcessing.ReadModelWorker.Tests.csproj b/OrderProcessing.ReadModelWorker.Tests/OrderProcessing.ReadModelWorker.Tests.csproj new file mode 100644 index 0000000..e3da6c0 --- /dev/null +++ b/OrderProcessing.ReadModelWorker.Tests/OrderProcessing.ReadModelWorker.Tests.csproj @@ -0,0 +1,25 @@ + + + + net10.0 + enable + enable + false + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/OrderProcessing.ReadModelWorker.Tests/ReadModelWorker.cs b/OrderProcessing.ReadModelWorker.Tests/ReadModelWorker.cs new file mode 100644 index 0000000..acdcafd --- /dev/null +++ b/OrderProcessing.ReadModelWorker.Tests/ReadModelWorker.cs @@ -0,0 +1,263 @@ +using System.Text.Json; +using Microsoft.Extensions.Logging.Abstractions; +using OrderProcessing.Contracts.Orders; +using OrderProcessing.ReadModelWorker.Messaging; +using OrderProcessing.ReadModelWorker.Persistence; +using OrderProcessing.ReadModelWorker.ReadModels; + +namespace OrderProcessing.ReadModelWorker.Tests; + +public sealed class OrderEventProjectionHandlerTests +{ + [Fact] + public async Task HandleAsync_ForOrderCreatedEvent_CreatesPendingReadModel() + { + // Arrange + var repository = new FakeOrderReadModelRepository(); + + var handler = new OrderEventProjectionHandler(repository, NullLogger.Instance); + + var occurredAtUtc = DateTime.UtcNow; + + var integrationEvent = + new OrderCreatedIntegrationEvent( + MessageId: Guid.NewGuid(), + OccurredAtUtc: occurredAtUtc, + OrderId: 123, + CustomerId: 456, + CustomerName: "John Smith", + CustomerEmail: "john@example.com", + TotalAmount: 199.98m, + CreatedAtUtc: occurredAtUtc, + Items: + [ + new OrderItemIntegrationModel( + ProductId: 10, + ProductName: "Keyboard", + Quantity: 2, + UnitPrice: 99.99m, + LineTotal: 199.98m) + ]); + + var body = Serialize(integrationEvent); + + // Act + await handler.HandleAsync( + typeof(OrderCreatedIntegrationEvent).FullName!, + body, + CancellationToken.None); + + // Assert + Assert.NotNull(repository.CreatedOrder); + + Assert.Equal( + integrationEvent.OrderId, + repository.CreatedOrder.OrderId); + + Assert.Equal( + integrationEvent.CustomerId, + repository.CreatedOrder.CustomerId); + + Assert.Equal( + integrationEvent.CustomerName, + repository.CreatedOrder.CustomerName); + + Assert.Equal( + "Pending", + repository.CreatedOrder.Status); + + Assert.Equal( + integrationEvent.TotalAmount, + repository.CreatedOrder.TotalAmount); + + Assert.Equal( + occurredAtUtc, + repository.CreatedOrder.LastUpdatedAtUtc); + + var item = + Assert.Single(repository.CreatedOrder.Items); + + Assert.Equal(10, item.ProductId); + Assert.Equal("Keyboard", item.ProductName); + Assert.Equal(2, item.Quantity); + Assert.Equal(99.99m, item.UnitPrice); + } + + [Fact] + public async Task HandleAsync_ForOrderCompletedEvent_MarksOrderCompleted() + { + // Arrange + var repository = new FakeOrderReadModelRepository(); + + var handler = new OrderEventProjectionHandler( + repository, + NullLogger.Instance); + + var occurredAtUtc = DateTime.UtcNow; + var completedAtUtc = occurredAtUtc; + + var integrationEvent = + new OrderCompletedIntegrationEvent( + MessageId: Guid.NewGuid(), + OccurredAtUtc: occurredAtUtc, + OrderId: 123, + CustomerId: 456, + CustomerName: "John Smith", + CustomerEmail: "john@example.com", + TotalAmount: 100m, + CompletedAtUtc: completedAtUtc); + + // Act + await handler.HandleAsync( + typeof(OrderCompletedIntegrationEvent).FullName!, + Serialize(integrationEvent), + CancellationToken.None); + + // Assert + Assert.Equal( + 123, + repository.CompletedOrderId); + + Assert.Equal( + completedAtUtc, + repository.CompletedAtUtc); + + Assert.Equal( + occurredAtUtc, + repository.CompletedEventOccurredAtUtc); + } + + [Fact] + public async Task HandleAsync_ForOrderCancelledEvent_MarksOrderCancelled() + { + // Arrange + var repository = new FakeOrderReadModelRepository(); + + var handler = new OrderEventProjectionHandler( + repository, + NullLogger.Instance); + + var occurredAtUtc = DateTime.UtcNow; + var cancelledAtUtc = occurredAtUtc; + + var integrationEvent = + new OrderCancelledIntegrationEvent( + MessageId: Guid.NewGuid(), + OccurredAtUtc: occurredAtUtc, + OrderId: 123, + CustomerId: 456, + CustomerName: "John Smith", + CustomerEmail: "john@example.com", + TotalAmount: 100m, + CancelledAtUtc: cancelledAtUtc); + + // Act + await handler.HandleAsync( + typeof(OrderCancelledIntegrationEvent).FullName!, + Serialize(integrationEvent), + CancellationToken.None); + + // Assert + Assert.Equal( + 123, + repository.CancelledOrderId); + + Assert.Equal( + cancelledAtUtc, + repository.CancelledAtUtc); + + Assert.Equal( + occurredAtUtc, + repository.CancelledEventOccurredAtUtc); + } + + [Fact] + public async Task HandleAsync_ForUnknownEventType_ThrowsException() + { + // Arrange + var handler = new OrderEventProjectionHandler( + new FakeOrderReadModelRepository(), + NullLogger.Instance); + + // Act + var action = () => handler.HandleAsync( + "UnknownIntegrationEvent", + [], + CancellationToken.None); + + // Assert + await Assert.ThrowsAsync( + action); + } + + private static byte[] Serialize(T value) + { + return JsonSerializer.SerializeToUtf8Bytes( + value, + new JsonSerializerOptions( + JsonSerializerDefaults.Web)); + } + + private sealed class FakeOrderReadModelRepository + : IOrderReadModelRepository + { + public OrderReadModel? CreatedOrder { get; private set; } + + public int? CompletedOrderId { get; private set; } + + public DateTime? CompletedAtUtc { get; private set; } + + public DateTime? CompletedEventOccurredAtUtc + { + get; + private set; + } + + public int? CancelledOrderId { get; private set; } + + public DateTime? CancelledAtUtc { get; private set; } + + public DateTime? CancelledEventOccurredAtUtc + { + get; + private set; + } + + public Task CreateIfMissingAsync( + OrderReadModel order, + CancellationToken cancellationToken) + { + CreatedOrder = order; + + return Task.CompletedTask; + } + + public Task MarkCompletedAsync( + int orderId, + DateTime completedAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken) + { + CompletedOrderId = orderId; + CompletedAtUtc = completedAtUtc; + CompletedEventOccurredAtUtc = + eventOccurredAtUtc; + + return Task.CompletedTask; + } + + public Task MarkCancelledAsync( + int orderId, + DateTime cancelledAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken) + { + CancelledOrderId = orderId; + CancelledAtUtc = cancelledAtUtc; + CancelledEventOccurredAtUtc = + eventOccurredAtUtc; + + return Task.CompletedTask; + } + } +} \ No newline at end of file diff --git a/OrderProcessing.ReadModelWorker/Configuration/RabbitMqOptions.cs b/OrderProcessing.ReadModelWorker/Configuration/RabbitMqOptions.cs new file mode 100644 index 0000000..08fd27b --- /dev/null +++ b/OrderProcessing.ReadModelWorker/Configuration/RabbitMqOptions.cs @@ -0,0 +1,39 @@ +namespace OrderProcessing.ReadModelWorker.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 ReadModelQueueName { get; set; } = "order-processing.read-model"; + + public string DeadLetterExchangeName { get; set; } = "order-processing.dead-letter"; + + public string DeadLetterQueueName { get; set; } = "order-processing.read-model.dead-letter"; + + public string DeadLetterRoutingKey { get; set; } = "read-model.failed"; + + public int DeliveryLimit { get; set; } = 5; + + public int RetryMinDelayMilliseconds { get; set; } = 5000; + + public int RetryMaxDelayMilliseconds { get; set; } = 30000; + + public string ClientProvidedName { get; set; } = + "order-processing-read-model-worker"; + + public int NetworkRecoveryIntervalSeconds { get; set; } = 5; + + public ushort PrefetchCount { get; set; } = 1; +} \ No newline at end of file diff --git a/OrderProcessing.ReadModelWorker/Messaging/OrderEventProjectionHandler.cs b/OrderProcessing.ReadModelWorker/Messaging/OrderEventProjectionHandler.cs new file mode 100644 index 0000000..b45f6c3 --- /dev/null +++ b/OrderProcessing.ReadModelWorker/Messaging/OrderEventProjectionHandler.cs @@ -0,0 +1,108 @@ +using System.Text.Json; +using OrderProcessing.Contracts.Orders; +using OrderProcessing.ReadModelWorker.Persistence; +using OrderProcessing.ReadModelWorker.ReadModels; + +namespace OrderProcessing.ReadModelWorker.Messaging; + +public sealed class OrderEventProjectionHandler +{ + private static readonly JsonSerializerOptions SerializerOptions = new(JsonSerializerDefaults.Web); + + private readonly IOrderReadModelRepository _repository; + private readonly ILogger _logger; + + public OrderEventProjectionHandler(IOrderReadModelRepository repository, ILogger logger) + { + _repository = repository; + _logger = logger; + } + + public async Task HandleAsync(string eventType, byte[] body, CancellationToken cancellationToken) + { + if (eventType == typeof(OrderCreatedIntegrationEvent).FullName) + { + var integrationEvent = Deserialize(body); + + await HandleCreatedAsync(integrationEvent, cancellationToken); + + return; + } + + if (eventType == typeof(OrderCompletedIntegrationEvent).FullName) + { + var integrationEvent = Deserialize(body); + + await _repository.MarkCompletedAsync( + integrationEvent.OrderId, + integrationEvent.CompletedAtUtc, + integrationEvent.OccurredAtUtc, + cancellationToken); + + return; + } + + if (eventType == typeof(OrderCancelledIntegrationEvent).FullName) + { + var integrationEvent = Deserialize(body); + + await _repository.MarkCancelledAsync( + integrationEvent.OrderId, + integrationEvent.CancelledAtUtc, + integrationEvent.OccurredAtUtc, + cancellationToken); + + return; + } + + throw new InvalidOperationException($"Unsupported integration event type '{eventType}'."); + } + + private async Task HandleCreatedAsync(OrderCreatedIntegrationEvent integrationEvent, CancellationToken cancellationToken) + { + var order = new OrderReadModel + { + OrderId = integrationEvent.OrderId, + CustomerId = integrationEvent.CustomerId, + CustomerName = + integrationEvent.CustomerName, + Status = "Pending", + TotalAmount = + integrationEvent.TotalAmount, + CreatedAtUtc = + integrationEvent.CreatedAtUtc, + CompletedAtUtc = null, + CancelledAtUtc = null, + LastUpdatedAtUtc = + integrationEvent.OccurredAtUtc, + + Items = integrationEvent.Items + .Select(item => + new OrderItemReadModel + { + ProductId = item.ProductId, + ProductName = item.ProductName, + Quantity = item.Quantity, + UnitPrice = item.UnitPrice, + LineTotal = item.LineTotal + }).ToList() + }; + + await _repository.CreateIfMissingAsync(order, cancellationToken); + + _logger.LogInformation( + "Projected order-created event into MongoDB " + + "for order {OrderId}", + 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.ReadModelWorker/Messaging/UnsupportedIntegrationEventException.cs b/OrderProcessing.ReadModelWorker/Messaging/UnsupportedIntegrationEventException.cs new file mode 100644 index 0000000..3580dc7 --- /dev/null +++ b/OrderProcessing.ReadModelWorker/Messaging/UnsupportedIntegrationEventException.cs @@ -0,0 +1,10 @@ +namespace OrderProcessing.ReadModelWorker.Messaging +{ + public sealed class UnsupportedIntegrationEventException : Exception + { + public UnsupportedIntegrationEventException(string eventType) : + base($"Integration event type '{eventType}' is not supported by the email worker.") + { + } + } +} diff --git a/OrderProcessing.ReadModelWorker/OrderProcessing.ReadModelWorker.csproj b/OrderProcessing.ReadModelWorker/OrderProcessing.ReadModelWorker.csproj index 57d3099..4f6f369 100644 --- a/OrderProcessing.ReadModelWorker/OrderProcessing.ReadModelWorker.csproj +++ b/OrderProcessing.ReadModelWorker/OrderProcessing.ReadModelWorker.csproj @@ -10,6 +10,7 @@ + diff --git a/OrderProcessing.ReadModelWorker/Persistence/IOrderReadModelRepository.cs b/OrderProcessing.ReadModelWorker/Persistence/IOrderReadModelRepository.cs new file mode 100644 index 0000000..01336d9 --- /dev/null +++ b/OrderProcessing.ReadModelWorker/Persistence/IOrderReadModelRepository.cs @@ -0,0 +1,20 @@ +using OrderProcessing.ReadModelWorker.ReadModels; + +namespace OrderProcessing.ReadModelWorker.Persistence; + +public interface IOrderReadModelRepository +{ + Task CreateIfMissingAsync(OrderReadModel order, CancellationToken cancellationToken); + + Task MarkCompletedAsync( + int orderId, + DateTime completedAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken); + + Task MarkCancelledAsync( + int orderId, + DateTime cancelledAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/OrderProcessing.ReadModelWorker/Persistence/MongoOrderReadModelRepository.cs b/OrderProcessing.ReadModelWorker/Persistence/MongoOrderReadModelRepository.cs new file mode 100644 index 0000000..25115cc --- /dev/null +++ b/OrderProcessing.ReadModelWorker/Persistence/MongoOrderReadModelRepository.cs @@ -0,0 +1,123 @@ +using MongoDB.Driver; +using OrderProcessing.ReadModelWorker.ReadModels; + +namespace OrderProcessing.ReadModelWorker.Persistence; + +public sealed class MongoOrderReadModelRepository : IOrderReadModelRepository +{ + private readonly OrderReadModelStore _store; + private readonly ILogger _logger; + + public MongoOrderReadModelRepository(OrderReadModelStore store, ILogger logger) + { + _store = store; + _logger = logger; + } + + public async Task CreateIfMissingAsync(OrderReadModel order, CancellationToken cancellationToken) + { + var filter = Builders.Filter.Eq(existing => existing.OrderId, order.OrderId); + + var update = Builders.Update + .SetOnInsert(existing => existing.OrderId, order.OrderId) + .SetOnInsert(existing => existing.CustomerId, order.CustomerId) + .SetOnInsert(existing => existing.CustomerName, order.CustomerName) + .SetOnInsert(existing => existing.Status, order.Status) + .SetOnInsert(existing => existing.TotalAmount, order.TotalAmount) + .SetOnInsert(existing => existing.CreatedAtUtc, order.CreatedAtUtc) + .SetOnInsert(existing => existing.CompletedAtUtc, order.CompletedAtUtc) + .SetOnInsert(existing => existing.CancelledAtUtc, order.CancelledAtUtc) + .SetOnInsert(existing => existing.Items, order.Items) + .SetOnInsert(existing => existing.LastUpdatedAtUtc, order.LastUpdatedAtUtc); + + await _store.Orders.UpdateOneAsync(filter, update, + new UpdateOptions + { + IsUpsert = true + }, + cancellationToken); + } + + public Task MarkCompletedAsync( + int orderId, + DateTime completedAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken) + { + return UpdateStatusAsync( + orderId, + status: "Completed", + completedAtUtc: completedAtUtc, + cancelledAtUtc: null, + eventOccurredAtUtc, + cancellationToken); + } + + public Task MarkCancelledAsync( + int orderId, + DateTime cancelledAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken) + { + return UpdateStatusAsync( + orderId, + status: "Cancelled", + completedAtUtc: null, + cancelledAtUtc: cancelledAtUtc, + eventOccurredAtUtc, + cancellationToken); + } + + private async Task UpdateStatusAsync( + int orderId, + string status, + DateTime? completedAtUtc, + DateTime? cancelledAtUtc, + DateTime eventOccurredAtUtc, + CancellationToken cancellationToken) + { + var existing = await _store.Orders.Find(order => order.OrderId == orderId) + .FirstOrDefaultAsync(cancellationToken); + + if (existing is null) + { + throw new InvalidOperationException($"Order read model {orderId} does not exist yet."); + } + + if (existing.LastUpdatedAtUtc > eventOccurredAtUtc) + { + _logger.LogInformation( + "Ignoring stale read-model event for order " + + "{OrderId}. Event occurred at {OccurredAtUtc}", + orderId, + eventOccurredAtUtc); + + return; + } + + var filter = Builders.Filter.And( + Builders.Filter.Eq( + order => order.OrderId, + orderId), + + Builders.Filter.Lte( + order => order.LastUpdatedAtUtc, + eventOccurredAtUtc)); + + var update = Builders.Update + .Set(order => order.Status, status) + .Set(order => order.LastUpdatedAtUtc, eventOccurredAtUtc); + + if (completedAtUtc.HasValue) + { + update = update.Set(order => order.CompletedAtUtc, completedAtUtc.Value); + } + + if (cancelledAtUtc.HasValue) + { + update = update.Set(order => order.CancelledAtUtc, cancelledAtUtc.Value); + } + + await _store.Orders.UpdateOneAsync(filter, update, cancellationToken: cancellationToken); + } +} \ No newline at end of file diff --git a/OrderProcessing.ReadModelWorker/Program.cs b/OrderProcessing.ReadModelWorker/Program.cs index cd1e697..a495023 100644 --- a/OrderProcessing.ReadModelWorker/Program.cs +++ b/OrderProcessing.ReadModelWorker/Program.cs @@ -1,6 +1,8 @@ using Microsoft.Extensions.Options; using MongoDB.Driver; +using OrderProcessing.ReadModelWorker; using OrderProcessing.ReadModelWorker.Configuration; +using OrderProcessing.ReadModelWorker.Messaging; using OrderProcessing.ReadModelWorker.Persistence; var builder = Host.CreateApplicationBuilder(args); @@ -26,6 +28,20 @@ builder.Services.AddSingleton(); +builder.Services + .AddOptions() + .Bind(builder.Configuration.GetSection(RabbitMqOptions.SectionName)) + .Validate(options => !string.IsNullOrWhiteSpace(options.HostName), "RabbitMQ host is required.") + .Validate(options => !string.IsNullOrWhiteSpace(options.ReadModelQueueName), "RabbitMQ read-model queue name is required.") + .ValidateOnStart(); + + +builder.Services.AddSingleton(); + +builder.Services.AddScoped(); + +builder.Services.AddHostedService(); + var host = builder.Build(); await host.RunAsync(); \ No newline at end of file diff --git a/OrderProcessing.ReadModelWorker/RabbitMqReadModelConsumer.cs b/OrderProcessing.ReadModelWorker/RabbitMqReadModelConsumer.cs new file mode 100644 index 0000000..478729e --- /dev/null +++ b/OrderProcessing.ReadModelWorker/RabbitMqReadModelConsumer.cs @@ -0,0 +1,338 @@ +using Microsoft.Extensions.Options; +using OrderProcessing.Contracts.Orders; +using OrderProcessing.ReadModelWorker.Configuration; +using OrderProcessing.ReadModelWorker.Messaging; +using RabbitMQ.Client; +using RabbitMQ.Client.Events; +using System; +using System.Collections.Generic; +using System.Text; +using System.Text.Json; + +namespace OrderProcessing.ReadModelWorker +{ + public class RabbitMqReadModelConsumer : BackgroundService + { + private readonly IServiceScopeFactory _scopeFactory; + private readonly RabbitMqOptions _options; + private readonly ILogger _logger; + + private IConnection? _connection; + private IChannel? _channel; + private string? _consumerTag; + + public RabbitMqReadModelConsumer( + 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 Read 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.ReadModelQueueName, + autoAck: false, + consumer: consumer, + cancellationToken: stoppingToken); + + _logger.LogInformation( + "RabbitMQ Read consumer is consuming queue {QueueName} " + + "on RabbitMQ {HostName}:{Port} " + + "with consumer tag {ConsumerTag}", + _options.ReadModelQueueName, + _options.HostName, + _options.Port, + _consumerTag); + + try + { + await Task.Delay(Timeout.InfiniteTimeSpan, stoppingToken); + } + catch (OperationCanceledException) + when (stoppingToken.IsCancellationRequested) + { + _logger.LogInformation( + "RabbitMQ Read 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; + var messageIdText = eventArgs.BasicProperties.MessageId; + + try + { + if (string.IsNullOrWhiteSpace(eventType)) + { + throw new UnsupportedIntegrationEventException(""); + } + + if (!Guid.TryParse(messageIdText, out var messageId)) + { + throw new UnsupportedIntegrationEventException($"Invalid message ID '{messageIdText}'."); + } + + await using var scope = _scopeFactory.CreateAsyncScope(); + + var handler = scope.ServiceProvider.GetRequiredService(); + + await handler.HandleAsync(eventType, body, CancellationToken.None); + + await channel.BasicAckAsync(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 retried", + eventArgs.BasicProperties.MessageId); + + await channel.BasicRejectAsync( + deliveryTag: eventArgs.DeliveryTag, + requeue: true); + } + } + + private async Task RejectPermanentFailureAsync( + IChannel channel, + BasicDeliverEventArgs eventArgs, + Exception exception) + { + _logger.LogError(exception, + "Permanent failure processing RabbitMQ " + + "message {MessageId}. " + + "The message will be dead-lettered", + eventArgs.BasicProperties.MessageId); + + await channel.BasicRejectAsync( + deliveryTag: eventArgs.DeliveryTag, + requeue: false); + } + + private async Task DeclareTopologyAsync(IChannel channel, CancellationToken cancellationToken) + { + var emailQueueArguments = new Dictionary + { + ["x-queue-type"] = "quorum", + + ["x-delivery-limit"] = _options.DeliveryLimit, + + ["x-delayed-retry-type"] = "failed", + + ["x-delayed-retry-min"] = _options.RetryMinDelayMilliseconds, + + ["x-delayed-retry-max"] = _options.RetryMaxDelayMilliseconds, + + ["x-dead-letter-exchange"] = _options.DeadLetterExchangeName, + + ["x-dead-letter-routing-key"] = _options.DeadLetterRoutingKey + }; + + await channel.ExchangeDeclareAsync( + exchange: _options.DeadLetterExchangeName, + type: ExchangeType.Direct, + durable: true, + autoDelete: false, + cancellationToken: cancellationToken); + + var deadLetterQueueArguments = new Dictionary + { + ["x-queue-type"] = "quorum" + }; + + await channel.QueueDeclareAsync( + queue: _options.DeadLetterQueueName, + durable: true, + exclusive: false, + autoDelete: false, + arguments: deadLetterQueueArguments, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.DeadLetterQueueName, + exchange: _options.DeadLetterExchangeName, + routingKey: _options.DeadLetterRoutingKey, + cancellationToken: cancellationToken); + + await channel.ExchangeDeclareAsync( + exchange: _options.ExchangeName, + type: ExchangeType.Topic, + durable: true, + autoDelete: false, + cancellationToken: cancellationToken); + + await channel.QueueDeclareAsync( + queue: _options.ReadModelQueueName, + durable: true, + exclusive: false, + autoDelete: false, + arguments: emailQueueArguments, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.ReadModelQueueName, + exchange: _options.ExchangeName, + routingKey: OrderEventRoutingKeys.Created, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.ReadModelQueueName, + exchange: _options.ExchangeName, + routingKey: OrderEventRoutingKeys.Completed, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.ReadModelQueueName, + 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(); + } + } + } +} diff --git a/OrderProcessing.ReadModelWorker/appsettings.json b/OrderProcessing.ReadModelWorker/appsettings.json index 47c79fc..14230f8 100644 --- a/OrderProcessing.ReadModelWorker/appsettings.json +++ b/OrderProcessing.ReadModelWorker/appsettings.json @@ -9,5 +9,23 @@ "ConnectionString": "mongodb://localhost:27017", "DatabaseName": "OrderProcessingReadDb", "OrdersCollectionName": "orders" + }, + "RabbitMq": { + "HostName": "localhost", + "Port": 5672, + "UserName": "guest", + "Password": "guest", + "VirtualHost": "/", + "ExchangeName": "order-processing.events", + "ReadModelQueueName": "order-processing.read-model", + "DeadLetterExchangeName": "order-processing.dead-letter", + "DeadLetterQueueName": "order-processing.read-model.dead-letter", + "DeadLetterRoutingKey": "read-model.failed", + "DeliveryLimit": 5, + "RetryMinDelayMilliseconds": 5000, + "RetryMaxDelayMilliseconds": 30000, + "ClientProvidedName": "order-processing-read-model-worker", + "NetworkRecoveryIntervalSeconds": 5, + "PrefetchCount": 1 } } diff --git a/OrderProcessingPlatform.slnx b/OrderProcessingPlatform.slnx index 92c5ba9..6f9cc13 100644 --- a/OrderProcessingPlatform.slnx +++ b/OrderProcessingPlatform.slnx @@ -2,6 +2,7 @@ +