using Dpz.Core.MessageQueue.Abstractions;
using Dpz.Core.MessageQueue.Models;
using Dpz.Core.MessageQueue.RabbitMQ;
using Medallion.Threading.FileSystem;
using Microsoft.Extensions.Logging.Abstractions;
using ZiggyCreatures.Caching.Fusion;

namespace Dpz.Core.MessageQueue.Test.RabbitMQ;

public class MessageOutboxRetryBackgroundServiceTests
{
    private const string StateCacheKey = "Dpz.Core.MessageQueue.OutboxRetry.LastState";

    [Fact]
    public async Task ExecuteAsync_ShouldThrottle_WhenNoMessagesNeedRetry()
    {
        var cache = new FusionCache(new FusionCacheOptions());
        var retryService = new CountingRetryService(new MessageOutboxRetryResult(0, 0, 0, 0));
        var tempDir = CreateTempDir();
        var lockProvider = new FileDistributedSynchronizationProvider(new DirectoryInfo(tempDir));
        var sut = CreateService(retryService, cache, lockProvider);

        try
        {
            await sut.StartAsync(CancellationToken.None);
            await WaitUntilAsync(async () =>
            {
                var state = await cache.TryGetAsync<MessageOutboxRetryState>(StateCacheKey);
                return retryService.CallCount == 1 && state.HasValue;
            });

            await Task.Delay(TimeSpan.FromMilliseconds(200));

            Assert.Equal(1, retryService.CallCount);
        }
        finally
        {
            await sut.StopAsync(CancellationToken.None);
            sut.Dispose();
            CleanupTempDir(tempDir);
        }
    }

    [Fact]
    public async Task ExecuteAsync_ShouldRunOnce_WhenTwoWorkersStartTogether()
    {
        var cache = new FusionCache(new FusionCacheOptions());
        var retryService = new CountingRetryService(
            new MessageOutboxRetryResult(1, 0, 0, 0),
            TimeSpan.FromMilliseconds(150)
        );
        var tempDir = CreateTempDir();
        var lockProvider = new FileDistributedSynchronizationProvider(new DirectoryInfo(tempDir));
        var first = CreateService(retryService, cache, lockProvider);
        var second = CreateService(retryService, cache, lockProvider);

        try
        {
            await first.StartAsync(CancellationToken.None);
            await second.StartAsync(CancellationToken.None);

            await WaitUntilAsync(async () =>
            {
                var state = await cache.TryGetAsync<MessageOutboxRetryState>(StateCacheKey);
                return retryService.CallCount == 1 && state.HasValue;
            });

            await Task.Delay(TimeSpan.FromMilliseconds(300));

            Assert.Equal(1, retryService.CallCount);
        }
        finally
        {
            await first.StopAsync(CancellationToken.None);
            await second.StopAsync(CancellationToken.None);
            first.Dispose();
            second.Dispose();
            CleanupTempDir(tempDir);
        }
    }

    private static MessageOutboxRetryBackgroundService CreateService(
        IMessageOutboxRetryService retryService,
        IFusionCache cache,
        FileDistributedSynchronizationProvider lockProvider
    )
    {
        return new MessageOutboxRetryBackgroundService(
            retryService,
            cache,
            lockProvider,
            NullLogger<MessageOutboxRetryBackgroundService>.Instance
        );
    }

    private static string CreateTempDir()
    {
        var tempDir = Path.Combine(
            Path.GetTempPath(),
            "dpz.core.mq.retry.test",
            Guid.NewGuid().ToString("N")
        );
        Directory.CreateDirectory(tempDir);
        return tempDir;
    }

    private static async Task WaitUntilAsync(Func<Task<bool>> predicate)
    {
        using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5));
        while (!timeout.IsCancellationRequested)
        {
            if (await predicate())
            {
                return;
            }

            await Task.Delay(TimeSpan.FromMilliseconds(25), timeout.Token);
        }

        throw new TimeoutException("The expected background service state was not reached.");
    }

    private static void CleanupTempDir(string tempDir)
    {
        if (Directory.Exists(tempDir))
        {
            Directory.Delete(tempDir, true);
        }
    }

    private sealed class CountingRetryService(
        MessageOutboxRetryResult result,
        TimeSpan? delay = null
    ) : IMessageOutboxRetryService
    {
        private int _callCount;

        public int CallCount => _callCount;

        public Task<MessageOutboxRetryResult> RetryAsync(
            int batchSize,
            CancellationToken cancellationToken = default
        )
        {
            Interlocked.Increment(ref _callCount);
            if (delay.HasValue)
            {
                return DelayAndReturnAsync(delay.Value, result, cancellationToken);
            }

            return Task.FromResult(result);
        }

        private static async Task<MessageOutboxRetryResult> DelayAndReturnAsync(
            TimeSpan delay,
            MessageOutboxRetryResult result,
            CancellationToken cancellationToken
        )
        {
            await Task.Delay(delay, cancellationToken);
            return result;
        }
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

这是一个针对 MessageOutboxRetryBackgroundService 类的单元测试文件。该服务用于重试消息队列中失败的消息发送操作。让我详细解释代码的各个部分:

测试类概述

这个测试类主要验证消息重试后台服务的两个关键行为:

  1. 当没有消息需要重试时的节流机制
  2. 多个工作进程同时启动时的并发控制

核心组件

依赖项

  • FusionCache: 用于缓存消息重试状态
  • FileDistributedSynchronizationProvider: 基于文件的分布式锁,确保多实例间的同步
  • CountingRetryService: 测试用的模拟重试服务,用于统计调用次数

测试方法

1. ExecuteAsync_ShouldThrottle_WhenNoMessagesNeedRetry

// 测试节流机制
var retryService = new CountingRetryService(new MessageOutboxRetryResult(0, 0, 0, 0));
  • 目的: 验证当没有消息需要重试时,服务应该实施节流,避免频繁执行
  • 验证点: 重试服务只被调用一次,然后进入等待状态
  • 关键断言: Assert.Equal(1, retryService.CallCount) - 确保只调用了一次

2. ExecuteAsync_ShouldRunOnce_WhenTwoWorkersStartTogether

// 测试并发控制
var first = CreateService(retryService, cache, lockProvider);
var second = CreateService(retryService, cache, lockProvider);
  • 目的: 验证多个服务实例同时运行时,分布式锁能确保重试逻辑只执行一次
  • 验证点: 即使启动两个服务实例,重试服务仍然只被调用一次
  • 关键断言: Assert.Equal(1, retryService.CallCount) - 确保分布式锁生效

辅助方法

CreateService

创建待测试的服务实例,注入所有必要的依赖项。

CreateTempDirCleanupTempDir

管理测试用的临时目录,用于文件锁的存储。

WaitUntilAsync

异步等待工具方法,用于等待后台服务达到预期状态:

var state = await cache.TryGetAsync<MessageOutboxRetryState>(StateCacheKey);
return retryService.CallCount == 1 && state.HasValue;

测试用模拟类 - CountingRetryService

private sealed class CountingRetryService : IMessageOutboxRetryService
{
    private int _callCount;
    public int CallCount => _callCount;
    
    public Task<MessageOutboxRetryResult> RetryAsync(...)
    {
        Interlocked.Increment(ref _callCount); // 线程安全的计数器
        // 可选延迟模拟真实处理时间
    }
}

测试策略

这些测试采用了以下策略:

  1. 隔离性: 每个测试使用独立的临时目录和缓存实例
  2. 可观察性: 通过计数器和状态缓存验证服务行为
  3. 超时保护: 使用 5 秒超时避免测试无限等待
  4. 资源清理: 确保测试后清理临时文件

业务价值

这些测试确保了:

  • 性能优化: 无消息时不会过度消耗资源
  • 数据一致性: 多实例部署时不会重复处理消息
  • 可靠性: 分布式锁机制正常工作

总的来说,这是一个设计良好的测试类,全面验证了消息重试后台服务的核心功能和边界条件。

评论加载中...