using Admin.NET.Core.Service; using Admin.NET.Plugin.AiDOP.FinishedWarehouse; using Admin.NET.Plugin.AiDOP.Manufacturing; using Admin.NET.Plugin.AiDOP.MaterialWarehouse; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// /// S5/S6/S7 T8 KPI 手动刷新统一编排(看板"刷新"按钮真实跑批入口)。 /// 单次刷新语义:抢全局刷新锁 → T8 基表入站(硬门,共享 6 表)→ 本模块 RunFull(MANUAL) → 汇总逐指标状态 → 释放锁。 /// 与 Cron job_s5_t8_kpi_refresh 共享同一把锁(),保证手动刷新与定时跑批互斥、并发手动刷新排他。 /// 入站为 S5/S6/S7 共用(T8BaseInboundMdpSyncService 6 张基表无法按模块拆),故刷任一模块均触发一次全量入站。 /// S5_L1_002 / S5_L1_004 仍 legacy 直连 T8(本轮不改),随 S5 RunFull 一并执行,逐指标标 LEGACY_SUCCESS。 /// /// 锁设计(加固):用持有者无关的硬 TTL 键(SysCacheService.Set 带过期),而非 BeginCacheLock。 /// 原因:BeginCacheLock 在持有者线程存活期间会续租,一旦某次运行意外挂死(如源库慢查询无超时)会永久占锁、 /// 阻塞后续所有刷新。改为硬 TTL 后:正常路径 finally 主动 Remove;即使挂死,键到 也自动过期。 /// 配合 QueryT8Async 的 T8 命令超时(60s)与主库命令超时(30s),单次运行时长恒 < TTL,正常不会误过期。 /// public sealed class AidopT8KpiManualRefreshService : ITransient { /// 手动刷新 + Cron 共享的全局互斥锁 key(覆盖 inbound + transform 整段)。 public const string RefreshLockKey = "aidop:t8kpi:refresh:lock"; /// 刷新锁硬 TTL(秒,持有者无关):到点自动过期,即使某次运行意外挂死也不会永久阻塞后续刷新。 /// 取值须大于加固后最坏运行时长(T8 查询各 ≤60s、主库命令 ≤30s,全程数分钟内),10 分钟为安全上限。 public const int RefreshLockTtlSeconds = 600; // legacy 直连 T8 的 KPI(仅剩 S5_L1_004 TVF;S5_L1_002 已中台化),逐指标状态标 LEGACY_SUCCESS 以示区分。 private static readonly HashSet LegacyKpiCodes = new(StringComparer.OrdinalIgnoreCase) { "S5_L1_004" }; private readonly T8BaseInboundMdpSyncService _inbound; private readonly S5MdpSyncTransformService _s5; private readonly S6MdpSyncTransformService _s6; private readonly S7MdpSyncTransformService _s7; private readonly SysCacheService _cache; private readonly ILogger _logger; public AidopT8KpiManualRefreshService( T8BaseInboundMdpSyncService inbound, S5MdpSyncTransformService s5, S6MdpSyncTransformService s6, S7MdpSyncTransformService s7, SysCacheService cache, ILogger logger) { _inbound = inbound; _s5 = s5; _s6 = s6; _s7 = s7; _cache = cache; _logger = logger; } /// 执行一次模块级手动刷新(S5/S6/S7)。返回统一契约,绝不抛业务异常(取消除外)。 public async Task RunModuleRefreshAsync(string moduleCode, CancellationToken cancellationToken) { var mc = string.IsNullOrWhiteSpace(moduleCode) ? string.Empty : moduleCode.Trim().ToUpperInvariant(); var startedAt = DateTime.Now; var result = new AidopKpiRefreshResult { ModuleCode = mc, StartedAt = startedAt.ToString("yyyy-MM-dd HH:mm:ss") }; if (mc is not ("S5" or "S6" or "S7")) { result.Ok = false; result.OverallStatus = "SKIPPED"; result.Message = $"模块 {mc} 不支持 T8 KPI 手动刷新"; Finish(result, startedAt); return result; } // 硬 TTL 锁:ExistKey 命中即判定已有刷新在跑(手动或定时);否则 Set(带 TTL) 占锁。 // 单实例 + 人工点击/每日 5 次定时下,ExistKey→Set 的微小竞态概率可忽略;即便未主动释放,TTL 到点自动过期。 if (_cache.ExistKey(RefreshLockKey)) { result.Ok = false; result.OverallStatus = "REFRESHING"; result.Message = "T8 KPI 刷新正在进行中(手动或定时任务占用),请稍后重试"; Finish(result, startedAt); return result; } _cache.Set(RefreshLockKey, $"MANUAL:{mc}:{startedAt:yyyyMMddHHmmss}", TimeSpan.FromSeconds(RefreshLockTtlSeconds)); try { // ① 入站硬门:共享 6 张 T8 基表全量入站;失败即中断,不用旧 std 伪装成功。 try { var inbound = await _inbound.RunInboundAsync(0, true, null, cancellationToken); result.Inbound.Ok = true; result.Inbound.StageRows = inbound.Tables.Sum(t => t.StgRows); result.Inbound.StandardRows = inbound.Tables.Sum(t => t.StdRows); } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } catch (Exception ex) { _logger.LogError(ex, "[手动刷新] {Module} T8 基表入站失败", mc); result.Ok = false; result.OverallStatus = "FAILED"; result.Inbound.Ok = false; result.Inbound.Message = ex.Message; result.Message = "源数据入站失败,未更新 KPI(保留原有看板数据)"; Finish(result, startedAt); return result; } // ② 本模块 KPI 转换(RunFull 为全有全无:任一 KPI 异常即整体抛错)。 string batchId; int dwdRows; int kpiRows; Dictionary perKpiRows; List denomStatus; try { switch (mc) { case "S5": var r5 = await _s5.RunFullAsync(cancellationToken, "MANUAL"); batchId = r5.BatchId; dwdRows = r5.DwdRows; kpiRows = r5.KpiRows; perKpiRows = r5.PerKpiKpiRows; denomStatus = r5.KpiDenominatorStatus; break; case "S6": var r6 = await _s6.RunFullAsync(cancellationToken, "MANUAL"); batchId = r6.BatchId; dwdRows = r6.DwdRows; kpiRows = r6.KpiRows; perKpiRows = r6.PerKpiKpiRows; denomStatus = r6.KpiDenominatorStatus; break; default: var r7 = await _s7.RunFullAsync(cancellationToken, "MANUAL"); batchId = r7.BatchId; dwdRows = r7.DwdRows; kpiRows = r7.KpiRows; perKpiRows = r7.PerKpiKpiRows; denomStatus = r7.KpiDenominatorStatus; break; } } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } catch (Exception ex) { _logger.LogError(ex, "[手动刷新] {Module} KPI 转换失败", mc); result.Ok = false; result.OverallStatus = "FAILED"; result.Transform.Ok = false; result.Transform.Message = ex.Message; result.Message = "KPI 计算失败(源数据已入站,看板保留原有数据)"; Finish(result, startedAt); return result; } result.BatchId = batchId; result.Transform.Ok = true; result.Transform.DwdRows = dwdRows; result.Transform.KpiRows = kpiRows; // ③ 逐指标状态:denom=OK → SUCCESS(legacy→LEGACY_SUCCESS);denom=NO_* → NO_DATA(真实无数据,非伪造 0)。 var denomMap = ParseDenominator(denomStatus); foreach (var kv in perKpiRows.OrderBy(k => k.Key, StringComparer.Ordinal)) { var code = kv.Key; var denom = denomMap.TryGetValue(code, out var d) ? d : "OK"; var ok = string.Equals(denom, "OK", StringComparison.OrdinalIgnoreCase); var status = ok ? (LegacyKpiCodes.Contains(code) ? "LEGACY_SUCCESS" : "SUCCESS") : "NO_DATA"; result.PerKpi.Add(new AidopKpiRefreshPerKpi { MetricCode = code, Status = status, Rows = kv.Value, Message = ok ? null : denom }); } result.SuccessCount = result.PerKpi.Count(p => p.Status is "SUCCESS" or "LEGACY_SUCCESS"); result.NoDataCount = result.PerKpi.Count(p => p.Status == "NO_DATA"); result.FailedCount = result.PerKpi.Count(p => p.Status is "FAILED" or "LEGACY_FAILED"); result.Ok = true; // RunFull 为全有全无,成功路径下 FailedCount 恒为 0;PARTIAL_SUCCESS 仅为契约保留态。 result.OverallStatus = result.FailedCount > 0 ? "PARTIAL_SUCCESS" : "SUCCESS"; Finish(result, startedAt); result.Message = $"刷新完成:成功 {result.SuccessCount},无数据 {result.NoDataCount},失败 {result.FailedCount},耗时 {result.DurationMs / 1000.0:0.0} 秒"; return result; } finally { // 无论成功/失败/取消,主动释放锁(硬 TTL 是兜底,正常路径这里即时释放)。 _cache.Remove(RefreshLockKey); } } /// 把 "S5_L1_001:OK" 形式的分母状态列表解析为 metricCode → status 字典。 private static Dictionary ParseDenominator(List denom) { var map = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var item in denom) { if (string.IsNullOrWhiteSpace(item)) continue; var idx = item.IndexOf(':'); if (idx <= 0) continue; map[item.Substring(0, idx).Trim()] = item.Substring(idx + 1).Trim(); } return map; } private static void Finish(AidopKpiRefreshResult result, DateTime startedAt) { var finishedAt = DateTime.Now; result.FinishedAt = finishedAt.ToString("yyyy-MM-dd HH:mm:ss"); result.DurationMs = (long)(finishedAt - startedAt).TotalMilliseconds; } } /// 手动刷新统一返回契约(PascalCase 属性经框架序列化为 camelCase:Ok→ok、BatchId→batchId)。 public sealed class AidopKpiRefreshResult { public bool Ok { get; set; } public string ModuleCode { get; set; } = string.Empty; /// 整体状态:SUCCESS / PARTIAL_SUCCESS / FAILED / REFRESHING / SKIPPED。 public string OverallStatus { get; set; } = string.Empty; public string? BatchId { get; set; } public string StartedAt { get; set; } = string.Empty; public string FinishedAt { get; set; } = string.Empty; public long DurationMs { get; set; } public AidopKpiRefreshInbound Inbound { get; set; } = new(); public AidopKpiRefreshTransform Transform { get; set; } = new(); public int SuccessCount { get; set; } public int NoDataCount { get; set; } public int FailedCount { get; set; } public List PerKpi { get; set; } = new(); public string Message { get; set; } = string.Empty; } public sealed class AidopKpiRefreshInbound { public bool Ok { get; set; } public int StageRows { get; set; } public int StandardRows { get; set; } public string? Message { get; set; } } public sealed class AidopKpiRefreshTransform { public bool Ok { get; set; } public int DwdRows { get; set; } public int KpiRows { get; set; } public string? Message { get; set; } } /// 单指标刷新状态:SUCCESS / LEGACY_SUCCESS / NO_DATA / FAILED / LEGACY_FAILED。 public sealed class AidopKpiRefreshPerKpi { public string MetricCode { get; set; } = string.Empty; public string Status { get; set; } = string.Empty; public int Rows { get; set; } public string? Message { get; set; } }