| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329 |
- 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<InvalidOperationException>(() => S1MdpRunScope.Create(0, 1));
- Assert.Throws<InvalidOperationException>(() => 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<AdoS1DashboardRebuildJob> 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<AdoS1DashboardRebuildJob> InsertQueuedAsync(AdoS1DashboardRebuildJob row, CancellationToken ct = default)
- {
- lock (_gate)
- {
- row.Id = _id++;
- Rows.Add(row);
- return Task.FromResult(row);
- }
- }
- public Task<AdoS1DashboardRebuildJob> 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<AdoS1DashboardRebuildJob> 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<AdoS1DashboardRebuildJob> 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<AdoS1DashboardRebuildJob> 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<AdoS1DashboardRebuildJob>(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;
- }
- }
- }
|