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 { private static void SeedSuccess(MemoryJobStore store, string module, long tenant, DateTime finishedAt) { store.Rows.Add(new AdoModuleDashboardRebuildJob { Id = 9000 + store.Rows.Count, ModuleCode = module, TenantId = tenant, FactoryId = 1, Status = ModuleRebuildStatus.Success, FinishedAt = finishedAt }); } [Fact] public async Task AutoEnqueue_WithinCooldown_ReturnsCooldown() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddHours(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var before = store.Rows.Count; var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO"); Assert.Equal(409, result.StatusCode); Assert.Equal("COOLDOWN", result.Body.Status); Assert.Equal(before, store.Rows.Count); } [Fact] public async Task AutoEnqueue_AfterCooldown_IsAccepted() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddHours(-7)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO"); Assert.Equal(202, result.StatusCode); } [Fact] public async Task ManualEnqueue_IsNeverCooledDown() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var result = await svc.EnqueueAsync("S2", 9, 1, null); Assert.Equal(202, result.StatusCode); } [Fact] public async Task ManualEnqueue_WithRequestedBy_IsNeverCooledDown() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var result = await svc.EnqueueAsync("S2", 9, 1, 123, "AUTO"); Assert.Equal(202, result.StatusCode); } [Fact] public async Task NightlyEnqueue_BypassesCooldown() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY"); Assert.Equal(202, result.StatusCode); } [Fact] public async Task BootstrapEnqueue_IsCooledDown() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var result = await svc.EnqueueAsync("S2", 9, 1, null, "BOOTSTRAP"); Assert.Equal(409, result.StatusCode); Assert.Equal("COOLDOWN", result.Body.Status); } [Fact] public async Task CooldownZero_DisablesWindow() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new ConfigurableCapability { AutoMinIntervalHours = 0 }); var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO"); Assert.Equal(202, result.StatusCode); } [Fact] public async Task RunningJob_StillBlocksBeforeCooldownCheck() { var store = new MemoryJobStore(); SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1)); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, null); var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO"); Assert.Equal(409, result.StatusCode); Assert.Equal(ModuleRebuildStatus.Queued, result.Body.Status); } /// /// 夜间兜底任务在非执行机上永远领不到(ClaimNextQueuedAsync 的白名单), /// 若它还能挡住人工入队,页面「数据重算」按钮就被永久卡死。人工请求必须能取代它。 /// [Fact] public async Task ManualEnqueue_SupersedesQueuedAutoJob() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); var nightly = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY"); var manual = await svc.EnqueueAsync("S2", 9, 1, 7); Assert.Equal(202, manual.StatusCode); Assert.NotEqual(nightly.Body.JobId, manual.Body.JobId); var superseded = store.Rows.Single(x => x.Id == nightly.Body.JobId); Assert.Equal(ModuleRebuildStatus.Cancelled, superseded.Status); Assert.Contains("AUTO_NIGHTLY", superseded.ErrorMessage); } /// 已经开跑的任务不得被掀掉:取代只针对 QUEUED。 [Fact] public async Task ManualEnqueue_DoesNotSupersedeRunningAutoJob() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability()); await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY"); var running = await store.ClaimNextQueuedAsync(new[] { "S2" }); var manual = await svc.EnqueueAsync("S2", 9, 1, 7); Assert.Equal(409, manual.StatusCode); Assert.Equal(ModuleRebuildStatus.Running, store.Rows.Single(x => x.Id == running.Id).Status); } /// 取代资格只给人工请求;自动触发器之间仍按原样互斥,避免夜间批次自相残杀。 [Fact] public async Task AutoEnqueue_DoesNotSupersedeQueuedAutoJob() { var store = new MemoryJobStore(); var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new ConfigurableCapability { AutoMinIntervalHours = 0 }); var nightly = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY"); var auto = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO"); Assert.Equal(409, auto.StatusCode); Assert.Equal(nightly.Body.JobId, auto.Body.JobId); Assert.Equal(ModuleRebuildStatus.Queued, store.Rows.Single(x => x.Id == nightly.Body.JobId).Status); } [Fact] public void S1Service_DelegatesToUnifiedPath() { var source = File.ReadAllText(RepoFile("server", "Plugins", "Admin.NET.Plugin.AiDOP", "DataPlatform", "S1Refresh", "S1DashboardRebuildService.cs")); var unified = source.IndexOf("if (UseUnified)", StringComparison.Ordinal); Assert.True(unified >= 0, "S1 统一路径分支丢失"); var body = source[unified..source.IndexOf("var scope = S1MdpRunScope.Create", unified, StringComparison.Ordinal)]; Assert.Contains("_moduleRebuild.EnqueueAsync(", body); } [Fact] public void NightlyJob_UsesNonCollidingCron() { var source = File.ReadAllText(RepoFile("server", "Plugins", "Admin.NET.Plugin.AiDOP", "Job", "MdpNightlyFullRebuildJob.cs")); Assert.Contains("Cron(\"40 3 * * *\"", source); Assert.DoesNotContain("\"15 3\"", source); Assert.DoesNotContain("\"0 2\"", source); } [Theory] [InlineData("AUTO", false)] [InlineData("auto", false)] [InlineData("MANUAL", true)] [InlineData("AUTO_NIGHTLY", true)] [InlineData("BOOTSTRAP", true)] [InlineData("", true)] public void ShouldPullFull_OnlyAutoIsIncremental(string triggerType, bool expected) => Assert.Equal(expected, ModuleRebuildTriggerType.ShouldPullFull(triggerType)); private static string RepoFile(params string[] parts) { var dir = new DirectoryInfo(AppContext.BaseDirectory); while (dir != null && !File.Exists(Path.Combine(dir.FullName, "AGENTS.md"))) dir = dir.Parent; Assert.NotNull(dir); return Path.Combine(new[] { dir!.FullName }.Concat(parts).ToArray()); } [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); } /// /// mdp_source_table_registry 是全局字典表,只有 (id, source_table, source_system, remark)。 /// 注入器过去把语句里每个 WHERE 都改写,给它补出 r.tenant_id=@TenantId, /// S3 的 mdp_std_supplier_item 因此每轮都报 Unknown column 'r.tenant_id' 并整轮 FAILED; /// 同样的子查询在 S4 只是碰巧被旧的 240 字符窗口盖住才没炸。 /// [Fact] public void SqlInject_SkipsGlobalRegistryLookup() { var sql = """ INSERT INTO mdp_std_supplier_item (source_system, item_code) SELECT COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = source_table LIMIT 1), '') FROM mdp_stg_source_list WHERE source_table='srm_purchase' """; var injected = MdpSqlScope.InjectTenantFactory(sql); Assert.Contains("WHERE r.source_table = source_table", injected); Assert.DoesNotContain("r.tenant_id", injected); Assert.DoesNotContain("r.factory_id", injected); Assert.Contains("WHERE tenant_id=@TenantId", injected); } /// /// 子查询自带的租户条件不能替外层背书。旧的 240 字符窗口会把它算进来, /// 外层于是整条漏掉租户过滤 —— 跨租户串数据,且不会报错。 /// [Fact] public void SqlInject_IgnoresTenantParamInNestedSubquery() { var sql = """ SELECT 1 FROM dwd_ship_trans d WHERE d.biz_date = (SELECT MAX(biz_date) FROM dwd_ship_trans WHERE tenant_id=@TenantId) """; var injected = MdpSqlScope.InjectTenantFactory(sql); Assert.Contains("WHERE d.tenant_id=@TenantId", injected); } /// /// 反过来:手写的租户条件即使排在 240 字符之外,也必须被认出来,不能重复注入。 /// [Fact] public void SqlInject_HonorsTenantPredicateBeyondOldWindow() { var filler = string.Join("\n", Enumerable.Range(0, 12) .Select(i => $" AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.f{i}')),'') <> ''")); var sql = $"SELECT 1 FROM mdp_std_so s\nWHERE s.order_no <> ''\n{filler}\n AND s.tenant_id=@TenantId"; var injected = MdpSqlScope.InjectTenantFactory(sql); Assert.Equal(sql, 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; public int GlobalMaxParallelScopes => 2; public int AutoMinIntervalHours => 6; } private sealed class ConfigurableCapability : IModuleRebuildCapability { public int AutoMinIntervalHours { get; init; } = 6; public bool IsEnabled(string moduleCode) => true; public int MaxParallelScopes => 2; public int GlobalMaxParallelScopes => 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 Task ClaimNextQueuedAsync(IReadOnlyCollection enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default) => Task.FromResult(null!); public Task RequestCancelAsync(long jobId, CancellationToken ct = default) { lock (_gate) { var row = Rows.FirstOrDefault(x => x.Id == jobId && x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running); if (row == null) return Task.FromResult(false); row.CancelRequestedFlag = true; return Task.FromResult(true); } } public Task IsCancelRequestedAsync(long jobId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.FirstOrDefault(x => x.Id == jobId)?.CancelRequestedFlag ?? false); } public Task FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter) => Task.FromResult(0L); public Task FindLastSuccessAsync(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 == ModuleRebuildStatus.Success) .OrderByDescending(x => x.FinishedAt) .FirstOrDefault()); } 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 TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default) { lock (_gate) { var row = Rows.FirstOrDefault(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Queued); if (row == null) return Task.FromResult(false); row.Status = ModuleRebuildStatus.Cancelled; row.CurrentStage = ModuleRebuildStages.Cancelled; row.FinishedAt = DateTime.Now; row.ErrorMessage = reason; return Task.FromResult(true); } } public Task CancelUnstartedByTriggerAsync(string triggerType, TimeSpan olderThan, string reason, CancellationToken ct = default) { lock (_gate) { var trigger = triggerType?.Trim().ToUpperInvariant(); var cutoff = DateTime.Now - olderThan; var hits = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued && string.Equals(x.TriggerType, trigger, StringComparison.Ordinal) && x.RequestedBy == null && (olderThan <= TimeSpan.Zero || x.SubmittedAt < cutoff)) .ToList(); foreach (var row in hits) { row.Status = ModuleRebuildStatus.Cancelled; row.CurrentStage = ModuleRebuildStages.Cancelled; row.FinishedAt = DateTime.Now; row.ErrorMessage = reason; } return Task.FromResult(hits.Count); } } public Task FailStaleQueuedAsync(TimeSpan staleAfter, string reason, CancellationToken ct = default) { lock (_gate) { var cutoff = DateTime.Now - staleAfter; var hits = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued && x.SubmittedAt < cutoff).ToList(); foreach (var row in hits) { row.Status = ModuleRebuildStatus.Failed; row.CurrentStage = ModuleRebuildStages.Failed; row.FinishedAt = DateTime.Now; row.ErrorMessage = reason; } return Task.FromResult(hits.Count); } } 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; } } } }