Add RabbitMQ integration tests and URL parsing logic
Enhanced test coverage for RabbitMQ integration by adding: - New package references for dependency injection, logging, and RabbitMQ. - `EmailSenderCollection` and `EmailSenderFixture` for managing static client lifecycle in integration tests. - Integration tests for `EmailSender` and `OutgoingEmailPublisher` to verify connection management, message publishing, and serialization. - `EmailSenderUrlOverloadTests` to validate URL-based connection overload behavior. - `RabbitMqTestConfig` for centralized RabbitMQ test configuration. - Unit tests for URL parsing logic in `EmailSenderUrlParsingTests`. These changes improve reliability, maintainability, and test coverage for the messaging service.
This commit is contained in:
@@ -0,0 +1,110 @@
|
||||
using DigitalData.MessagingService.Client.DependencyInjection;
|
||||
using DigitalData.MessagingService.Publisher.Abstraction;
|
||||
|
||||
namespace DigitalData.MessagingService.Tests.Integration;
|
||||
|
||||
/// <summary>
|
||||
/// xUnit collection that serializes all <see cref="EmailSenderTests"/> so they share
|
||||
/// the same process-wide static state of <see cref="EmailSender.LazyProvider"/>.
|
||||
/// </summary>
|
||||
[CollectionDefinition(Name)]
|
||||
public sealed class EmailSenderCollection : ICollectionFixture<EmailSenderFixture>
|
||||
{
|
||||
public const string Name = "EmailSender";
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Fixture that connects <see cref="EmailSender"/> once before all tests in the collection run.
|
||||
/// </summary>
|
||||
public sealed class EmailSenderFixture
|
||||
{
|
||||
public EmailSenderFixture()
|
||||
{
|
||||
// EmailSender is a static class with a Lazy<IServiceProvider>.
|
||||
// ConnectRabbitMq can only be called successfully once per process.
|
||||
if (!EmailSender.IsConnected)
|
||||
{
|
||||
EmailSender.ConnectRabbitMq(cfg =>
|
||||
{
|
||||
cfg.HostName = RabbitMqTestConfig.HostName;
|
||||
cfg.Port = RabbitMqTestConfig.Port;
|
||||
cfg.UserName = RabbitMqTestConfig.UserName;
|
||||
cfg.Password = RabbitMqTestConfig.Password;
|
||||
cfg.VirtualHost = RabbitMqTestConfig.VirtualHost;
|
||||
cfg.QueueName = RabbitMqTestConfig.QueueName;
|
||||
cfg.ExchangeName = RabbitMqTestConfig.ExchangeName;
|
||||
cfg.RoutingKey = RabbitMqTestConfig.RoutingKey;
|
||||
cfg.DlqQueueName = RabbitMqTestConfig.DlqQueueName;
|
||||
cfg.DlqExchangeName = RabbitMqTestConfig.DlqExchangeName;
|
||||
cfg.DlqRoutingKey = RabbitMqTestConfig.DlqRoutingKey;
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Integration tests for the <see cref="EmailSender"/> static client.
|
||||
/// All tests run inside <see cref="EmailSenderCollection"/> to share the single
|
||||
/// static connection established by <see cref="EmailSenderFixture"/>.
|
||||
/// </summary>
|
||||
[Collection(EmailSenderCollection.Name)]
|
||||
public sealed class EmailSenderTests
|
||||
{
|
||||
[Fact]
|
||||
public void IsConnected_AfterConnectRabbitMq_ReturnsTrue()
|
||||
{
|
||||
Assert.True(EmailSender.IsConnected);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void ConnectRabbitMq_WhenAlreadyConnected_WithThrowException_ThrowsInvalidOperationException()
|
||||
{
|
||||
Assert.Throws<InvalidOperationException>(() =>
|
||||
EmailSender.ConnectRabbitMq(cfg => { }, OnReconnect.ThrowException));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void ConnectRabbitMq_WhenAlreadyConnected_WithIgnore_DoesNotThrow()
|
||||
{
|
||||
var exception = Record.Exception(() =>
|
||||
EmailSender.ConnectRabbitMq(cfg => { }, OnReconnect.Ignore));
|
||||
|
||||
Assert.Null(exception);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void Send_WithValidEmail_DoesNotThrow()
|
||||
{
|
||||
var email = new OutgoingEmailEvent
|
||||
{
|
||||
Id = Guid.NewGuid(),
|
||||
Recipient = "hakanttek@gmail.com",
|
||||
Subject = "EmailSender.Send Integration Test",
|
||||
Body = "<p>Sent via EmailSender static client.</p>",
|
||||
IsHtml = true,
|
||||
QueuedAt = DateTime.Now
|
||||
};
|
||||
|
||||
var exception = Record.Exception(() => EmailSender.Send(email));
|
||||
|
||||
Assert.Null(exception);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void Send_WithPlainTextBody_DoesNotThrow()
|
||||
{
|
||||
var email = new OutgoingEmailEvent
|
||||
{
|
||||
Id = Guid.NewGuid(),
|
||||
Recipient = "hakanttek@gmail.com",
|
||||
Subject = "Plain Text Test",
|
||||
Body = "This is a plain text email.",
|
||||
IsHtml = false,
|
||||
QueuedAt = DateTime.Now
|
||||
};
|
||||
|
||||
var exception = Record.Exception(() => EmailSender.Send(email));
|
||||
|
||||
Assert.Null(exception);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
using DigitalData.MessagingService.Client.DependencyInjection;
|
||||
using DigitalData.MessagingService.Publisher.Abstraction;
|
||||
|
||||
namespace DigitalData.MessagingService.Tests.Integration;
|
||||
|
||||
/// <summary>
|
||||
/// Integration tests for the URL-based <see cref="EmailSender.ConnectRabbitMq(string, string, string, OnReconnect)"/>
|
||||
/// overload. All tests run inside <see cref="EmailSenderCollection"/> so they share the already-established
|
||||
/// static connection. The URL overload is exercised via <see cref="OnReconnect.Ignore"/> and
|
||||
/// <see cref="OnReconnect.ThrowException"/> to verify its delegation and guard behavior.
|
||||
/// </summary>
|
||||
[Collection(EmailSenderCollection.Name)]
|
||||
public sealed class EmailSenderUrlOverloadTests
|
||||
{
|
||||
private static readonly string ValidUrl = $"amqp://{RabbitMqTestConfig.HostName}:{RabbitMqTestConfig.Port}";
|
||||
private static readonly string ValidUrlWithVHost = $"amqp://{RabbitMqTestConfig.HostName}:{RabbitMqTestConfig.Port}/myvhost";
|
||||
|
||||
[Fact]
|
||||
public void ConnectRabbitMq_UrlOverload_WhenAlreadyConnected_WithIgnore_DoesNotThrow()
|
||||
{
|
||||
var exception = Record.Exception(() =>
|
||||
EmailSender.ConnectRabbitMq(
|
||||
ValidUrl,
|
||||
RabbitMqTestConfig.UserName,
|
||||
RabbitMqTestConfig.Password,
|
||||
OnReconnect.Ignore));
|
||||
|
||||
Assert.Null(exception);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void ConnectRabbitMq_UrlOverload_WhenAlreadyConnected_WithThrowException_Throws()
|
||||
{
|
||||
Assert.Throws<InvalidOperationException>(() =>
|
||||
EmailSender.ConnectRabbitMq(
|
||||
ValidUrl,
|
||||
RabbitMqTestConfig.UserName,
|
||||
RabbitMqTestConfig.Password,
|
||||
OnReconnect.ThrowException));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void ConnectRabbitMq_UrlOverload_WithVHostInPath_WithIgnore_DoesNotThrow()
|
||||
{
|
||||
var exception = Record.Exception(() =>
|
||||
EmailSender.ConnectRabbitMq(
|
||||
ValidUrlWithVHost,
|
||||
RabbitMqTestConfig.UserName,
|
||||
RabbitMqTestConfig.Password,
|
||||
OnReconnect.Ignore));
|
||||
|
||||
Assert.Null(exception);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void Send_AfterUrlOverloadConnection_DoesNotThrow()
|
||||
{
|
||||
var email = new OutgoingEmailEvent
|
||||
{
|
||||
Id = Guid.NewGuid(),
|
||||
Recipient = "url-overload-test@example.com",
|
||||
Subject = "URL Overload Integration Test",
|
||||
Body = "<p>Sent after URL-based connection.</p>",
|
||||
IsHtml = true,
|
||||
QueuedAt = DateTime.Now
|
||||
};
|
||||
|
||||
// Connection was established via Action<> overload in the fixture;
|
||||
// Send should work regardless of which overload was used to connect.
|
||||
var exception = Record.Exception(() => EmailSender.Send(email));
|
||||
|
||||
Assert.Null(exception);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,195 @@
|
||||
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(),
|
||||
Recipient = "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(),
|
||||
Recipient = $"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(),
|
||||
Recipient = "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,
|
||||
Recipient = "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.Recipient);
|
||||
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();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user