You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 06:54:23