IqcInventoryEventConsumer.cs 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188
  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.Receipt;
  4. namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.L1;
  5. public enum IqcConsumeOutcome
  6. {
  7. /// <summary>首次处理并成功(ack)。</summary>
  8. Processed = 1,
  9. /// <summary>既有已完成(POSTED,或崩溃恢复检出完成锚)→ 幂等 ack、不重复执行。</summary>
  10. AlreadyProcessed = 2,
  11. /// <summary>另一路处理中且 lease 未过期(或抢占竞争失败)→ 不起第二套事务。</summary>
  12. Skipped = 3,
  13. /// <summary>业务失败/状态未就绪/domain 未解析 → 事件保持可重试。</summary>
  14. Retryable = 4,
  15. /// <summary>无库存动作(ReturnOnly/未匹配)→ 幂等完成。</summary>
  16. NoAction = 5,
  17. }
  18. public sealed class IqcConsumeResult
  19. {
  20. public IqcConsumeOutcome Outcome { get; set; }
  21. public string Reason { get; set; }
  22. }
  23. /// <summary>
  24. /// L1 IQC 合格事件消费入口。**只负责触发编排**,不含 MRB 数量公式 / 9 守卫 / C6 计算 / 库存 delta(均归已有组件)。
  25. ///
  26. /// 幂等/可靠性(Phase 4A + 7C)+ 接线(5B-3):
  27. /// 0. domain 解析:经 <see cref="IqcDomainResolver"/> 校验非空(无 8010 fallback),空→Retryable。
  28. /// 1. 幂等认领四态:POSTED→ack;PROCESSING→(先查完成锚)已提交则修复 POSTED+ack,否则原子抢占 stale 否则 skip;FirstProcess/FailedRetryable→继续。
  29. /// 2. commit→MarkPosted 崩溃窗口:完成锚与业务同事务提交;重投检出锚即修复 POSTED、不重复过账。
  30. /// 3. Runtime Gate 事务内新鲜度 + freshness 水位:标准层落后事件→Retryable。
  31. ///
  32. /// ack 边界:业务 Commit 成功后才 MarkPosted(ack);失败/异常 MarkFailed→retryable。端到端幂等锚 = posting/completion (tenant,domain,FBILLNO) UNIQUE。
  33. /// </summary>
  34. public sealed class IqcInventoryEventConsumer
  35. {
  36. private static readonly TimeSpan DefaultLease = TimeSpan.FromMinutes(5);
  37. private readonly IqcPostingIdempotencyService _idempotency;
  38. private readonly IBusinessCompletionStore _completion;
  39. private readonly IIqcReceiptStateLoader _loader;
  40. private readonly IqcReceiptOrchestrator _orchestrator;
  41. private readonly IqcMrbSelectionTaskService _mrbSelection;
  42. private readonly IqcDomainResolver _domainResolver;
  43. private readonly TimeSpan _lease;
  44. public IqcInventoryEventConsumer(
  45. IqcPostingIdempotencyService idempotency, IBusinessCompletionStore completion, IIqcReceiptStateLoader loader,
  46. IqcReceiptOrchestrator orchestrator, IqcMrbSelectionTaskService mrbSelection,
  47. TimeSpan? processingLease = null, IqcDomainResolver domainResolver = null)
  48. {
  49. _idempotency = idempotency ?? throw new ArgumentNullException(nameof(idempotency));
  50. _completion = completion ?? throw new ArgumentNullException(nameof(completion));
  51. _loader = loader ?? throw new ArgumentNullException(nameof(loader));
  52. _orchestrator = orchestrator ?? throw new ArgumentNullException(nameof(orchestrator));
  53. _mrbSelection = mrbSelection ?? throw new ArgumentNullException(nameof(mrbSelection));
  54. _domainResolver = domainResolver;
  55. _lease = processingLease ?? DefaultLease;
  56. }
  57. public async Task<IqcConsumeResult> ConsumeAsync(IqcInventoryEvent e)
  58. {
  59. if (e == null) throw new ArgumentNullException(nameof(e));
  60. if (string.IsNullOrWhiteSpace(e.Fbillno)) throw new ArgumentException("FBILLNO 不能为空", nameof(e));
  61. // 0. domain 解析(无 8010 fallback;空→Retryable)
  62. string domain;
  63. try { domain = _domainResolver != null ? _domainResolver.Resolve(e.TenantId, e.DomainCode) : (e.DomainCode ?? string.Empty); }
  64. catch (Exception ex) { return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = ex.Message }; }
  65. // 1. Phase 4A 幂等认领(四态)
  66. var claim = await _idempotency.ClaimAsync(e.TenantId, domain, e.Fbillno, p =>
  67. {
  68. p.Receiver = e.Receiver; p.RctQcNbr = e.RctQcNbr;
  69. });
  70. if (claim.Outcome == IqcPostingClaimOutcome.AlreadyProcessed)
  71. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.AlreadyProcessed }; // POSTED → ack
  72. if (claim.Outcome != IqcPostingClaimOutcome.FirstProcess)
  73. {
  74. // A3:完成锚存在 ⇔ 业务事务已完整提交 → 修复 POSTED + ack、不重复过账
  75. if (await _completion.IsCommittedAsync(e.TenantId, domain, e.Fbillno))
  76. {
  77. await _idempotency.MarkPostedAsync(claim.Posting);
  78. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.AlreadyProcessed, Reason = "崩溃恢复:检出完成锚" };
  79. }
  80. // A2:仅 PROCESSING 需抢占;lease 未过期或竞争失败 → 不起第二套
  81. if (claim.Outcome == IqcPostingClaimOutcome.Processing)
  82. {
  83. var recovered = await _idempotency.TryRecoverStaleAsync(e.TenantId, domain, e.Fbillno, _lease);
  84. if (!recovered)
  85. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Skipped, Reason = "PROCESSING 未过期或已被他路抢占" };
  86. }
  87. }
  88. var posting = claim.Posting;
  89. try
  90. {
  91. // 2. C1 路由映射
  92. var map = IqcResultMapping.Map(new IqcResultMappingInput { Pd = e.Pd, Clfs = e.Clfs, Dhsl = e.Dhsl, Bhgsl = e.Bhgsl });
  93. if (!map.Matched || map.RoutingKind == IqcRoutingKind.Unmatched || map.RoutingKind == IqcRoutingKind.ReturnOnly)
  94. {
  95. await _idempotency.MarkPostedAsync(posting);
  96. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.NoAction, Reason = map.BranchSemantics };
  97. }
  98. // 3. 加载收货态(routing/identity 基线;最终数量/门禁在事务内 reloadInTx 复核)
  99. var st = await _loader.LoadAsync(e.TenantId, domain, e.Fbillno);
  100. if (!st.Found)
  101. {
  102. await _idempotency.MarkFailedAsync(posting);
  103. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = "收货状态未就绪" };
  104. }
  105. // B8 freshness(5B-3 §17 fail-closed):水位缺失(null=未知) 或 落后事件 → 不过账、保持 retryable。
  106. if (!st.SourceSyncedAt.HasValue || st.SourceSyncedAt.Value < e.EventTime)
  107. {
  108. await _idempotency.MarkFailedAsync(posting);
  109. return new IqcConsumeResult
  110. {
  111. Outcome = IqcConsumeOutcome.Retryable,
  112. Reason = st.SourceSyncedAt.HasValue ? "标准层落后于事件(freshness)" : "标准层水位未知(freshness fail-closed)",
  113. };
  114. }
  115. if (map.RoutingKind == IqcRoutingKind.FullReceiptPosting)
  116. {
  117. var gate = new PurchaseReceiptGateInput
  118. {
  119. HasPendingReceipt = st.HasPendingReceipt, AcceptedQty = map.QtyAcc, HasReceivableDetails = st.HasReceivableDetails,
  120. };
  121. var receipt = new IqcAcceptedReceipt
  122. {
  123. TenantId = e.TenantId, DomainCode = domain,
  124. PostingId = posting.Id, TransactionGroupId = posting.Id,
  125. RctNbr = st.RctNbr, ItemNum = st.ItemNum, Location = st.Location, LotSerial = st.LotSerial ?? "", Refs = st.Refs ?? "",
  126. AcceptedQty = map.QtyAcc, SampleQty = st.SampleQty, UmConversion = 1m,
  127. PurOrd = st.PurOrd, PurLine = st.PurLine, Potype = st.Potype,
  128. OrderedQty = st.OrderedQty, BeforeReceivedQty = st.BeforeReceivedQty, ReturnedQty = st.ReturnedQty,
  129. RctQcNbr = e.RctQcNbr, Receiver = e.Receiver, Fbillno = e.Fbillno, User = e.User,
  130. };
  131. var r = await _orchestrator.PostAsync(receipt, gate, null, () => ReloadForPostingAsync(e.TenantId, domain, e.Fbillno));
  132. return await FinalizeAsync(posting, r.Outcome == IqcReceiptOutcome.Success, r.Reason ?? r.LegacyReturnMsg);
  133. }
  134. else // MrbChooseMobileTask
  135. {
  136. var mrb = new IqcMrbSelectionCommand
  137. {
  138. TenantId = e.TenantId, DomainCode = domain, PostingId = posting.Id,
  139. QcNbr = e.RctQcNbr, ItemNum = st.ItemNum, Fbillno = e.Fbillno, Receiver = e.Receiver, User = e.User,
  140. };
  141. var r = await _mrbSelection.PostAsync(mrb);
  142. return await FinalizeAsync(posting, r.Outcome == IqcReceiptOutcome.Success, r.Reason ?? r.LegacyReturnMsg);
  143. }
  144. }
  145. catch (Exception ex)
  146. {
  147. try { await _idempotency.MarkFailedAsync(posting); } catch { }
  148. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = ex.Message };
  149. }
  150. }
  151. /// <summary>A5:事务内重读可变收货态 → Receipt 层刷新对象(orchestrator 在事务内调用)。</summary>
  152. private async Task<ReceiptPostingRefresh> ReloadForPostingAsync(long tenantId, string domain, string fbillno)
  153. {
  154. var f = await _loader.LoadAsync(tenantId, domain, fbillno);
  155. return new ReceiptPostingRefresh
  156. {
  157. Found = f.Found, HasPendingReceipt = f.HasPendingReceipt, HasReceivableDetails = f.HasReceivableDetails,
  158. OrderedQty = f.OrderedQty, BeforeReceivedQty = f.BeforeReceivedQty, ReturnedQty = f.ReturnedQty, SampleQty = f.SampleQty,
  159. };
  160. }
  161. private async Task<IqcConsumeResult> FinalizeAsync(AdoIqcInventoryPosting posting, bool success, string reason)
  162. {
  163. if (success)
  164. {
  165. await _idempotency.MarkPostedAsync(posting); // 业务 Commit 成功 → ack(完成锚已在业务事务内落地)
  166. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Processed };
  167. }
  168. await _idempotency.MarkFailedAsync(posting);
  169. return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = reason };
  170. }
  171. }