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 测试了两种类型的消息消费者:

  1. 普通消费者 (RabbitMQConsumerBackgroundService<T>)
  2. 带返回值的消费者 (RabbitMQConsumerBackgroundServiceWithResult<T, TResult>)

主要测试场景

普通消费者测试

  1. 成功处理消息 (ConsumerCore_ShouldAck_WhenHandlerSucceeds)

    • 处理成功时应发送 ACK 确认
    • 消息 ID 被记录到已消费列表
  2. 处理失败但未达最大重试次数 (ConsumerCore_ShouldRepublishAndAck_WhenHandlerFailsBeforeMaxRetry)

    • 重新发布消息到队列(重试计数 +1)
    • 发送 ACK 确认原消息
    • 记录消费尝试失败
  3. 达到最大重试次数 (ConsumerCore_ShouldNackWithoutRequeue_WhenHandlerFailsAtMaxRetry)

    • 发送 NACK(不重新入队)
    • 将消息标记为最终失败
  4. 处理过程抛异常 (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: 标记为最终失败(丢弃)

这套测试确保了消息消费的可靠性和容错能力。

评论加载中...