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; }
}