using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess; using Microsoft.Extensions.Logging; using SqlSugar; using System.Text.RegularExpressions; namespace Admin.NET.Plugin.AiDOP.Service.S8.Rules.Health; /// /// 的数据中台适配实现: /// 从 mdp_transform_run_log 采集生产运行事实,必要时再校验快照表的发布不变量, /// 然后交给纯函数 判定。 /// /// 本类是「临时适配」:中台一旦提供数据集级 freshness 接口, /// 换掉本类即可,端口与调用方不动。 /// public class S8MdpAuthorityHealthResolver : IS8AuthorityHealthResolver, ITransient { /// 合法标识符(表名 / 列名)。声明值虽由开发者书写,仍不直接拼进 SQL 而先做白名单校验。 private static readonly Regex IdentifierPattern = new("^[A-Za-z_][A-Za-z0-9_]*$", RegexOptions.Compiled); private readonly ISqlSugarClient _db; private readonly IS8DatasetCatalog _datasetCatalog; private readonly ILogger _logger; public S8MdpAuthorityHealthResolver( ISqlSugarClient db, IS8DatasetCatalog datasetCatalog, ILogger logger) { _db = db; _datasetCatalog = datasetCatalog; _logger = logger; } public async Task ResolveAsync(long tenantId, string datasetCode) { var observedAt = DateTime.Now; try { var spec = _datasetCatalog.Find(datasetCode)?.AuthoritySpec; // 未声明 AuthoritySpec / 声明不完整 → 交给纯函数判 UNKNOWN,不在这里编造默认值。 if (spec is null || string.IsNullOrWhiteSpace(spec.ProducerJobCode)) { return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation { DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt, SpecFound = false }); } // 生产者不可信时不必查库 —— 结论与运行日志无关。 if (!spec.ProducerTrusted) { return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation { DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt, SpecFound = true, ProducerTrusted = false, UntrustedReasonCode = spec.UntrustedReasonCode, StaleWindow = spec.StaleWindow }); } var lookbackFrom = observedAt - spec.LookbackWindow; var lastSuccess = await QueryLastSuccessAsync(tenantId, spec.ProducerJobCode, lookbackFrom); // 锚点:有成功则以其结束时刻为界,否则以回溯窗口起点为界。 var anchor = lastSuccess?.EndTime ?? lookbackFrom; var newerFailure = await QueryNewerTerminalFailureAsync(tenantId, spec.ProducerJobCode, anchor); var newerInFlight = await QueryNewerInFlightAsync(tenantId, spec.ProducerJobCode, anchor); var (invariant, detail) = await ResolveSnapshotInvariantAsync(tenantId, spec, lastSuccess); return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation { DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt, SpecFound = true, ProducerTrusted = true, LastSuccessEndAt = lastSuccess?.EndTime, HasNewerTerminalFailure = newerFailure is not null, NewerTerminalFailureStatus = newerFailure?.Status, OldestNewerInFlightStartAt = newerInFlight, SnapshotInvariant = invariant, SnapshotDetail = detail, StaleWindow = spec.StaleWindow }); } catch (Exception ex) { // fail-safe:探测自身失败一律判 UNKNOWN、拦截恢复。绝不因为查不动库就放行。 _logger.LogWarning(ex, "authority_health_probe_failed tenant={Tenant} dataset={Dataset}", tenantId, datasetCode); return new S8AuthorityHealthResult { DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt, State = S8AuthorityHealthState.Unknown, ReasonCode = S8AuthorityHealthReason.ResolverFailed, Reason = $"Authority 健康探测失败:{ex.GetType().Name}" }; } } /// /// 最近一次终态成功的生产运行。 /// /// 三处刻意为之,改动前先读完 /// /// 租户条件是严格相等,绝不能写成 (tenant_id = @TenantId OR tenant_id = 0)。 /// 仓内 MdpMonitorService.BuildMdpRunLogTenantWhere 就是后者,本类不得复用: /// 该表确实存在 tenant_id = 0 的平台行(实测 S1 有 111 条、 /// 且 S5_PURCHASE_RECEIPT_MDP_SYNC 有 325 条平台级 SUCCESS), /// 一条这样的行会让所有租户同时判成健康。 /// 必须显式过滤终态。该表有 298 条永不回收的 RUNNING 行(S1 占 42 条,最老 570 小时), /// 裸 ORDER BY start_time DESC LIMIT 1 会取到孤儿行,其 end_time 为 NULL, /// 后续新鲜度运算全部失效。 /// start_time 排序(走 idx_job_start 的反向索引扫描,无 filesort), /// 但新鲜度用取到那行的 end_time 计算 —— 那才是 Authority 真正变成当前态的时刻。 /// /// private async Task QueryLastSuccessAsync(long tenantId, string jobCode, DateTime lookbackFrom) { const string sql = """ SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime, JSON_UNQUOTE(JSON_EXTRACT(r.summary_json, '$.publish.batchId')) AS PublishedBatchId, JSON_EXTRACT(r.summary_json, '$.publish.currentRows') AS PublishedCurrentRows FROM mdp_transform_run_log r WHERE r.job_code = @JobCode AND r.tenant_id = @TenantId AND r.status = 'SUCCESS' AND r.end_time IS NOT NULL AND r.start_time >= @LookbackFrom ORDER BY r.start_time DESC, r.id DESC LIMIT 1 """; var rows = await _db.Ado.SqlQueryAsync(sql, new SugarParameter[] { new("@JobCode", jobCode), new("@TenantId", tenantId), new("@LookbackFrom", lookbackFrom) }); return rows.FirstOrDefault(); } /// /// 锚点之后是否出现终态非成功的运行。 /// status <> 'SUCCESS' 而非 IN ('FAILED', ...): /// 将来若新增未知的终态状态,它会被归为「失败」而不是被静默当成健康 —— 失败方向必须保守。 /// private async Task QueryNewerTerminalFailureAsync(long tenantId, string jobCode, DateTime anchor) { const string sql = """ SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime FROM mdp_transform_run_log r WHERE r.job_code = @JobCode AND r.tenant_id = @TenantId AND r.end_time IS NOT NULL AND r.status <> 'SUCCESS' AND r.start_time > @Anchor ORDER BY r.start_time DESC, r.id DESC LIMIT 1 """; var rows = await _db.Ado.SqlQueryAsync(sql, new SugarParameter[] { new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor) }); return rows.FirstOrDefault(); } /// /// 锚点之后是否存在尚未终态的运行,取最早那条的开始时刻。 /// end_time IS NULLstatus = 'RUNNING' 在该表上实测完全等价(298/298 双向成立), /// 用前者是因为它对将来新增的非终态状态同样成立。 /// private async Task QueryNewerInFlightAsync(long tenantId, string jobCode, DateTime anchor) { const string sql = """ SELECT MIN(r.start_time) AS StartTime FROM mdp_transform_run_log r WHERE r.job_code = @JobCode AND r.tenant_id = @TenantId AND r.end_time IS NULL AND r.start_time > @Anchor """; var rows = await _db.Ado.SqlQueryAsync(sql, new SugarParameter[] { new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor) }); return rows.FirstOrDefault()?.StartTime; } /// /// 快照型 Authority 的发布不变量。 /// /// 当前批次数 = 0 时返回 而非「健康的空」: /// 生产侧的发布语句在本轮零行时命中 0 行、在库里不留任何 batch marker, /// 于是「源侧本轮确实没有数据」与「发布压根没跑成」字面等同。 /// 这一点必须保守 —— 猜成健康会让真实预警被整体判成已恢复。 /// 消除它需要生产侧补一个发布标记,属下一批的工作。 /// /// 作用域只按租户。快照的真实作用域是 (租户, 工厂),但运行日志表没有 factory 列, /// 无法把两侧对齐;若某租户将来出现多工厂,本查询会看到多个当前批次并判 VIOLATED —— /// 是拦截而非放行,失败方向安全。 /// private async Task<(string Invariant, string? Detail)> ResolveSnapshotInvariantAsync( long tenantId, S8AuthoritySpec spec, RunRow? lastSuccess) { var lastSuccessBatchId = lastSuccess?.BatchId; if (!string.Equals(spec.AuthorityKind, S8AuthorityKind.PublishedSnapshot, StringComparison.Ordinal)) return (S8SnapshotInvariant.NotApplicable, null); if (string.IsNullOrWhiteSpace(spec.SnapshotTable)) return (S8SnapshotInvariant.Violated, "AuthorityKind 声明为 PUBLISHED_SNAPSHOT 但未声明 SnapshotTable"); var table = RequireIdentifier(spec.SnapshotTable, nameof(spec.SnapshotTable)); var batchCol = RequireIdentifier(spec.SnapshotBatchColumn, nameof(spec.SnapshotBatchColumn)); var flagCol = RequireIdentifier(spec.SnapshotCurrentFlagColumn, nameof(spec.SnapshotCurrentFlagColumn)); var batchSql = $""" SELECT COUNT(DISTINCT s.`{batchCol}`) AS CurrentBatches, MIN(s.`{batchCol}`) AS BatchId FROM `{table}` s WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1 """; var batchRows = await _db.Ado.SqlQueryAsync( batchSql, new SugarParameter[] { new("@TenantId", tenantId) }); var batch = batchRows.FirstOrDefault(); var currentBatches = batch?.CurrentBatches ?? 0; if (currentBatches == 0) { // 零个当前批次有两种截然不同的成因,DWD 表本身分不出来: // ① 源侧本轮确实没有数据 —— 发布语句命中 0 行,成功发布了一个空快照(合法) // ② 发布压根没跑成 —— 上一批次被退休了,新批次没被翻牌(故障) // 生产侧的发布证据是唯一能区分二者的东西,且它与 status='SUCCESS' 写在同一条 UPDATE 里, // 不存在「成功了但没有证据」的中间态。 var publishedBatchId = lastSuccess?.PublishedBatchId; var publishedRows = lastSuccess?.PublishedCurrentRows; if (!string.IsNullOrWhiteSpace(publishedBatchId) && string.Equals(publishedBatchId, lastSuccessBatchId, StringComparison.Ordinal) && publishedRows == 0) { return (S8SnapshotInvariant.Satisfied, null); } // 没有发布证据(历史行,或生产侧尚未升级)→ 维持保守判定,与升级前行为完全一致。 var why = string.IsNullOrWhiteSpace(publishedBatchId) ? "最近一次成功运行没有留下发布证据" : $"发布证据指向批次 {publishedBatchId}(当前行数 {publishedRows?.ToString() ?? "未知"})," + $"与最近一次成功运行的批次 {lastSuccessBatchId} 不一致"; return (S8SnapshotInvariant.Unconfirmed, $"{table} 在租户 {tenantId} 下当前批次数为 0,且{why}" + "(合法空快照与发布失败在库里字面等同,不猜)"); } if (currentBatches > 1) return (S8SnapshotInvariant.Violated, $"{table} 在租户 {tenantId} 下同时存在 {currentBatches} 个当前批次"); if (!string.IsNullOrWhiteSpace(lastSuccessBatchId) && !string.Equals(batch!.BatchId, lastSuccessBatchId, StringComparison.Ordinal)) { return (S8SnapshotInvariant.Violated, $"{table} 的当前批次 {batch.BatchId} 与最近一次成功运行的批次 {lastSuccessBatchId} 不一致"); } // 必填列校验用绝对计数,空快照上天然真空成立。 foreach (var column in spec.SnapshotRequiredColumns) { var col = RequireIdentifier(column, nameof(spec.SnapshotRequiredColumns)); var nullSql = $""" SELECT COUNT(*) AS NullCount FROM `{table}` s WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1 AND s.`{col}` IS NULL """; var nullRows = await _db.Ado.SqlQueryAsync( nullSql, new SugarParameter[] { new("@TenantId", tenantId) }); var nullCount = nullRows.FirstOrDefault()?.NullCount ?? 0; if (nullCount > 0) return (S8SnapshotInvariant.Violated, $"{table} 当前快照中 {col} 为空的行数 = {nullCount}"); } return (S8SnapshotInvariant.Satisfied, null); } private static string RequireIdentifier(string value, string field) { if (!IdentifierPattern.IsMatch(value)) throw new InvalidOperationException($"AuthoritySpec.{field} 不是合法标识符:{value}"); return value; } /// 运行日志投影行。属性名与 SELECT 别名一一对应。 private sealed class RunRow { public string? BatchId { get; set; } public string? Status { get; set; } public DateTime? StartTime { get; set; } public DateTime? EndTime { get; set; } /// 生产侧写入的发布证据;历史行没有该字段,取出为 null。 public string? PublishedBatchId { get; set; } /// 发布后的当前行数。0 是合法值(成功发布了一个空快照)。 public int? PublishedCurrentRows { get; set; } } /// 快照校验投影行。 private sealed class SnapshotRow { public int CurrentBatches { get; set; } public string? BatchId { get; set; } public int NullCount { get; set; } } }