使用MassTransit Test Harness的单元测试在构建服务器上随机失败
测试代码结构
using MassTransit.Testing; namespace g2fp.PosApi.UnitTests; public class SampleTest : IAsyncLifetime { public async Task InitializeAsync() { var serviceCollection = new ServiceCollection(); serviceCollection.AddTransient<ISomeService, SomeService>(); serviceCollection.AddMassTransitTestHarness(cfg => { cfg.AddConsumer<SomeConsumer>(); }); Services = serviceCollection.BuildServiceProvider(); Harness = Services.GetRequiredService<ITestHarness>(); ConsumerHarness = Services.GetRequiredService<IConsumerTestHarness<SomeConsumer>>(); await Harness.Start(); } [Fact] public async Task Consume() { await Harness.Bus.Publish(new SomeMessage("something")); // await Harness.InactivityTask; (await Harness.Consumed.SelectAsync<SomeMessage>().Any()).Should().BeTrue(); (await ConsumerHarness.Consumed.SelectAsync<SomeMessage>().Any()).Should().BeTrue(); } public IConsumerTestHarness<SomeConsumer> ConsumerHarness { get; set; } public ITestHarness Harness { get; set; } public ServiceProvider Services { get; set; } public async Task DisposeAsync() { await Harness.Stop(); } } public record SomeMessage(string Message); public class SomeConsumer : IConsumer<SomeMessage> { private readonly ISomeService _someService; public SomeConsumer(ISomeService someService) { _someService = someService; } public async Task Consume(ConsumeContext<SomeMessage> context) { await _someService.DoSomething(); } } public interface ISomeService { Task DoSomething(); } public class SomeService : ISomeService { public async Task DoSomething() { await Task.Delay(500); } }
问题描述
本地运行该测试1000次均正常,但在GitLab流水线中,测试会随机在(await Harness.Consumed.SelectAsync<SomeMessage>().Any()).Should().BeTrue();这一行失败。无论调整SomeService中的延迟时长,测试都能等待其执行完成,但断言仍会随机失败。
已尝试的方案
- 调整等待顺序,先创建等待任务再发布消息:
var waitTask = Harness.Consumed.SelectAsync<SomeMessage>().Any(); await Harness.Bus.Publish(new SomeMessage("something")); await waitTask.Should().BeTrue();
该方式本地仍正常运行,但不确定是否能减少流水线中的随机性。
- 设置10秒无活动超时:
serviceCollection.AddMassTransitTestHarness(cfg => { cfg.SetTestTimeouts(testInactivityTimeout: TimeSpan.FromSeconds(10)); cfg.SetKebabCaseEndpointNameFormatter(); ConfigureTestHarness(cfg); });
但测试仍以相同方式随机失败,且未等待完整的超时时间。
补充日志信息(带时间戳)
2023-06-06T09:30:35.5251966+00:00 - Information - 0 - MassTransit - Configured endpoint my-message-endpoint, Consumer: MyConsumer 2023-06-06T09:30:35.5270254+00:00 - Debug - 0 - MassTransit.Transports.BusDepot - Starting bus instances: IBus 2023-06-06T09:30:35.5270426+00:00 - Debug - 0 - MassTransit - Starting bus: loopback://localhost/ 2023-06-06T09:30:39.5021627+00:00 - Debug - 0 - MassTransit - Endpoint Ready: loopback://localhost/my-message-endpoint 2023-06-06T09:30:39.5049111+00:00 - Debug - 0 - MassTransit - Endpoint Ready: loopback://localhost/runnerwspzdzwlproject21928941concurrent0_testhost_bus_9hboyyfcnrbrfprdbdpschfqfb 2023-06-06T09:30:39.5050096+00:00 - Information - 0 - MassTransit - Bus started: loopback://localhost/ 2023-06-06T09:30:39.5363258+00:00 - Debug - 0 - MassTransit - Create send transport: loopback://localhost/urn:message:MyMessage 2023-06-06T09:30:39.5365853+00:00 - Debug - 0 - MassTransit.Messages - SEND loopback://localhost/urn:message:Events:MyMessage ff030000-ac11-0242-a321-08db6670b105 Events.MyMessage 2023-06-06T09:30:41.6678433+00:00 - Information - 0 - MyConsumer - Updating for Product 42 <-- 消费者代码日志 2023-06-06T09:30:41.7407242+00:00 - Information - 0 - MyConsumer - 0 modified for Product 42 <-- 消费者代码日志 2023-06-06T09:30:41.7719252+00:00 - Debug - 0 - MassTransit.Messages - RECEIVE loopback://localhost/my-message-endpoint ff030000-ac11-0242-a321-08db6670b105 MyMessage MyConsumer(00:00:00.0734534) 2023-06-06T09:30:42.0018748+00:00 - Debug - 0 - MassTransit.Transports.BusDepot - Stopping bus instances: IBus 2023-06-06T09:30:42.0021543+00:00 - Debug - 0 - MassTransit - Stopping bus: loopback://localhost/ 2023-06-06T09:30:42.0023811+00:00 - Debug - 0 - MassTransit - Endpoint Stopping: loopback://localhost/my-message-endpoint 2023-06-06T09:30:42.0026290+00:00 - Debug - 0 - MassTransit - Endpoint Completed: loopback://localhost/my-message-endpoint 2023-06-06T09:30:42.0026475+00:00 - Debug - 0 - MassTransit.Messages - Consumer Completed: loopback://localhost/my-message-endpoint: 1 received, 1 concurrent 2023-06-06T09:30:42.0026759+00:00 - Debug - 0 - MassTransit - Endpoint Stopping: loopback://localhost/runnerwspzdzwlproject21928941concurrent0_testhost_bus_9hboyyfcnrbrfprdbdpschfqfb 2023-06-06T09:30:42.0027023+00:00 - Debug - 0 - MassTransit - Endpoint Completed: loopback://localhost/runnerwspzdzwlproject21928941concurrent0_testhost_bus_9hboyyfcnrbrfprdbdpschfqfb 2023-06-06T09:30:42.0027102+00:00 - Debug - 0 - MassTransit.Messages - Consumer Completed: loopback://localhost/runnerwspzdzwlproject21928941concurrent0_testhost_bus_9hboyyfcnrbrfprdbdpschfqfb: 0 received, 0 concurrent 2023-06-06T09:30:42.0027993+00:00 - Information - 0 - MassTransit - Bus stopped: loopback://localhost/
错误详情
错误发生在以下代码行:
await Harness.Bus.Publish(message); (await Harness.Consumed.SelectAsync<TMessage>().Any()).Should().BeTrue(); // <- 此处报错 (await ConsumerHarness.Consumed.SelectAsync<TMessage>().Any()).Should().BeTrue();
错误信息:Expected (Harness.Consumed.SelectAsync().Any()) to be true, but found False.
可行解决思路
使用Test Harness提供的
WaitUntilReceived方法主动等待消息被消费,替代直接调用Any(),该方法会阻塞直到目标消息被接收或超时:await Harness.Consumed.WaitUntilReceived<SomeMessage>(TimeSpan.FromSeconds(10)); await ConsumerHarness.Consumed.WaitUntilReceived<SomeMessage>(TimeSpan.FromSeconds(10));优先断言特定消费者的消费记录,
ConsumerHarness.Consumed是针对当前测试消费者的精准统计,比全局的Harness.Consumed更可靠,可以先确保这一步断言通过,再验证全局记录:await Harness.Bus.Publish(new SomeMessage("something")); var consumerConsumed = await ConsumerHarness.Consumed.SelectAsync<SomeMessage>().Any(); consumerConsumed.Should().BeTrue(); // 若需要验证全局消费,可在此处增加短暂延迟或等待全局记录更新 await Task.Delay(100); (await Harness.Consumed.SelectAsync<SomeMessage>().Any()).Should().BeTrue();检查GitLab流水线的资源限制,若流水线环境CPU/内存不足,会导致消息处理和Test Harness的统计更新延迟,可适当增加测试的超时时间,同时确保xUnit等测试框架的全局超时设置足够长。
避免在测试中同时依赖
Harness.Consumed和ConsumerHarness.Consumed,如果只需要验证消费者是否处理了消息,仅使用ConsumerHarness的断言即可,减少不必要的全局状态依赖。发布消息后,先等待总线确认消息发送完成,再等待消费者处理:
var publishTask = Harness.Bus.Publish(new SomeMessage("something")); var consumeTask = ConsumerHarness.Consumed.SelectAsync<SomeMessage>().FirstOrDefault(); await publishTask.ConfigureAwait(false); var consumedMessage = await consumeTask.ConfigureAwait(false); consumedMessage.Should().NotBeNull();
内容的提问来源于stack exchange,提问作者Krzysztof Skowronek

