using System.Text; using System.Text.Json; using DigitalData.MessagingService.Publisher; using DigitalData.MessagingService.Abstraction; using DigitalData.MessagingService.RabbitMQ; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using RabbitMQ.Client; namespace DigitalData.MessagingService.Tests.Integration; /// /// Integration tests for 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). /// public sealed class SendingEmailPublisherTests : IAsyncDisposable { private readonly ServiceProvider _serviceProvider; private readonly ISendingEmailPublisher _publisher; private readonly RabbitMqConnectionFactory _factory; public SendingEmailPublisherTests() { var services = new ServiceCollection(); services.AddLogging(); services.AddMessagingServicePublisher(RabbitMqTestConfig.Apply); _serviceProvider = services.BuildServiceProvider(); _publisher = _serviceProvider.GetRequiredService(); _factory = _serviceProvider.GetRequiredService(); } [Fact] public async Task EnqueueAsync_PublishesMessage_MessageArrivesInQueue() { var email = new SendingEmailEvent { Id = Guid.NewGuid(), Mail = new EmailContext { Sender = new EmailAccountDto { Username = "test@example.com", Password = "password", SmtpServer = "smtp.example.com" }, Recipients = ["test@example.com"], Subject = "Integration Test - EnqueueAsync", Body = "

Hello from integration test.

", 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 SendingEmailEvent { Id = Guid.NewGuid(), Mail = new EmailContext { Sender = new EmailAccountDto { Username = "test@example.com", Password = "password", SmtpServer = "smtp.example.com" }, Recipients = new List { $"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 SendingEmailEvent { Id = Guid.NewGuid(), Mail = new EmailContext { Sender = new EmailAccountDto { Username = "test@example.com", Password = "password", SmtpServer = "smtp.example.com" }, 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 SendingEmailEvent { Id = id, Mail = new EmailContext { Sender = new EmailAccountDto { Username = "test@example.com", Password = "password", SmtpServer = "smtp.example.com" }, Recipients = new List { "serialize@example.com" }, Subject = "Serialization Test", Body = "Bold", 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.Mail.Recipients); Assert.Equal("Serialization Test", received.Mail.Subject); Assert.Equal("Bold", received.Mail.Body); Assert.True(received.Mail.IsHtml); } /// /// Reads a single message from the queue without acknowledging it (peek via nack+requeue). /// private async Task 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(json); } /// /// Scans queue messages (up to a limit) to find a message matching the given . /// All messages are re-queued after inspection. /// private async Task 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)>(); SendingEmailEvent? 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(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(); } }