S1DashboardRebuildServiceTests.cs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh;
  2. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  3. using Admin.NET.Plugin.AiDOP.Order;
  4. using Xunit;
  5. namespace Admin.NET.Plugin.AiDOP.Tests.DataPlatform;
  6. public class S1DashboardRebuildServiceTests
  7. {
  8. [Fact]
  9. public async Task SecondEnqueue_SameScope_ReturnsConflict()
  10. {
  11. var store = new MemoryJobStore();
  12. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  13. var first = await svc.EnqueueAsync(9, 1, 1);
  14. var second = await svc.EnqueueAsync(9, 1, 1);
  15. Assert.Equal(202, first.StatusCode);
  16. Assert.Equal(409, second.StatusCode);
  17. Assert.Equal(first.Body.JobId, second.Body.JobId);
  18. Assert.Equal(1, store.ActiveCount());
  19. }
  20. [Fact]
  21. public async Task OtherTenantEnqueue_NotBlocked()
  22. {
  23. var store = new MemoryJobStore();
  24. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  25. var a = await svc.EnqueueAsync(9, 1, 1);
  26. var b = await svc.EnqueueAsync(10, 1, 2);
  27. Assert.Equal(202, a.StatusCode);
  28. Assert.Equal(202, b.StatusCode);
  29. Assert.NotEqual(a.Body.JobId, b.Body.JobId);
  30. Assert.Equal(2, store.ActiveCount());
  31. }
  32. [Fact]
  33. public async Task ConcurrentEnqueue_AtMostOneActivePerScope()
  34. {
  35. var store = new MemoryJobStore();
  36. var queue = new S1DashboardRebuildQueue();
  37. var tasks = Enumerable.Range(0, 20).Select(_ =>
  38. new S1DashboardRebuildService(store, queue).EnqueueAsync(9, 1, 1));
  39. var results = await Task.WhenAll(tasks);
  40. Assert.Equal(1, results.Count(x => x.StatusCode == 202));
  41. Assert.Equal(19, results.Count(x => x.StatusCode == 409));
  42. Assert.Equal(1, store.ActiveCount(9, 1));
  43. }
  44. [Fact]
  45. public async Task GetById_OtherTenant_NotFound()
  46. {
  47. var store = new MemoryJobStore();
  48. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  49. var created = await svc.EnqueueAsync(9, 1, 1);
  50. var row = await store.GetByIdAsync(created.Body.JobId!.Value, 10, 1);
  51. Assert.Null(row);
  52. }
  53. [Fact]
  54. public async Task Latest_IsScopeIsolated()
  55. {
  56. var store = new MemoryJobStore();
  57. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  58. await svc.EnqueueAsync(9, 1, 1);
  59. await svc.EnqueueAsync(10, 1, 2);
  60. var latestA = await svc.LatestAsync(9, 1);
  61. var latestB = await svc.LatestAsync(10, 1);
  62. Assert.Equal(9, store.Rows.First(x => x.Id == latestA.JobId).TenantId);
  63. Assert.Equal(10, store.Rows.First(x => x.Id == latestB.JobId).TenantId);
  64. }
  65. [Fact]
  66. public async Task NewJob_StartsQueuedAtZero()
  67. {
  68. var store = new MemoryJobStore();
  69. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  70. await svc.EnqueueAsync(9, 1, 1);
  71. var row = store.Latest(9, 1);
  72. Assert.Equal(S1DashboardRebuildStatus.Queued, row.Status);
  73. Assert.Equal(S1MdpRebuildStage.Queued, row.CurrentStage);
  74. Assert.Equal(0, row.ProgressPercent);
  75. Assert.Equal(5, row.StageTotal);
  76. }
  77. [Fact]
  78. public async Task AlreadyRunning_RequeuesJob()
  79. {
  80. var store = new MemoryJobStore();
  81. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  82. await svc.EnqueueAsync(9, 1, 1);
  83. await svc.RunNextAsync((_, _, _, _) => throw new S1MdpAlreadyRunningException(), CancellationToken.None);
  84. Assert.Equal(S1DashboardRebuildStatus.Queued, store.Latest(9, 1).Status);
  85. Assert.Equal(2, store.Latest(9, 1).ProgressPercent);
  86. }
  87. [Fact]
  88. public async Task WorkerTokenCancel_MarksCancelled()
  89. {
  90. var store = new MemoryJobStore();
  91. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  92. await svc.EnqueueAsync(9, 1, 1);
  93. using var cts = new CancellationTokenSource();
  94. cts.Cancel();
  95. await svc.RunNextAsync(async (_, _, _, ct) =>
  96. {
  97. ct.ThrowIfCancellationRequested();
  98. return new S1MdpSyncTransformResult();
  99. }, cts.Token);
  100. Assert.Equal(S1DashboardRebuildStatus.Cancelled, store.Latest(9, 1).Status);
  101. }
  102. [Fact]
  103. public async Task Success_WritesHundredPercent()
  104. {
  105. var store = new MemoryJobStore();
  106. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  107. await svc.EnqueueAsync(9, 1, 1);
  108. await svc.RunNextAsync((_, _, _, _) => Task.FromResult(new S1MdpSyncTransformResult { BatchId = "b1", StageRows = 1 }), CancellationToken.None);
  109. var row = store.Latest(9, 1);
  110. Assert.Equal(S1DashboardRebuildStatus.Success, row.Status);
  111. Assert.Equal(100, row.ProgressPercent);
  112. Assert.Equal(S1MdpRebuildStage.Success, row.CurrentStage);
  113. Assert.Equal("b1", row.BatchId);
  114. }
  115. [Fact]
  116. public async Task Failed_KeepsFailedStage()
  117. {
  118. var store = new MemoryJobStore();
  119. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  120. await svc.EnqueueAsync(9, 1, 1);
  121. await svc.RunNextAsync(async (_, _, report, _) =>
  122. {
  123. await report(new S1MdpProgressUpdate(S1MdpRebuildStage.Staging, 1, 10, "正在拉取源数据"));
  124. throw new InvalidOperationException("stg boom");
  125. }, CancellationToken.None);
  126. var row = store.Latest(9, 1);
  127. Assert.Equal(S1DashboardRebuildStatus.Failed, row.Status);
  128. Assert.Equal(S1MdpRebuildStage.Staging, row.FailedStage);
  129. Assert.True(row.ProgressPercent < 100);
  130. }
  131. [Fact]
  132. public async Task StaleRunning_MarkedFailed()
  133. {
  134. var store = new MemoryJobStore();
  135. store.Rows.Add(new AdoS1DashboardRebuildJob
  136. {
  137. Id = 8,
  138. TenantId = 9,
  139. FactoryId = 1,
  140. Status = S1DashboardRebuildStatus.Running,
  141. HeartbeatAt = DateTime.Now.AddMinutes(-40),
  142. SubmittedAt = DateTime.Now.AddMinutes(-41),
  143. CreateTime = DateTime.Now,
  144. UpdateTime = DateTime.Now
  145. });
  146. var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue());
  147. await svc.FailStaleAsync();
  148. Assert.Equal(S1DashboardRebuildStatus.Failed, store.Latest(9, 1).Status);
  149. }
  150. [Fact]
  151. public async Task InMemoryLock_AllowsDifferentScopes()
  152. {
  153. var lck = new InMemoryS1MdpFullRunLock();
  154. var a = await lck.TryAcquireAsync("S1_MDP_FULL:9:1", "a");
  155. var b = await lck.TryAcquireAsync("S1_MDP_FULL:10:1", "b");
  156. var a2 = await lck.TryAcquireAsync("S1_MDP_FULL:9:1", "a2");
  157. Assert.NotNull(a);
  158. Assert.NotNull(b);
  159. Assert.Null(a2);
  160. await a.DisposeAsync();
  161. var a3 = await lck.TryAcquireAsync("S1_MDP_FULL:9:1", "a3");
  162. Assert.NotNull(a3);
  163. await b.DisposeAsync();
  164. await a3.DisposeAsync();
  165. }
  166. [Fact]
  167. public void RunScope_RejectsZeroTenant()
  168. {
  169. Assert.Throws<InvalidOperationException>(() => S1MdpRunScope.Create(0, 1));
  170. Assert.Throws<InvalidOperationException>(() => S1MdpRunScope.Create(9, 0));
  171. }
  172. [Fact]
  173. public void ProgressPercent_NeverExceedsHundred()
  174. {
  175. Assert.Equal(100, 100);
  176. Assert.True(S1MdpRebuildStage.ToStageIndex(S1MdpRebuildStage.Success) <= 5);
  177. }
  178. private sealed class MemoryJobStore : IS1DashboardRebuildJobStore
  179. {
  180. public List<AdoS1DashboardRebuildJob> Rows { get; } = new();
  181. private long _id = 1;
  182. private readonly object _gate = new();
  183. public int ActiveCount()
  184. {
  185. lock (_gate)
  186. return Rows.Count(x => x.Status is S1DashboardRebuildStatus.Queued or S1DashboardRebuildStatus.Running);
  187. }
  188. public int ActiveCount(long tenantId, long factoryId)
  189. {
  190. lock (_gate)
  191. return Rows.Count(x => x.TenantId == tenantId && x.FactoryId == factoryId
  192. && x.Status is S1DashboardRebuildStatus.Queued or S1DashboardRebuildStatus.Running);
  193. }
  194. public AdoS1DashboardRebuildJob Latest(long tenantId, long factoryId)
  195. {
  196. lock (_gate)
  197. return Rows.Where(x => x.TenantId == tenantId && x.FactoryId == factoryId).OrderByDescending(x => x.Id).First();
  198. }
  199. public Task<AdoS1DashboardRebuildJob> InsertQueuedAsync(AdoS1DashboardRebuildJob row, CancellationToken ct = default)
  200. {
  201. lock (_gate)
  202. {
  203. row.Id = _id++;
  204. Rows.Add(row);
  205. return Task.FromResult(row);
  206. }
  207. }
  208. public Task<AdoS1DashboardRebuildJob> FindActiveAsync(long tenantId, long factoryId, CancellationToken ct = default)
  209. {
  210. lock (_gate)
  211. return Task.FromResult(Rows.Where(x => x.TenantId == tenantId && x.FactoryId == factoryId
  212. && x.Status is S1DashboardRebuildStatus.Queued or S1DashboardRebuildStatus.Running)
  213. .OrderBy(x => x.Id).FirstOrDefault());
  214. }
  215. public Task<AdoS1DashboardRebuildJob> GetByIdAsync(long id, long tenantId, long factoryId, CancellationToken ct = default)
  216. {
  217. lock (_gate)
  218. return Task.FromResult(Rows.FirstOrDefault(x => x.Id == id && x.TenantId == tenantId && x.FactoryId == factoryId));
  219. }
  220. public Task<AdoS1DashboardRebuildJob> GetLatestAsync(long tenantId, long factoryId, CancellationToken ct = default)
  221. {
  222. lock (_gate)
  223. return Task.FromResult(Rows.Where(x => x.TenantId == tenantId && x.FactoryId == factoryId)
  224. .OrderByDescending(x => x.Id).FirstOrDefault());
  225. }
  226. public Task<AdoS1DashboardRebuildJob> ClaimNextQueuedAsync(CancellationToken ct = default)
  227. {
  228. lock (_gate)
  229. {
  230. var row = Rows.Where(x => x.Status == S1DashboardRebuildStatus.Queued)
  231. .Where(x => !Rows.Any(y => y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == S1DashboardRebuildStatus.Running))
  232. .OrderBy(x => x.Id).FirstOrDefault();
  233. if (row == null) return Task.FromResult<AdoS1DashboardRebuildJob>(null);
  234. row.Status = S1DashboardRebuildStatus.Running;
  235. row.StartedAt = DateTime.Now;
  236. row.HeartbeatAt = DateTime.Now;
  237. return Task.FromResult(row);
  238. }
  239. }
  240. public Task UpdateAsync(AdoS1DashboardRebuildJob row, CancellationToken ct = default) => Task.CompletedTask;
  241. public Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default)
  242. {
  243. lock (_gate)
  244. {
  245. var row = Rows.FirstOrDefault(x => x.Id == jobId);
  246. if (row == null) return Task.CompletedTask;
  247. row.CurrentStage = currentStage;
  248. row.StageIndex = stageIndex;
  249. row.ProgressPercent = progressPercent;
  250. row.ProgressMessage = message;
  251. row.LastProgressAt = now;
  252. row.HeartbeatAt = now;
  253. }
  254. return Task.CompletedTask;
  255. }
  256. public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, DateTime now, CancellationToken ct = default)
  257. {
  258. lock (_gate)
  259. {
  260. var row = Rows.FirstOrDefault(x => x.Id == jobId);
  261. if (row == null) return Task.CompletedTask;
  262. switch (completedStage)
  263. {
  264. case S1MdpRebuildStage.Staging: row.StageRows = rows; break;
  265. case S1MdpRebuildStage.Standard: row.StandardRows = rows; break;
  266. case S1MdpRebuildStage.Dwd: row.DwdRows = rows; break;
  267. case S1MdpRebuildStage.Kpi: row.KpiRows = rows; break;
  268. case S1MdpRebuildStage.Atomic: row.AtomicRows = rows; break;
  269. }
  270. row.ProgressPercent = progressPercent;
  271. row.ProgressMessage = nextMessage;
  272. row.LastProgressAt = now;
  273. }
  274. return Task.CompletedTask;
  275. }
  276. public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default)
  277. {
  278. lock (_gate)
  279. {
  280. var row = Rows.FirstOrDefault(x => x.Id == jobId && x.Status == S1DashboardRebuildStatus.Running);
  281. if (row != null) row.HeartbeatAt = now;
  282. }
  283. return Task.CompletedTask;
  284. }
  285. public Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default)
  286. {
  287. lock (_gate)
  288. {
  289. var cutoff = DateTime.Now - staleAfter;
  290. foreach (var row in Rows.Where(x => x.Status == S1DashboardRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff)))
  291. {
  292. row.Status = S1DashboardRebuildStatus.Failed;
  293. row.ErrorMessage = "服务中断,任务未正常结束";
  294. row.FinishedAt = DateTime.Now;
  295. }
  296. }
  297. return Task.CompletedTask;
  298. }
  299. }
  300. }