| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449 |
- 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<InvalidOperationException>(() => MdpRebuildScope.Create("S2", 0, 1));
- Assert.Throws<InvalidOperationException>(() => MdpRebuildScope.Create("S2", 9, 0));
- Assert.Throws<InvalidOperationException>(() => MdpRebuildScope.Create("S2", 9, 9));
- Assert.Throws<InvalidOperationException>(() => 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<ModuleRebuildResult> RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func<ModuleProgressUpdate, Task> report, CancellationToken cancellationToken) =>
- Task.FromResult(new ModuleRebuildResult { BatchId = "b1", StageRows = 1 });
- }
- private sealed class ThrowingHandler : IModuleRebuildHandler
- {
- public string ModuleCode => "S2";
- public async Task<ModuleRebuildResult> RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func<ModuleProgressUpdate, Task> 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<ModuleRebuildResult> RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func<ModuleProgressUpdate, Task> 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<AdoModuleDashboardRebuildJob> 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<AdoModuleDashboardRebuildJob> InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default)
- {
- lock (_gate)
- {
- row.Id = _id++;
- Rows.Add(row);
- return Task.FromResult(row);
- }
- }
- public Task<AdoModuleDashboardRebuildJob> 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<AdoModuleDashboardRebuildJob> 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<AdoModuleDashboardRebuildJob> 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<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(IReadOnlyCollection<string> 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<AdoModuleDashboardRebuildJob>(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;
- }
- }
- }
- }
|