ModuleRebuildServiceTests.cs 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  2. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  3. using Admin.NET.Plugin.AiDOP.Supply;
  4. using Xunit;
  5. namespace Admin.NET.Plugin.AiDOP.Tests.DataPlatform;
  6. public class ModuleRebuildServiceTests
  7. {
  8. private static void SeedSuccess(MemoryJobStore store, string module, long tenant, DateTime finishedAt)
  9. {
  10. store.Rows.Add(new AdoModuleDashboardRebuildJob
  11. {
  12. Id = 9000 + store.Rows.Count,
  13. ModuleCode = module,
  14. TenantId = tenant,
  15. FactoryId = 1,
  16. Status = ModuleRebuildStatus.Success,
  17. FinishedAt = finishedAt
  18. });
  19. }
  20. [Fact]
  21. public async Task AutoEnqueue_WithinCooldown_ReturnsCooldown()
  22. {
  23. var store = new MemoryJobStore();
  24. SeedSuccess(store, "S2", 9, DateTime.Now.AddHours(-1));
  25. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  26. var before = store.Rows.Count;
  27. var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
  28. Assert.Equal(409, result.StatusCode);
  29. Assert.Equal("COOLDOWN", result.Body.Status);
  30. Assert.Equal(before, store.Rows.Count);
  31. }
  32. [Fact]
  33. public async Task AutoEnqueue_AfterCooldown_IsAccepted()
  34. {
  35. var store = new MemoryJobStore();
  36. SeedSuccess(store, "S2", 9, DateTime.Now.AddHours(-7));
  37. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  38. var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
  39. Assert.Equal(202, result.StatusCode);
  40. }
  41. [Fact]
  42. public async Task ManualEnqueue_IsNeverCooledDown()
  43. {
  44. var store = new MemoryJobStore();
  45. SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
  46. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  47. var result = await svc.EnqueueAsync("S2", 9, 1, null);
  48. Assert.Equal(202, result.StatusCode);
  49. }
  50. [Fact]
  51. public async Task ManualEnqueue_WithRequestedBy_IsNeverCooledDown()
  52. {
  53. var store = new MemoryJobStore();
  54. SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
  55. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  56. var result = await svc.EnqueueAsync("S2", 9, 1, 123, "AUTO");
  57. Assert.Equal(202, result.StatusCode);
  58. }
  59. [Fact]
  60. public async Task NightlyEnqueue_BypassesCooldown()
  61. {
  62. var store = new MemoryJobStore();
  63. SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
  64. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  65. var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY");
  66. Assert.Equal(202, result.StatusCode);
  67. }
  68. [Fact]
  69. public async Task BootstrapEnqueue_IsCooledDown()
  70. {
  71. var store = new MemoryJobStore();
  72. SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
  73. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  74. var result = await svc.EnqueueAsync("S2", 9, 1, null, "BOOTSTRAP");
  75. Assert.Equal(409, result.StatusCode);
  76. Assert.Equal("COOLDOWN", result.Body.Status);
  77. }
  78. [Fact]
  79. public async Task CooldownZero_DisablesWindow()
  80. {
  81. var store = new MemoryJobStore();
  82. SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
  83. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new ConfigurableCapability { AutoMinIntervalHours = 0 });
  84. var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
  85. Assert.Equal(202, result.StatusCode);
  86. }
  87. [Fact]
  88. public async Task RunningJob_StillBlocksBeforeCooldownCheck()
  89. {
  90. var store = new MemoryJobStore();
  91. SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
  92. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  93. await svc.EnqueueAsync("S2", 9, 1, null);
  94. var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
  95. Assert.Equal(409, result.StatusCode);
  96. Assert.Equal(ModuleRebuildStatus.Queued, result.Body.Status);
  97. }
  98. [Fact]
  99. public void S1Service_DelegatesToUnifiedPath()
  100. {
  101. var source = File.ReadAllText(RepoFile("server", "Plugins", "Admin.NET.Plugin.AiDOP", "DataPlatform", "S1Refresh", "S1DashboardRebuildService.cs"));
  102. var unified = source.IndexOf("if (UseUnified)", StringComparison.Ordinal);
  103. Assert.True(unified >= 0, "S1 统一路径分支丢失");
  104. var body = source[unified..source.IndexOf("var scope = S1MdpRunScope.Create", unified, StringComparison.Ordinal)];
  105. Assert.Contains("_moduleRebuild.EnqueueAsync(", body);
  106. }
  107. [Fact]
  108. public void NightlyJob_UsesNonCollidingCron()
  109. {
  110. var source = File.ReadAllText(RepoFile("server", "Plugins", "Admin.NET.Plugin.AiDOP", "Job", "MdpNightlyFullRebuildJob.cs"));
  111. Assert.Contains("Cron(\"40 3 * * *\"", source);
  112. Assert.DoesNotContain("\"15 3\"", source);
  113. Assert.DoesNotContain("\"0 2\"", source);
  114. }
  115. [Theory]
  116. [InlineData("AUTO", false)]
  117. [InlineData("auto", false)]
  118. [InlineData("MANUAL", true)]
  119. [InlineData("AUTO_NIGHTLY", true)]
  120. [InlineData("BOOTSTRAP", true)]
  121. [InlineData("", true)]
  122. public void ShouldPullFull_OnlyAutoIsIncremental(string triggerType, bool expected)
  123. => Assert.Equal(expected, ModuleRebuildTriggerType.ShouldPullFull(triggerType));
  124. private static string RepoFile(params string[] parts)
  125. {
  126. var dir = new DirectoryInfo(AppContext.BaseDirectory);
  127. while (dir != null && !File.Exists(Path.Combine(dir.FullName, "AGENTS.md")))
  128. dir = dir.Parent;
  129. Assert.NotNull(dir);
  130. return Path.Combine(new[] { dir!.FullName }.Concat(parts).ToArray());
  131. }
  132. [Fact]
  133. public async Task SecondEnqueue_SameScope_ReturnsConflict()
  134. {
  135. var store = new MemoryJobStore();
  136. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  137. var first = await svc.EnqueueAsync("S2", 9, 1, 1);
  138. var second = await svc.EnqueueAsync("S2", 9, 1, 1);
  139. Assert.Equal(202, first.StatusCode);
  140. Assert.Equal(409, second.StatusCode);
  141. Assert.Equal(first.Body.JobId, second.Body.JobId);
  142. }
  143. [Fact]
  144. public async Task OtherTenantEnqueue_NotBlocked()
  145. {
  146. var store = new MemoryJobStore();
  147. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  148. var a = await svc.EnqueueAsync("S2", 9, 1, 1);
  149. var b = await svc.EnqueueAsync("S2", 10, 1, 2);
  150. Assert.Equal(202, a.StatusCode);
  151. Assert.Equal(202, b.StatusCode);
  152. Assert.NotEqual(a.Body.JobId, b.Body.JobId);
  153. }
  154. [Fact]
  155. public async Task DifferentModule_SameTenant_NotBlocked()
  156. {
  157. var store = new MemoryJobStore();
  158. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  159. var a = await svc.EnqueueAsync("S2", 9, 1, 1);
  160. var b = await svc.EnqueueAsync("S3", 9, 1, 1);
  161. Assert.Equal(202, a.StatusCode);
  162. Assert.Equal(202, b.StatusCode);
  163. }
  164. [Fact]
  165. public async Task GetById_OtherTenant_NotFound()
  166. {
  167. var store = new MemoryJobStore();
  168. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  169. var created = await svc.EnqueueAsync("S2", 9, 1, 1);
  170. var row = await store.GetByIdAsync("S2", created.Body.JobId!.Value, 10, 1);
  171. Assert.Null(row);
  172. }
  173. [Fact]
  174. public void Scope_RejectsZeroTenant()
  175. {
  176. Assert.Throws<InvalidOperationException>(() => MdpRebuildScope.Create("S2", 0, 1));
  177. Assert.Throws<InvalidOperationException>(() => MdpRebuildScope.Create("S2", 9, 0));
  178. Assert.Throws<InvalidOperationException>(() => MdpRebuildScope.Create("S2", 9, 9));
  179. Assert.Throws<InvalidOperationException>(() => MdpRebuildScope.Create("S9", 9, 1));
  180. }
  181. [Fact]
  182. public void SqlInject_AddsTenantToEveryWhere()
  183. {
  184. var sql = """
  185. SELECT 1 FROM t WHERE a=1
  186. UNION ALL
  187. SELECT 1 FROM t WHERE b=2
  188. """;
  189. var injected = MdpSqlScope.InjectTenantFactory(sql);
  190. Assert.Equal(2, injected.Split("@TenantId", StringSplitOptions.None).Length - 1);
  191. }
  192. [Fact]
  193. public void TenantOnlySqlInject_DoesNotRequireFactoryColumn()
  194. {
  195. var injected = MdpSqlScope.InjectTenantOnly("SELECT 1 FROM mdp_std_delivery_schedule ds WHERE ds.tenant_id > 0");
  196. Assert.Contains("ds.tenant_id=@TenantId", injected);
  197. Assert.DoesNotContain("factory_id", injected, StringComparison.OrdinalIgnoreCase);
  198. Assert.DoesNotContain("@FactoryId", injected);
  199. }
  200. [Fact]
  201. public void TenantOnlySqlInject_QualifiesRootAndNestedJoinScopes()
  202. {
  203. var sql = """
  204. SELECT po.id
  205. FROM mdp_std_purchase_order po
  206. LEFT JOIN (
  207. SELECT tenant_id, po_no
  208. FROM mdp_std_delivery_schedule
  209. WHERE po_no <> ''
  210. ) ds ON po.tenant_id=ds.tenant_id
  211. WHERE po.po_no <> ''
  212. """;
  213. var injected = MdpSqlScope.InjectTenantOnly(sql);
  214. Assert.Contains("FROM mdp_std_delivery_schedule\n WHERE tenant_id=@TenantId", injected);
  215. Assert.Contains("WHERE po.tenant_id=@TenantId", injected);
  216. }
  217. /// <summary>
  218. /// mdp_source_table_registry 是全局字典表,只有 (id, source_table, source_system, remark)。
  219. /// 注入器过去把语句里每个 WHERE 都改写,给它补出 r.tenant_id=@TenantId,
  220. /// S3 的 mdp_std_supplier_item 因此每轮都报 Unknown column 'r.tenant_id' 并整轮 FAILED;
  221. /// 同样的子查询在 S4 只是碰巧被旧的 240 字符窗口盖住才没炸。
  222. /// </summary>
  223. [Fact]
  224. public void SqlInject_SkipsGlobalRegistryLookup()
  225. {
  226. var sql = """
  227. INSERT INTO mdp_std_supplier_item (source_system, item_code)
  228. SELECT COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = source_table LIMIT 1), '')
  229. FROM mdp_stg_source_list
  230. WHERE source_table='srm_purchase'
  231. """;
  232. var injected = MdpSqlScope.InjectTenantFactory(sql);
  233. Assert.Contains("WHERE r.source_table = source_table", injected);
  234. Assert.DoesNotContain("r.tenant_id", injected);
  235. Assert.DoesNotContain("r.factory_id", injected);
  236. Assert.Contains("WHERE tenant_id=@TenantId", injected);
  237. }
  238. /// <summary>
  239. /// 子查询自带的租户条件不能替外层背书。旧的 240 字符窗口会把它算进来,
  240. /// 外层于是整条漏掉租户过滤 —— 跨租户串数据,且不会报错。
  241. /// </summary>
  242. [Fact]
  243. public void SqlInject_IgnoresTenantParamInNestedSubquery()
  244. {
  245. var sql = """
  246. SELECT 1 FROM dwd_ship_trans d
  247. WHERE d.biz_date = (SELECT MAX(biz_date) FROM dwd_ship_trans WHERE tenant_id=@TenantId)
  248. """;
  249. var injected = MdpSqlScope.InjectTenantFactory(sql);
  250. Assert.Contains("WHERE d.tenant_id=@TenantId", injected);
  251. }
  252. /// <summary>
  253. /// 反过来:手写的租户条件即使排在 240 字符之外,也必须被认出来,不能重复注入。
  254. /// </summary>
  255. [Fact]
  256. public void SqlInject_HonorsTenantPredicateBeyondOldWindow()
  257. {
  258. var filler = string.Join("\n", Enumerable.Range(0, 12)
  259. .Select(i => $" AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.f{i}')),'') <> ''"));
  260. var sql = $"SELECT 1 FROM mdp_std_so s\nWHERE s.order_no <> ''\n{filler}\n AND s.tenant_id=@TenantId";
  261. var injected = MdpSqlScope.InjectTenantFactory(sql);
  262. Assert.Equal(sql, injected);
  263. }
  264. [Fact]
  265. public void S3ScopeDiscovery_ReadsNormalizedFactoryColumn()
  266. {
  267. var sql = MdpRebuildScopeCatalog.GetFactScopeSql("S3");
  268. Assert.NotNull(sql);
  269. Assert.Contains("COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId", sql);
  270. }
  271. [Fact]
  272. public void S3DefaultInbound_ExcludesArchivedInvMaster()
  273. {
  274. Assert.DoesNotContain(S3MdpEntityConfig.All, x =>
  275. x.EntityCode == "S3_INVENTORY" ||
  276. x.SourceTable.Equals("InvMaster", StringComparison.OrdinalIgnoreCase));
  277. }
  278. [Fact]
  279. public async Task Success_WritesHundredPercent()
  280. {
  281. var store = new MemoryJobStore();
  282. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  283. await svc.EnqueueAsync("S2", 9, 1, 1);
  284. var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
  285. var handler = new FakeHandler();
  286. await svc.RunClaimedAsync(job, handler, new InMemoryModuleRebuildLock(), CancellationToken.None);
  287. var row = store.Latest("S2", 9, 1);
  288. Assert.Equal(ModuleRebuildStatus.Success, row.Status);
  289. Assert.Equal(100, row.ProgressPercent);
  290. }
  291. [Fact]
  292. public async Task AlreadyRunning_RequeuesJob()
  293. {
  294. var store = new MemoryJobStore();
  295. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  296. await svc.EnqueueAsync("S2", 9, 1, 1);
  297. var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
  298. var lck = new InMemoryModuleRebuildLock();
  299. await using var held = await lck.TryAcquireAsync(MdpRebuildScope.Create("S2", 9, 1), "holder");
  300. await svc.RunClaimedAsync(job, new FakeHandler(), lck, CancellationToken.None);
  301. Assert.Equal(ModuleRebuildStatus.Queued, store.Latest("S2", 9, 1).Status);
  302. }
  303. [Fact]
  304. public async Task Enqueue_InitialQueuedZeroPercent()
  305. {
  306. var store = new MemoryJobStore();
  307. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  308. var accepted = await svc.EnqueueAsync("S4", 9, 1, 1, "AUTO");
  309. var row = store.Latest("S4", 9, 1);
  310. Assert.Equal(202, accepted.StatusCode);
  311. Assert.Equal(ModuleRebuildStatus.Queued, row.Status);
  312. Assert.Equal(ModuleRebuildStages.Queued, row.CurrentStage);
  313. Assert.Equal(0, row.ProgressPercent);
  314. Assert.Equal(0, row.StageIndex);
  315. Assert.Equal(4, row.StageTotal);
  316. Assert.Equal("AUTO", row.TriggerType);
  317. }
  318. [Fact]
  319. public async Task ConcurrentEnqueue_SameScope_OnlyOneActive()
  320. {
  321. var store = new MemoryJobStore();
  322. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  323. var results = await Task.WhenAll(Enumerable.Range(0, 20).Select(_ => svc.EnqueueAsync("S2", 9, 1, 1)));
  324. Assert.Equal(1, results.Count(x => x.StatusCode == 202));
  325. Assert.Equal(19, results.Count(x => x.StatusCode == 409));
  326. Assert.Single(store.Rows.Where(x => x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running));
  327. Assert.True(results.Where(x => x.StatusCode == 409).All(x => x.Body.JobId == results.First(y => y.StatusCode == 202).Body.JobId));
  328. }
  329. [Fact]
  330. public async Task DifferentFactory_NotBlocked()
  331. {
  332. var store = new MemoryJobStore();
  333. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  334. var a = await svc.EnqueueAsync("S2", 9, 1, 1);
  335. var b = await svc.EnqueueAsync("S2", 9, 2, 1);
  336. Assert.Equal(202, a.StatusCode);
  337. Assert.Equal(202, b.StatusCode);
  338. Assert.NotEqual(a.Body.JobId, b.Body.JobId);
  339. }
  340. [Fact]
  341. public async Task DisabledModule_Returns404()
  342. {
  343. var store = new MemoryJobStore();
  344. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new DisabledModuleRebuildCapability());
  345. var result = await svc.EnqueueAsync("S2", 9, 1, 1);
  346. Assert.Equal(404, result.StatusCode);
  347. Assert.Empty(store.Rows);
  348. }
  349. [Fact]
  350. public async Task Failed_KeepsStageAndPercent()
  351. {
  352. var store = new MemoryJobStore();
  353. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  354. await svc.EnqueueAsync("S2", 9, 1, 1);
  355. var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
  356. await svc.RunClaimedAsync(job, new ThrowingHandler(), new InMemoryModuleRebuildLock(), CancellationToken.None);
  357. var row = store.Latest("S2", 9, 1);
  358. Assert.Equal(ModuleRebuildStatus.Failed, row.Status);
  359. Assert.Equal(ModuleRebuildStages.Kpi, row.FailedStage);
  360. Assert.Equal(65, row.ProgressPercent);
  361. Assert.NotEqual(100, row.ProgressPercent);
  362. }
  363. [Fact]
  364. public async Task Cancellation_MarksCancelled()
  365. {
  366. var store = new MemoryJobStore();
  367. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  368. await svc.EnqueueAsync("S2", 9, 1, 1);
  369. var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
  370. using var cts = new CancellationTokenSource();
  371. var run = svc.RunClaimedAsync(job, new HangingHandler(), new InMemoryModuleRebuildLock(), cts.Token);
  372. await Task.Delay(200);
  373. cts.Cancel();
  374. await run;
  375. Assert.Equal(ModuleRebuildStatus.Cancelled, store.Latest("S2", 9, 1).Status);
  376. }
  377. [Fact]
  378. public async Task FailStale_ClosesRunningWithoutHeartbeat()
  379. {
  380. var store = new MemoryJobStore();
  381. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  382. await svc.EnqueueAsync("S2", 9, 1, 1);
  383. var job = store.Latest("S2", 9, 1);
  384. job.Status = ModuleRebuildStatus.Running;
  385. job.HeartbeatAt = DateTime.Now.AddHours(-1);
  386. await svc.FailStaleAsync();
  387. Assert.Equal(ModuleRebuildStatus.Failed, store.Latest("S2", 9, 1).Status);
  388. Assert.Contains("服务中断", store.Latest("S2", 9, 1).ErrorMessage);
  389. }
  390. [Fact]
  391. public async Task Progress_DoesNotRegress()
  392. {
  393. var store = new MemoryJobStore();
  394. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  395. await svc.EnqueueAsync("S2", 9, 1, 1);
  396. var job = store.Latest("S2", 9, 1);
  397. await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 20, "stg"));
  398. await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 40, "std"));
  399. await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 60, "dwd"));
  400. var percents = store.ProgressHistory.Select(x => x.Percent).ToArray();
  401. for (var i = 1; i < percents.Length; i++)
  402. Assert.True(percents[i] >= percents[i - 1], $"percent regress {percents[i - 1]} -> {percents[i]}");
  403. Assert.Equal(60, job.ProgressPercent);
  404. }
  405. [Fact]
  406. public void StageIndex_IsMonotonicPerModule()
  407. {
  408. Assert.True(ModuleRebuildStages.ToStageIndex("S2", ModuleRebuildStages.Staging)
  409. < ModuleRebuildStages.ToStageIndex("S2", ModuleRebuildStages.Standard));
  410. Assert.True(ModuleRebuildStages.ToStageIndex("S4", ModuleRebuildStages.Dwd)
  411. < ModuleRebuildStages.ToStageIndex("S4", ModuleRebuildStages.Kpi));
  412. Assert.True(ModuleRebuildStages.ToStageIndex("S5", ModuleRebuildStages.T8Inbound)
  413. < ModuleRebuildStages.ToStageIndex("S5", ModuleRebuildStages.KpiCalculating));
  414. Assert.Equal(5, ModuleRebuildStages.StageTotal("S1"));
  415. Assert.Equal(4, ModuleRebuildStages.StageTotal("S4"));
  416. Assert.Equal(4, ModuleRebuildStages.StageTotal("S7"));
  417. }
  418. [Fact]
  419. public void LockName_IsolatesTenantFactoryModule()
  420. {
  421. var a = MdpRebuildScope.Create("S2", 9, 1);
  422. var b = MdpRebuildScope.Create("S2", 10, 1);
  423. var c = MdpRebuildScope.Create("S3", 9, 1);
  424. Assert.Equal("S2_MDP_FULL:9:1", a.LockName);
  425. Assert.NotEqual(a.LockName, b.LockName);
  426. Assert.NotEqual(a.LockName, c.LockName);
  427. }
  428. [Fact]
  429. public void SqlInject_DoesNotDoubleInject()
  430. {
  431. var sql = "SELECT 1 FROM t w WHERE w.tenant_id=@TenantId AND a=1";
  432. var injected = MdpSqlScope.InjectTenantFactory(sql);
  433. Assert.Equal(1, injected.Split("@TenantId", StringSplitOptions.None).Length - 1);
  434. }
  435. [Fact]
  436. public async Task JobRunner_EnqueuesEachScope()
  437. {
  438. var store = new MemoryJobStore();
  439. var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
  440. var a = await svc.EnqueueAsync("S2", 9, 1, null, "BOOTSTRAP");
  441. var b = await svc.EnqueueAsync("S2", 10, 1, null, "BOOTSTRAP");
  442. Assert.Equal(202, a.StatusCode);
  443. Assert.Equal(202, b.StatusCode);
  444. Assert.Equal(2, store.Rows.Count);
  445. Assert.All(store.Rows, x => Assert.Equal("BOOTSTRAP", x.TriggerType));
  446. }
  447. private sealed class FakeHandler : IModuleRebuildHandler
  448. {
  449. public string ModuleCode => "S2";
  450. public Task<ModuleRebuildResult> RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func<ModuleProgressUpdate, Task> report, CancellationToken cancellationToken) =>
  451. Task.FromResult(new ModuleRebuildResult { BatchId = "b1", StageRows = 1 });
  452. }
  453. private sealed class ThrowingHandler : IModuleRebuildHandler
  454. {
  455. public string ModuleCode => "S2";
  456. public async Task<ModuleRebuildResult> RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func<ModuleProgressUpdate, Task> report, CancellationToken cancellationToken)
  457. {
  458. await report(new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 65, "正在重算 KPI"));
  459. throw new InvalidOperationException("kpi failed");
  460. }
  461. }
  462. private sealed class HangingHandler : IModuleRebuildHandler
  463. {
  464. public string ModuleCode => "S2";
  465. public async Task<ModuleRebuildResult> RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func<ModuleProgressUpdate, Task> report, CancellationToken cancellationToken)
  466. {
  467. await Task.Delay(Timeout.Infinite, cancellationToken);
  468. return new ModuleRebuildResult();
  469. }
  470. }
  471. private sealed class DisabledModuleRebuildCapability : IModuleRebuildCapability
  472. {
  473. public bool IsEnabled(string moduleCode) => false;
  474. public int MaxParallelScopes => 2;
  475. public int GlobalMaxParallelScopes => 2;
  476. public int AutoMinIntervalHours => 6;
  477. }
  478. private sealed class ConfigurableCapability : IModuleRebuildCapability
  479. {
  480. public int AutoMinIntervalHours { get; init; } = 6;
  481. public bool IsEnabled(string moduleCode) => true;
  482. public int MaxParallelScopes => 2;
  483. public int GlobalMaxParallelScopes => 2;
  484. }
  485. private sealed class MemoryJobStore : IModuleRebuildJobStore
  486. {
  487. public List<AdoModuleDashboardRebuildJob> Rows { get; } = new();
  488. public List<(long JobId, int Percent, string Stage)> ProgressHistory { get; } = new();
  489. private long _id = 1;
  490. private readonly object _gate = new();
  491. public Task<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(IReadOnlyCollection<string> enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default)
  492. => Task.FromResult<AdoModuleDashboardRebuildJob>(null!);
  493. public Task<bool> RequestCancelAsync(long jobId, CancellationToken ct = default)
  494. {
  495. lock (_gate)
  496. {
  497. var row = Rows.FirstOrDefault(x => x.Id == jobId
  498. && x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running);
  499. if (row == null) return Task.FromResult(false);
  500. row.CancelRequestedFlag = true;
  501. return Task.FromResult(true);
  502. }
  503. }
  504. public Task<bool> IsCancelRequestedAsync(long jobId, CancellationToken ct = default)
  505. {
  506. lock (_gate)
  507. return Task.FromResult(Rows.FirstOrDefault(x => x.Id == jobId)?.CancelRequestedFlag ?? false);
  508. }
  509. public Task<long> FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter)
  510. => Task.FromResult(0L);
  511. public Task<AdoModuleDashboardRebuildJob> FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
  512. {
  513. lock (_gate)
  514. return Task.FromResult(Rows
  515. .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
  516. && x.Status == ModuleRebuildStatus.Success)
  517. .OrderByDescending(x => x.FinishedAt)
  518. .FirstOrDefault());
  519. }
  520. public AdoModuleDashboardRebuildJob Latest(string module, long tenantId, long factoryId)
  521. {
  522. lock (_gate)
  523. return Rows.Where(x => x.ModuleCode == module && x.TenantId == tenantId && x.FactoryId == factoryId).OrderByDescending(x => x.Id).First();
  524. }
  525. public Task<AdoModuleDashboardRebuildJob> InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default)
  526. {
  527. lock (_gate)
  528. {
  529. row.Id = _id++;
  530. Rows.Add(row);
  531. return Task.FromResult(row);
  532. }
  533. }
  534. public Task<AdoModuleDashboardRebuildJob> FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
  535. {
  536. lock (_gate)
  537. return Task.FromResult(Rows.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
  538. && x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running)
  539. .OrderBy(x => x.Id).FirstOrDefault());
  540. }
  541. public Task<AdoModuleDashboardRebuildJob> GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default)
  542. {
  543. lock (_gate)
  544. return Task.FromResult(Rows.FirstOrDefault(x => x.Id == id && x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId));
  545. }
  546. public Task<AdoModuleDashboardRebuildJob> GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
  547. {
  548. lock (_gate)
  549. return Task.FromResult(Rows.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId)
  550. .OrderByDescending(x => x.Id).FirstOrDefault());
  551. }
  552. public Task<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(IReadOnlyCollection<string> enabledModules, CancellationToken ct = default)
  553. {
  554. lock (_gate)
  555. {
  556. var row = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued && enabledModules.Contains(x.ModuleCode))
  557. .Where(x => !Rows.Any(y => y.ModuleCode == x.ModuleCode && y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == ModuleRebuildStatus.Running))
  558. .OrderBy(x => x.Id).FirstOrDefault();
  559. if (row == null) return Task.FromResult<AdoModuleDashboardRebuildJob>(null);
  560. row.Status = ModuleRebuildStatus.Running;
  561. row.StartedAt = DateTime.Now;
  562. row.HeartbeatAt = DateTime.Now;
  563. return Task.FromResult(row);
  564. }
  565. }
  566. public Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) => Task.CompletedTask;
  567. public Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default)
  568. {
  569. lock (_gate)
  570. {
  571. var row = Rows.FirstOrDefault(x => x.Id == jobId);
  572. if (row == null) return Task.CompletedTask;
  573. row.CurrentStage = currentStage;
  574. row.StageIndex = stageIndex;
  575. row.ProgressPercent = progressPercent;
  576. row.ProgressMessage = message;
  577. row.LastProgressAt = now;
  578. row.HeartbeatAt = now;
  579. row.UpdateTime = now;
  580. ProgressHistory.Add((jobId, progressPercent, currentStage));
  581. return Task.CompletedTask;
  582. }
  583. }
  584. public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default) =>
  585. UpdateProgressAsync(jobId, completedStage, 0, progressPercent, nextMessage, now, ct);
  586. public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default)
  587. {
  588. lock (_gate)
  589. {
  590. var row = Rows.FirstOrDefault(x => x.Id == jobId);
  591. if (row != null) row.HeartbeatAt = now;
  592. return Task.CompletedTask;
  593. }
  594. }
  595. public Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default)
  596. {
  597. lock (_gate)
  598. {
  599. var cutoff = DateTime.Now - staleAfter;
  600. foreach (var row in Rows.Where(x => x.Status == ModuleRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff)))
  601. {
  602. row.Status = ModuleRebuildStatus.Failed;
  603. row.CurrentStage = ModuleRebuildStages.Failed;
  604. row.FinishedAt = DateTime.Now;
  605. row.LastProgressAt = DateTime.Now;
  606. row.UpdateTime = DateTime.Now;
  607. row.ErrorMessage = "服务中断,任务未正常结束";
  608. }
  609. return Task.CompletedTask;
  610. }
  611. }
  612. }
  613. }