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