Refactored the `Email` and `OutgoingEmailEvent` classes to replace the `Recipient` property with a `Recipients` collection, enabling support for multiple recipients. Updated all related test cases, including `EmailSenderTests`, `EmailSenderUrlOverloadTests`, and `OutgoingEmailPublisherTests`, to reflect this change. Moved the `EmailAccountDto` class and its references from the `DigitalData.MessagingService.Application.Common.Dtos` namespace to the `DigitalData.MessagingService.Publisher.Abstraction` namespace for better code organization. Updated `using` directives across affected files. Removed unused `using` directives and updated the `Email` class's `ToEvent` method to map the new `Recipients` property. Adjusted test assertions to validate collections instead of single recipient strings.
196 lines
6.6 KiB
C#
196 lines
6.6 KiB
C#
using System.Text;
|
|
using System.Text.Json;
|
|
using DigitalData.MessagingService.Publisher;
|
|
using DigitalData.MessagingService.Publisher.Abstraction;
|
|
using DigitalData.MessagingService.RabbitMQ;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Microsoft.Extensions.Logging;
|
|
using RabbitMQ.Client;
|
|
|
|
namespace DigitalData.MessagingService.Tests.Integration;
|
|
|
|
/// <summary>
|
|
/// Integration tests for <see cref="OutgoingEmailPublisher"/> against the real RabbitMQ broker.
|
|
/// Each test publishes a message and immediately reads it back via BasicGetAsync to verify
|
|
/// the full round-trip without starting the consumer (which requires Limilabs Mail.dll).
|
|
/// </summary>
|
|
public sealed class OutgoingEmailPublisherTests : IAsyncDisposable
|
|
{
|
|
private readonly ServiceProvider _serviceProvider;
|
|
private readonly IOutgoingEmailPublisher _publisher;
|
|
private readonly RabbitMqConnectionFactory _factory;
|
|
|
|
public OutgoingEmailPublisherTests()
|
|
{
|
|
var services = new ServiceCollection();
|
|
services.AddLogging();
|
|
services.AddMessagingServicePublisher(RabbitMqTestConfig.Apply);
|
|
|
|
_serviceProvider = services.BuildServiceProvider();
|
|
_publisher = _serviceProvider.GetRequiredService<IOutgoingEmailPublisher>();
|
|
_factory = _serviceProvider.GetRequiredService<RabbitMqConnectionFactory>();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task EnqueueAsync_PublishesMessage_MessageArrivesInQueue()
|
|
{
|
|
var email = new OutgoingEmailEvent
|
|
{
|
|
Id = Guid.NewGuid(),
|
|
Recipients = ["test@example.com"],
|
|
Subject = "Integration Test - EnqueueAsync",
|
|
Body = "<p>Hello from integration test.</p>",
|
|
IsHtml = true,
|
|
QueuedAt = DateTime.Now
|
|
};
|
|
|
|
var depthBefore = await _publisher.GetQueueDepthAsync();
|
|
|
|
await _publisher.EnqueueAsync(email);
|
|
|
|
await Task.Delay(300);
|
|
|
|
var depthAfter = await _publisher.GetQueueDepthAsync();
|
|
|
|
Assert.True(depthAfter >= depthBefore + 1,
|
|
$"Expected queue depth to increase by at least 1. Before: {depthBefore}, After: {depthAfter}.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task EnqueueAsync_MultipleMessages_AllArrivesInQueue()
|
|
{
|
|
var emails = Enumerable.Range(1, 3).Select(i => new OutgoingEmailEvent
|
|
{
|
|
Id = Guid.NewGuid(),
|
|
Recipients = new List<string> { $"recipient{i}@example.com" },
|
|
Subject = $"Integration Test - Batch #{i}",
|
|
Body = $"Batch message {i}",
|
|
IsHtml = false,
|
|
QueuedAt = DateTime.Now
|
|
}).ToList();
|
|
|
|
foreach (var email in emails)
|
|
await _publisher.EnqueueAsync(email);
|
|
|
|
await Task.Delay(500);
|
|
|
|
// Verify at least one message is present
|
|
var received = await PeekMessageAsync();
|
|
Assert.NotNull(received);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task GetQueueDepthAsync_AfterPublish_ReturnsPositiveDepth()
|
|
{
|
|
var email = new OutgoingEmailEvent
|
|
{
|
|
Id = Guid.NewGuid(),
|
|
Recipients = ["depth-test@example.com"],
|
|
Subject = "Integration Test - GetQueueDepth",
|
|
Body = "Queue depth test",
|
|
IsHtml = false,
|
|
QueuedAt = DateTime.Now
|
|
};
|
|
|
|
await _publisher.EnqueueAsync(email);
|
|
await Task.Delay(300);
|
|
|
|
var depth = await _publisher.GetQueueDepthAsync();
|
|
|
|
Assert.True(depth > 0, $"Expected queue depth > 0, but got {depth}.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task EnqueueAsync_SerializesAllFields_DeserializesCorrectly()
|
|
{
|
|
var id = Guid.NewGuid();
|
|
|
|
var email = new OutgoingEmailEvent
|
|
{
|
|
Id = id,
|
|
Recipients = new List<string> { "serialize@example.com" },
|
|
Subject = "Serialization Test",
|
|
Body = "<strong>Bold</strong>",
|
|
IsHtml = true,
|
|
QueuedAt = DateTime.Now
|
|
};
|
|
|
|
await _publisher.EnqueueAsync(email);
|
|
await Task.Delay(300);
|
|
|
|
// Read messages until we find the one we just published
|
|
var received = await FindMessageAsync(id);
|
|
|
|
Assert.NotNull(received);
|
|
Assert.Equal(id, received.Id);
|
|
Assert.Equal(["serialize@example.com"], received.Recipients);
|
|
Assert.Equal("Serialization Test", received.Subject);
|
|
Assert.Equal("<strong>Bold</strong>", received.Body);
|
|
Assert.True(received.IsHtml);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Reads a single message from the queue without acknowledging it (peek via nack+requeue).
|
|
/// </summary>
|
|
private async Task<OutgoingEmailEvent?> PeekMessageAsync()
|
|
{
|
|
var connection = await _factory.GetDefaultConnectionAsync();
|
|
await using var channel = await connection.CreateChannelAsync();
|
|
|
|
var result = await channel.BasicGetAsync(RabbitMqTestConfig.QueueName, autoAck: false);
|
|
|
|
if (result is null)
|
|
return null;
|
|
|
|
// Nack with requeue=true so the message stays in the queue for consumer
|
|
await channel.BasicNackAsync(result.DeliveryTag, multiple: false, requeue: true);
|
|
|
|
var json = Encoding.UTF8.GetString(result.Body.ToArray());
|
|
return JsonSerializer.Deserialize<OutgoingEmailEvent>(json);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Scans queue messages (up to a limit) to find a message matching the given <paramref name="id"/>.
|
|
/// All messages are re-queued after inspection.
|
|
/// </summary>
|
|
private async Task<OutgoingEmailEvent?> FindMessageAsync(Guid id, int maxMessages = 50)
|
|
{
|
|
var connection = await _factory.GetDefaultConnectionAsync();
|
|
await using var channel = await connection.CreateChannelAsync();
|
|
|
|
var requeue = new List<(ulong DeliveryTag, byte[] Body)>();
|
|
|
|
OutgoingEmailEvent? found = null;
|
|
|
|
for (int i = 0; i < maxMessages; i++)
|
|
{
|
|
var result = await channel.BasicGetAsync(RabbitMqTestConfig.QueueName, autoAck: false);
|
|
if (result is null)
|
|
break;
|
|
|
|
requeue.Add((result.DeliveryTag, result.Body.ToArray()));
|
|
|
|
var json = Encoding.UTF8.GetString(result.Body.ToArray());
|
|
var evt = JsonSerializer.Deserialize<OutgoingEmailEvent>(json);
|
|
|
|
if (evt?.Id == id)
|
|
{
|
|
found = evt;
|
|
break;
|
|
}
|
|
}
|
|
|
|
// Re-queue all inspected messages so the consumer can still process them
|
|
foreach (var (tag, _) in requeue)
|
|
await channel.BasicNackAsync(tag, multiple: false, requeue: true);
|
|
|
|
return found;
|
|
}
|
|
|
|
public async ValueTask DisposeAsync()
|
|
{
|
|
await _factory.DisposeAsync();
|
|
await _serviceProvider.DisposeAsync();
|
|
}
|
|
}
|