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

MassTransit SagaStateMachine单元测试未通过问题求助

MassTransit SagaStateMachine单元测试失败排查

我参照MassTransit容器测试示例编写SagaStateMachine单元测试,但测试始终失败。Test Harness可以正常消费StripeSubscriptionCreated事件,但Saga Harness的断言始终无法通过。

我的SagaStateMachine实现

public class SubscriptionCreatedStateMachine :
    MassTransitStateMachine<SubscriptionCreatedState>
{
    public SubscriptionCreatedStateMachine()
    {
        InstanceState(x => x.CurrentState);

        Event(() => StripeSubscriptionCreated, x =>
            x.CorrelateById(m => m.Message.EventId));

        Event(() => FinbuckleTenantCreated, x => 
            x.CorrelateById(m => m.Message.CorrelationId));

        Initially(
            When(StripeSubscriptionCreated)
                .Then(context =>
                {
                    context.Saga.TenantId = context.Message.TenantId;
                    context.Saga.Identifier = context.Message.TenantId; // 新租户的标识符始终是租户ID(GUID)
                    context.Saga.SubscriptionId = context.Message.SubscriptionId;
                    context.Saga.CustomerId = context.Message.CustomerId;
                    context.Saga.CustomerEmail = context.Message.CustomerEmail;
                    context.Saga.Description = context.Message.Description;
                    context.Saga.CurrentPeriodStart = context.Message.CurrentPeriodStart;
                    context.Saga.CurrentPeriodEnd = context.Message.CurrentPeriodEnd;
                    context.Saga.Created  = context.Message.Created;
                })
                .TransitionTo(SubscriptionCreated)
                .Publish(context => new CreateTenant()
                {
                    CorrelationId = context.Saga.CorrelationId,
                    TenantId = context.Saga.TenantId,
                    Identifier = context.Saga.Identifier, 
                    SubscriptionId = context.Saga.SubscriptionId,
                    CustomerId = context.Saga.CustomerId,
                    CustomerEmail = context.Saga.CustomerEmail,
                    Description = context.Saga.Description,
                    CurrentPeriodStart = context.Saga.CurrentPeriodStart,
                    CurrentPeriodEnd = context.Saga.CurrentPeriodEnd,
                    Created = context.Saga.Created
                })
        );

        During(SubscriptionCreated,
                    When(FinbuckleTenantCreated)
                        .TransitionTo(TenantCreated)
                        .Publish(context => new SendTenantCreatedEmail()
                        {
                            CorrelationId = context.Saga.CorrelationId,
                            Identifier = context.Saga.Identifier,
                            CustomerEmail = context.Saga.CustomerEmail,
                            CurrentPeriodStart = context.Saga.CurrentPeriodStart,
                            CurrentPeriodEnd = context.Saga.CurrentPeriodEnd
                        })
        );
    }

    public State SubscriptionCreated { get; } = null!;
    public State TenantCreated { get; } = null!;
    public Event<StripeSubscriptionCreated> StripeSubscriptionCreated { get; } = null!;
    public Event<TenantCreated> FinbuckleTenantCreated { get;  }= null!;
}

单元测试代码

[Fact]
public async Task Can_consume_stripe_subscription_created_event()
{
    await using var provider = new ServiceCollection()
        .AddMassTransitTestHarness(cfg =>
        {
            cfg.AddSagaStateMachine<SubscriptionCreatedStateMachine, SubscriptionCreatedState>()
                .InMemoryRepository();
        })
        .BuildServiceProvider(true);

    var harness = provider.GetTestHarness();
    await harness.Start();

    var subscriptionEvent = new StripeSubscriptionCreated()
    {
        SubscriptionId = "abc",
        CustomerId = "abc",
        CustomerEmail = "customer@email.com",
        Description = "description",
        CurrentPeriodStart = DateTime.UtcNow,
        CurrentPeriodEnd = DateTime.UtcNow.AddMonths(1),
        Created = DateTime.UtcNow,
        TenantId = Guid.NewGuid().ToString()
    };

    await harness.Bus.Publish(subscriptionEvent);

    (await harness.Consumed.Any<StripeSubscriptionCreated>()).ShouldBeTrue();

    var sagaHarness = harness.GetSagaStateMachineHarness<SubscriptionCreatedStateMachine, SubscriptionCreatedState>();  
    (await sagaHarness.Consumed.Any<StripeSubscriptionCreated>()).ShouldBeTrue();
    (await sagaHarness.Created.Any(x => x.CorrelationId == subscriptionEvent.EventId)).ShouldBeTrue();
}

问题分析与解决方案

测试失败的核心原因如下:

  • 未设置StripeSubscriptionCreated事件的EventId
    Saga配置中,StripeSubscriptionCreated事件通过x.CorrelateById(m => m.Message.EventId)关联Saga实例,但测试中创建的subscriptionEvent未赋值EventId,导致Saga无法创建或匹配对应的实例,断言自然失败。

  • Saga实例关联逻辑依赖事件EventId
    首次处理事件时,Saga会将事件的EventId作为自身实例的CorrelationId,缺失该值会直接导致Saga实例无法正常创建。

修复后的测试代码

修改事件创建逻辑,添加EventId赋值,并优化断言等待逻辑:

[Fact]
public async Task Can_consume_stripe_subscription_created_event()
{
    await using var provider = new ServiceCollection()
        .AddMassTransitTestHarness(cfg =>
        {
            cfg.AddSagaStateMachine<SubscriptionCreatedStateMachine, SubscriptionCreatedState>()
                .InMemoryRepository();
        })
        .BuildServiceProvider(true);

    var harness = provider.GetTestHarness();
    await harness.Start();

    var eventId = Guid.NewGuid();
    var subscriptionEvent = new StripeSubscriptionCreated()
    {
        EventId = eventId, // 新增:设置EventId
        SubscriptionId = "abc",
        CustomerId = "abc",
        CustomerEmail = "customer@email.com",
        Description = "description",
        CurrentPeriodStart = DateTime.UtcNow,
        CurrentPeriodEnd = DateTime.UtcNow.AddMonths(1),
        Created = DateTime.UtcNow,
        TenantId = Guid.NewGuid().ToString()
    };

    await harness.Bus.Publish(subscriptionEvent);

    (await harness.Consumed.Any<StripeSubscriptionCreated>()).ShouldBeTrue();

    var sagaHarness = harness.GetSagaStateMachineHarness<SubscriptionCreatedStateMachine, SubscriptionCreatedState>();  
    // 等待Saga完成事件处理
    await sagaHarness.Consumed.Any<StripeSubscriptionCreated>(x => x.Context.Message.EventId == eventId);
    // 断言Saga实例创建成功
    (await sagaHarness.Created.Any(x => x.CorrelationId == eventId)).ShouldBeTrue();
}

内容的提问来源于stack exchange,提问作者user351479

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 01:59:51