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 };
}
}