diff --git a/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs b/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs index 1e767c6..db06eaa 100644 --- a/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs +++ b/OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs @@ -1,100 +1,183 @@ using System.Text.Json; +using Microsoft.Data.Sqlite; +using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Logging.Abstractions; using OrderProcessing.Contracts.Orders; using OrderProcessing.EmailWorker.Emailing; using OrderProcessing.EmailWorker.Messaging; +using OrderProcessing.EmailWorker.Persistence; namespace OrderProcessing.EmailWorker.Tests; -public sealed class OrderEventEmailHandlerTests +public sealed class IdempotentEmailMessageProcessorTests { [Fact] - public async Task HandleAsync_ForOrderCreatedEvent_SendsCreatedEmail() + public async Task ProcessAsync_WhenMessageIsDeliveredTwice_SendsEmailOnce() { // Arrange - var sender = new TestEmailSender(); + await using var connection = + new SqliteConnection("Data Source=:memory:"); - var handler = new OrderEventEmailHandler( - sender, - NullLogger.Instance); + await connection.OpenAsync(); + + var options = + new DbContextOptionsBuilder< + EmailWorkerDbContext>() + .UseSqlite(connection) + .Options; + + await using var dbContext = + new EmailWorkerDbContext(options); + + await dbContext.Database.EnsureCreatedAsync(); + + var emailSender = new TestEmailSender(); + + var eventHandler = + new OrderEventEmailHandler( + emailSender, + NullLogger.Instance); + + var processor = + new IdempotentEmailMessageProcessor( + dbContext, + eventHandler, + NullLogger< + IdempotentEmailMessageProcessor>.Instance); + + var messageId = Guid.NewGuid(); var integrationEvent = new OrderCreatedIntegrationEvent( - MessageId: Guid.NewGuid(), + MessageId: messageId, OccurredAtUtc: DateTime.UtcNow, - OrderId: 1050, - CustomerId: 1001, - CustomerName: "John Smith", - CustomerEmail: "john.smith@example.com", - TotalAmount: 99.99m, + OrderId: 123, + CustomerId: 456, + CustomerName: "Test Customer", + CustomerEmail: "test@example.com", + TotalAmount: 100m, 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)); + Items: []); + + var body = + JsonSerializer.SerializeToUtf8Bytes( + integrationEvent, + new JsonSerializerOptions( + JsonSerializerDefaults.Web)); + + var eventType = + typeof(OrderCreatedIntegrationEvent) + .FullName!; // Act - await handler.HandleAsync( - typeof(OrderCreatedIntegrationEvent).FullName!, + await processor.ProcessAsync( + messageId, + eventType, + body, + CancellationToken.None); + + await processor.ProcessAsync( + messageId, + eventType, body, CancellationToken.None); // Assert - var email = Assert.Single(sender.Messages); + Assert.Single(emailSender.Messages); - Assert.Equal( - "john.smith@example.com", - email.Recipient); + var processedMessages = + await dbContext.ProcessedMessages + .ToListAsync(); - Assert.Contains( - "1050", - email.Subject); + var processed = + Assert.Single(processedMessages); - Assert.Contains( - "created", - email.Subject, - StringComparison.OrdinalIgnoreCase); + Assert.Equal( + messageId, + processed.MessageId); } [Fact] - public async Task HandleAsync_ForUnknownEvent_ThrowsException() + public async Task ProcessAsync_WhenEmailSendingFails_DoesNotRecordMessage() { - // Arrange - var handler = new OrderEventEmailHandler( - new TestEmailSender(), - NullLogger.Instance); + await using var connection = + new SqliteConnection("Data Source=:memory:"); - // Act - var action = () => handler.HandleAsync( - "UnknownIntegrationEvent", - [], - CancellationToken.None); + await connection.OpenAsync(); - // Assert - await Assert.ThrowsAsync< - UnsupportedIntegrationEventException>( - action); + var options = + new DbContextOptionsBuilder< + EmailWorkerDbContext>() + .UseSqlite(connection) + .Options; + + await using var dbContext = + new EmailWorkerDbContext(options); + + await dbContext.Database.EnsureCreatedAsync(); + + var eventHandler = + new OrderEventEmailHandler( + new FailingEmailSender(), + NullLogger.Instance); + + var processor = + new IdempotentEmailMessageProcessor( + dbContext, + eventHandler, + NullLogger< + IdempotentEmailMessageProcessor>.Instance); + + var messageId = Guid.NewGuid(); + + var integrationEvent = + new OrderCompletedIntegrationEvent( + MessageId: messageId, + OccurredAtUtc: DateTime.UtcNow, + OrderId: 123, + CustomerId: 456, + CustomerName: "Test Customer", + CustomerEmail: "test@example.com", + TotalAmount: 100m, + CompletedAtUtc: DateTime.UtcNow); + + var body = + JsonSerializer.SerializeToUtf8Bytes( + integrationEvent, + new JsonSerializerOptions( + JsonSerializerDefaults.Web)); + + await Assert.ThrowsAsync( + () => processor.ProcessAsync( + messageId, + typeof(OrderCompletedIntegrationEvent) + .FullName!, + body, + CancellationToken.None)); + + Assert.Empty( + await dbContext.ProcessedMessages.ToListAsync()); } private sealed class TestEmailSender : IEmailSender { public List Messages { get; } = []; - public Task SendAsync(EmailMessage message, CancellationToken cancellationToken) + public Task SendAsync( + EmailMessage message, + CancellationToken cancellationToken) { Messages.Add(message); return Task.CompletedTask; } } + + private sealed class FailingEmailSender : IEmailSender + { + public Task SendAsync(EmailMessage message, CancellationToken cancellationToken) + { + throw new InvalidOperationException( "Simulated email failure."); + } + } } \ No newline at end of file diff --git a/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj b/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj index cc30b6d..db4df1a 100644 --- a/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj +++ b/OrderProcessing.EmailWorker.Tests/OrderProcessing.EmailWorker.Tests.csproj @@ -9,6 +9,7 @@ + diff --git a/OrderProcessing.EmailWorker/Messaging/IdempotentEmailMessageProcessor.cs b/OrderProcessing.EmailWorker/Messaging/IdempotentEmailMessageProcessor.cs new file mode 100644 index 0000000..73f1c7f --- /dev/null +++ b/OrderProcessing.EmailWorker/Messaging/IdempotentEmailMessageProcessor.cs @@ -0,0 +1,55 @@ +using Microsoft.EntityFrameworkCore; +using OrderProcessing.EmailWorker.Persistence; +using OrderProcessing.EmailWorker.Persistence.Entities; + +namespace OrderProcessing.EmailWorker.Messaging; + +public sealed class IdempotentEmailMessageProcessor +{ + private readonly EmailWorkerDbContext _dbContext; + private readonly OrderEventEmailHandler _eventHandler; + private readonly ILogger + _logger; + + public IdempotentEmailMessageProcessor(EmailWorkerDbContext dbContext, OrderEventEmailHandler eventHandler, ILogger logger) + { + _dbContext = dbContext; + _eventHandler = eventHandler; + _logger = logger; + } + + public async Task ProcessAsync(Guid messageId, string eventType, byte[] body, CancellationToken cancellationToken) + { + var alreadyProcessed = await _dbContext.ProcessedMessages + .AsNoTracking() + .AnyAsync(message => message.MessageId == messageId, cancellationToken); + + if (alreadyProcessed) + { + _logger.LogInformation( + "Skipping duplicate integration event " + + "{MessageId} of type {EventType}", + messageId, + eventType); + + return; + } + + await _eventHandler.HandleAsync(eventType, body, cancellationToken); + + _dbContext.ProcessedMessages.Add( + new ProcessedMessage + { + MessageId = messageId, + EventType = eventType, + ProcessedAtUtc = DateTime.UtcNow + }); + + await _dbContext.SaveChangesAsync(cancellationToken); + + _logger.LogInformation( + "Integration event {MessageId} recorded " + + "as successfully processed", + messageId); + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs b/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs index 7eecce5..9ef69f6 100644 --- a/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs +++ b/OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs @@ -62,9 +62,7 @@ private async Task HandleCreatedAsync(OrderCreatedIntegrationEvent integrationEv $"was created with a total of " + $"{integrationEvent.TotalAmount:C}."); - await _emailSender.SendAsync( - email, - cancellationToken); + await _emailSender.SendAsync(email, cancellationToken); _logger.LogInformation( "Handled order-created event {MessageId} " + @@ -73,9 +71,7 @@ await _emailSender.SendAsync( integrationEvent.OrderId); } - private async Task HandleCompletedAsync( - OrderCompletedIntegrationEvent integrationEvent, - CancellationToken cancellationToken) + private async Task HandleCompletedAsync(OrderCompletedIntegrationEvent integrationEvent, CancellationToken cancellationToken) { var email = new EmailMessage( Recipient: integrationEvent.CustomerEmail, @@ -86,9 +82,7 @@ private async Task HandleCompletedAsync( $"your order #{integrationEvent.OrderId} " + "has been completed."); - await _emailSender.SendAsync( - email, - cancellationToken); + await _emailSender.SendAsync(email, cancellationToken); _logger.LogInformation( "Handled order-completed event {MessageId} " + @@ -110,9 +104,7 @@ private async Task HandleCancelledAsync( $"your order #{integrationEvent.OrderId} " + "has been cancelled."); - await _emailSender.SendAsync( - email, - cancellationToken); + await _emailSender.SendAsync(email, cancellationToken); _logger.LogInformation( "Handled order-cancelled event {MessageId} " + diff --git a/OrderProcessing.EmailWorker/Migrations/20260810183537_InitialEmailWorkerInbox.Designer.cs b/OrderProcessing.EmailWorker/Migrations/20260810183537_InitialEmailWorkerInbox.Designer.cs new file mode 100644 index 0000000..55648cc --- /dev/null +++ b/OrderProcessing.EmailWorker/Migrations/20260810183537_InitialEmailWorkerInbox.Designer.cs @@ -0,0 +1,50 @@ +// +using System; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using OrderProcessing.EmailWorker.Persistence; + +#nullable disable + +namespace OrderProcessing.EmailWorker.Migrations +{ + [DbContext(typeof(EmailWorkerDbContext))] + [Migration("20260810183537_InitialEmailWorkerInbox")] + partial class InitialEmailWorkerInbox + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.9") + .HasAnnotation("Relational:MaxIdentifierLength", 128); + + SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + + modelBuilder.Entity("OrderProcessing.EmailWorker.Persistence.Entities.ProcessedMessage", b => + { + b.Property("MessageId") + .HasColumnType("uniqueidentifier"); + + b.Property("EventType") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("nvarchar(256)"); + + b.Property("ProcessedAtUtc") + .HasColumnType("datetime2"); + + b.HasKey("MessageId"); + + b.HasIndex("ProcessedAtUtc"); + + b.ToTable("ProcessedMessages", "email"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/OrderProcessing.EmailWorker/Migrations/20260810183537_InitialEmailWorkerInbox.cs b/OrderProcessing.EmailWorker/Migrations/20260810183537_InitialEmailWorkerInbox.cs new file mode 100644 index 0000000..f5f081c --- /dev/null +++ b/OrderProcessing.EmailWorker/Migrations/20260810183537_InitialEmailWorkerInbox.cs @@ -0,0 +1,46 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace OrderProcessing.EmailWorker.Migrations +{ + /// + public partial class InitialEmailWorkerInbox : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.EnsureSchema( + name: "email"); + + migrationBuilder.CreateTable( + name: "ProcessedMessages", + schema: "email", + columns: table => new + { + MessageId = table.Column(type: "uniqueidentifier", nullable: false), + EventType = table.Column(type: "nvarchar(256)", maxLength: 256, nullable: false), + ProcessedAtUtc = table.Column(type: "datetime2", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_ProcessedMessages", x => x.MessageId); + }); + + migrationBuilder.CreateIndex( + name: "IX_ProcessedMessages_ProcessedAtUtc", + schema: "email", + table: "ProcessedMessages", + column: "ProcessedAtUtc"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "ProcessedMessages", + schema: "email"); + } + } +} diff --git a/OrderProcessing.EmailWorker/Migrations/EmailWorkerDbContextModelSnapshot.cs b/OrderProcessing.EmailWorker/Migrations/EmailWorkerDbContextModelSnapshot.cs new file mode 100644 index 0000000..77d681d --- /dev/null +++ b/OrderProcessing.EmailWorker/Migrations/EmailWorkerDbContextModelSnapshot.cs @@ -0,0 +1,47 @@ +// +using System; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using OrderProcessing.EmailWorker.Persistence; + +#nullable disable + +namespace OrderProcessing.EmailWorker.Migrations +{ + [DbContext(typeof(EmailWorkerDbContext))] + partial class EmailWorkerDbContextModelSnapshot : ModelSnapshot + { + protected override void BuildModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.9") + .HasAnnotation("Relational:MaxIdentifierLength", 128); + + SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + + modelBuilder.Entity("OrderProcessing.EmailWorker.Persistence.Entities.ProcessedMessage", b => + { + b.Property("MessageId") + .HasColumnType("uniqueidentifier"); + + b.Property("EventType") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("nvarchar(256)"); + + b.Property("ProcessedAtUtc") + .HasColumnType("datetime2"); + + b.HasKey("MessageId"); + + b.HasIndex("ProcessedAtUtc"); + + b.ToTable("ProcessedMessages", "email"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj b/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj index 7ebe150..86eb105 100644 --- a/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj +++ b/OrderProcessing.EmailWorker/OrderProcessing.EmailWorker.csproj @@ -8,6 +8,12 @@ + + + all + runtime; build; native; contentfiles; analyzers; buildtransitive + + diff --git a/OrderProcessing.EmailWorker/Persistence/EmailWorkerDbContext.cs b/OrderProcessing.EmailWorker/Persistence/EmailWorkerDbContext.cs new file mode 100644 index 0000000..1b8ca6c --- /dev/null +++ b/OrderProcessing.EmailWorker/Persistence/EmailWorkerDbContext.cs @@ -0,0 +1,41 @@ +using Microsoft.EntityFrameworkCore; +using OrderProcessing.EmailWorker.Persistence.Entities; + +namespace OrderProcessing.EmailWorker.Persistence; + +public sealed class EmailWorkerDbContext + : DbContext +{ + public EmailWorkerDbContext(DbContextOptions options) : base(options) + { + } + + public DbSet ProcessedMessages => Set(); + + protected override void OnModelCreating(ModelBuilder modelBuilder) + { + base.OnModelCreating(modelBuilder); + + modelBuilder.Entity( + builder => + { + builder.ToTable( + "ProcessedMessages", + schema: "email"); + + builder.HasKey(message =>message.MessageId); + + builder.Property(message => message.MessageId) + .ValueGeneratedNever(); + + builder.Property(message => message.EventType) + .HasMaxLength(256) + .IsRequired(); + + builder.Property(message => message.ProcessedAtUtc) + .IsRequired(); + + builder.HasIndex(message => message.ProcessedAtUtc); + }); + } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Persistence/Entities/ProcessedMessage.cs b/OrderProcessing.EmailWorker/Persistence/Entities/ProcessedMessage.cs new file mode 100644 index 0000000..3fcc629 --- /dev/null +++ b/OrderProcessing.EmailWorker/Persistence/Entities/ProcessedMessage.cs @@ -0,0 +1,10 @@ +namespace OrderProcessing.EmailWorker.Persistence.Entities; + +public sealed class ProcessedMessage +{ + public Guid MessageId { get; set; } + + public required string EventType { get; set; } + + public DateTime ProcessedAtUtc { get; set; } +} \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/Program.cs b/OrderProcessing.EmailWorker/Program.cs index 1f4f08d..9856f12 100644 --- a/OrderProcessing.EmailWorker/Program.cs +++ b/OrderProcessing.EmailWorker/Program.cs @@ -1,10 +1,26 @@ +using Microsoft.EntityFrameworkCore; using OrderProcessing.EmailWorker; using OrderProcessing.EmailWorker.Configuration; using OrderProcessing.EmailWorker.Emailing; using OrderProcessing.EmailWorker.Messaging; +using OrderProcessing.EmailWorker.Persistence; var builder = Host.CreateApplicationBuilder(args); +var connectionString = builder.Configuration.GetConnectionString("EmailWorkerConnection") ?? + throw new InvalidOperationException( "EmailWorker database connection string is missing."); + +builder.Services.AddDbContext( + options => + options.UseSqlServer( + connectionString, + sqlOptions => + { + sqlOptions.MigrationsHistoryTable("__EFMigrationsHistory", "email"); + })); + +builder.Services.AddScoped(); + builder.Services .AddOptions() .Bind( @@ -38,15 +54,12 @@ "RabbitMQ prefetch count must be greater than zero.") .ValidateOnStart(); -builder.Services.AddSingleton< - IEmailSender, - LoggingEmailSender>(); -builder.Services.AddScoped< - OrderEventEmailHandler>(); +builder.Services.AddSingleton(); + +builder.Services.AddScoped(); -builder.Services.AddHostedService< - RabbitMqEmailConsumer>(); +builder.Services.AddHostedService(); var host = builder.Build(); diff --git a/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs b/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs index 2a8199d..d00d4c9 100644 --- a/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs +++ b/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs @@ -130,8 +130,7 @@ await Task.Delay( } } - private async Task HandleDeliveryAsync( - BasicDeliverEventArgs eventArgs) + private async Task HandleDeliveryAsync(BasicDeliverEventArgs eventArgs) { var channel = _channel ?? throw new InvalidOperationException( @@ -142,7 +141,8 @@ private async Task HandleDeliveryAsync( var body = eventArgs.Body.ToArray(); var eventType = eventArgs.BasicProperties.Type; - + var messageIdText = eventArgs.BasicProperties.MessageId; + try { if (string.IsNullOrWhiteSpace(eventType)) @@ -150,19 +150,18 @@ private async Task HandleDeliveryAsync( 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(); + var processor = scope.ServiceProvider.GetRequiredService(); - await handler.HandleAsync( - eventType, - body, - CancellationToken.None); + await processor.ProcessAsync(messageId, eventType, body, CancellationToken.None); - await channel.BasicAckAsync( - deliveryTag: eventArgs.DeliveryTag, - multiple: false); + await channel.BasicAckAsync(deliveryTag: eventArgs.DeliveryTag, multiple: false); _logger.LogInformation( "Acknowledged RabbitMQ message {MessageId} " + diff --git a/OrderProcessing.EmailWorker/appsettings.json b/OrderProcessing.EmailWorker/appsettings.json index 4c8cdd7..c5e9bb7 100644 --- a/OrderProcessing.EmailWorker/appsettings.json +++ b/OrderProcessing.EmailWorker/appsettings.json @@ -22,5 +22,8 @@ "ClientProvidedName": "order-processing-email-worker", "NetworkRecoveryIntervalSeconds": 5, "PrefetchCount": 1 + }, + "ConnectionStrings": { + "EmailWorkerConnection": "Server=(localdb)\\mssqllocaldb;Database=OrderProcessingDb;Trusted_Connection=True;TrustServerCertificate=True;" } }