using System.Threading.Channels; using Admin.NET.Plugin.AiDOP.Entity.SmartOps; namespace Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh; public sealed class S1DashboardRebuildQueue : 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 S1RebuildJobAccepted { public bool Ok { get; set; } public long? JobId { get; set; } public string Status { get; set; } public string Message { get; set; } } public sealed class S1RebuildJobDto { 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; } = S1MdpRebuildStage.StageTotal; 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 ErrorMessage { get; set; } } public interface IS1DashboardRebuildJobStore { Task InsertQueuedAsync(AdoS1DashboardRebuildJob row, CancellationToken ct = default); Task FindActiveAsync(long tenantId, long factoryId, CancellationToken ct = default); Task GetByIdAsync(long id, long tenantId, long factoryId, CancellationToken ct = default); Task GetLatestAsync(long tenantId, long factoryId, CancellationToken ct = default); Task ClaimNextQueuedAsync(CancellationToken ct = default); Task UpdateAsync(AdoS1DashboardRebuildJob 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, DateTime now, CancellationToken ct = default); Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default); Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default); } public sealed class S1DashboardRebuildJobStore : IS1DashboardRebuildJobStore, ITransient { private readonly ISqlSugarClient _db; public S1DashboardRebuildJobStore(ISqlSugarClient db) => _db = db; public async Task InsertQueuedAsync(AdoS1DashboardRebuildJob row, CancellationToken ct = default) { var id = await _db.Insertable(row).ExecuteReturnIdentityAsync(); row.Id = id; return row; } public Task FindActiveAsync(long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId && (x.Status == S1DashboardRebuildStatus.Queued || x.Status == S1DashboardRebuildStatus.Running)) .OrderBy(x => x.Id) .FirstAsync(ct); public Task GetByIdAsync(long id, long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .FirstAsync(x => x.Id == id && x.TenantId == tenantId && x.FactoryId == factoryId, ct); public Task GetLatestAsync(long tenantId, long factoryId, CancellationToken ct = default) => _db.Queryable() .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId) .OrderBy(x => x.Id, OrderByType.Desc) .FirstAsync(ct); public async Task ClaimNextQueuedAsync(CancellationToken ct = default) { var row = await _db.Queryable() .Where(x => x.Status == S1DashboardRebuildStatus.Queued) .Where(x => SqlFunc.Subqueryable() .Where(y => y.TenantId == x.TenantId && y.FactoryId == x.FactoryId && y.Status == S1DashboardRebuildStatus.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 AdoS1DashboardRebuildJob { Status = S1DashboardRebuildStatus.Running, CurrentStage = S1MdpRebuildStage.AcquiringLock, StageIndex = 0, ProgressPercent = 2, ProgressMessage = "等待现有 S1 全量任务完成", LastProgressAt = now, StartedAt = now, HeartbeatAt = now, UpdateTime = now }) .Where(x => x.Id == row.Id && x.Status == S1DashboardRebuildStatus.Queued) .ExecuteCommandAsync(ct); if (n <= 0) return null; row.Status = S1DashboardRebuildStatus.Running; row.CurrentStage = S1MdpRebuildStage.AcquiringLock; row.ProgressPercent = 2; row.StartedAt = now; row.HeartbeatAt = now; row.UpdateTime = now; return row; } public Task UpdateAsync(AdoS1DashboardRebuildJob 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 AdoS1DashboardRebuildJob { CurrentStage = currentStage, StageIndex = stageIndex, ProgressPercent = progressPercent, ProgressMessage = message, LastProgressAt = now, HeartbeatAt = now, UpdateTime = now }) .Where(x => x.Id == jobId && x.Status == S1DashboardRebuildStatus.Running) .ExecuteCommandAsync(ct); public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, DateTime now, CancellationToken ct = default) { var column = completedStage switch { S1MdpRebuildStage.Staging => "stage_rows", S1MdpRebuildStage.Standard => "standard_rows", S1MdpRebuildStage.Dwd => "dwd_rows", S1MdpRebuildStage.Kpi => "kpi_rows", S1MdpRebuildStage.Atomic => "atomic_rows", _ => null }; if (column == null) return UpdateProgressAsync(jobId, completedStage, S1MdpRebuildStage.ToStageIndex(completedStage), progressPercent, nextMessage, now, ct); return _db.Ado.ExecuteCommandAsync( $""" UPDATE ado_s1_dashboard_rebuild_job SET `{column}`=@Rows, progress_percent=@Pct, progress_message=@Msg, 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("@Now", now), new SugarParameter("@Id", jobId)); } public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default) => _db.Updateable() .SetColumns(x => new AdoS1DashboardRebuildJob { HeartbeatAt = now, UpdateTime = now }) .Where(x => x.Id == jobId && x.Status == S1DashboardRebuildStatus.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 AdoS1DashboardRebuildJob { Status = S1DashboardRebuildStatus.Failed, CurrentStage = S1MdpRebuildStage.Failed, FinishedAt = now, LastProgressAt = now, UpdateTime = now, ErrorMessage = "服务中断,任务未正常结束" }) .Where(x => x.Status == S1DashboardRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff)) .ExecuteCommandAsync(ct); } }