using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
using Admin.NET.Plugin.AiDOP.Supply;
using Xunit;
namespace Admin.NET.Plugin.AiDOP.Tests.DataPlatform;
public class ModuleRebuildServiceTests
{
private static void SeedSuccess(MemoryJobStore store, string module, long tenant, DateTime finishedAt)
{
store.Rows.Add(new AdoModuleDashboardRebuildJob
{
Id = 9000 + store.Rows.Count,
ModuleCode = module,
TenantId = tenant,
FactoryId = 1,
Status = ModuleRebuildStatus.Success,
FinishedAt = finishedAt
});
}
[Fact]
public async Task AutoEnqueue_WithinCooldown_ReturnsCooldown()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddHours(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var before = store.Rows.Count;
var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
Assert.Equal(409, result.StatusCode);
Assert.Equal("COOLDOWN", result.Body.Status);
Assert.Equal(before, store.Rows.Count);
}
[Fact]
public async Task AutoEnqueue_AfterCooldown_IsAccepted()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddHours(-7));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
Assert.Equal(202, result.StatusCode);
}
[Fact]
public async Task ManualEnqueue_IsNeverCooledDown()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var result = await svc.EnqueueAsync("S2", 9, 1, null);
Assert.Equal(202, result.StatusCode);
}
[Fact]
public async Task ManualEnqueue_WithRequestedBy_IsNeverCooledDown()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var result = await svc.EnqueueAsync("S2", 9, 1, 123, "AUTO");
Assert.Equal(202, result.StatusCode);
}
[Fact]
public async Task NightlyEnqueue_BypassesCooldown()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY");
Assert.Equal(202, result.StatusCode);
}
[Fact]
public async Task BootstrapEnqueue_IsCooledDown()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var result = await svc.EnqueueAsync("S2", 9, 1, null, "BOOTSTRAP");
Assert.Equal(409, result.StatusCode);
Assert.Equal("COOLDOWN", result.Body.Status);
}
[Fact]
public async Task CooldownZero_DisablesWindow()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new ConfigurableCapability { AutoMinIntervalHours = 0 });
var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
Assert.Equal(202, result.StatusCode);
}
[Fact]
public async Task RunningJob_StillBlocksBeforeCooldownCheck()
{
var store = new MemoryJobStore();
SeedSuccess(store, "S2", 9, DateTime.Now.AddMinutes(-1));
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, null);
var result = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
Assert.Equal(409, result.StatusCode);
Assert.Equal(ModuleRebuildStatus.Queued, result.Body.Status);
}
///
/// 夜间兜底任务在非执行机上永远领不到(ClaimNextQueuedAsync 的白名单),
/// 若它还能挡住人工入队,页面「数据重算」按钮就被永久卡死。人工请求必须能取代它。
///
[Fact]
public async Task ManualEnqueue_SupersedesQueuedAutoJob()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var nightly = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY");
var manual = await svc.EnqueueAsync("S2", 9, 1, 7);
Assert.Equal(202, manual.StatusCode);
Assert.NotEqual(nightly.Body.JobId, manual.Body.JobId);
var superseded = store.Rows.Single(x => x.Id == nightly.Body.JobId);
Assert.Equal(ModuleRebuildStatus.Cancelled, superseded.Status);
Assert.Contains("AUTO_NIGHTLY", superseded.ErrorMessage);
}
/// 已经开跑的任务不得被掀掉:取代只针对 QUEUED。
[Fact]
public async Task ManualEnqueue_DoesNotSupersedeRunningAutoJob()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY");
var running = await store.ClaimNextQueuedAsync(new[] { "S2" });
var manual = await svc.EnqueueAsync("S2", 9, 1, 7);
Assert.Equal(409, manual.StatusCode);
Assert.Equal(ModuleRebuildStatus.Running, store.Rows.Single(x => x.Id == running.Id).Status);
}
/// 取代资格只给人工请求;自动触发器之间仍按原样互斥,避免夜间批次自相残杀。
[Fact]
public async Task AutoEnqueue_DoesNotSupersedeQueuedAutoJob()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new ConfigurableCapability { AutoMinIntervalHours = 0 });
var nightly = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO_NIGHTLY");
var auto = await svc.EnqueueAsync("S2", 9, 1, null, "AUTO");
Assert.Equal(409, auto.StatusCode);
Assert.Equal(nightly.Body.JobId, auto.Body.JobId);
Assert.Equal(ModuleRebuildStatus.Queued, store.Rows.Single(x => x.Id == nightly.Body.JobId).Status);
}
[Fact]
public void S1Service_DelegatesToUnifiedPath()
{
var source = File.ReadAllText(RepoFile("server", "Plugins", "Admin.NET.Plugin.AiDOP", "DataPlatform", "S1Refresh", "S1DashboardRebuildService.cs"));
var unified = source.IndexOf("if (UseUnified)", StringComparison.Ordinal);
Assert.True(unified >= 0, "S1 统一路径分支丢失");
var body = source[unified..source.IndexOf("var scope = S1MdpRunScope.Create", unified, StringComparison.Ordinal)];
Assert.Contains("_moduleRebuild.EnqueueAsync(", body);
}
[Fact]
public void NightlyJob_UsesNonCollidingCron()
{
var source = File.ReadAllText(RepoFile("server", "Plugins", "Admin.NET.Plugin.AiDOP", "Job", "MdpNightlyFullRebuildJob.cs"));
Assert.Contains("Cron(\"40 3 * * *\"", source);
Assert.DoesNotContain("\"15 3\"", source);
Assert.DoesNotContain("\"0 2\"", source);
}
[Theory]
[InlineData("AUTO", false)]
[InlineData("auto", false)]
[InlineData("MANUAL", true)]
[InlineData("AUTO_NIGHTLY", true)]
[InlineData("BOOTSTRAP", true)]
[InlineData("", true)]
public void ShouldPullFull_OnlyAutoIsIncremental(string triggerType, bool expected)
=> Assert.Equal(expected, ModuleRebuildTriggerType.ShouldPullFull(triggerType));
private static string RepoFile(params string[] parts)
{
var dir = new DirectoryInfo(AppContext.BaseDirectory);
while (dir != null && !File.Exists(Path.Combine(dir.FullName, "AGENTS.md")))
dir = dir.Parent;
Assert.NotNull(dir);
return Path.Combine(new[] { dir!.FullName }.Concat(parts).ToArray());
}
[Fact]
public async Task SecondEnqueue_SameScope_ReturnsConflict()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var first = await svc.EnqueueAsync("S2", 9, 1, 1);
var second = await svc.EnqueueAsync("S2", 9, 1, 1);
Assert.Equal(202, first.StatusCode);
Assert.Equal(409, second.StatusCode);
Assert.Equal(first.Body.JobId, second.Body.JobId);
}
[Fact]
public async Task OtherTenantEnqueue_NotBlocked()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var a = await svc.EnqueueAsync("S2", 9, 1, 1);
var b = await svc.EnqueueAsync("S2", 10, 1, 2);
Assert.Equal(202, a.StatusCode);
Assert.Equal(202, b.StatusCode);
Assert.NotEqual(a.Body.JobId, b.Body.JobId);
}
[Fact]
public async Task DifferentModule_SameTenant_NotBlocked()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var a = await svc.EnqueueAsync("S2", 9, 1, 1);
var b = await svc.EnqueueAsync("S3", 9, 1, 1);
Assert.Equal(202, a.StatusCode);
Assert.Equal(202, b.StatusCode);
}
[Fact]
public async Task GetById_OtherTenant_NotFound()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var created = await svc.EnqueueAsync("S2", 9, 1, 1);
var row = await store.GetByIdAsync("S2", created.Body.JobId!.Value, 10, 1);
Assert.Null(row);
}
[Fact]
public void Scope_RejectsZeroTenant()
{
Assert.Throws(() => MdpRebuildScope.Create("S2", 0, 1));
Assert.Throws(() => MdpRebuildScope.Create("S2", 9, 0));
Assert.Throws(() => MdpRebuildScope.Create("S2", 9, 9));
Assert.Throws(() => MdpRebuildScope.Create("S9", 9, 1));
}
[Fact]
public void SqlInject_AddsTenantToEveryWhere()
{
var sql = """
SELECT 1 FROM t WHERE a=1
UNION ALL
SELECT 1 FROM t WHERE b=2
""";
var injected = MdpSqlScope.InjectTenantFactory(sql);
Assert.Equal(2, injected.Split("@TenantId", StringSplitOptions.None).Length - 1);
}
[Fact]
public void TenantOnlySqlInject_DoesNotRequireFactoryColumn()
{
var injected = MdpSqlScope.InjectTenantOnly("SELECT 1 FROM mdp_std_delivery_schedule ds WHERE ds.tenant_id > 0");
Assert.Contains("ds.tenant_id=@TenantId", injected);
Assert.DoesNotContain("factory_id", injected, StringComparison.OrdinalIgnoreCase);
Assert.DoesNotContain("@FactoryId", injected);
}
[Fact]
public void TenantOnlySqlInject_QualifiesRootAndNestedJoinScopes()
{
var sql = """
SELECT po.id
FROM mdp_std_purchase_order po
LEFT JOIN (
SELECT tenant_id, po_no
FROM mdp_std_delivery_schedule
WHERE po_no <> ''
) ds ON po.tenant_id=ds.tenant_id
WHERE po.po_no <> ''
""";
var injected = MdpSqlScope.InjectTenantOnly(sql);
Assert.Contains("FROM mdp_std_delivery_schedule\n WHERE tenant_id=@TenantId", injected);
Assert.Contains("WHERE po.tenant_id=@TenantId", injected);
}
///
/// mdp_source_table_registry 是全局字典表,只有 (id, source_table, source_system, remark)。
/// 注入器过去把语句里每个 WHERE 都改写,给它补出 r.tenant_id=@TenantId,
/// S3 的 mdp_std_supplier_item 因此每轮都报 Unknown column 'r.tenant_id' 并整轮 FAILED;
/// 同样的子查询在 S4 只是碰巧被旧的 240 字符窗口盖住才没炸。
///
[Fact]
public void SqlInject_SkipsGlobalRegistryLookup()
{
var sql = """
INSERT INTO mdp_std_supplier_item (source_system, item_code)
SELECT COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = source_table LIMIT 1), '')
FROM mdp_stg_source_list
WHERE source_table='srm_purchase'
""";
var injected = MdpSqlScope.InjectTenantFactory(sql);
Assert.Contains("WHERE r.source_table = source_table", injected);
Assert.DoesNotContain("r.tenant_id", injected);
Assert.DoesNotContain("r.factory_id", injected);
Assert.Contains("WHERE tenant_id=@TenantId", injected);
}
///
/// 子查询自带的租户条件不能替外层背书。旧的 240 字符窗口会把它算进来,
/// 外层于是整条漏掉租户过滤 —— 跨租户串数据,且不会报错。
///
[Fact]
public void SqlInject_IgnoresTenantParamInNestedSubquery()
{
var sql = """
SELECT 1 FROM dwd_ship_trans d
WHERE d.biz_date = (SELECT MAX(biz_date) FROM dwd_ship_trans WHERE tenant_id=@TenantId)
""";
var injected = MdpSqlScope.InjectTenantFactory(sql);
Assert.Contains("WHERE d.tenant_id=@TenantId", injected);
}
///
/// 反过来:手写的租户条件即使排在 240 字符之外,也必须被认出来,不能重复注入。
///
[Fact]
public void SqlInject_HonorsTenantPredicateBeyondOldWindow()
{
var filler = string.Join("\n", Enumerable.Range(0, 12)
.Select(i => $" AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.f{i}')),'') <> ''"));
var sql = $"SELECT 1 FROM mdp_std_so s\nWHERE s.order_no <> ''\n{filler}\n AND s.tenant_id=@TenantId";
var injected = MdpSqlScope.InjectTenantFactory(sql);
Assert.Equal(sql, injected);
}
[Fact]
public void S3ScopeDiscovery_ReadsNormalizedFactoryColumn()
{
var sql = MdpRebuildScopeCatalog.GetFactScopeSql("S3");
Assert.NotNull(sql);
Assert.Contains("COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId", sql);
}
[Fact]
public void S3DefaultInbound_ExcludesArchivedInvMaster()
{
Assert.DoesNotContain(S3MdpEntityConfig.All, x =>
x.EntityCode == "S3_INVENTORY" ||
x.SourceTable.Equals("InvMaster", StringComparison.OrdinalIgnoreCase));
}
[Fact]
public async Task Success_WritesHundredPercent()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, 1);
var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
var handler = new FakeHandler();
await svc.RunClaimedAsync(job, handler, new InMemoryModuleRebuildLock(), CancellationToken.None);
var row = store.Latest("S2", 9, 1);
Assert.Equal(ModuleRebuildStatus.Success, row.Status);
Assert.Equal(100, row.ProgressPercent);
}
[Fact]
public async Task AlreadyRunning_RequeuesJob()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, 1);
var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
var lck = new InMemoryModuleRebuildLock();
await using var held = await lck.TryAcquireAsync(MdpRebuildScope.Create("S2", 9, 1), "holder");
await svc.RunClaimedAsync(job, new FakeHandler(), lck, CancellationToken.None);
Assert.Equal(ModuleRebuildStatus.Queued, store.Latest("S2", 9, 1).Status);
}
[Fact]
public async Task Enqueue_InitialQueuedZeroPercent()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var accepted = await svc.EnqueueAsync("S4", 9, 1, 1, "AUTO");
var row = store.Latest("S4", 9, 1);
Assert.Equal(202, accepted.StatusCode);
Assert.Equal(ModuleRebuildStatus.Queued, row.Status);
Assert.Equal(ModuleRebuildStages.Queued, row.CurrentStage);
Assert.Equal(0, row.ProgressPercent);
Assert.Equal(0, row.StageIndex);
Assert.Equal(4, row.StageTotal);
Assert.Equal("AUTO", row.TriggerType);
}
[Fact]
public async Task ConcurrentEnqueue_SameScope_OnlyOneActive()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var results = await Task.WhenAll(Enumerable.Range(0, 20).Select(_ => svc.EnqueueAsync("S2", 9, 1, 1)));
Assert.Equal(1, results.Count(x => x.StatusCode == 202));
Assert.Equal(19, results.Count(x => x.StatusCode == 409));
Assert.Single(store.Rows.Where(x => x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running));
Assert.True(results.Where(x => x.StatusCode == 409).All(x => x.Body.JobId == results.First(y => y.StatusCode == 202).Body.JobId));
}
[Fact]
public async Task DifferentFactory_NotBlocked()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var a = await svc.EnqueueAsync("S2", 9, 1, 1);
var b = await svc.EnqueueAsync("S2", 9, 2, 1);
Assert.Equal(202, a.StatusCode);
Assert.Equal(202, b.StatusCode);
Assert.NotEqual(a.Body.JobId, b.Body.JobId);
}
[Fact]
public async Task DisabledModule_Returns404()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new DisabledModuleRebuildCapability());
var result = await svc.EnqueueAsync("S2", 9, 1, 1);
Assert.Equal(404, result.StatusCode);
Assert.Empty(store.Rows);
}
[Fact]
public async Task Failed_KeepsStageAndPercent()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, 1);
var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
await svc.RunClaimedAsync(job, new ThrowingHandler(), new InMemoryModuleRebuildLock(), CancellationToken.None);
var row = store.Latest("S2", 9, 1);
Assert.Equal(ModuleRebuildStatus.Failed, row.Status);
Assert.Equal(ModuleRebuildStages.Kpi, row.FailedStage);
Assert.Equal(65, row.ProgressPercent);
Assert.NotEqual(100, row.ProgressPercent);
}
[Fact]
public async Task Cancellation_MarksCancelled()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, 1);
var job = await store.ClaimNextQueuedAsync(new[] { "S2" });
using var cts = new CancellationTokenSource();
var run = svc.RunClaimedAsync(job, new HangingHandler(), new InMemoryModuleRebuildLock(), cts.Token);
await Task.Delay(200);
cts.Cancel();
await run;
Assert.Equal(ModuleRebuildStatus.Cancelled, store.Latest("S2", 9, 1).Status);
}
[Fact]
public async Task FailStale_ClosesRunningWithoutHeartbeat()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, 1);
var job = store.Latest("S2", 9, 1);
job.Status = ModuleRebuildStatus.Running;
job.HeartbeatAt = DateTime.Now.AddHours(-1);
await svc.FailStaleAsync();
Assert.Equal(ModuleRebuildStatus.Failed, store.Latest("S2", 9, 1).Status);
Assert.Contains("服务中断", store.Latest("S2", 9, 1).ErrorMessage);
}
[Fact]
public async Task Progress_DoesNotRegress()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
await svc.EnqueueAsync("S2", 9, 1, 1);
var job = store.Latest("S2", 9, 1);
await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 20, "stg"));
await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 40, "std"));
await svc.ApplyProgressAsync(job.Id, "S2", new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 60, "dwd"));
var percents = store.ProgressHistory.Select(x => x.Percent).ToArray();
for (var i = 1; i < percents.Length; i++)
Assert.True(percents[i] >= percents[i - 1], $"percent regress {percents[i - 1]} -> {percents[i]}");
Assert.Equal(60, job.ProgressPercent);
}
[Fact]
public void StageIndex_IsMonotonicPerModule()
{
Assert.True(ModuleRebuildStages.ToStageIndex("S2", ModuleRebuildStages.Staging)
< ModuleRebuildStages.ToStageIndex("S2", ModuleRebuildStages.Standard));
Assert.True(ModuleRebuildStages.ToStageIndex("S4", ModuleRebuildStages.Dwd)
< ModuleRebuildStages.ToStageIndex("S4", ModuleRebuildStages.Kpi));
Assert.True(ModuleRebuildStages.ToStageIndex("S5", ModuleRebuildStages.T8Inbound)
< ModuleRebuildStages.ToStageIndex("S5", ModuleRebuildStages.KpiCalculating));
Assert.Equal(5, ModuleRebuildStages.StageTotal("S1"));
Assert.Equal(4, ModuleRebuildStages.StageTotal("S4"));
Assert.Equal(4, ModuleRebuildStages.StageTotal("S7"));
}
[Fact]
public void LockName_IsolatesTenantFactoryModule()
{
var a = MdpRebuildScope.Create("S2", 9, 1);
var b = MdpRebuildScope.Create("S2", 10, 1);
var c = MdpRebuildScope.Create("S3", 9, 1);
Assert.Equal("S2_MDP_FULL:9:1", a.LockName);
Assert.NotEqual(a.LockName, b.LockName);
Assert.NotEqual(a.LockName, c.LockName);
}
[Fact]
public void SqlInject_DoesNotDoubleInject()
{
var sql = "SELECT 1 FROM t w WHERE w.tenant_id=@TenantId AND a=1";
var injected = MdpSqlScope.InjectTenantFactory(sql);
Assert.Equal(1, injected.Split("@TenantId", StringSplitOptions.None).Length - 1);
}
[Fact]
public async Task JobRunner_EnqueuesEachScope()
{
var store = new MemoryJobStore();
var svc = new ModuleRebuildService(store, new ModuleRebuildQueue(), new AlwaysOnModuleRebuildCapability());
var a = await svc.EnqueueAsync("S2", 9, 1, null, "BOOTSTRAP");
var b = await svc.EnqueueAsync("S2", 10, 1, null, "BOOTSTRAP");
Assert.Equal(202, a.StatusCode);
Assert.Equal(202, b.StatusCode);
Assert.Equal(2, store.Rows.Count);
Assert.All(store.Rows, x => Assert.Equal("BOOTSTRAP", x.TriggerType));
}
private sealed class FakeHandler : IModuleRebuildHandler
{
public string ModuleCode => "S2";
public Task RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func report, CancellationToken cancellationToken) =>
Task.FromResult(new ModuleRebuildResult { BatchId = "b1", StageRows = 1 });
}
private sealed class ThrowingHandler : IModuleRebuildHandler
{
public string ModuleCode => "S2";
public async Task RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func report, CancellationToken cancellationToken)
{
await report(new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 65, "正在重算 KPI"));
throw new InvalidOperationException("kpi failed");
}
}
private sealed class HangingHandler : IModuleRebuildHandler
{
public string ModuleCode => "S2";
public async Task RunAsync(MdpRebuildScope scope, string triggerType, long jobId, Func report, CancellationToken cancellationToken)
{
await Task.Delay(Timeout.Infinite, cancellationToken);
return new ModuleRebuildResult();
}
}
private sealed class DisabledModuleRebuildCapability : IModuleRebuildCapability
{
public bool IsEnabled(string moduleCode) => false;
public int MaxParallelScopes => 2;
public int GlobalMaxParallelScopes => 2;
public int AutoMinIntervalHours => 6;
}
private sealed class ConfigurableCapability : IModuleRebuildCapability
{
public int AutoMinIntervalHours { get; init; } = 6;
public bool IsEnabled(string moduleCode) => true;
public int MaxParallelScopes => 2;
public int GlobalMaxParallelScopes => 2;
}
private sealed class MemoryJobStore : IModuleRebuildJobStore
{
public List Rows { get; } = new();
public List<(long JobId, int Percent, string Stage)> ProgressHistory { get; } = new();
private long _id = 1;
private readonly object _gate = new();
public Task ClaimNextQueuedAsync(IReadOnlyCollection enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default)
=> Task.FromResult(null!);
public Task RequestCancelAsync(long jobId, CancellationToken ct = default)
{
lock (_gate)
{
var row = Rows.FirstOrDefault(x => x.Id == jobId
&& x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running);
if (row == null) return Task.FromResult(false);
row.CancelRequestedFlag = true;
return Task.FromResult(true);
}
}
public Task IsCancelRequestedAsync(long jobId, CancellationToken ct = default)
{
lock (_gate)
return Task.FromResult(Rows.FirstOrDefault(x => x.Id == jobId)?.CancelRequestedFlag ?? false);
}
public Task FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter)
=> Task.FromResult(0L);
public Task FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
{
lock (_gate)
return Task.FromResult(Rows
.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
&& x.Status == ModuleRebuildStatus.Success)
.OrderByDescending(x => x.FinishedAt)
.FirstOrDefault());
}
public AdoModuleDashboardRebuildJob Latest(string module, long tenantId, long factoryId)
{
lock (_gate)
return Rows.Where(x => x.ModuleCode == module && x.TenantId == tenantId && x.FactoryId == factoryId).OrderByDescending(x => x.Id).First();
}
public Task InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default)
{
lock (_gate)
{
row.Id = _id++;
Rows.Add(row);
return Task.FromResult(row);
}
}
public Task FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
{
lock (_gate)
return Task.FromResult(Rows.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
&& x.Status is ModuleRebuildStatus.Queued or ModuleRebuildStatus.Running)
.OrderBy(x => x.Id).FirstOrDefault());
}
public Task TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default)
{
lock (_gate)
{
var row = Rows.FirstOrDefault(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Queued);
if (row == null) return Task.FromResult(false);
row.Status = ModuleRebuildStatus.Cancelled;
row.CurrentStage = ModuleRebuildStages.Cancelled;
row.FinishedAt = DateTime.Now;
row.ErrorMessage = reason;
return Task.FromResult(true);
}
}
public Task CancelUnstartedByTriggerAsync(string triggerType, TimeSpan olderThan, string reason, CancellationToken ct = default)
{
lock (_gate)
{
var trigger = triggerType?.Trim().ToUpperInvariant();
var cutoff = DateTime.Now - olderThan;
var hits = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued
&& string.Equals(x.TriggerType, trigger, StringComparison.Ordinal)
&& x.RequestedBy == null
&& (olderThan <= TimeSpan.Zero || x.SubmittedAt < cutoff))
.ToList();
foreach (var row in hits)
{
row.Status = ModuleRebuildStatus.Cancelled;
row.CurrentStage = ModuleRebuildStages.Cancelled;
row.FinishedAt = DateTime.Now;
row.ErrorMessage = reason;
}
return Task.FromResult(hits.Count);
}
}
public Task FailStaleQueuedAsync(TimeSpan staleAfter, string reason, CancellationToken ct = default)
{
lock (_gate)
{
var cutoff = DateTime.Now - staleAfter;
var hits = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued && x.SubmittedAt < cutoff).ToList();
foreach (var row in hits)
{
row.Status = ModuleRebuildStatus.Failed;
row.CurrentStage = ModuleRebuildStages.Failed;
row.FinishedAt = DateTime.Now;
row.ErrorMessage = reason;
}
return Task.FromResult(hits.Count);
}
}
public Task GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default)
{
lock (_gate)
return Task.FromResult(Rows.FirstOrDefault(x => x.Id == id && x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId));
}
public Task GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
{
lock (_gate)
return Task.FromResult(Rows.Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId)
.OrderByDescending(x => x.Id).FirstOrDefault());
}
public Task ClaimNextQueuedAsync(IReadOnlyCollection enabledModules, CancellationToken ct = default)
{
lock (_gate)
{
var row = Rows.Where(x => x.Status == ModuleRebuildStatus.Queued && enabledModules.Contains(x.ModuleCode))
.Where(x => !Rows.Any(y => y.ModuleCode == x.ModuleCode && y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == ModuleRebuildStatus.Running))
.OrderBy(x => x.Id).FirstOrDefault();
if (row == null) return Task.FromResult(null);
row.Status = ModuleRebuildStatus.Running;
row.StartedAt = DateTime.Now;
row.HeartbeatAt = DateTime.Now;
return Task.FromResult(row);
}
}
public Task UpdateAsync(AdoModuleDashboardRebuildJob 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;
row.UpdateTime = now;
ProgressHistory.Add((jobId, progressPercent, currentStage));
return Task.CompletedTask;
}
}
public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default) =>
UpdateProgressAsync(jobId, completedStage, 0, progressPercent, nextMessage, now, ct);
public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default)
{
lock (_gate)
{
var row = Rows.FirstOrDefault(x => x.Id == jobId);
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 == ModuleRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff)))
{
row.Status = ModuleRebuildStatus.Failed;
row.CurrentStage = ModuleRebuildStages.Failed;
row.FinishedAt = DateTime.Now;
row.LastProgressAt = DateTime.Now;
row.UpdateTime = DateTime.Now;
row.ErrorMessage = "服务中断,任务未正常结束";
}
return Task.CompletedTask;
}
}
}
}