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); /// /// 把一条仍在 QUEUED 的任务收口为 CANCELLED,给人工重算腾位。 /// /// 条件更新(WHERE status='QUEUED'):若 Worker 已把它抢成 RUNNING 则返回 false, /// 调用方必须据此退回 409,不得把正在跑的任务掀掉。 /// Task TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default); /// /// 批量把某个自动触发类型下**仍未开跑**的任务收口为 CANCELLED,返回收口行数。 /// /// 只碰 QUEUED 且 requested_by IS NULL 的行: /// 条件里带 trigger_type 是为了让调用方只清自己那一类, /// 带 requested_by IS NULL 是为了绝不误伤页面手工任务 /// (ModuleRebuildTriggerType.IsManualTrigger 把带 requested_by 的行也算人工)。 /// RUNNING 一律不碰。 /// /// 刻意跨租户:调用方是夜间全量作业,它本来就一次性扇出全部启用 scope, /// 清的也正是自己上一轮那一批。不要给这里加 tenant 过滤,否则只能清掉一个租户。 /// /// 要清理的触发类型,如 AUTO_NIGHTLY。 /// 仅清理入队时间早于 now - olderThan 的行;TimeSpan.Zero 表示不限。 Task CancelUnstartedByTriggerAsync(string triggerType, TimeSpan olderThan, string reason, CancellationToken ct = default); /// /// 把长期滞留在 QUEUED、已无消费可能的任务收口为 FAILED,返回收口行数。 /// /// 为什么 不够:它只管 RUNNING。 /// QUEUED 的滞留有独立成因——非执行机只领手工任务,所以指派一丢, /// 自动任务在队列里没有任何消费者,却会一直挡着同 scope 的后续入队。 /// /// 阈值必须给足:GlobalMaxParallelScopes = 1 且单 scope 跑 10~20 分钟, /// 夜间一次扇出 28 个 scope 时,最后一个**合法**等待可达 7 小时左右。 /// 阈值定短了会把正常排队的任务误杀,见 ModuleRebuildLock.QueuedStaleAfter。 /// Task FailStaleQueuedAsync(TimeSpan staleAfter, string reason, 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 async Task TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default) { var now = DateTime.Now; var n = await _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { Status = ModuleRebuildStatus.Cancelled, CurrentStage = ModuleRebuildStages.Cancelled, FinishedAt = now, ErrorMessage = reason, UpdateTime = now }) .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Queued) .ExecuteCommandAsync(ct); return n > 0; } public async Task CancelUnstartedByTriggerAsync(string triggerType, TimeSpan olderThan, string reason, CancellationToken ct = default) { if (string.IsNullOrWhiteSpace(triggerType)) return 0; var trigger = triggerType.Trim().ToUpperInvariant(); var now = DateTime.Now; var cutoff = now - olderThan; return await _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { Status = ModuleRebuildStatus.Cancelled, CurrentStage = ModuleRebuildStages.Cancelled, FinishedAt = now, ErrorMessage = reason, UpdateTime = now }) // requested_by IS NULL 这一支不可省:带 requested_by 的行按 // ModuleRebuildTriggerType.IsManualTrigger 算人工触发,清了就是误伤页面按钮。 .Where(x => x.Status == ModuleRebuildStatus.Queued && x.TriggerType == trigger && x.RequestedBy == null) .WhereIF(olderThan > TimeSpan.Zero, x => x.SubmittedAt < cutoff) .ExecuteCommandAsync(ct); } public async Task FailStaleQueuedAsync(TimeSpan staleAfter, string reason, CancellationToken ct = default) { var now = DateTime.Now; var cutoff = now - staleAfter; return await _db.Updateable() .SetColumns(x => new AdoModuleDashboardRebuildJob { Status = ModuleRebuildStatus.Failed, CurrentStage = ModuleRebuildStages.Failed, FinishedAt = now, ErrorMessage = reason, UpdateTime = now }) .Where(x => x.Status == ModuleRebuildStatus.Queued && x.SubmittedAt < cutoff) .ExecuteCommandAsync(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); } }