using Dpz.Core.MessageQueue.Abstractions;
using Dpz.Core.MessageQueue.Models;
using Dpz.Core.MessageQueue.RabbitMQ;
using Microsoft.Extensions.Logging.Abstractions;

namespace Dpz.Core.MessageQueue.Test.RabbitMQ;

public class MessageOutboxRetryServiceTests
{
    [Fact]
    public async Task RetryAsync_ShouldRepublishPendingEntriesAndMarkSent()
    {
        var outboxStore = new InMemoryOutboxStore
        {
            PendingPublishEntries = [Entry("publish-1"), Entry("publish-2")],
            PendingConsumeEntries = [Entry("consume-1")],
        };
        var retryPublisher = new InMemoryRetryPublisher();
        var sut = new MessageOutboxRetryService(
            outboxStore,
            retryPublisher,
            NullLogger<MessageOutboxRetryService>.Instance
        );

        var result = await sut.RetryAsync(50);

        Assert.Equal(2, result.PublishRetryCount);
        Assert.Equal(0, result.PublishFailureCount);
        Assert.Equal(1, result.ConsumeRetryCount);
        Assert.Equal(0, result.ConsumeFailureCount);
        Assert.Equal(["publish-1", "publish-2", "consume-1"], outboxStore.SentMessageIds);
        Assert.Equal(["publish-1", "publish-2", "consume-1"], retryPublisher.PublishedMessageIds);
    }

    [Fact]
    public async Task RetryAsync_ShouldMarkPublishFailed_WhenPublishRetryThrows()
    {
        var outboxStore = new InMemoryOutboxStore { PendingPublishEntries = [Entry("publish-1")] };
        var retryPublisher = new InMemoryRetryPublisher { ThrowForMessageIds = ["publish-1"] };
        var sut = new MessageOutboxRetryService(
            outboxStore,
            retryPublisher,
            NullLogger<MessageOutboxRetryService>.Instance
        );

        var result = await sut.RetryAsync(50);

        Assert.Equal(1, result.PublishRetryCount);
        Assert.Equal(1, result.PublishFailureCount);
        Assert.Empty(outboxStore.SentMessageIds);
        Assert.Equal(["publish-1"], outboxStore.PublishFailedMessageIds);
    }

    [Fact]
    public async Task RetryAsync_ShouldMarkConsumeFailed_WhenConsumeRetryThrows()
    {
        var outboxStore = new InMemoryOutboxStore { PendingConsumeEntries = [Entry("consume-1")] };
        var retryPublisher = new InMemoryRetryPublisher { ThrowForMessageIds = ["consume-1"] };
        var sut = new MessageOutboxRetryService(
            outboxStore,
            retryPublisher,
            NullLogger<MessageOutboxRetryService>.Instance
        );

        var result = await sut.RetryAsync(50);

        Assert.Equal(1, result.ConsumeRetryCount);
        Assert.Equal(1, result.ConsumeFailureCount);
        Assert.Empty(outboxStore.SentMessageIds);
        Assert.Equal(["consume-1"], outboxStore.ConsumeFailedMessageIds);
    }

    private static MessageOutboxEntry Entry(string messageId) =>
        new(messageId, "test.exchange", "test.routing", "{}", 1, 1);

    private sealed class InMemoryRetryPublisher : IMessageOutboxRetryPublisher
    {
        public List<string> PublishedMessageIds { get; } = [];

        public HashSet<string> ThrowForMessageIds { get; init; } = [];

        public Task PublishRawAsync(
            string exchange,
            string routingKey,
            string messageId,
            string jsonPayload,
            CancellationToken cancellationToken = default
        )
        {
            if (ThrowForMessageIds.Contains(messageId))
            {
                throw new InvalidOperationException("retry failed");
            }

            PublishedMessageIds.Add(messageId);
            return Task.CompletedTask;
        }
    }

    private sealed class InMemoryOutboxStore : IMessageOutboxStore
    {
        public IReadOnlyList<MessageOutboxEntry> PendingPublishEntries { get; init; } = [];

        public IReadOnlyList<MessageOutboxEntry> PendingConsumeEntries { get; init; } = [];

        public List<string> SentMessageIds { get; } = [];

        public List<string> PublishFailedMessageIds { get; } = [];

        public List<string> ConsumeFailedMessageIds { 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)
        {
            SentMessageIds.Add(messageId);
            return Task.CompletedTask;
        }

        public Task MarkPublishFailedAsync(
            string messageId,
            string error,
            CancellationToken cancellationToken = default
        )
        {
            PublishFailedMessageIds.Add(messageId);
            return Task.CompletedTask;
        }

        public Task MarkConsumedAsync(
            string messageId,
            CancellationToken cancellationToken = default
        ) => Task.CompletedTask;

        public Task RecordConsumeAttemptFailedAsync(
            string messageId,
            string error,
            CancellationToken cancellationToken = default
        ) => Task.CompletedTask;

        public Task MarkConsumeFailedAsync(
            string messageId,
            string error,
            CancellationToken cancellationToken = default
        )
        {
            ConsumeFailedMessageIds.Add(messageId);
            return Task.CompletedTask;
        }

        public Task<IReadOnlyList<MessageOutboxEntry>> GetPendingPublishRetryAsync(
            int batchSize,
            CancellationToken cancellationToken = default
        ) => Task.FromResult(PendingPublishEntries);

        public Task<IReadOnlyList<MessageOutboxEntry>> GetPendingConsumeRetryAsync(
            int batchSize,
            CancellationToken cancellationToken = default
        ) => Task.FromResult(PendingConsumeEntries);
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

代码说明

这是一个针对 MessageOutboxRetryService 的单元测试类,用于测试消息重试服务的核心功能。该服务实现了消息发布和消费的失败重试机制(Outbox 模式)。

整体架构

测试使用了内存模拟实现来替代真实的依赖项:

  • InMemoryOutboxStore: 模拟消息持久化存储
  • InMemoryRetryPublisher: 模拟消息重新发布器

三个核心测试用例

1. RetryAsync_ShouldRepublishPendingEntriesAndMarkSent

测试场景: 正常重试流程

准备数据:
- 2条待重试的发布消息 (publish-1, publish-2)
- 1条待重试的消费消息 (consume-1)

验证结果:
✓ 发布重试计数 = 2
✓ 消费重试计数 = 1
✓ 失败计数均为 0
✓ 所有消息被标记为已发送
✓ 所有消息被重新发布

2. RetryAsync_ShouldMarkPublishFailed_WhenPublishRetryThrows

测试场景: 发布重试失败

模拟异常:
- publish-1 在重试时抛出异常

验证结果:
✓ 发布重试计数 = 1
✓ 发布失败计数 = 1
✓ 消息未被标记为已发送
✓ 消息被标记为发布失败

3. RetryAsync_ShouldMarkConsumeFailed_WhenConsumeRetryThrows

测试场景: 消费重试失败

模拟异常:
- consume-1 在重试时抛出异常

验证结果:
✓ 消费重试计数 = 1
✓ 消费失败计数 = 1
✓ 消息未被标记为已发送
✓ 消息被标记为消费失败

Mock 实现说明

InMemoryRetryPublisher

  • 记录所有成功发布的消息ID
  • 可配置特定消息ID抛出异常(模拟发布失败)

InMemoryOutboxStore

实现了完整的 IMessageOutboxStore 接口,提供:

  • 可配置的待重试消息列表
  • 记录各种状态变更(已发送、发布失败、消费失败)
  • 返回预设的待重试消息

测试覆盖的业务逻辑

  1. 成功路径: 批量重试并标记成功
  2. 异常处理: 发布/消费失败时的错误处理
  3. 状态管理: 正确更新消息状态(Sent/PublishFailed/ConsumeFailed)
  4. 统计信息: 返回准确的重试和失败计数

这是一个典型的事务性消息模式测试,确保消息最终一致性。

评论加载中...