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 类的单元测试文件。该服务用于重试消息队列中失败的消息发送操作。让我详细解释代码的各个部分:
测试类概述
这个测试类主要验证消息重试后台服务的两个关键行为:
- 当没有消息需要重试时的节流机制
- 多个工作进程同时启动时的并发控制
核心组件
依赖项
- 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
创建待测试的服务实例,注入所有必要的依赖项。
CreateTempDir 和 CleanupTempDir
管理测试用的临时目录,用于文件锁的存储。
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); // 线程安全的计数器
// 可选延迟模拟真实处理时间
}
}
测试策略
这些测试采用了以下策略:
- 隔离性: 每个测试使用独立的临时目录和缓存实例
- 可观察性: 通过计数器和状态缓存验证服务行为
- 超时保护: 使用 5 秒超时避免测试无限等待
- 资源清理: 确保测试后清理临时文件
业务价值
这些测试确保了:
- 性能优化: 无消息时不会过度消耗资源
- 数据一致性: 多实例部署时不会重复处理消息
- 可靠性: 分布式锁机制正常工作
总的来说,这是一个设计良好的测试类,全面验证了消息重试后台服务的核心功能和边界条件。
AI 正在分析代码…
评论加载中...