Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
191 changes: 137 additions & 54 deletions OrderProcessing.EmailWorker.Tests/OrderEventEmailHandlerTests.cs
Original file line number Diff line number Diff line change
@@ -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<OrderEventEmailHandler>.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<OrderEventEmailHandler>.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<OrderEventEmailHandler>.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<OrderEventEmailHandler>.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<InvalidOperationException>(
() => processor.ProcessAsync(
messageId,
typeof(OrderCompletedIntegrationEvent)
.FullName!,
body,
CancellationToken.None));

Assert.Empty(
await dbContext.ProcessedMessages.ToListAsync());
}

private sealed class TestEmailSender : IEmailSender
{
public List<EmailMessage> 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.");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

<ItemGroup>
<PackageReference Include="coverlet.collector" Version="6.0.4" />
<PackageReference Include="Microsoft.EntityFrameworkCore.Sqlite" Version="10.0.10" />
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.14.1" />
<PackageReference Include="xunit" Version="2.9.3" />
<PackageReference Include="xunit.runner.visualstudio" Version="3.1.4" />
Expand Down
Original file line number Diff line number Diff line change
@@ -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<IdempotentEmailMessageProcessor>
_logger;

public IdempotentEmailMessageProcessor(EmailWorkerDbContext dbContext, OrderEventEmailHandler eventHandler, ILogger<IdempotentEmailMessageProcessor> 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);
}
}
16 changes: 4 additions & 12 deletions OrderProcessing.EmailWorker/Messaging/OrderEventEmailHandler.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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} " +
Expand All @@ -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,
Expand All @@ -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} " +
Expand All @@ -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} " +
Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading