using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild; using Admin.NET.Plugin.AiDOP.Entity.SmartOps; using Admin.NET.Plugin.AiDOP.Supply; using Xunit; namespace Admin.NET.Plugin.AiDOP.Tests.DataPlatform; public class ModuleRebuildServiceTests { [Fact] public async Task SecondEnqueue_SameScope_ReturnsConflict() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var first = await svc.EnqueueAsync("S2", 9, 1, 1); var second = await svc.EnqueueAsync("S2", 9, 1, 1); Assert.Equal(202, first.StatusCode); Assert.Equal(409, second.StatusCode); Assert.Equal(first.Body.JobId, second.Body.JobId); } [Fact] public async Task OtherTenantEnqueue_NotBlocked() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var a = await svc.EnqueueAsync("S2", 9, 1, 1); var b = await svc.EnqueueAsync("S2", 10, 1, 2); Assert.Equal(202, a.StatusCode); Assert.Equal(202, b.StatusCode); Assert.NotEqual(a.Body.JobId, b.Body.JobId); } [Fact] public async Task DifferentModule_SameTenant_NotBlocked() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var a = await svc.EnqueueAsync("S2", 9, 1, 1); var b = await svc.EnqueueAsync("S3", 9, 1, 1); Assert.Equal(202, a.StatusCode); Assert.Equal(202, b.StatusCode); } [Fact] public async Task GetById_OtherTenant_NotFound() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var created = await svc.EnqueueAsync("S2", 9, 1, 1); var row = await store.GetByIdAsync("S2", created.Body.JobId!.Value, 10, 1); Assert.Null(row); } [Fact] public void Scope_RejectsZeroTenant() { Assert.Throws(() => MdpRebuildScope.Create("S2", 0, 1)); Assert.Throws(() => MdpRebuildScope.Create("S2", 9, 0)); Assert.Throws(() => MdpRebuildScope.Create("S2", 9, 9)); Assert.Throws(() => MdpRebuildScope.Create("S9", 9, 1)); } [Fact] public void SqlInject_AddsTenantToEveryWhere() { var sql = """ SELECT 1 FROM t WHERE a=1 UNION ALL SELECT 1 FROM t WHERE b=2 """; var injected = MdpSqlScope.InjectTenantFactory(sql); Assert.Equal(2, injected.Split("@TenantId", StringSplitOptions.None).Length - 1); } [Fact] public void TenantOnlySqlInject_DoesNotRequireFactoryColumn() { var injected = MdpSqlScope.InjectTenantOnly("SELECT 1 FROM mdp_std_delivery_schedule ds WHERE ds.tenant_id > 0"); Assert.Contains("ds.tenant_id=@TenantId", injected); Assert.DoesNotContain("factory_id", injected, StringComparison.OrdinalIgnoreCase); Assert.DoesNotContain("@FactoryId", injected); } [Fact] public void TenantOnlySqlInject_QualifiesRootAndNestedJoinScopes() { var sql = """ SELECT po.id FROM mdp_std_purchase_order po LEFT JOIN ( SELECT tenant_id, po_no FROM mdp_std_delivery_schedule WHERE po_no <> '' ) ds ON po.tenant_id=ds.tenant_id WHERE po.po_no <> '' """; var injected = MdpSqlScope.InjectTenantOnly(sql); Assert.Contains("FROM mdp_std_delivery_schedule\n WHERE tenant_id=@TenantId", injected); Assert.Contains("WHERE po.tenant_id=@TenantId", injected); } [Fact] public void S3ScopeDiscovery_ReadsNormalizedFactoryColumn() { var sql = MdpRebuildScopeCatalog.GetFactScopeSql("S3"); Assert.NotNull(sql); Assert.Contains("COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId", sql); } [Fact] public void S3DefaultInbound_ExcludesArchivedInvMaster() { Assert.DoesNotContain(S3MdpEntityConfig.All, x => x.EntityCode == "S3_INVENTORY" || x.SourceTable.Equals("InvMaster", StringComparison.OrdinalIgnoreCase)); } [Fact] public async Task Success_WritesHundredPercent() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, 1); var job = await store.ClaimNextQueuedAsync(new[] { "S2" }); var handler = new FakeHandler(); await svc.RunClaimedAsync(job, handler, new InMemoryModuleRebuildLock(), CancellationToken.None); var row = store.Latest("S2", 9, 1); Assert.Equal(ModuleRebuildStatus.Success, row.Status); Assert.Equal(100, row.ProgressPercent); } [Fact] public async Task AlreadyRunning_RequeuesJob() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, 1); var job = await store.ClaimNextQueuedAsync(new[] { "S2" }); var lck = new InMemoryModuleRebuildLock(); await using var held = await lck.TryAcquireAsync(MdpRebuildScope.Create("S2", 9, 1), "holder"); await svc.RunClaimedAsync(job, new FakeHandler(), lck, CancellationToken.None); Assert.Equal(ModuleRebuildStatus.Queued, store.Latest("S2", 9, 1).Status); } [Fact] public async Task Enqueue_InitialQueuedZeroPercent() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var accepted = await svc.EnqueueAsync("S4", 9, 1, 1, "AUTO"); var row = store.Latest("S4", 9, 1); Assert.Equal(202, accepted.StatusCode); Assert.Equal(ModuleRebuildStatus.Queued, row.Status); Assert.Equal(ModuleRebuildStages.Queued, row.CurrentStage); Assert.Equal(0, row.ProgressPercent); Assert.Equal(0, row.StageIndex); Assert.Equal(4, row.StageTotal); Assert.Equal("AUTO", row.TriggerType); } [Fact] public async Task ConcurrentEnqueue_SameScope_OnlyOneActive() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var results = await Task.WhenAll(Enumerable.Range(0, 20).Select(_ => svc.EnqueueAsync("S2", 9, 1, 1))); Assert.Equal(1, results.Count(x => x.StatusCode == 202)); Assert.Equal(19, results.Count(x => x.StatusCode == 409)); Assert.Single(store.Rows.Where(x => x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running)); Assert.True(results.Where(x => x.StatusCode == 409).All(x => x.Body.JobId == results.First(y => y.StatusCode == 202).Body.JobId)); } [Fact] public async Task DifferentFactory_NotBlocked() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var a = await svc.EnqueueAsync("S2", 9, 1, 1); var b = await svc.EnqueueAsync("S2", 9, 2, 1); Assert.Equal(202, a.StatusCode); Assert.Equal(202, b.StatusCode); Assert.NotEqual(a.Body.JobId, b.Body.JobId); } [Fact] public async Task DisabledModule_Returns404() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new DisabledModuleRebuildCapability()); var result = await svc.EnqueueAsync("S2", 9, 1, 1); Assert.Equal(404, result.StatusCode); Assert.Empty(store.Rows); } [Fact] public async Task Failed_KeepsStageAndPercent() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, 1); var job = await store.ClaimNextQueuedAsync(new[] { "S2" }); await svc.RunClaimedAsync(job, new ThrowingHandler(), new InMemoryModuleRebuildLock(), CancellationToken.None); var row = store.Latest("S2", 9, 1); Assert.Equal(ModuleRebuildStatus.Failed, row.Status); Assert.Equal(ModuleRebuildStages.Kpi, row.FailedStage); Assert.Equal(65, row.ProgressPercent); Assert.NotEqual(100, row.ProgressPercent); } [Fact] public async Task Cancellation_MarksCancelled() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, 1); var job = await store.ClaimNextQueuedAsync(new[] { "S2" }); using var cts = new CancellationTokenSource(); var run = svc.RunClaimedAsync(job, new HangingHandler(), new InMemoryModuleRebuildLock(), cts.Token); await Task.Delay(200); cts.Cancel(); await run; Assert.Equal(ModuleRebuildStatus.Cancelled, store.Latest("S2", 9, 1).Status); } [Fact] public async Task FailStale_ClosesRunningWithoutHeartbeat() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, 1); var job = store.Latest("S2", 9, 1); job.Status = ModuleRebuildStatus.Running; job.HeartbeatAt = DateTime.Now.AddHours(-1); await svc.FailStaleAsync(); Assert.Equal(ModuleRebuildStatus.Failed, store.Latest("S2", 9, 1).Status); Assert.Contains("服务中断", store.Latest("S2", 9, 1).ErrorMessage); } [Fact] public async Task Progress_DoesNotRegress() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, 1); var job = store.Latest("S2", 9, 1); await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 20, "stg")); await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 40, "std")); await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 60, "dwd")); var percents = store.ProgressHistory.Select(x => x.Percent).ToArray(); for (var i = 1; i < percents.Length; i++) Assert.True(percents[i] >= percents[i - 1], $"percent regress {percents[i - 1]} -> {percents[i]}"); Assert.Equal(60, job.ProgressPercent); } [Fact] public void StageIndex_IsMonotonicPerModule() { Assert.True(ModuleRebuildStages.ToStageIndex("S2", ModuleRebuildStages.Staging) < ModuleRebuildStages.ToStageIndex("S2", ModuleRebuildStages.Standard)); Assert.True(ModuleRebuildStages.ToStageIndex("S4", ModuleRebuildStages.Dwd) < ModuleRebuildStages.ToStageIndex("S4", ModuleRebuildStages.Kpi)); Assert.True(ModuleRebuildStages.ToStageIndex("S5", ModuleRebuildStages.T8Inbound) < ModuleRebuildStages.ToStageIndex("S5", ModuleRebuildStages.KpiCalculating)); Assert.Equal(5, ModuleRebuildStages.StageTotal("S1")); Assert.Equal(4, ModuleRebuildStages.StageTotal("S4")); Assert.Equal(4, ModuleRebuildStages.StageTotal("S7")); } [Fact] public void LockName_IsolatesTenantFactoryModule() { var a = MdpRebuildScope.Create("S2", 9, 1); var b = MdpRebuildScope.Create("S2", 10, 1); var c = MdpRebuildScope.Create("S3", 9, 1); Assert.Equal("S2_MDP_FULL:9:1", a.LockName); Assert.NotEqual(a.LockName, b.LockName); Assert.NotEqual(a.LockName, c.LockName); } [Fact] public void SqlInject_DoesNotDoubleInject() { var sql = "SELECT 1 FROM t w WHERE w.tenant_id=@TenantId AND a=1"; var injected = MdpSqlScope.InjectTenantFactory(sql); Assert.Equal(1, injected.Split("@TenantId", StringSplitOptions.None).Length - 1); } [Fact] public async Task JobRunner_EnqueuesEachScope() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var a = await svc.EnqueueAsync("S2", 9, 1, null, "BOOTSTRAP"); var b = await svc.EnqueueAsync("S2", 10, 1, null, "BOOTSTRAP"); Assert.Equal(202, a.StatusCode); Assert.Equal(202, b.StatusCode); Assert.Equal(2, store.Rows.Count); Assert.All(store.Rows, x => Assert.Equal("BOOTSTRAP", x.TriggerType)); } private sealed class FakeHandler : IModuleRebuildHandler { public string ModuleCode => "S2"; public Task RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func report, CancellationToken cancellationToken) => Task.FromResult(new ModuleRebuildResult { BatchId = "b1", StageRows = 1 }); } private sealed class ThrowingHandler : IModuleRebuildHandler { public string ModuleCode => "S2"; public async Task RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func report, CancellationToken cancellationToken) { await report(new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 65, "正在重算 KPI")); throw new InvalidOperationException("kpi failed"); } } private sealed class HangingHandler : IModuleRebuildHandler { public string ModuleCode => "S2"; public async Task RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func report, CancellationToken cancellationToken) { await Task.Delay(Timeout.Infinite, cancellationToken); return new ModuleRebuildResult(); } } private sealed class DisabledModuleRebuildCapability : IModuleRebuildCapability { public bool IsEnabled(string moduleCode) => false; public int MaxParallelScopes => 2; } private sealed class MemoryJobStore : IModuleRebuildJobStore { public List Rows { get; } = new(); public List<(long JobId, int Percent, string Stage)> ProgressHistory { get; } = new(); private long _id = 1; private readonly object _gate = new(); public AdoModuleDashboardRebuildJob Latest(string module, long tenantId, long factoryId) { lock (_gate) return Rows.Where(x => x.ModuleCode == module && x.TenantId == tenantId && x.FactoryId == factoryId).OrderByDescending(x => x.Id).First(); } public Task InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) { lock (_gate) { row.Id = _id++; Rows.Add(row); return Task.FromResult(row); } } public Task FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId && x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running) .OrderBy(x => x.Id).FirstOrDefault()); } public Task GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.FirstOrDefault(x => x.Id == id && x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId)); } public Task GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId) .OrderByDescending(x => x.Id).FirstOrDefault()); } public Task ClaimNextQueuedAsync(IReadOnlyCollection enabledModules, CancellationToken ct = default) { lock (_gate) { var row = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued && enabledModules.Contains(x.ModuleCode)) .Where(x => !Rows.Any(y => y.ModuleCode == x.ModuleCode && y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == ModuleRebuildStatus.Running)) .OrderBy(x => x.Id).FirstOrDefault(); if (row == null) return Task.FromResult(null); row.Status = ModuleRebuildStatus.Running; row.StartedAt = DateTime.Now; row.HeartbeatAt = DateTime.Now; return Task.FromResult(row); } } public Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) => Task.CompletedTask; public Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default) { lock (_gate) { var row = Rows.FirstOrDefault(x => x.Id == jobId); if (row == null) return Task.CompletedTask; row.CurrentStage = currentStage; row.StageIndex = stageIndex; row.ProgressPercent = progressPercent; row.ProgressMessage = message; row.LastProgressAt = now; row.HeartbeatAt = now; row.UpdateTime = now; ProgressHistory.Add((jobId, progressPercent, currentStage)); return Task.CompletedTask; } } public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default) => UpdateProgressAsync(jobId, completedStage, 0, progressPercent, nextMessage, now, ct); public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default) { lock (_gate) { var row = Rows.FirstOrDefault(x => x.Id == jobId); if (row != null) row.HeartbeatAt = now; return Task.CompletedTask; } } public Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default) { lock (_gate) { var cutoff = DateTime.Now - staleAfter; foreach (var row in Rows.Where(x => x.Status == ModuleRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff))) { row.Status = ModuleRebuildStatus.Failed; row.CurrentStage = ModuleRebuildStages.Failed; row.FinishedAt = DateTime.Now; row.LastProgressAt = DateTime.Now; row.UpdateTime = DateTime.Now; row.ErrorMessage = "服务中断,任务未正常结束"; } return Task.CompletedTask; } } } }