IqcInventoryEventConsumerTests.cs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291
  1. using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Entity;
  2. using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Idempotency;
  3. using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Inventory;
  4. using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.L1;
  5. using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Receipt;
  6. using Xunit;
  7. namespace Admin.NET.Plugin.AiDOP.Tests.S5.MaterialWarehouse;
  8. /// <summary>
  9. /// 可配置 L1 收货态 loader(测试用)。<see cref="State"/> = 首次(routing)读取;
  10. /// <see cref="ReloadState"/> 若设置则第 2 次起(orchestrator 事务内 reload)返回它——用于 A5 新鲜度对照。
  11. /// </summary>
  12. public sealed class FakeIqcReceiptStateLoader : IIqcReceiptStateLoader
  13. {
  14. public IqcReceiptState State { get; set; }
  15. public IqcReceiptState ReloadState { get; set; }
  16. private int _calls;
  17. public Task<IqcReceiptState> LoadAsync(long tenantId, string domainCode, string fbillno)
  18. {
  19. var i = System.Threading.Interlocked.Increment(ref _calls);
  20. return Task.FromResult(i == 1 ? State : (ReloadState ?? State));
  21. }
  22. }
  23. /// <summary>完成锚只读实现(测试用,over FakeInventoryDatabase.Completions)。</summary>
  24. public sealed class InMemoryBusinessCompletionStore : IBusinessCompletionStore
  25. {
  26. private readonly FakeInventoryDatabase _db;
  27. public InMemoryBusinessCompletionStore(FakeInventoryDatabase db) => _db = db;
  28. public Task<bool> IsCommittedAsync(long tenantId, string domainCode, string fbillNo)
  29. => Task.FromResult(_db.Completions.ContainsKey($"{tenantId}|{domainCode ?? ""}|{fbillNo}"));
  30. }
  31. /// <summary>
  32. /// Phase 7-2 + 7C:L1 IQC 事件消费——四态、端到端幂等(重复投递)、retry/ack、路由,
  33. /// 以及 7C 可靠性:stale Processing 恢复、commit→MarkPosted 崩溃自愈、Runtime Gate 事务内新鲜度。
  34. /// </summary>
  35. public class IqcInventoryEventConsumerTests
  36. {
  37. private const long Tenant = 797403760988229L;
  38. private static IqcInventoryEvent Event(int pd = 0, string clfs = "", decimal dhsl = 100m, decimal bhgsl = 0m, string fbillno = "FB-1")
  39. => new IqcInventoryEvent
  40. {
  41. TenantId = Tenant, DomainCode = "8010", Fbillno = fbillno, Receiver = "RCV-1", RctQcNbr = "QC-1",
  42. Pd = pd, Clfs = clfs, Dhsl = dhsl, Bhgsl = bhgsl, User = "t", EventTime = new DateTime(2026, 8, 9),
  43. };
  44. private static IqcReceiptState Receivable()
  45. => new IqcReceiptState
  46. {
  47. Found = true, HasPendingReceipt = true, HasReceivableDetails = true,
  48. ItemNum = "IT-1", Location = "L1", LotSerial = "", Refs = "", RctNbr = "RCT-1",
  49. PurOrd = "PO-1", PurLine = 1, Potype = "po", OrderedQty = 100m, BeforeReceivedQty = 0m, ReturnedQty = 0m, SampleQty = 0m,
  50. SourceSyncedAt = new DateTime(2026, 8, 9), // 与事件同刻→非 stale;满足 fail-closed 水位非空
  51. };
  52. private sealed class Rig
  53. {
  54. public FakeInventoryDatabase Db;
  55. public InMemoryIqcPostingStore Store;
  56. public FakeIqcReceiptStateLoader Loader;
  57. public IqcInventoryEventConsumer Consumer;
  58. }
  59. private static Rig Build(IqcReceiptState state, FakeReceiptFailPhase fail = FakeReceiptFailPhase.None,
  60. IqcReceiptState reload = null, TimeSpan? lease = null)
  61. {
  62. var db = new FakeInventoryDatabase();
  63. db.SeedMaster(Tenant, "8010", "IT-1", "L1");
  64. var store = new InMemoryIqcPostingStore();
  65. var idem = new IqcPostingIdempotencyService(store);
  66. var factory = new FakeIqcReceiptUnitOfWorkFactory(db, fail);
  67. var loader = new FakeIqcReceiptStateLoader { State = state, ReloadState = reload };
  68. var consumer = new IqcInventoryEventConsumer(
  69. idem, new InMemoryBusinessCompletionStore(db), loader,
  70. new IqcReceiptOrchestrator(factory), new IqcMrbSelectionTaskService(factory), lease);
  71. return new Rig { Db = db, Store = store, Loader = loader, Consumer = consumer };
  72. }
  73. // ═══ Phase 7-2 基线 ═══
  74. [Fact]
  75. public async Task FirstProcess_FullReceipt_Processed_AndPosted()
  76. {
  77. var rig = Build(Receivable());
  78. var r = await rig.Consumer.ConsumeAsync(Event());
  79. Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
  80. Assert.Equal(1, rig.Db.InvMaster.Count);
  81. Assert.Equal(1, rig.Db.MobileTasks.Count);
  82. Assert.Equal(1, rig.Db.Completions.Count); // 完成锚已随业务事务落地
  83. Assert.Equal("POSTED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus);
  84. }
  85. [Fact]
  86. public async Task Duplicate_SecondConsume_AlreadyProcessed_NoDoublePost()
  87. {
  88. var rig = Build(Receivable());
  89. await rig.Consumer.ConsumeAsync(Event());
  90. var r2 = await rig.Consumer.ConsumeAsync(Event());
  91. Assert.Equal(IqcConsumeOutcome.AlreadyProcessed, r2.Outcome);
  92. Assert.Equal(1, rig.Db.InvMaster.Count);
  93. Assert.Equal(100m, rig.Db.InvMaster[FakeInventoryDatabase.InvK(Tenant, "8010", "IT-1", "L1")].QtyOnHand);
  94. }
  95. [Fact]
  96. public async Task FailedRetryable_RetriesAndSucceeds()
  97. {
  98. var rig = Build(Receivable());
  99. await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting { TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "FAILED" });
  100. var r = await rig.Consumer.ConsumeAsync(Event());
  101. Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
  102. Assert.Equal(1, rig.Db.InvMaster.Count);
  103. }
  104. [Fact]
  105. public async Task Duplicate1000_ExactlyOneProcessed_SingleSetPersisted()
  106. {
  107. var rig = Build(Receivable());
  108. var outcomes = new List<IqcConsumeOutcome>();
  109. for (var i = 0; i < 1000; i++)
  110. outcomes.Add((await rig.Consumer.ConsumeAsync(Event())).Outcome);
  111. Assert.Equal(1, outcomes.Count(o => o == IqcConsumeOutcome.Processed));
  112. Assert.Equal(999, outcomes.Count(o => o == IqcConsumeOutcome.AlreadyProcessed));
  113. Assert.Equal(1, rig.Db.InvMaster.Count);
  114. Assert.Equal(1, rig.Db.MobileTasks.Count);
  115. Assert.Equal(1, rig.Db.ReceiptWriteback.Count);
  116. Assert.Equal(1, rig.Db.Completions.Count);
  117. }
  118. [Fact]
  119. public async Task BusinessFail_EventRetryable_NoPersist()
  120. {
  121. var rig = Build(Receivable(), FakeReceiptFailPhase.C5);
  122. var r = await rig.Consumer.ConsumeAsync(Event());
  123. Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
  124. Assert.Equal(0, rig.Db.TotalPersisted);
  125. Assert.Equal(0, rig.Db.Completions.Count); // 失败→完成锚未落地
  126. Assert.Equal("FAILED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus);
  127. }
  128. [Fact]
  129. public async Task Freshness_LoaderNotReceivable_Retryable_NotProcessed()
  130. {
  131. var state = Receivable();
  132. state.HasReceivableDetails = false; // 预门禁即拦截
  133. var rig = Build(state);
  134. var r = await rig.Consumer.ConsumeAsync(Event());
  135. Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
  136. Assert.Equal(0, rig.Db.TotalPersisted);
  137. }
  138. [Fact]
  139. public async Task Route_MrbChooseEvent_CreatesTask_NoInventory()
  140. {
  141. var rig = Build(Receivable());
  142. var r = await rig.Consumer.ConsumeAsync(Event(pd: 0, clfs: "3", bhgsl: 3m));
  143. Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
  144. Assert.Equal(1, rig.Db.MobileTasks.Count);
  145. Assert.Equal("PurOrdRctByMRBChoose", rig.Db.MobileTasks[$"{Tenant}|8010|QC-1"].Name);
  146. Assert.Equal(0, rig.Db.InvMaster.Count);
  147. Assert.Equal(1, rig.Db.Completions.Count); // MRB 亦落完成锚
  148. }
  149. [Fact]
  150. public async Task Route_ReturnOnlyEvent_NoAction()
  151. {
  152. var rig = Build(Receivable());
  153. var r = await rig.Consumer.ConsumeAsync(Event(pd: 1, clfs: "0", bhgsl: 3m));
  154. Assert.Equal(IqcConsumeOutcome.NoAction, r.Outcome);
  155. Assert.Equal(0, rig.Db.TotalPersisted);
  156. Assert.Equal("POSTED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus);
  157. }
  158. [Fact]
  159. public async Task StateNotFound_Retryable()
  160. {
  161. var rig = Build(new IqcReceiptState { Found = false });
  162. var r = await rig.Consumer.ConsumeAsync(Event());
  163. Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
  164. }
  165. // ═══ Phase 7C 可靠性 ═══
  166. // R1:PROCESSING 未过期 → skip、不进 C2
  167. [Fact]
  168. public async Task R1_Processing_LeaseNotExpired_Skipped()
  169. {
  170. var rig = Build(Receivable());
  171. await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
  172. {
  173. TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
  174. CreateTime = DateTime.Now, UpdateTime = DateTime.Now, // 心跳新鲜
  175. });
  176. var r = await rig.Consumer.ConsumeAsync(Event());
  177. Assert.Equal(IqcConsumeOutcome.Skipped, r.Outcome);
  178. Assert.Equal(0, rig.Db.TotalPersisted);
  179. }
  180. // R2:PROCESSING 已过期 + 业务未执行(无完成锚)→ 抢占恢复 → 重新执行成功
  181. [Fact]
  182. public async Task R2_Processing_Stale_NoCompletion_RecoversAndRuns()
  183. {
  184. var rig = Build(Receivable());
  185. await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
  186. {
  187. TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
  188. CreateTime = DateTime.Now.AddMinutes(-10), UpdateTime = DateTime.Now.AddMinutes(-10), // 过期
  189. });
  190. var r = await rig.Consumer.ConsumeAsync(Event());
  191. Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
  192. Assert.Equal(1, rig.Db.InvMaster.Count);
  193. Assert.Equal(1, rig.Db.Completions.Count);
  194. }
  195. // R3:业务 Commit 成功(完成锚已落)但 MarkPosted 前崩溃 → 重投检出完成锚 → 修复 POSTED、不重复业务
  196. [Fact]
  197. public async Task R3_CommitThenCrashBeforeMarkPosted_RecoversFromCompletion_NoDoubleInventory()
  198. {
  199. var rig = Build(Receivable());
  200. // 模拟“上一轮业务已提交”:库存已到位 + 完成锚已落,但 posting 仍卡 PROCESSING
  201. rig.Db.SeedStock(Tenant, "8010", "IT-1", "L1", 100m);
  202. rig.Db.Completions[$"{Tenant}|8010|FB-1"] = new IqcCompletionMarker { TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingId = 1 };
  203. await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
  204. {
  205. TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
  206. CreateTime = DateTime.Now.AddMinutes(-10), UpdateTime = DateTime.Now.AddMinutes(-10),
  207. });
  208. var r = await rig.Consumer.ConsumeAsync(Event());
  209. Assert.Equal(IqcConsumeOutcome.AlreadyProcessed, r.Outcome); // 崩溃恢复
  210. Assert.Empty(rig.Db.CallOrder); // 未再跑 C6/C5/C7
  211. Assert.Equal(100m, rig.Db.InvMaster[FakeInventoryDatabase.InvK(Tenant, "8010", "IT-1", "L1")].QtyOnHand); // 库存未二次变化
  212. Assert.Equal("POSTED", (await rig.Store.GetAsync(Tenant, "8010", "FB-1")).PostingStatus); // 已修复
  213. }
  214. // R4:两个消费者同时抢 stale PROCESSING → 只有一个获得恢复 claim(一 Processed / 一 Skipped)
  215. [Fact]
  216. public async Task R4_TwoConsumersRaceStale_OnlyOneRecovers()
  217. {
  218. var rig = Build(Receivable());
  219. await rig.Store.TryInsertAsync(new AdoIqcInventoryPosting
  220. {
  221. TenantId = Tenant, DomainCode = "8010", FbillNo = "FB-1", PostingStatus = "PROCESSING",
  222. CreateTime = DateTime.Now.AddMinutes(-10), UpdateTime = DateTime.Now.AddMinutes(-10),
  223. });
  224. var results = await Task.WhenAll(
  225. Task.Run(() => rig.Consumer.ConsumeAsync(Event())),
  226. Task.Run(() => rig.Consumer.ConsumeAsync(Event())));
  227. // 不变式:恰好一路真正执行(Processed);另一路是幂等 no-op——
  228. // Skipped(在赢家提交完成锚前判定) 或 AlreadyProcessed(在其后判定),二者都对,不重复过账。
  229. Assert.Equal(1, results.Count(x => x.Outcome == IqcConsumeOutcome.Processed));
  230. Assert.All(results.Where(x => x.Outcome != IqcConsumeOutcome.Processed),
  231. x => Assert.Contains(x.Outcome, new[] { IqcConsumeOutcome.Skipped, IqcConsumeOutcome.AlreadyProcessed }));
  232. Assert.Equal(1, rig.Db.InvMaster.Count); // 仅一套库存
  233. Assert.Equal(1, rig.Db.Completions.Count);
  234. }
  235. // R5:事务内 reload 报“当前不可收” → 0 C6/C5/C7/库存写 → Retryable(用事务内新鲜态否决)
  236. [Fact]
  237. public async Task R5_InTxReloadNotReceivable_NoWrites_Retryable()
  238. {
  239. var reload = Receivable();
  240. reload.HasReceivableDetails = false; // 事务内最新:不可收
  241. var rig = Build(Receivable(), reload: reload); // 事务外 snapshot 仍“可收”(预门禁放行)
  242. var r = await rig.Consumer.ConsumeAsync(Event());
  243. Assert.Equal(IqcConsumeOutcome.Retryable, r.Outcome);
  244. Assert.Equal(0, rig.Db.TotalPersisted);
  245. Assert.Empty(rig.Db.CallOrder); // reload 否决在 C6 之前
  246. Assert.Equal(0, rig.Db.Completions.Count);
  247. }
  248. // R6:C6 对账数量用事务内最新 BeforeReceived(非事件/事务外 snapshot)
  249. [Fact]
  250. public async Task R6_C6UsesInTxFreshQuantities()
  251. {
  252. var reload = Receivable();
  253. reload.BeforeReceivedQty = 20m; // 事务内最新已收 20(事务外 snapshot=0)
  254. var rig = Build(Receivable(), reload: reload);
  255. var r = await rig.Consumer.ConsumeAsync(Event());
  256. Assert.Equal(IqcConsumeOutcome.Processed, r.Outcome);
  257. var c6 = rig.Db.ReceiptWriteback.Values.Single();
  258. Assert.Equal(20m, c6.BeforeReceivedQty); // 事务内新鲜值
  259. Assert.Equal(120m, c6.AfterReceivedQty); // 20 + 合格100
  260. }
  261. }