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) { } public ModuleRebuildService( IModuleRebuildJobStore jobs, ModuleRebuildQueue queue, IModuleRebuildCapability capability, ILoggerFactory loggerFactory) { _jobs = jobs; _queue = queue; _capability = capability; _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 active = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct); if (active != null) { return (409, Conflict(scope.ModuleCode, active)); } 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 = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant(), 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); } public async Task FailStaleAsync(CancellationToken ct = default) => await _jobs.FailStaleRunningAsync(ModuleRebuildLock.StaleAfter, ct); 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 { var result = await handler.RunAsync( scope, job.TriggerType, job.Id, async update => { failedStage = update.Stage; await ApplyProgressAsync(job.Id, scope.ModuleCode, update, CancellationToken.None); }, 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 (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]; }