MassTransit测试Saga状态机报ProcessInvoiceReceived未配置错误
错误原因
ProcessInvoiceReceived was not specified配置错误的核心原因是状态机初始事件未配置Saga实例关联规则。
MassTransit Automatonymous状态机强制要求所有声明的事件必须配置与Saga实例的匹配逻辑:
- 启动新Saga实例的初始事件,除匹配规则外,还需要明确新实例
CorrelationId的生成规则 - 关联已有实例的后续事件,需要指定消息中用于匹配Saga实例的字段
当前代码仅通过Event(() => ProcessInvoiceReceived)声明了事件,未添加任何关联配置,状态机构建阶段的校验逻辑就会抛出该错误。
除此之外代码还有两处会导致测试运行异常的问题:
- 重复注册Saga仓储:手动注入了一次InMemorySagaRepository,配置Saga时又调用
InMemoryRepository()重复注册,会引发依赖注入冲突 - 后续事件关联ID传递错误:Mock消费者回传事件时用
context.MessageId作为CorrelationId,该值是消息自身的唯一标识,和Saga实例ID无关,会导致后续事件无法匹配到已创建的Saga实例,状态流转中断。
修复方案
- 修正Saga事件配置
为所有事件显式配置关联规则,初始事件补充实例ID生成逻辑,修改状态机构造函数中的事件声明部分:
public ProcessEntryInvoiceSaga() { InstanceState(x => x.CurrentState); // 初始事件配置:按InvoiceId关联,新实例自动生成Guid作为CorrelationId Event(() => ProcessInvoiceReceived, e => e .CorrelateBy((instance, ctx) => instance.InvoiceId == ctx.Message.InvoiceId) .SetIdFrom(_ => NewId.NextGuid()) ); // 后续事件配置:按消息携带的CorrelationId匹配已有Saga实例 Event(() => ProductsAddedEvent, e => e.CorrelateById(ctx => ctx.Message.CorrelationId)); Event(() => DuplicatesRegisteredEvent, e => e.CorrelateById(ctx => ctx.Message.CorrelationId)); Initially( When(ProcessInvoiceReceived) .Then(x => x.Saga.InvoiceId = x.Message.InvoiceId) .Activity(x=> x.OfType<ProcessEntryInvoiceStartActivity>()) .TransitionTo(Started) ); During(Started, When(ProductsAddedEvent) .TransitionTo(Processing), When(DuplicatesRegisteredEvent) .TransitionTo(Processing)); During(Processing, When(ProductsAddedEvent) .TransitionTo(Closed), When(DuplicatesRegisteredEvent) .TransitionTo(Closed)); WhenEnter(Closed, binder => binder .Activity(x=>x.OfType<ProcessEntryInvoiceFinishActivity>()) ); SetCompleted(async instance => { State currentState = await this.GetState(instance); return Closed.Equals(currentState); }); }
原有状态流转逻辑无需改动,但要确保初始事件对应的ProcessEntryInvoiceStartActivity中发布下游消息时,将Saga实例的CorrelationId赋值给下游消息的CorrelationId字段,供后续事件关联使用。
- 清理测试项目的重复依赖注册
删除测试Setup方法中手动注册Saga仓储的代码,保留AddSagaStateMachine后的.InMemoryRepository()调用即可,修正后的服务注册部分:
var provider = new ServiceCollection() .AddSingleton(_entryInvoiceRepository.Object) // 删除手动注册Saga仓储的冗余代码 .AddMassTransitTestHarness(cfg => { cfg.AddSagaStateMachine<ProcessEntryInvoiceSaga, ProcessEntryInvoiceSagaState>().InMemoryRepository(); cfg.AddConsumer<MockSagaStepConsumers>(); }) .BuildServiceProvider(true);
- 修正Mock消费者的CorrelationId传递逻辑
下游消息发布时必须传递Saga实例的CorrelationId,不能使用消息本身的MessageId。如果ProcessEntryInvoiceStartActivity中发布AddProductsToInventory、RegisterDuplicatesOfInvoice消息时已经正确携带了Saga的CorrelationId,Mock消费者直接透传该值即可:
public class MockSagaStepConsumers : IConsumer<AddProductsToInventory>, IConsumer<RegisterDuplicatesOfInvoice> { public async Task Consume(ConsumeContext<AddProductsToInventory> context) { await context.Publish(new ProductsAdded() { // 直接透传上游消息携带的Saga CorrelationId,不要使用context.MessageId CorrelationId = context.Message.CorrelationId }); } public async Task Consume(ConsumeContext<RegisterDuplicatesOfInvoice> context) { await context.Publish(new DuplicatesRegistered() { // 直接透传上游消息携带的Saga CorrelationId,不要使用context.MessageId CorrelationId = context.Message.CorrelationId, InvoiceId = context.Message.InvoiceId }); } }
验证
完成以上修改后,状态机配置校验可正常通过,测试执行时Saga实例可以正常创建、响应后续事件并流转到Closed状态。
内容的提问来源于stack exchange,提问作者Lelis718
相关产品推荐
相关产品推荐

