using Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh; using Admin.NET.Plugin.AiDOP.Entity.SmartOps; using Admin.NET.Plugin.AiDOP.Order; using Xunit; namespace Admin.NET.Plugin.AiDOP.Tests.DataPlatform; public class S1DashboardRebuildServiceTests { [Fact] public async Task SecondEnqueue_SameScope_ReturnsConflict() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); var first = await svc.EnqueueAsync(9, 1, 1); var second = await svc.EnqueueAsync(9, 1, 1); Assert.Equal(202, first.StatusCode); Assert.Equal(409, second.StatusCode); Assert.Equal(first.Body.JobId, second.Body.JobId); Assert.Equal(1, store.ActiveCount()); } [Fact] public async Task OtherTenantEnqueue_NotBlocked() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); var a = await svc.EnqueueAsync(9, 1, 1); var b = await svc.EnqueueAsync(10, 1, 2); Assert.Equal(202, a.StatusCode); Assert.Equal(202, b.StatusCode); Assert.NotEqual(a.Body.JobId, b.Body.JobId); Assert.Equal(2, store.ActiveCount()); } [Fact] public async Task ConcurrentEnqueue_AtMostOneActivePerScope() { var store = new MemoryJobStore(); var queue = new S1DashboardRebuildQueue(); var tasks = Enumerable.Range(0, 20).Select(_ => new S1DashboardRebuildService(store, queue).EnqueueAsync(9, 1, 1)); var results = await Task.WhenAll(tasks); Assert.Equal(1, results.Count(x => x.StatusCode == 202)); Assert.Equal(19, results.Count(x => x.StatusCode == 409)); Assert.Equal(1, store.ActiveCount(9, 1)); } [Fact] public async Task GetById_OtherTenant_NotFound() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); var created = await svc.EnqueueAsync(9, 1, 1); var row = await store.GetByIdAsync(created.Body.JobId!.Value, 10, 1); Assert.Null(row); } [Fact] public async Task Latest_IsScopeIsolated() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.EnqueueAsync(9, 1, 1); await svc.EnqueueAsync(10, 1, 2); var latestA = await svc.LatestAsync(9, 1); var latestB = await svc.LatestAsync(10, 1); Assert.Equal(9, store.Rows.First(x => x.Id == latestA.JobId).TenantId); Assert.Equal(10, store.Rows.First(x => x.Id == latestB.JobId).TenantId); } [Fact] public async Task NewJob_StartsQueuedAtZero() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.EnqueueAsync(9, 1, 1); var row = store.Latest(9, 1); Assert.Equal(S1DashboardRebuildStatus.Queued, row.Status); Assert.Equal(S1MdpRebuildStage.Queued, row.CurrentStage); Assert.Equal(0, row.ProgressPercent); Assert.Equal(5, row.StageTotal); } [Fact] public async Task AlreadyRunning_RequeuesJob() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.EnqueueAsync(9, 1, 1); await svc.RunNextAsync((_, _, _, _) => throw new S1MdpAlreadyRunningException(), CancellationToken.None); Assert.Equal(S1DashboardRebuildStatus.Queued, store.Latest(9, 1).Status); Assert.Equal(2, store.Latest(9, 1).ProgressPercent); } [Fact] public async Task WorkerTokenCancel_MarksCancelled() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.EnqueueAsync(9, 1, 1); using var cts = new CancellationTokenSource(); cts.Cancel(); await svc.RunNextAsync(async (_, _, _, ct) => { ct.ThrowIfCancellationRequested(); return new S1MdpSyncTransformResult(); }, cts.Token); Assert.Equal(S1DashboardRebuildStatus.Cancelled, store.Latest(9, 1).Status); } [Fact] public async Task Success_WritesHundredPercent() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.EnqueueAsync(9, 1, 1); await svc.RunNextAsync((_, _, _, _) => Task.FromResult(new S1MdpSyncTransformResult { BatchId = "b1", StageRows = 1 }), CancellationToken.None); var row = store.Latest(9, 1); Assert.Equal(S1DashboardRebuildStatus.Success, row.Status); Assert.Equal(100, row.ProgressPercent); Assert.Equal(S1MdpRebuildStage.Success, row.CurrentStage); Assert.Equal("b1", row.BatchId); } [Fact] public async Task Failed_KeepsFailedStage() { var store = new MemoryJobStore(); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.EnqueueAsync(9, 1, 1); await svc.RunNextAsync(async (_, _, report, _) => { await report(new S1MdpProgressUpdate(S1MdpRebuildStage.Staging, 1, 10, "正在拉取源数据")); throw new InvalidOperationException("stg boom"); }, CancellationToken.None); var row = store.Latest(9, 1); Assert.Equal(S1DashboardRebuildStatus.Failed, row.Status); Assert.Equal(S1MdpRebuildStage.Staging, row.FailedStage); Assert.True(row.ProgressPercent < 100); } [Fact] public async Task StaleRunning_MarkedFailed() { var store = new MemoryJobStore(); store.Rows.Add(new AdoS1DashboardRebuildJob { Id = 8, TenantId = 9, FactoryId = 1, Status = S1DashboardRebuildStatus.Running, HeartbeatAt = DateTime.Now.AddMinutes(-40), SubmittedAt = DateTime.Now.AddMinutes(-41), CreateTime = DateTime.Now, UpdateTime = DateTime.Now }); var svc = new S1DashboardRebuildService(store, new S1DashboardRebuildQueue()); await svc.FailStaleAsync(); Assert.Equal(S1DashboardRebuildStatus.Failed, store.Latest(9, 1).Status); } [Fact] public async Task InMemoryLock_AllowsDifferentScopes() { var lck = new InMemoryS1MdpFullRunLock(); var a = await lck.TryAcquireAsync("S1_MDP_FULL:9:1", "a"); var b = await lck.TryAcquireAsync("S1_MDP_FULL:10:1", "b"); var a2 = await lck.TryAcquireAsync("S1_MDP_FULL:9:1", "a2"); Assert.NotNull(a); Assert.NotNull(b); Assert.Null(a2); await a.DisposeAsync(); var a3 = await lck.TryAcquireAsync("S1_MDP_FULL:9:1", "a3"); Assert.NotNull(a3); await b.DisposeAsync(); await a3.DisposeAsync(); } [Fact] public void RunScope_RejectsZeroTenant() { Assert.Throws(() => S1MdpRunScope.Create(0, 1)); Assert.Throws(() => S1MdpRunScope.Create(9, 0)); } [Fact] public void ProgressPercent_NeverExceedsHundred() { Assert.Equal(100, 100); Assert.True(S1MdpRebuildStage.ToStageIndex(S1MdpRebuildStage.Success) <= 5); } private sealed class MemoryJobStore : IS1DashboardRebuildJobStore { public List Rows { get; } = new(); private long _id = 1; private readonly object _gate = new(); public int ActiveCount() { lock (_gate) return Rows.Count(x => x.Status is S1DashboardRebuildStatus.Queued or S1DashboardRebuildStatus.Running); } public int ActiveCount(long tenantId, long factoryId) { lock (_gate) return Rows.Count(x => x.TenantId == tenantId && x.FactoryId == factoryId && x.Status is S1DashboardRebuildStatus.Queued or S1DashboardRebuildStatus.Running); } public AdoS1DashboardRebuildJob Latest(long tenantId, long factoryId) { lock (_gate) return Rows.Where(x => x.TenantId == tenantId && x.FactoryId == factoryId).OrderByDescending(x => x.Id).First(); } public Task InsertQueuedAsync(AdoS1DashboardRebuildJob row, CancellationToken ct = default) { lock (_gate) { row.Id = _id++; Rows.Add(row); return Task.FromResult(row); } } public Task FindActiveAsync(long tenantId, long factoryId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.Where(x => x.TenantId == tenantId && x.FactoryId == factoryId && x.Status is S1DashboardRebuildStatus.Queued or S1DashboardRebuildStatus.Running) .OrderBy(x => x.Id).FirstOrDefault()); } public Task GetByIdAsync(long id, long tenantId, long factoryId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.FirstOrDefault(x => x.Id == id && x.TenantId == tenantId && x.FactoryId == factoryId)); } public Task GetLatestAsync(long tenantId, long factoryId, CancellationToken ct = default) { lock (_gate) return Task.FromResult(Rows.Where(x => x.TenantId == tenantId && x.FactoryId == factoryId) .OrderByDescending(x => x.Id).FirstOrDefault()); } public Task ClaimNextQueuedAsync(CancellationToken ct = default) { lock (_gate) { var row = Rows.Where(x => x.Status == S1DashboardRebuildStatus.Queued) .Where(x => !Rows.Any(y => y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == S1DashboardRebuildStatus.Running)) .OrderBy(x => x.Id).FirstOrDefault(); if (row == null) return Task.FromResult(null); row.Status = S1DashboardRebuildStatus.Running; row.StartedAt = DateTime.Now; row.HeartbeatAt = DateTime.Now; return Task.FromResult(row); } } public Task UpdateAsync(AdoS1DashboardRebuildJob 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; } return Task.CompletedTask; } public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, DateTime now, CancellationToken ct = default) { lock (_gate) { var row = Rows.FirstOrDefault(x => x.Id == jobId); if (row == null) return Task.CompletedTask; switch (completedStage) { case S1MdpRebuildStage.Staging: row.StageRows = rows; break; case S1MdpRebuildStage.Standard: row.StandardRows = rows; break; case S1MdpRebuildStage.Dwd: row.DwdRows = rows; break; case S1MdpRebuildStage.Kpi: row.KpiRows = rows; break; case S1MdpRebuildStage.Atomic: row.AtomicRows = rows; break; } row.ProgressPercent = progressPercent; row.ProgressMessage = nextMessage; row.LastProgressAt = now; } return Task.CompletedTask; } public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default) { lock (_gate) { var row = Rows.FirstOrDefault(x => x.Id == jobId && x.Status == S1DashboardRebuildStatus.Running); 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 == S1DashboardRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff))) { row.Status = S1DashboardRebuildStatus.Failed; row.ErrorMessage = "服务中断,任务未正常结束"; row.FinishedAt = DateTime.Now; } } return Task.CompletedTask; } } }