using System.Text;
using System.Text.Json;
using Dpz.Core.Entity.Base;
using Dpz.Core.MessageQueue.Abstractions;
using Dpz.Core.MessageQueue.Models;
using Dpz.Core.MessageQueue.RabbitMQ;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
namespace Dpz.Core.MessageQueue.Test.RabbitMQ;
public class RabbitMQConsumerBackgroundServiceTests
{
[Fact]
public async Task ConsumerCore_ShouldAck_WhenHandlerSucceeds()
{
var handler = new DelegatingMessageHandler(
(_, _) => Task.FromResult(MessageHandlerResult.Ok())
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildRegularConsumerProvider(handler, outboxStore);
var consumer = provider.GetRequiredService<TestableRabbitMQConsumerBackgroundService>();
var channel = BuildChannelMock();
var message = new TestMessage { MessageId = "m-success", RetryCount = 0 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(x => x.BasicAckAsync(100, false, It.IsAny<CancellationToken>()), Times.Once);
channel.Verify(
x =>
x.BasicNackAsync(
It.IsAny<ulong>(),
It.IsAny<bool>(),
It.IsAny<bool>(),
It.IsAny<CancellationToken>()
),
Times.Never
);
channel.Verify(
x =>
x.BasicPublishAsync(
It.IsAny<string>(),
It.IsAny<string>(),
It.IsAny<bool>(),
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
),
Times.Never
);
Assert.Equal(["m-success"], outboxStore.ConsumedMessageIds);
Assert.Empty(outboxStore.ConsumeFailedMessageIds);
}
[Fact]
public async Task ConsumerCore_ShouldRepublishAndAck_WhenHandlerFailsBeforeMaxRetry()
{
var handler = new DelegatingMessageHandler(
(_, _) => Task.FromResult(MessageHandlerResult.Fail("regular failed"))
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildRegularConsumerProvider(handler, outboxStore);
var consumer = provider.GetRequiredService<TestableRabbitMQConsumerBackgroundService>();
var channel = BuildChannelMock();
ReadOnlyMemory<byte> republishedBody = default;
channel
.Setup(x =>
x.BasicPublishAsync(
It.IsAny<string>(),
It.IsAny<string>(),
It.IsAny<bool>(),
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
)
)
.Callback<
string,
string,
bool,
BasicProperties,
ReadOnlyMemory<byte>,
CancellationToken
>((_, _, _, _, body, _) => republishedBody = body)
.Returns(ValueTask.CompletedTask);
var message = new TestMessage { MessageId = "m-retry", RetryCount = 0 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(
x =>
x.BasicPublishAsync(
"test-exchange",
"test-routing-key",
false,
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
),
Times.Once
);
channel.Verify(x => x.BasicAckAsync(100, false, It.IsAny<CancellationToken>()), Times.Once);
var republishedMessage = JsonSerializer.Deserialize<TestMessage>(republishedBody.Span);
Assert.NotNull(republishedMessage);
Assert.Equal(1, republishedMessage.RetryCount);
Assert.Equal(message.MessageId, republishedMessage.MessageId);
Assert.Equal(["m-retry"], outboxStore.ConsumeAttemptFailedMessageIds);
Assert.Equal(["regular failed"], outboxStore.ConsumeAttemptFailedErrors);
Assert.Empty(outboxStore.ConsumeFailedMessageIds);
Assert.Empty(outboxStore.ConsumedMessageIds);
}
[Fact]
public async Task ConsumerCore_ShouldNackWithoutRequeue_WhenHandlerFailsAtMaxRetry()
{
var handler = new DelegatingMessageHandler(
(_, _) => Task.FromResult(MessageHandlerResult.Fail("regular final failed"))
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildRegularConsumerProvider(handler, outboxStore);
var consumer = provider.GetRequiredService<TestableRabbitMQConsumerBackgroundService>();
var channel = BuildChannelMock();
var message = new TestMessage { MessageId = "m-drop", RetryCount = 3 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(
x => x.BasicNackAsync(100, false, false, It.IsAny<CancellationToken>()),
Times.Once
);
channel.Verify(
x =>
x.BasicAckAsync(It.IsAny<ulong>(), It.IsAny<bool>(), It.IsAny<CancellationToken>()),
Times.Never
);
channel.Verify(
x =>
x.BasicPublishAsync(
It.IsAny<string>(),
It.IsAny<string>(),
It.IsAny<bool>(),
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
),
Times.Never
);
Assert.Equal(["m-drop"], outboxStore.ConsumeFailedMessageIds);
Assert.Equal(["regular final failed"], outboxStore.ConsumeFailedErrors);
Assert.Empty(outboxStore.ConsumedMessageIds);
}
[Fact]
public async Task ConsumerCore_ShouldConvertHandlerExceptionToFailedResult()
{
var handler = new DelegatingMessageHandler(
(_, _) => throw new InvalidOperationException("regular boom")
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildRegularConsumerProvider(handler, outboxStore);
var consumer = provider.GetRequiredService<TestableRabbitMQConsumerBackgroundService>();
var channel = BuildChannelMock();
var message = new TestMessage { MessageId = "m-exception", RetryCount = 3 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(
x => x.BasicNackAsync(100, false, false, It.IsAny<CancellationToken>()),
Times.Once
);
Assert.Equal(["m-exception"], outboxStore.ConsumeFailedMessageIds);
Assert.Equal(["regular boom"], outboxStore.ConsumeFailedErrors);
}
[Fact]
public async Task ConsumerWithResultCore_ShouldAck_WhenResultSuccess()
{
var handler = new DelegatingResultHandler(
(_, _) => Task.FromResult(MessageHandlerResult<string>.Ok("ok"))
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildResultConsumerProvider(handler, outboxStore);
var consumer =
provider.GetRequiredService<TestableRabbitMQConsumerBackgroundServiceWithResult>();
var channel = BuildChannelMock();
var message = new TestMessage { MessageId = "r-success", RetryCount = 0 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(x => x.BasicAckAsync(100, false, It.IsAny<CancellationToken>()), Times.Once);
channel.Verify(
x =>
x.BasicNackAsync(
It.IsAny<ulong>(),
It.IsAny<bool>(),
It.IsAny<bool>(),
It.IsAny<CancellationToken>()
),
Times.Never
);
Assert.Equal(["r-success"], outboxStore.ConsumedMessageIds);
Assert.Empty(outboxStore.ConsumeFailedMessageIds);
}
[Fact]
public async Task ConsumerWithResultCore_ShouldRepublishAndAck_WhenResultFailedBeforeMaxRetry()
{
var handler = new DelegatingResultHandler(
(_, _) => Task.FromResult(MessageHandlerResult<string>.Fail("failed"))
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildResultConsumerProvider(handler, outboxStore);
var consumer =
provider.GetRequiredService<TestableRabbitMQConsumerBackgroundServiceWithResult>();
var channel = BuildChannelMock();
ReadOnlyMemory<byte> republishedBody = default;
channel
.Setup(x =>
x.BasicPublishAsync(
It.IsAny<string>(),
It.IsAny<string>(),
It.IsAny<bool>(),
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
)
)
.Callback<
string,
string,
bool,
BasicProperties,
ReadOnlyMemory<byte>,
CancellationToken
>((_, _, _, _, body, _) => republishedBody = body)
.Returns(ValueTask.CompletedTask);
var message = new TestMessage { MessageId = "r-retry", RetryCount = 0 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(
x =>
x.BasicPublishAsync(
"test-exchange",
"test-routing-key",
false,
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
),
Times.Once
);
channel.Verify(x => x.BasicAckAsync(100, false, It.IsAny<CancellationToken>()), Times.Once);
var republishedMessage = JsonSerializer.Deserialize<TestMessage>(republishedBody.Span);
Assert.NotNull(republishedMessage);
Assert.Equal(1, republishedMessage.RetryCount);
Assert.Equal(message.MessageId, republishedMessage.MessageId);
Assert.Equal(["r-retry"], outboxStore.ConsumeAttemptFailedMessageIds);
Assert.Equal(["failed"], outboxStore.ConsumeAttemptFailedErrors);
Assert.Empty(outboxStore.ConsumeFailedMessageIds);
Assert.Empty(outboxStore.ConsumedMessageIds);
}
[Fact]
public async Task ConsumerWithResultCore_ShouldNackWithoutRequeue_WhenResultFailedAtMaxRetry()
{
var handler = new DelegatingResultHandler(
(_, _) => Task.FromResult(MessageHandlerResult<string>.Fail("failed"))
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildResultConsumerProvider(handler, outboxStore);
var consumer =
provider.GetRequiredService<TestableRabbitMQConsumerBackgroundServiceWithResult>();
var channel = BuildChannelMock();
var message = new TestMessage { MessageId = "r-drop", RetryCount = 3 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(
x => x.BasicNackAsync(100, false, false, It.IsAny<CancellationToken>()),
Times.Once
);
channel.Verify(
x =>
x.BasicAckAsync(It.IsAny<ulong>(), It.IsAny<bool>(), It.IsAny<CancellationToken>()),
Times.Never
);
Assert.Equal(["r-drop"], outboxStore.ConsumeFailedMessageIds);
Assert.Equal(["failed"], outboxStore.ConsumeFailedErrors);
Assert.Empty(outboxStore.ConsumedMessageIds);
}
[Fact]
public async Task ConsumerWithResultCore_ShouldConvertHandlerExceptionToFailedResult()
{
var handler = new DelegatingResultHandler(
(_, _) => throw new InvalidOperationException("result boom")
);
var outboxStore = new InMemoryOutboxStore();
await using var provider = BuildResultConsumerProvider(handler, outboxStore);
var consumer =
provider.GetRequiredService<TestableRabbitMQConsumerBackgroundServiceWithResult>();
var channel = BuildChannelMock();
var message = new TestMessage { MessageId = "r-exception", RetryCount = 3 };
var eventArgs = CreateEventArgs(message.MessageId);
await consumer.InvokeHandleMessageAsyncCore(
message,
eventArgs,
channel.Object,
CancellationToken.None
);
channel.Verify(
x => x.BasicNackAsync(100, false, false, It.IsAny<CancellationToken>()),
Times.Once
);
Assert.Equal(["r-exception"], outboxStore.ConsumeFailedMessageIds);
Assert.Equal(["result boom"], outboxStore.ConsumeFailedErrors);
}
private static ServiceProvider BuildRegularConsumerProvider(
DelegatingMessageHandler handler,
IMessageOutboxStore? outboxStore = null
)
{
var routing = BuildRoutingConventionMock();
var factory = new Mock<IRabbitMQConnectionFactory>();
var services = new ServiceCollection();
services.AddSingleton(factory.Object);
services.AddSingleton(routing.Object);
services.AddSingleton<ILoggerFactory>(NullLoggerFactory.Instance);
services.AddSingleton(outboxStore ?? NullMessageOutboxStore.Instance.Value);
services.AddSingleton(handler);
services.AddScoped<IMessageHandler<TestMessage>>(_ => handler);
services.AddTransient<TestableRabbitMQConsumerBackgroundService>();
return services.BuildServiceProvider();
}
private static ServiceProvider BuildResultConsumerProvider(
DelegatingResultHandler handler,
IMessageOutboxStore? outboxStore = null
)
{
var routing = BuildRoutingConventionMock();
var factory = new Mock<IRabbitMQConnectionFactory>();
var services = new ServiceCollection();
services.AddSingleton(factory.Object);
services.AddSingleton(routing.Object);
services.AddSingleton<ILoggerFactory>(NullLoggerFactory.Instance);
services.AddSingleton(outboxStore ?? NullMessageOutboxStore.Instance.Value);
services.AddSingleton(handler);
services.AddScoped<IMessageHandler<TestMessage, MessageHandlerResult<string>>>(_ =>
handler
);
services.AddTransient<TestableRabbitMQConsumerBackgroundServiceWithResult>();
return services.BuildServiceProvider();
}
private static Mock<IMessageRoutingConvention> BuildRoutingConventionMock()
{
var routing = new Mock<IMessageRoutingConvention>();
routing.Setup(x => x.GetExchangeName<TestMessage>()).Returns("test-exchange");
routing.Setup(x => x.GetQueueName<TestMessage>()).Returns("test-queue");
routing.Setup(x => x.GetRoutingKey<TestMessage>()).Returns("test-routing-key");
routing.Setup(x => x.GetExchangeType<TestMessage>()).Returns(Enums.ExchangeType.Direct);
return routing;
}
private static Mock<IChannel> BuildChannelMock()
{
var channel = new Mock<IChannel>();
channel
.Setup(x =>
x.BasicAckAsync(It.IsAny<ulong>(), It.IsAny<bool>(), It.IsAny<CancellationToken>())
)
.Returns(ValueTask.CompletedTask);
channel
.Setup(x =>
x.BasicNackAsync(
It.IsAny<ulong>(),
It.IsAny<bool>(),
It.IsAny<bool>(),
It.IsAny<CancellationToken>()
)
)
.Returns(ValueTask.CompletedTask);
channel
.Setup(x =>
x.BasicPublishAsync(
It.IsAny<string>(),
It.IsAny<string>(),
It.IsAny<bool>(),
It.IsAny<BasicProperties>(),
It.IsAny<ReadOnlyMemory<byte>>(),
It.IsAny<CancellationToken>()
)
)
.Returns(ValueTask.CompletedTask);
return channel;
}
private static BasicDeliverEventArgs CreateEventArgs(string messageId)
{
return new BasicDeliverEventArgs(
consumerTag: "consumer-1",
deliveryTag: 100,
redelivered: false,
exchange: "test-exchange",
routingKey: "test-routing-key",
properties: new BasicProperties { MessageId = messageId },
body: Encoding.UTF8.GetBytes("{}"),
cancellationToken: CancellationToken.None
);
}
private sealed class TestableRabbitMQConsumerBackgroundService(
IRabbitMQConnectionFactory connectionFactory,
IMessageRoutingConvention routingConvention,
IServiceProvider serviceProvider,
IMessageOutboxStore outboxStore,
ILoggerFactory loggerFactory
)
: RabbitMQConsumerBackgroundService<TestMessage>(
connectionFactory,
routingConvention,
serviceProvider,
outboxStore,
loggerFactory
)
{
public Task InvokeHandleMessageAsyncCore(
TestMessage message,
BasicDeliverEventArgs eventArgs,
IChannel channel,
CancellationToken cancellationToken
) => HandleMessageAsyncCore(message, eventArgs, channel, cancellationToken);
}
private sealed class TestableRabbitMQConsumerBackgroundServiceWithResult(
IRabbitMQConnectionFactory connectionFactory,
IMessageRoutingConvention routingConvention,
IServiceProvider serviceProvider,
IMessageOutboxStore outboxStore,
ILoggerFactory loggerFactory
)
: RabbitMQConsumerBackgroundServiceWithResult<TestMessage, MessageHandlerResult<string>>(
connectionFactory,
routingConvention,
serviceProvider,
outboxStore,
loggerFactory
)
{
public Task InvokeHandleMessageAsyncCore(
TestMessage message,
BasicDeliverEventArgs eventArgs,
IChannel channel,
CancellationToken cancellationToken
) => HandleMessageAsyncCore(message, eventArgs, channel, cancellationToken);
}
private sealed class DelegatingMessageHandler(
Func<TestMessage, CancellationToken, Task<MessageHandlerResult>> handler
) : IMessageHandler<TestMessage>
{
public Task<MessageHandlerResult> HandleAsync(
TestMessage message,
CancellationToken cancellationToken = default
)
{
return handler(message, cancellationToken);
}
}
private sealed class DelegatingResultHandler(
Func<TestMessage, CancellationToken, Task<MessageHandlerResult<string>>> handler
) : IMessageHandler<TestMessage, MessageHandlerResult<string>>
{
public Task<MessageHandlerResult<string>> HandleAsync(
TestMessage message,
CancellationToken cancellationToken = default
)
{
return handler(message, cancellationToken);
}
}
private sealed class InMemoryOutboxStore : IMessageOutboxStore
{
public List<string> ConsumedMessageIds { get; } = [];
public List<string> ConsumeAttemptFailedMessageIds { get; } = [];
public List<string> ConsumeAttemptFailedErrors { get; } = [];
public List<string> ConsumeFailedMessageIds { get; } = [];
public List<string> ConsumeFailedErrors { get; } = [];
public Task CreateAsync(
string messageId,
string messageType,
string exchange,
string routingKey,
string payload,
string? source,
CancellationToken cancellationToken = default
) => Task.CompletedTask;
public Task MarkSentAsync(
string messageId,
CancellationToken cancellationToken = default
) => Task.CompletedTask;
public Task MarkPublishFailedAsync(
string messageId,
string error,
CancellationToken cancellationToken = default
) => Task.CompletedTask;
public Task MarkConsumedAsync(
string messageId,
CancellationToken cancellationToken = default
)
{
ConsumedMessageIds.Add(messageId);
return Task.CompletedTask;
}
public Task RecordConsumeAttemptFailedAsync(
string messageId,
string error,
CancellationToken cancellationToken = default
)
{
ConsumeAttemptFailedMessageIds.Add(messageId);
ConsumeAttemptFailedErrors.Add(error);
return Task.CompletedTask;
}
public Task MarkConsumeFailedAsync(
string messageId,
string error,
CancellationToken cancellationToken = default
)
{
ConsumeFailedMessageIds.Add(messageId);
ConsumeFailedErrors.Add(error);
return Task.CompletedTask;
}
public Task<IReadOnlyList<MessageOutboxEntry>> GetPendingPublishRetryAsync(
int batchSize,
CancellationToken cancellationToken = default
) => Task.FromResult<IReadOnlyList<MessageOutboxEntry>>([]);
public Task<IReadOnlyList<MessageOutboxEntry>> GetPendingConsumeRetryAsync(
int batchSize,
CancellationToken cancellationToken = default
) => Task.FromResult<IReadOnlyList<MessageOutboxEntry>>([]);
}
private sealed class TestMessage : MessageBase
{
public string Payload { get; set; } = string.Empty;
}
}
⚠⚠ 以下内容为AI分析的结果,请根据实际情况进行判断。
代码说明
这是一个针对 RabbitMQ 消息队列消费者后台服务的单元测试文件,主要测试消息处理的各种场景。
核心功能
该测试类 RabbitMQConsumerBackgroundServiceTests 测试了两种类型的消息消费者:
- 普通消费者 (
RabbitMQConsumerBackgroundService<T>) - 带返回值的消费者 (
RabbitMQConsumerBackgroundServiceWithResult<T, TResult>)
主要测试场景
普通消费者测试
成功处理消息 (
ConsumerCore_ShouldAck_WhenHandlerSucceeds)- 处理成功时应发送 ACK 确认
- 消息 ID 被记录到已消费列表
处理失败但未达最大重试次数 (
ConsumerCore_ShouldRepublishAndAck_WhenHandlerFailsBeforeMaxRetry)- 重新发布消息到队列(重试计数 +1)
- 发送 ACK 确认原消息
- 记录消费尝试失败
达到最大重试次数 (
ConsumerCore_ShouldNackWithoutRequeue_WhenHandlerFailsAtMaxRetry)- 发送 NACK(不重新入队)
- 将消息标记为最终失败
处理过程抛异常 (
ConsumerCore_ShouldConvertHandlerExceptionToFailedResult)- 异常被转换为失败结果
- 达到最大重试时发送 NACK
带返回值消费者测试
对应上述场景进行了相同的测试(前缀为 ConsumerWithResultCore_)
辅助组件
测试用类
TestableRabbitMQConsumerBackgroundService: 可测试的消费者实现,暴露了内部的HandleMessageAsyncCore方法DelegatingMessageHandler/DelegatingResultHandler: 委托型处理器,允许通过 lambda 自定义处理逻辑InMemoryOutboxStore: 内存实现的消息出站箱存储,用于验证消息状态变更TestMessage: 测试用消息类型
Mock 构建方法
BuildRegularConsumerProvider: 构建普通消费者的 DI 容器BuildResultConsumerProvider: 构建带返回值消费者的 DI 容器BuildRoutingConventionMock: Mock 路由约定(exchange、queue、routing key)BuildChannelMock: Mock RabbitMQ 通道CreateEventArgs: 创建消息投递事件参数
验证逻辑
使用 Moq 框架验证:
BasicAckAsync: 消息确认BasicNackAsync: 消息拒绝BasicPublishAsync: 消息重新发布
通过 InMemoryOutboxStore 验证消息状态:
ConsumedMessageIds: 成功消费ConsumeFailedMessageIds: 最终失败ConsumeAttemptFailedMessageIds: 重试失败
重试机制
测试验证了基于 RetryCount 的重试逻辑:
RetryCount < 3: 重新发布消息(重试)RetryCount >= 3: 标记为最终失败(丢弃)
这套测试确保了消息消费的可靠性和容错能力。
AI 正在分析代码…
评论加载中...