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; /// /// 可配置 L1 收货态 loader(测试用)。 = 首次(routing)读取; /// 若设置则第 2 次起(orchestrator 事务内 reload)返回它——用于 A5 新鲜度对照。 /// public sealed class FakeIqcReceiptStateLoader : IIqcReceiptStateLoader { public IqcReceiptState State { get; set; } public IqcReceiptState ReloadState { get; set; } private int _calls; public Task LoadAsync(long tenantId, string domainCode, string fbillno) { var i = System.Threading.Interlocked.Increment(ref _calls); return Task.FromResult(i == 1 ? State : (ReloadState ?? State)); } } /// 完成锚只读实现(测试用,over FakeInventoryDatabase.Completions)。 public sealed class InMemoryBusinessCompletionStore : IBusinessCompletionStore { private readonly FakeInventoryDatabase _db; public InMemoryBusinessCompletionStore(FakeInventoryDatabase db) => _db = db; public Task IsCommittedAsync(long tenantId, string domainCode, string fbillNo) => Task.FromResult(_db.Completions.ContainsKey($"{tenantId}|{domainCode ?? ""}|{fbillNo}")); } /// /// Phase 7-2 + 7C:L1 IQC 事件消费——四态、端到端幂等(重复投递)、retry/ack、路由, /// 以及 7C 可靠性:stale Processing 恢复、commit→MarkPosted 崩溃自愈、Runtime Gate 事务内新鲜度。 /// 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(); 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 } }