using System.Threading.Channels; using Admin.NET.Plugin.AiDOP.Entity.SmartOps; using Admin.NET.Plugin.AiDOP.MaterialWarehouse; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild; public sealed class ModuleRebuildQueue : ISingleton { private readonly Channel _channel = Channel.CreateBounded(new BoundedChannelOptions(8) { FullMode = BoundedChannelFullMode.DropOldest, SingleReader = true, SingleWriter = false }); public void Pulse() => _channel.Writer.TryWrite(0); public ChannelReader Reader => _channel.Reader; } public sealed class ModuleRebuildJobAccepted { public bool Ok { get; set; } public string ModuleCode { get; set; } = string.Empty; public long? JobId { get; set; } public string Status { get; set; } public string Message { get; set; } } public sealed class ModuleRebuildJobDto { public bool Ok { get; set; } = true; public string ModuleCode { get; set; } = string.Empty; public long JobId { get; set; } public string Status { get; set; } public string CurrentStage { get; set; } public int StageIndex { get; set; } public int StageTotal { get; set; } public int ProgressPercent { get; set; } public string ProgressMessage { get; set; } public DateTime? LastProgressAt { get; set; } public DateTime? HeartbeatAt { get; set; } public string FailedStage { get; set; } public DateTime SubmittedAt { get; set; } public DateTime? StartedAt { get; set; } public DateTime? FinishedAt { get; set; } public int? DurationMs { get; set; } public string BatchId { get; set; } public int StageRows { get; set; } public int StandardRows { get; set; } public int DwdRows { get; set; } public int KpiRows { get; set; } public int AtomicRows { get; set; } public string DetailJson { get; set; } public string ErrorMessage { get; set; } public string Message { get; set; } } public interface IModuleRebuildJobStore { Task InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default); Task FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default); /// 该 scope 最近一条 SUCCESS,用于 AUTO 冷却判定。没有则返回 null。 Task FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default); Task GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default); Task GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default); /// 本实例启用的模块码集合,空集合直接返回 null。 /// /// 全库并发上限:统计**整个库**处于 RUNNING 的 job 条数,达到上限即不再抢占。 /// 与单进程的 MaxParallelScopes 不同,它跨实例生效。 /// /// /// 本实例是否为 ETL 执行机(AidopJobGate.IsRunner)。 /// false 时只领取页面手工任务,不领 AUTO / BOOTSTRAP。 /// 不要改成「非执行机干脆不启动 Worker」:队列同时承载页面「数据重算」按钮, /// 停掉 Worker 会让开发机上的手工重算永远停在 QUEUED。 /// /// 取消令牌。 Task ClaimNextQueuedAsync(IReadOnlyCollection enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default); Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default); Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default); Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default); Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default); Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default); /// 置协作式取消标志。只对仍在 QUEUED / RUNNING 的行生效,返回是否改到行。 Task RequestCancelAsync(long jobId, CancellationToken ct = default); /// 阶段边界查询取消标志。查不到行视为未取消。 Task IsCancelRequestedAsync(long jobId, CancellationToken ct = default); /// /// 找该模块、该租户在本次重算开始后仍为 RUNNING 的 transform run log。 /// job_code 约定为 {module}_MDP_SYNC_TRANSFORM。找不到返回 0。 /// Task FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter); } public sealed class ModuleRebuildJobStore : IModuleRebuildJobStore, ITransient { private readonly ISqlSugarClient _db; private readonly ILogger _logger; public ModuleRebuildJobStore(ISqlSugarClient db, ILoggerFactory loggerFactory) { _db = db; _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildJobStore)); } public async Task InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) { var id = await _db.Insertable(row).ExecuteReturnIdentityAsync(); row.Id = id; return row; } public Task FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId && (x.Status == ModuleRebuildStatus.Queued || x.Status == ModuleRebuildStatus.Running)) .OrderBy(x => x.Id) .FirstAsync(ct); public Task FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId && x.Status == ModuleRebuildStatus.Success) .OrderBy(x => x.FinishedAt, OrderByType.Desc) .FirstAsync(ct); public Task GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .FirstAsync(x => x.Id == id && x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId, ct); public Task GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId) .OrderBy(x => x.Id, OrderByType.Desc) .FirstAsync(ct); public async Task RequestCancelAsync(long jobId, CancellationToken ct = default) { var affected = await _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { CancelRequestedFlag = true, UpdateTime = DateTime.Now }) .Where(x => x.Id == jobId && (x.Status == ModuleRebuildStatus.Queued || x.Status == ModuleRebuildStatus.Running)) .ExecuteCommandAsync(ct); return affected > 0; } public async Task FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter) { var jobCode = $"{moduleCode}_MDP_SYNC_TRANSFORM"; var id = await _db.Ado.GetLongAsync( """ SELECT id FROM mdp_transform_run_log WHERE job_code=@JobCode AND tenant_id=@TenantId AND status='RUNNING' AND start_time>=@StartedAfter ORDER BY id DESC LIMIT 1 """, new List { new("@JobCode", jobCode), new("@TenantId", tenantId), new("@StartedAfter", startedAfter.AddMinutes(-1)) }); return id; } public async Task IsCancelRequestedAsync(long jobId, CancellationToken ct = default) { var flag = await _db.Queryable() .Where(x => x.Id == jobId) .Select(x => x.CancelRequestedFlag) .FirstAsync(ct); return flag; } /// 跨实例抢占互斥锁键。只保护「数 RUNNING + 置 RUNNING」这一小段,不覆盖重算本体。 public const string ClaimLockKey = "aidop:mdp:rebuild:claim"; public async Task ClaimNextQueuedAsync( IReadOnlyCollection enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default) { if (enabledModules == null || enabledModules.Count == 0) return null; // 计数与抢占必须在同一把跨实例锁内:否则 N 个实例会在同一瞬间都读到 // 「RUNNING 还没到上限」,各自抢走一个 job,上限形同虚设。 // 非阻塞取锁——抢不到说明别的实例正在抢占,直接让出,5 秒后下一拍再来。 await using var claimGate = await InventoryInboundLockGuard.TryAcquireAsync( _db, _logger, ClaimLockKey, ct); if (!claimGate.Acquired) return null; var running = await _db.Queryable() .Where(x => x.Status == ModuleRebuildStatus.Running) .CountAsync(ct); if (running >= globalMaxParallelScopes) return null; var enabled = enabledModules.ToList(); var row = await _db.Queryable() .Where(x => x.Status == ModuleRebuildStatus.Queued && enabled.Contains(x.ModuleCode)) // 非执行机只领页面手工任务。用**白名单**而非「排除 AUTO/BOOTSTRAP」: // 后续新增的自动触发类型(如夜间兜底的 AUTO_NIGHTLY)天然落在白名单外, // 不会因为有人忘了往黑名单里补一项就被开发机领去跑全量。 .WhereIF(!runner, x => x.TriggerType == ModuleRebuildTriggerType.Manual || x.RequestedBy != null) .Where(x => SqlFunc.Subqueryable() .Where(y => y.ModuleCode == x.ModuleCode && y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == ModuleRebuildStatus.Running) .NotAny()) .OrderBy(x => x.Id) .FirstAsync(ct); if (row == null) return null; var now = DateTime.Now; var n = await _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { Status = ModuleRebuildStatus.Running, CurrentStage = ModuleRebuildStages.AcquiringLock, StageIndex = 0, ProgressPercent = 2, ProgressMessage = $"等待现有 {row.ModuleCode} 全量任务完成", LastProgressAt = now, StartedAt = now, HeartbeatAt = now, UpdateTime = now }) .Where(x => x.Id == row.Id && x.Status == ModuleRebuildStatus.Queued) .ExecuteCommandAsync(ct); if (n <= 0) return null; row.Status = ModuleRebuildStatus.Running; row.CurrentStage = ModuleRebuildStages.AcquiringLock; row.ProgressPercent = 2; row.StartedAt = now; row.HeartbeatAt = now; row.UpdateTime = now; return row; } public Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) => _db.Updateable(row).ExecuteCommandAsync(ct); public Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default) => _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { CurrentStage = currentStage, StageIndex = stageIndex, ProgressPercent = progressPercent, ProgressMessage = message, LastProgressAt = now, HeartbeatAt = now, UpdateTime = now }) .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Running) .ExecuteCommandAsync(ct); public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default) { var column = completedStage switch { ModuleRebuildStages.Staging or ModuleRebuildStages.T8Inbound => "stage_rows", ModuleRebuildStages.Standard or ModuleRebuildStages.KpiPreparing => "standard_rows", ModuleRebuildStages.Dwd => "dwd_rows", ModuleRebuildStages.Kpi or ModuleRebuildStages.KpiCalculating => "kpi_rows", ModuleRebuildStages.Atomic => "atomic_rows", _ => null }; if (column == null) return UpdateProgressAsync(jobId, completedStage, 0, progressPercent, nextMessage, now, ct); return _db.Ado.ExecuteCommandAsync( $""" UPDATE ado_module_dashboard_rebuild_job SET `{column}`=@Rows, progress_percent=@Pct, progress_message=@Msg, detail_json=IFNULL(@Detail, detail_json), last_progress_at=@Now, heartbeat_at=@Now, update_time=@Now WHERE id=@Id AND status='RUNNING' """, new SugarParameter("@Rows", rows), new SugarParameter("@Pct", progressPercent), new SugarParameter("@Msg", nextMessage), new SugarParameter("@Detail", (object?)detailJson ?? DBNull.Value), new SugarParameter("@Now", now), new SugarParameter("@Id", jobId)); } public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default) => _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { HeartbeatAt = now, UpdateTime = now }) .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Running) .ExecuteCommandAsync(ct); public Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default) { var cutoff = DateTime.Now - staleAfter; var now = DateTime.Now; return _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { Status = ModuleRebuildStatus.Failed, CurrentStage = ModuleRebuildStages.Failed, FinishedAt = now, LastProgressAt = now, UpdateTime = now, ErrorMessage = "服务中断,任务未正常结束" }) .Where(x => x.Status == ModuleRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff)) .ExecuteCommandAsync(ct); } }