using Admin.NET.Plugin.AiDOP.Entity.SmartOps; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild; public sealed class ModuleRebuildService : ITransient { private readonly IModuleRebuildJobStore _jobs; private readonly ModuleRebuildQueue _queue; private readonly IModuleRebuildCapability _capability; private readonly ILogger _logger; public ModuleRebuildService(IModuleRebuildJobStore jobs, ModuleRebuildQueue queue, IModuleRebuildCapability capability) : this(jobs, queue, capability, NullLoggerFactory.Instance) { } private readonly MdpNeutralSourceCleanup? _cleanup; private readonly TransformRunLogFinalizer? _runLogFinalizer; public ModuleRebuildService( IModuleRebuildJobStore jobs, ModuleRebuildQueue queue, IModuleRebuildCapability capability, ILoggerFactory loggerFactory) : this(jobs, queue, capability, loggerFactory, null) { } public ModuleRebuildService( IModuleRebuildJobStore jobs, ModuleRebuildQueue queue, IModuleRebuildCapability capability, ILoggerFactory loggerFactory, MdpNeutralSourceCleanup? cleanup, TransformRunLogFinalizer? runLogFinalizer = null) { _jobs = jobs; _queue = queue; _capability = capability; _cleanup = cleanup; _runLogFinalizer = runLogFinalizer; _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildService)); } public async Task<(int StatusCode, ModuleRebuildJobAccepted Body)> EnqueueAsync( string moduleCode, long tenantId, long factoryId, long? requestedBy, string triggerType = "MANUAL", CancellationToken ct = default) { var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId); if (!_capability.IsEnabled(scope.ModuleCode)) { return (404, new ModuleRebuildJobAccepted { Ok = false, ModuleCode = scope.ModuleCode, Status = "DISABLED", Message = $"{scope.ModuleCode} 数据重算尚未启用" }); } var normalizedTrigger = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant(); var manualRequest = ModuleRebuildTriggerType.IsManualTrigger(normalizedTrigger, requestedBy); var active = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct); if (active != null) { // 人工重算可以取代同 scope 里「还没开跑的自动任务」:二者都是同 scope 全量,取代不丢工作量。 // 不让位则会死锁——ClaimNextQueuedAsync 只让执行机领自动任务,非执行机上这条 QUEUED 永远 // 不会消失,页面「数据重算」按钮就被它永久挡住(执行机指派还会随进程重启失效)。 var supersedable = manualRequest && active.Status == ModuleRebuildStatus.Queued && !ModuleRebuildTriggerType.IsManualTrigger(active.TriggerType, active.RequestedBy); if (!supersedable || !await _jobs.TrySupersedeQueuedAsync(active.Id, $"已被人工重算取代(原 {active.TriggerType})", ct)) { return (409, Conflict(scope.ModuleCode, active)); } _logger.LogInformation( "[ModuleRebuild] {Module} tenant={Tenant} factory={Factory} 人工重算取代待执行任务 job={JobId} trigger={Trigger}", scope.ModuleCode, scope.TenantId, scope.FactoryId, active.Id, active.TriggerType); } // AUTO / BOOTSTRAP 才冷却。MANUAL 是页面「数据重算」按钮,AUTO_NIGHTLY 是全量兜底,两者都必须放行。 var cooldownHours = _capability.AutoMinIntervalHours; if (cooldownHours > 0 && !manualRequest && normalizedTrigger != ModuleRebuildTriggerType.AutoNightly) { var lastSuccess = await _jobs.FindLastSuccessAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct); if (lastSuccess?.FinishedAt is { } finishedAt && DateTime.Now - finishedAt < TimeSpan.FromHours(cooldownHours)) { return (409, new ModuleRebuildJobAccepted { Ok = false, ModuleCode = scope.ModuleCode, JobId = lastSuccess.Id, Status = "COOLDOWN", Message = $"{scope.ModuleCode} 距上次成功重算不足 {cooldownHours} 小时,本次跳过" }); } } var now = DateTime.Now; var created = await _jobs.InsertQueuedAsync(new AdoModuleDashboardRebuildJob { ModuleCode = scope.ModuleCode, TenantId = scope.TenantId, FactoryId = scope.FactoryId, Status = ModuleRebuildStatus.Queued, CurrentStage = ModuleRebuildStages.Queued, StageIndex = 0, StageTotal = ModuleRebuildStages.StageTotal(scope.ModuleCode), ProgressPercent = 0, ProgressMessage = "已入队", LastProgressAt = now, TriggerType = normalizedTrigger, RequestedBy = requestedBy, SubmittedAt = now, CreateTime = now, UpdateTime = now }, ct); var oldest = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct); if (oldest != null && oldest.Id != created.Id) { created.Status = ModuleRebuildStatus.Failed; created.CurrentStage = ModuleRebuildStages.Failed; created.FinishedAt = DateTime.Now; created.ErrorMessage = $"{scope.ModuleCode} 数据重算正在执行,请勿重复提交"; created.UpdateTime = DateTime.Now; await _jobs.UpdateAsync(created, ct); return (409, Conflict(scope.ModuleCode, oldest)); } _queue.Pulse(); return (202, new ModuleRebuildJobAccepted { Ok = true, ModuleCode = scope.ModuleCode, JobId = created.Id, Status = ModuleRebuildStatus.Queued, Message = $"{scope.ModuleCode} 数据重算已排队" }); } public async Task GetAsync(string moduleCode, long jobId, long tenantId, long factoryId, CancellationToken ct = default) { var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId); var row = await _jobs.GetByIdAsync(scope.ModuleCode, jobId, scope.TenantId, scope.FactoryId, ct); if (row == null) throw Oops.Oh("任务不存在"); return ToDto(row); } public async Task LatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) { var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId); var row = await _jobs.GetLatestAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct); return row == null ? null : ToDto(row); } /// /// 陈旧任务收口:RUNNING 无心跳、以及 QUEUED 长期无人消费,两类都要捞。 /// /// QUEUED 那一支是 2026-09-28 补的:非执行机只领手工任务,所以执行机指派一丢, /// 自动任务在队列里没有任何消费者,却会一直挡着同 scope 的后续入队。 /// 只扫 RUNNING 的旧实现看不见这种滞留。 /// public async Task FailStaleAsync(CancellationToken ct = default) { await _jobs.FailStaleRunningAsync(ModuleRebuildLock.StaleAfter, ct); var reaped = await _jobs.FailStaleQueuedAsync( ModuleRebuildLock.QueuedStaleAfter, $"排队超过 {ModuleRebuildLock.QueuedStaleAfter.TotalHours:0} 小时仍无人领取,已收口;" + "常见成因是 ETL 执行机未指派或指派已失效", ct); if (reaped > 0) _logger.LogWarning( "[ModuleRebuild] 收口长期滞留的 QUEUED 任务 {Count} 条(阈值 {Hours} 小时)。" + "请检查 MDP 运行监控页是否存在「零台存活执行机」", reaped, ModuleRebuildLock.QueuedStaleAfter.TotalHours); } public async Task TouchHeartbeatAsync(long jobId, CancellationToken ct = default) { try { await _jobs.TouchHeartbeatAsync(jobId, DateTime.Now, ct); } catch (Exception ex) { _logger.LogWarning(ex, "[MdpRebuild] heartbeat failed jobId={JobId}", jobId); } } public async Task ApplyProgressAsync(long jobId, string moduleCode, ModuleProgressUpdate update, CancellationToken ct = default) { try { var now = DateTime.Now; if (update.CompletedStage != null && update.Rows.HasValue) { await _jobs.UpdateStageResultAsync( jobId, update.CompletedStage, update.Rows.Value, update.ProgressPercent, update.Message, update.DetailJson, now, ct); } await _jobs.UpdateProgressAsync( jobId, update.Stage, update.StageIndex, update.ProgressPercent, update.Message, now, ct); } catch (Exception ex) { _logger.LogWarning(ex, "[MdpRebuild] progress write failed jobId={JobId} module={Module} stage={Stage}", jobId, moduleCode, update.Stage); } } public async Task RunClaimedAsync( AdoModuleDashboardRebuildJob job, IModuleRebuildHandler handler, IModuleRebuildLock runLock, CancellationToken stoppingToken) { var failedStage = job.CurrentStage; var scope = MdpRebuildScope.Create(job.ModuleCode, job.TenantId, job.FactoryId); try { await ApplyProgressAsync(job.Id, scope.ModuleCode, new ModuleProgressUpdate( ModuleRebuildStages.AcquiringLock, 0, 2, $"等待现有 {scope.ModuleCode} 全量任务完成"), CancellationToken.None); var holderId = $"{job.TriggerType}:{scope.ScopeKey}:{Environment.MachineName}:{Guid.NewGuid():N}"; var lease = await runLock.TryAcquireAsync(scope, holderId, job.Id, stoppingToken); if (lease == null) throw new ModuleRebuildAlreadyRunningException(scope.ModuleCode); await using (lease) { using var lockHeartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); var lockHeartbeat = KeepLeaseAliveAsync(lease, lockHeartbeatCts.Token); try { if (_cleanup != null) await _cleanup.PurgeModuleAsync(scope.TenantId, scope.ModuleCode, stoppingToken); var result = await handler.RunAsync( scope, job.TriggerType, job.Id, async update => { failedStage = update.Stage; await ApplyProgressAsync(job.Id, scope.ModuleCode, update, CancellationToken.None); // 阶段边界:每个 handler 每跨一个阶段都会走到这里,是唯一的通用取消点。 if (await _jobs.IsCancelRequestedAsync(job.Id, CancellationToken.None)) throw new ModuleRebuildCancelledException(); }, stoppingToken); var now = DateTime.Now; job.Status = ModuleRebuildStatus.Success; job.CurrentStage = ModuleRebuildStages.Success; job.StageIndex = ModuleRebuildStages.StageTotal(scope.ModuleCode); job.ProgressPercent = 100; job.ProgressMessage = $"{scope.ModuleCode} 数据重算已完成"; job.FailedStage = null; job.LastProgressAt = now; job.BatchId = result.BatchId; job.TransformRunLogId = result.RunLogId; job.StageRows = result.StageRows; job.StandardRows = result.StandardRows; job.DwdRows = result.DwdRows; job.KpiRows = result.KpiRows; job.AtomicRows = result.AtomicRows; job.DetailJson = result.DetailJson; job.FinishedAt = now; job.HeartbeatAt = now; job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null; job.UpdateTime = now; job.ErrorMessage = null; await _jobs.UpdateAsync(job, CancellationToken.None); } finally { lockHeartbeatCts.Cancel(); try { await lockHeartbeat; } catch (OperationCanceledException) { } } } } catch (ModuleRebuildAlreadyRunningException) { var now = DateTime.Now; job.Status = ModuleRebuildStatus.Queued; job.CurrentStage = ModuleRebuildStages.AcquiringLock; job.StageIndex = 0; job.ProgressPercent = 2; job.ProgressMessage = $"等待现有 {scope.ModuleCode} 全量任务完成"; job.StartedAt = null; job.HeartbeatAt = now; job.LastProgressAt = now; job.UpdateTime = now; job.ErrorMessage = null; await _jobs.UpdateAsync(job, CancellationToken.None); _queue.Pulse(); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { await FailJobAsync(job, "服务停止,任务已取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage); } catch (ModuleRebuildCancelledException) { await FailJobAsync(job, "已由管理员取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage); if (_runLogFinalizer != null) { var runLogId = await _jobs.FindOpenRunLogIdAsync(scope.ModuleCode, scope.TenantId, job.StartedAt ?? job.SubmittedAt); if (runLogId > 0) await _runLogFinalizer.FinalizeAsAdminCancelledAsync(runLogId, job.StartedAt ?? DateTime.Now); } } catch (Exception ex) { await FailJobAsync( job, Truncate(ex.Message, 2000), ModuleRebuildStatus.Failed, ModuleRebuildStages.Failed, failedStage); } } private async Task FailJobAsync(AdoModuleDashboardRebuildJob job, string message, string status, string currentStage, string failedStage) { var now = DateTime.Now; job.Status = status; job.CurrentStage = currentStage; job.FailedStage = status == ModuleRebuildStatus.Failed ? failedStage : job.FailedStage; job.ProgressMessage = status == ModuleRebuildStatus.Failed ? $"在【{ModuleRebuildStages.ToChinese(failedStage)}】失败" : message; job.FinishedAt = now; job.HeartbeatAt = now; job.LastProgressAt = now; job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null; job.ErrorMessage = message; job.UpdateTime = now; await _jobs.UpdateAsync(job, CancellationToken.None); } public static ModuleRebuildJobDto ToDto(AdoModuleDashboardRebuildJob row) => new() { Ok = true, ModuleCode = row.ModuleCode, JobId = row.Id, Status = row.Status, CurrentStage = row.CurrentStage, StageIndex = row.StageIndex, StageTotal = row.StageTotal <= 0 ? ModuleRebuildStages.StageTotal(row.ModuleCode) : row.StageTotal, ProgressPercent = row.ProgressPercent, ProgressMessage = row.ProgressMessage, LastProgressAt = row.LastProgressAt, HeartbeatAt = row.HeartbeatAt, FailedStage = row.FailedStage, SubmittedAt = row.SubmittedAt, StartedAt = row.StartedAt, FinishedAt = row.FinishedAt, DurationMs = row.DurationMs, BatchId = row.BatchId, StageRows = row.StageRows, StandardRows = row.StandardRows, DwdRows = row.DwdRows, KpiRows = row.KpiRows, AtomicRows = row.AtomicRows, DetailJson = row.DetailJson, ErrorMessage = row.ErrorMessage }; private static ModuleRebuildJobAccepted Conflict(string moduleCode, AdoModuleDashboardRebuildJob active) => new() { Ok = false, ModuleCode = moduleCode, JobId = active.Id, Status = active.Status, Message = $"{moduleCode} 数据重算正在执行,请勿重复提交" }; private static async Task KeepLeaseAliveAsync(IModuleRebuildLease lease, CancellationToken ct) { while (!ct.IsCancellationRequested) { try { await Task.Delay(TimeSpan.FromSeconds(30), ct); await lease.HeartbeatAsync(CancellationToken.None); } catch (OperationCanceledException) { return; } catch { // 单例 SqlSugar 连接被业务阶段占用时跳过本次锁心跳,避免打断 MDP。 } } } private static int ToDurationMs(TimeSpan elapsed) { var ms = elapsed.TotalMilliseconds; if (double.IsNaN(ms) || ms <= 0) return 0; return ms >= int.MaxValue ? int.MaxValue : (int)ms; } private static string Truncate(string value, int max) => string.IsNullOrEmpty(value) || value.Length <= max ? value : value[..max]; }