ModuleRebuildServiceTests.cs 33 KB

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