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 @@
+