| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291 |
- using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Entity;
- using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Idempotency;
- using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Inventory;
- using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.L1;
- using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Receipt;
- using Xunit;
- namespace Admin.NET.Plugin.AiDOP.Tests.S5.MaterialWarehouse;
- /// <summary>
- /// 可配置 L1 收货态 loader(测试用)。<see cref="State"/> = 首次(routing)读取;
- /// <see cref="ReloadState"/> 若设置则第 2 次起(orchestrator 事务内 reload)返回它——用于 A5 新鲜度对照。
- /// </summary>
- public sealed class FakeIqcReceiptStateLoader : IIqcReceiptStateLoader
- {
- public IqcReceiptState State { get; set; }
- public IqcReceiptState ReloadState { get; set; }
- private int _calls;
- public Task<IqcReceiptState> LoadAsync(long tenantId, string domainCode, string fbillno)
- {
- var i = System.Threading.Interlocked.Increment(ref _calls);
- return Task.FromResult(i == 1 ? State : (ReloadState ?? State));
- }
- }
- /// <summary>完成锚只读实现(测试用,over FakeInventoryDatabase.Completions)。</summary>
- public sealed class InMemoryBusinessCompletionStore : IBusinessCompletionStore
- {
- private readonly FakeInventoryDatabase _db;
- public InMemoryBusinessCompletionStore(FakeInventoryDatabase db) => _db = db;
- public Task<bool> IsCommittedAsync(long tenantId, string domainCode, string fbillNo)
- => Task.FromResult(_db.Completions.ContainsKey($"{tenantId}|{domainCode ?? ""}|{fbillNo}"));
- }
- /// <summary>
- /// Phase 7-2 + 7C:L1 IQC 事件消费——四态、端到端幂等(重复投递)、retry/ack、路由,
- /// 以及 7C 可靠性:stale Processing 恢复、commit→MarkPosted 崩溃自愈、Runtime Gate 事务内新鲜度。
- /// </summary>
- public class IqcInventoryEventConsumerTests
- {
- private const long Tenant = 797403760988229L;
- private static IqcInventoryEvent Event(int pd = 0, string clfs = "", decimal dhsl = 100m, decimal bhgsl = 0m, string fbillno = "FB-1")
- => new IqcInventoryEvent
- {
- TenantId = Tenant, DomainCode = "8010", Fbillno = fbillno, Receiver = "RCV-1", RctQcNbr = "QC-1",
- Pd = pd, Clfs = clfs, Dhsl = dhsl, Bhgsl = bhgsl, User = "t", EventTime = new DateTime(2026, 8, 9),
- };
- private static IqcReceiptState Receivable()
- => new IqcReceiptState
- {
- Found = true, HasPendingReceipt = true, HasReceivableDetails = true,
- ItemNum = "IT-1", Location = "L1", LotSerial = "", Refs = "", RctNbr = "RCT-1",
- PurOrd = "PO-1", PurLine = 1, Potype = "po", OrderedQty = 100m, BeforeReceivedQty = 0m, ReturnedQty = 0m, SampleQty = 0m,
- SourceSyncedAt = new DateTime(2026, 8, 9), // 与事件同刻→非 stale;满足 fail-closed 水位非空
- };
- private sealed class Rig
- {
- public FakeInventoryDatabase Db;
- public InMemoryIqcPostingStore Store;
- public FakeIqcReceiptStateLoader Loader;
- public IqcInventoryEventConsumer Consumer;
- }
- private static Rig Build(IqcReceiptState state, FakeReceiptFailPhase fail = FakeReceiptFailPhase.None,
- IqcReceiptState reload = null, TimeSpan? lease = null)
- {
- var db = new FakeInventoryDatabase();
- db.SeedMaster(Tenant, "8010", "IT-1", "L1");
- var store = new InMemoryIqcPostingStore();
- var idem = new IqcPostingIdempotencyService(store);
- var factory = new FakeIqcReceiptUnitOfWorkFactory(db, fail);
- var loader = new FakeIqcReceiptStateLoader { State = state, ReloadState = reload };
- var consumer = new IqcInventoryEventConsumer(
- idem, new InMemoryBusinessCompletionStore(db), loader,
- new IqcReceiptOrchestrator(factory), new IqcMrbSelectionTaskService(factory), lease);
- return new Rig { Db = db, Store = store, Loader = loader, Consumer = consumer };
- }
- // ═══ Phase 7-2 基线 ═══
- [Fact]
- public async Task FirstProcess_FullReceipt_Processed_AndPosted()
- {
- var rig = Build(Receivable());
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
- Assert.Equal(1, rig.Db.InvMaster.Count);
- Assert.Equal(1, rig.Db.MobileTasks.Count);
- Assert.Equal(1, rig.Db.Completions.Count); // 完成锚已随业务事务落地
- Assert.Equal("POSTED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus);
- }
- [Fact]
- public async Task Duplicate_SecondConsume_AlreadyProcessed_NoDoublePost()
- {
- var rig = Build(Receivable());
- await rig.Consumer.ConsumeAsync(Event());
- var r2 = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.AlreadyProcessed, r2.Outcome);
- Assert.Equal(1, rig.Db.InvMaster.Count);
- Assert.Equal(100m, rig.Db.InvMaster[FakeInventoryDatabase.InvK(Tenant, "8010", "IT-1", "L1")].QtyOnHand);
- }
- [Fact]
- public async Task FailedRetryable_RetriesAndSucceeds()
- {
- var rig = Build(Receivable());
- await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting { TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "FAILED" });
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
- Assert.Equal(1, rig.Db.InvMaster.Count);
- }
- [Fact]
- public async Task Duplicate1000_ExactlyOneProcessed_SingleSetPersisted()
- {
- var rig = Build(Receivable());
- var outcomes = new List<IqcConsumeOutcome>();
- for (var i = 0; i < 1000; i++)
- outcomes.Add((await rig.Consumer.ConsumeAsync(Event())).Outcome);
- Assert.Equal(1, outcomes.Count(o => o == IqcConsumeOutcome.Processed));
- Assert.Equal(999, outcomes.Count(o => o == IqcConsumeOutcome.AlreadyProcessed));
- Assert.Equal(1, rig.Db.InvMaster.Count);
- Assert.Equal(1, rig.Db.MobileTasks.Count);
- Assert.Equal(1, rig.Db.ReceiptWriteback.Count);
- Assert.Equal(1, rig.Db.Completions.Count);
- }
- [Fact]
- public async Task BusinessFail_EventRetryable_NoPersist()
- {
- var rig = Build(Receivable(), FakeReceiptFailPhase.C5);
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
- Assert.Equal(0, rig.Db.TotalPersisted);
- Assert.Equal(0, rig.Db.Completions.Count); // 失败→完成锚未落地
- Assert.Equal("FAILED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus);
- }
- [Fact]
- public async Task Freshness_LoaderNotReceivable_Retryable_NotProcessed()
- {
- var state = Receivable();
- state.HasReceivableDetails = false; // 预门禁即拦截
- var rig = Build(state);
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
- Assert.Equal(0, rig.Db.TotalPersisted);
- }
- [Fact]
- public async Task Route_MrbChooseEvent_CreatesTask_NoInventory()
- {
- var rig = Build(Receivable());
- var r = await rig.Consumer.ConsumeAsync(Event(pd: 0, clfs: "3", bhgsl: 3m));
- Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
- Assert.Equal(1, rig.Db.MobileTasks.Count);
- Assert.Equal("PurOrdRctByMRBChoose", rig.Db.MobileTasks[$"{Tenant}|8010|QC-1"].Name);
- Assert.Equal(0, rig.Db.InvMaster.Count);
- Assert.Equal(1, rig.Db.Completions.Count); // MRB 亦落完成锚
- }
- [Fact]
- public async Task Route_ReturnOnlyEvent_NoAction()
- {
- var rig = Build(Receivable());
- var r = await rig.Consumer.ConsumeAsync(Event(pd: 1, clfs: "0", bhgsl: 3m));
- Assert.Equal(IqcConsumeOutcome.NoAction, r.Outcome);
- Assert.Equal(0, rig.Db.TotalPersisted);
- Assert.Equal("POSTED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus);
- }
- [Fact]
- public async Task StateNotFound_Retryable()
- {
- var rig = Build(new IqcReceiptState { Found = false });
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
- }
- // ═══ Phase 7C 可靠性 ═══
- // R1:PROCESSING 未过期 → skip、不进 C2
- [Fact]
- public async Task R1_Processing_LeaseNotExpired_Skipped()
- {
- var rig = Build(Receivable());
- await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
- {
- TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
- CreateTime = DateTime.Now, UpdateTime = DateTime.Now, // 心跳新鲜
- });
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Skipped, r.Outcome);
- Assert.Equal(0, rig.Db.TotalPersisted);
- }
- // R2:PROCESSING 已过期 + 业务未执行(无完成锚)→ 抢占恢复 → 重新执行成功
- [Fact]
- public async Task R2_Processing_Stale_NoCompletion_RecoversAndRuns()
- {
- var rig = Build(Receivable());
- await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
- {
- TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
- CreateTime = DateTime.Now.AddMinutes(-10), UpdateTime = DateTime.Now.AddMinutes(-10), // 过期
- });
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
- Assert.Equal(1, rig.Db.InvMaster.Count);
- Assert.Equal(1, rig.Db.Completions.Count);
- }
- // R3:业务 Commit 成功(完成锚已落)但 MarkPosted 前崩溃 → 重投检出完成锚 → 修复 POSTED、不重复业务
- [Fact]
- public async Task R3_CommitThenCrashBeforeMarkPosted_RecoversFromCompletion_NoDoubleInventory()
- {
- var rig = Build(Receivable());
- // 模拟“上一轮业务已提交”:库存已到位 + 完成锚已落,但 posting 仍卡 PROCESSING
- rig.Db.SeedStock(Tenant, "8010", "IT-1", "L1", 100m);
- rig.Db.Completions[$"{Tenant}|8010|FB-1"] = new IqcCompletionMarker { TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingId = 1 };
- await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
- {
- TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
- CreateTime = DateTime.Now.AddMinutes(-10), UpdateTime = DateTime.Now.AddMinutes(-10),
- });
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.AlreadyProcessed, r.Outcome); // 崩溃恢复
- Assert.Empty(rig.Db.CallOrder); // 未再跑 C6/C5/C7
- Assert.Equal(100m, rig.Db.InvMaster[FakeInventoryDatabase.InvK(Tenant, "8010", "IT-1", "L1")].QtyOnHand); // 库存未二次变化
- Assert.Equal("POSTED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus); // 已修复
- }
- // R4:两个消费者同时抢 stale PROCESSING → 只有一个获得恢复 claim(一 Processed / 一 Skipped)
- [Fact]
- public async Task R4_TwoConsumersRaceStale_OnlyOneRecovers()
- {
- var rig = Build(Receivable());
- await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
- {
- TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
- CreateTime = DateTime.Now.AddMinutes(-10), UpdateTime = DateTime.Now.AddMinutes(-10),
- });
- var results = await Task.WhenAll(
- Task.Run(() => rig.Consumer.ConsumeAsync(Event())),
- Task.Run(() => rig.Consumer.ConsumeAsync(Event())));
- // 不变式:恰好一路真正执行(Processed);另一路是幂等 no-op——
- // Skipped(在赢家提交完成锚前判定) 或 AlreadyProcessed(在其后判定),二者都对,不重复过账。
- Assert.Equal(1, results.Count(x => x.Outcome == IqcConsumeOutcome.Processed));
- Assert.All(results.Where(x => x.Outcome != IqcConsumeOutcome.Processed),
- x => Assert.Contains(x.Outcome, new[] { IqcConsumeOutcome.Skipped, IqcConsumeOutcome.AlreadyProcessed }));
- Assert.Equal(1, rig.Db.InvMaster.Count); // 仅一套库存
- Assert.Equal(1, rig.Db.Completions.Count);
- }
- // R5:事务内 reload 报“当前不可收” → 0 C6/C5/C7/库存写 → Retryable(用事务内新鲜态否决)
- [Fact]
- public async Task R5_InTxReloadNotReceivable_NoWrites_Retryable()
- {
- var reload = Receivable();
- reload.HasReceivableDetails = false; // 事务内最新:不可收
- var rig = Build(Receivable(), reload: reload); // 事务外 snapshot 仍“可收”(预门禁放行)
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
- Assert.Equal(0, rig.Db.TotalPersisted);
- Assert.Empty(rig.Db.CallOrder); // reload 否决在 C6 之前
- Assert.Equal(0, rig.Db.Completions.Count);
- }
- // R6:C6 对账数量用事务内最新 BeforeReceived(非事件/事务外 snapshot)
- [Fact]
- public async Task R6_C6UsesInTxFreshQuantities()
- {
- var reload = Receivable();
- reload.BeforeReceivedQty = 20m; // 事务内最新已收 20(事务外 snapshot=0)
- var rig = Build(Receivable(), reload: reload);
- var r = await rig.Consumer.ConsumeAsync(Event());
- Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
- var c6 = rig.Db.ReceiptWriteback.Values.Single();
- Assert.Equal(20m, c6.BeforeReceivedQty); // 事务内新鲜值
- Assert.Equal(120m, c6.AfterReceivedQty); // 20 + 合格100
- }
- }
|