using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Entity; using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Idempotency; using Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.Receipt; namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse.InventoryPosting.L1; public enum IqcConsumeOutcome { /// 首次处理并成功(ack)。 Processed = 1, /// 既有已完成(POSTED,或崩溃恢复检出完成锚)→ 幂等 ack、不重复执行。 AlreadyProcessed = 2, /// 另一路处理中且 lease 未过期(或抢占竞争失败)→ 不起第二套事务。 Skipped = 3, /// 业务失败/状态未就绪/domain 未解析 → 事件保持可重试。 Retryable = 4, /// 无库存动作(ReturnOnly/未匹配)→ 幂等完成。 NoAction = 5, } public sealed class IqcConsumeResult { public IqcConsumeOutcome Outcome { get; set; } public string Reason { get; set; } } /// /// L1 IQC 合格事件消费入口。**只负责触发编排**,不含 MRB 数量公式 / 9 守卫 / C6 计算 / 库存 delta(均归已有组件)。 /// /// 幂等/可靠性(Phase 4A + 7C)+ 接线(5B-3): /// 0. domain 解析:经 校验非空(无 8010 fallback),空→Retryable。 /// 1. 幂等认领四态:POSTED→ack;PROCESSING→(先查完成锚)已提交则修复 POSTED+ack,否则原子抢占 stale 否则 skip;FirstProcess/FailedRetryable→继续。 /// 2. commit→MarkPosted 崩溃窗口:完成锚与业务同事务提交;重投检出锚即修复 POSTED、不重复过账。 /// 3. Runtime Gate 事务内新鲜度 + freshness 水位:标准层落后事件→Retryable。 /// /// ack 边界:业务 Commit 成功后才 MarkPosted(ack);失败/异常 MarkFailed→retryable。端到端幂等锚 = posting/completion (tenant,domain,FBILLNO) UNIQUE。 /// public sealed class IqcInventoryEventConsumer { private static readonly TimeSpan DefaultLease = TimeSpan.FromMinutes(5); private readonly IqcPostingIdempotencyService _idempotency; private readonly IBusinessCompletionStore _completion; private readonly IIqcReceiptStateLoader _loader; private readonly IqcReceiptOrchestrator _orchestrator; private readonly IqcMrbSelectionTaskService _mrbSelection; private readonly IqcDomainResolver _domainResolver; private readonly TimeSpan _lease; public IqcInventoryEventConsumer( IqcPostingIdempotencyService idempotency, IBusinessCompletionStore completion, IIqcReceiptStateLoader loader, IqcReceiptOrchestrator orchestrator, IqcMrbSelectionTaskService mrbSelection, TimeSpan? processingLease = null, IqcDomainResolver domainResolver = null) { _idempotency = idempotency ?? throw new ArgumentNullException(nameof(idempotency)); _completion = completion ?? throw new ArgumentNullException(nameof(completion)); _loader = loader ?? throw new ArgumentNullException(nameof(loader)); _orchestrator = orchestrator ?? throw new ArgumentNullException(nameof(orchestrator)); _mrbSelection = mrbSelection ?? throw new ArgumentNullException(nameof(mrbSelection)); _domainResolver = domainResolver; _lease = processingLease ?? DefaultLease; } public async Task ConsumeAsync(IqcInventoryEvent e) { if (e == null) throw new ArgumentNullException(nameof(e)); if (string.IsNullOrWhiteSpace(e.Fbillno)) throw new ArgumentException("FBILLNO 不能为空", nameof(e)); // 0. domain 解析(无 8010 fallback;空→Retryable) string domain; try { domain = _domainResolver != null ? _domainResolver.Resolve(e.TenantId, e.DomainCode) : (e.DomainCode ?? string.Empty); } catch (Exception ex) { return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = ex.Message }; } // 1. Phase 4A 幂等认领(四态) var claim = await _idempotency.ClaimAsync(e.TenantId, domain, e.Fbillno, p => { p.Receiver = e.Receiver; p.RctQcNbr = e.RctQcNbr; }); if (claim.Outcome == IqcPostingClaimOutcome.AlreadyProcessed) return new IqcConsumeResult { Outcome = IqcConsumeOutcome.AlreadyProcessed }; // POSTED → ack if (claim.Outcome != IqcPostingClaimOutcome.FirstProcess) { // A3:完成锚存在 ⇔ 业务事务已完整提交 → 修复 POSTED + ack、不重复过账 if (await _completion.IsCommittedAsync(e.TenantId, domain, e.Fbillno)) { await _idempotency.MarkPostedAsync(claim.Posting); return new IqcConsumeResult { Outcome = IqcConsumeOutcome.AlreadyProcessed, Reason = "崩溃恢复:检出完成锚" }; } // A2:仅 PROCESSING 需抢占;lease 未过期或竞争失败 → 不起第二套 if (claim.Outcome == IqcPostingClaimOutcome.Processing) { var recovered = await _idempotency.TryRecoverStaleAsync(e.TenantId, domain, e.Fbillno, _lease); if (!recovered) return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Skipped, Reason = "PROCESSING 未过期或已被他路抢占" }; } } var posting = claim.Posting; try { // 2. C1 路由映射 var map = IqcResultMapping.Map(new IqcResultMappingInput { Pd = e.Pd, Clfs = e.Clfs, Dhsl = e.Dhsl, Bhgsl = e.Bhgsl }); if (!map.Matched || map.RoutingKind == IqcRoutingKind.Unmatched || map.RoutingKind == IqcRoutingKind.ReturnOnly) { await _idempotency.MarkPostedAsync(posting); return new IqcConsumeResult { Outcome = IqcConsumeOutcome.NoAction, Reason = map.BranchSemantics }; } // 3. 加载收货态(routing/identity 基线;最终数量/门禁在事务内 reloadInTx 复核) var st = await _loader.LoadAsync(e.TenantId, domain, e.Fbillno); if (!st.Found) { await _idempotency.MarkFailedAsync(posting); return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = "收货状态未就绪" }; } // B8 freshness(5B-3 §17 fail-closed):水位缺失(null=未知) 或 落后事件 → 不过账、保持 retryable。 if (!st.SourceSyncedAt.HasValue || st.SourceSyncedAt.Value < e.EventTime) { await _idempotency.MarkFailedAsync(posting); return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = st.SourceSyncedAt.HasValue ? "标准层落后于事件(freshness)" : "标准层水位未知(freshness fail-closed)", }; } if (map.RoutingKind == IqcRoutingKind.FullReceiptPosting) { var gate = new PurchaseReceiptGateInput { HasPendingReceipt = st.HasPendingReceipt, AcceptedQty = map.QtyAcc, HasReceivableDetails = st.HasReceivableDetails, }; var receipt = new IqcAcceptedReceipt { TenantId = e.TenantId, DomainCode = domain, PostingId = posting.Id, TransactionGroupId = posting.Id, RctNbr = st.RctNbr, ItemNum = st.ItemNum, Location = st.Location, LotSerial = st.LotSerial ?? "", Refs = st.Refs ?? "", AcceptedQty = map.QtyAcc, SampleQty = st.SampleQty, UmConversion = 1m, PurOrd = st.PurOrd, PurLine = st.PurLine, Potype = st.Potype, OrderedQty = st.OrderedQty, BeforeReceivedQty = st.BeforeReceivedQty, ReturnedQty = st.ReturnedQty, RctQcNbr = e.RctQcNbr, Receiver = e.Receiver, Fbillno = e.Fbillno, User = e.User, }; var r = await _orchestrator.PostAsync(receipt, gate, null, () => ReloadForPostingAsync(e.TenantId, domain, e.Fbillno)); return await FinalizeAsync(posting, r.Outcome == IqcReceiptOutcome.Success, r.Reason ?? r.LegacyReturnMsg); } else // MrbChooseMobileTask { var mrb = new IqcMrbSelectionCommand { TenantId = e.TenantId, DomainCode = domain, PostingId = posting.Id, QcNbr = e.RctQcNbr, ItemNum = st.ItemNum, Fbillno = e.Fbillno, Receiver = e.Receiver, User = e.User, }; var r = await _mrbSelection.PostAsync(mrb); return await FinalizeAsync(posting, r.Outcome == IqcReceiptOutcome.Success, r.Reason ?? r.LegacyReturnMsg); } } catch (Exception ex) { try { await _idempotency.MarkFailedAsync(posting); } catch { } return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = ex.Message }; } } /// A5:事务内重读可变收货态 → Receipt 层刷新对象(orchestrator 在事务内调用)。 private async Task ReloadForPostingAsync(long tenantId, string domain, string fbillno) { var f = await _loader.LoadAsync(tenantId, domain, fbillno); return new ReceiptPostingRefresh { Found = f.Found, HasPendingReceipt = f.HasPendingReceipt, HasReceivableDetails = f.HasReceivableDetails, OrderedQty = f.OrderedQty, BeforeReceivedQty = f.BeforeReceivedQty, ReturnedQty = f.ReturnedQty, SampleQty = f.SampleQty, }; } private async Task FinalizeAsync(AdoIqcInventoryPosting posting, bool success, string reason) { if (success) { await _idempotency.MarkPostedAsync(posting); // 业务 Commit 成功 → ack(完成锚已在业务事务内落地) return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Processed }; } await _idempotency.MarkFailedAsync(posting); return new IqcConsumeResult { Outcome = IqcConsumeOutcome.Retryable, Reason = reason }; } }