S8MdpAuthorityHealthResolver.cs 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329
  1. using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess;
  2. using Microsoft.Extensions.Logging;
  3. using SqlSugar;
  4. using System.Text.RegularExpressions;
  5. namespace Admin.NET.Plugin.AiDOP.Service.S8.Rules.Health;
  6. /// <summary>
  7. /// <see cref="IS8AuthorityHealthResolver"/> 的数据中台适配实现:
  8. /// 从 <c>mdp_transform_run_log</c> 采集生产运行事实,必要时再校验快照表的发布不变量,
  9. /// 然后交给纯函数 <see cref="S8AuthorityHealthEvaluator"/> 判定。
  10. ///
  11. /// <para>本类是「临时适配」:中台一旦提供数据集级 freshness 接口,
  12. /// 换掉本类即可,端口与调用方不动。</para>
  13. /// </summary>
  14. public class S8MdpAuthorityHealthResolver : IS8AuthorityHealthResolver, ITransient
  15. {
  16. /// <summary>合法标识符(表名 / 列名)。声明值虽由开发者书写,仍不直接拼进 SQL 而先做白名单校验。</summary>
  17. private static readonly Regex IdentifierPattern = new("^[A-Za-z_][A-Za-z0-9_]*$", RegexOptions.Compiled);
  18. private readonly ISqlSugarClient _db;
  19. private readonly IS8DatasetCatalog _datasetCatalog;
  20. private readonly ILogger<S8MdpAuthorityHealthResolver> _logger;
  21. public S8MdpAuthorityHealthResolver(
  22. ISqlSugarClient db,
  23. IS8DatasetCatalog datasetCatalog,
  24. ILogger<S8MdpAuthorityHealthResolver> logger)
  25. {
  26. _db = db;
  27. _datasetCatalog = datasetCatalog;
  28. _logger = logger;
  29. }
  30. public async Task<S8AuthorityHealthResult> ResolveAsync(long tenantId, string datasetCode)
  31. {
  32. var observedAt = DateTime.Now;
  33. try
  34. {
  35. var spec = _datasetCatalog.Find(datasetCode)?.AuthoritySpec;
  36. // 未声明 AuthoritySpec / 声明不完整 → 交给纯函数判 UNKNOWN,不在这里编造默认值。
  37. if (spec is null || string.IsNullOrWhiteSpace(spec.ProducerJobCode))
  38. {
  39. return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation
  40. {
  41. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt, SpecFound = false
  42. });
  43. }
  44. // 生产者不可信时不必查库 —— 结论与运行日志无关。
  45. if (!spec.ProducerTrusted)
  46. {
  47. return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation
  48. {
  49. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt,
  50. SpecFound = true, ProducerTrusted = false,
  51. UntrustedReasonCode = spec.UntrustedReasonCode,
  52. StaleWindow = spec.StaleWindow
  53. });
  54. }
  55. var lookbackFrom = observedAt - spec.LookbackWindow;
  56. var lastSuccess = await QueryLastSuccessAsync(tenantId, spec.ProducerJobCode, lookbackFrom);
  57. // 锚点:有成功则以其结束时刻为界,否则以回溯窗口起点为界。
  58. var anchor = lastSuccess?.EndTime ?? lookbackFrom;
  59. var newerFailure = await QueryNewerTerminalFailureAsync(tenantId, spec.ProducerJobCode, anchor);
  60. var newerInFlight = await QueryNewerInFlightAsync(tenantId, spec.ProducerJobCode, anchor);
  61. var (invariant, detail) = await ResolveSnapshotInvariantAsync(tenantId, spec, lastSuccess);
  62. return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation
  63. {
  64. DatasetCode = datasetCode,
  65. TenantId = tenantId,
  66. ObservedAt = observedAt,
  67. SpecFound = true,
  68. ProducerTrusted = true,
  69. LastSuccessEndAt = lastSuccess?.EndTime,
  70. HasNewerTerminalFailure = newerFailure is not null,
  71. NewerTerminalFailureStatus = newerFailure?.Status,
  72. OldestNewerInFlightStartAt = newerInFlight,
  73. SnapshotInvariant = invariant,
  74. SnapshotDetail = detail,
  75. StaleWindow = spec.StaleWindow
  76. });
  77. }
  78. catch (Exception ex)
  79. {
  80. // fail-safe:探测自身失败一律判 UNKNOWN、拦截恢复。绝不因为查不动库就放行。
  81. _logger.LogWarning(ex,
  82. "authority_health_probe_failed tenant={Tenant} dataset={Dataset}", tenantId, datasetCode);
  83. return new S8AuthorityHealthResult
  84. {
  85. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt,
  86. State = S8AuthorityHealthState.Unknown,
  87. ReasonCode = S8AuthorityHealthReason.ResolverFailed,
  88. Reason = $"Authority 健康探测失败:{ex.GetType().Name}"
  89. };
  90. }
  91. }
  92. /// <summary>
  93. /// 最近一次<b>终态成功</b>的生产运行。
  94. ///
  95. /// <para><b>三处刻意为之,改动前先读完</b>:</para>
  96. /// <list type="number">
  97. /// <item>租户条件是严格相等,<b>绝不能</b>写成 <c>(tenant_id = @TenantId OR tenant_id = 0)</c>。
  98. /// 仓内 <c>MdpMonitorService.BuildMdpRunLogTenantWhere</c> 就是后者,本类不得复用:
  99. /// 该表确实存在 <c>tenant_id = 0</c> 的平台行(实测 S1 有 111 条、
  100. /// 且 <c>S5_PURCHASE_RECEIPT_MDP_SYNC</c> 有 325 条平台级 SUCCESS),
  101. /// 一条这样的行会让所有租户同时判成健康。</item>
  102. /// <item>必须显式过滤终态。该表有 298 条永不回收的 RUNNING 行(S1 占 42 条,最老 570 小时),
  103. /// 裸 <c>ORDER BY start_time DESC LIMIT 1</c> 会取到孤儿行,其 <c>end_time</c> 为 NULL,
  104. /// 后续新鲜度运算全部失效。</item>
  105. /// <item>按 <c>start_time</c> 排序(走 <c>idx_job_start</c> 的反向索引扫描,无 filesort),
  106. /// 但新鲜度用取到那行的 <c>end_time</c> 计算 —— 那才是 Authority 真正变成当前态的时刻。</item>
  107. /// </list>
  108. /// </summary>
  109. private async Task<RunRow?> QueryLastSuccessAsync(long tenantId, string jobCode, DateTime lookbackFrom)
  110. {
  111. const string sql =
  112. """
  113. SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime,
  114. JSON_UNQUOTE(JSON_EXTRACT(r.summary_json, '$.publish.batchId')) AS PublishedBatchId,
  115. JSON_EXTRACT(r.summary_json, '$.publish.currentRows') AS PublishedCurrentRows
  116. FROM mdp_transform_run_log r
  117. WHERE r.job_code = @JobCode
  118. AND r.tenant_id = @TenantId
  119. AND r.status = 'SUCCESS'
  120. AND r.end_time IS NOT NULL
  121. AND r.start_time >= @LookbackFrom
  122. ORDER BY r.start_time DESC, r.id DESC
  123. LIMIT 1
  124. """;
  125. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  126. {
  127. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@LookbackFrom", lookbackFrom)
  128. });
  129. return rows.FirstOrDefault();
  130. }
  131. /// <summary>
  132. /// 锚点之后是否出现终态非成功的运行。
  133. /// <para>用 <c>status &lt;&gt; 'SUCCESS'</c> 而非 <c>IN ('FAILED', ...)</c>:
  134. /// 将来若新增未知的终态状态,它会被归为「失败」而不是被静默当成健康 —— 失败方向必须保守。</para>
  135. /// </summary>
  136. private async Task<RunRow?> QueryNewerTerminalFailureAsync(long tenantId, string jobCode, DateTime anchor)
  137. {
  138. const string sql =
  139. """
  140. SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime
  141. FROM mdp_transform_run_log r
  142. WHERE r.job_code = @JobCode
  143. AND r.tenant_id = @TenantId
  144. AND r.end_time IS NOT NULL
  145. AND r.status <> 'SUCCESS'
  146. AND r.start_time > @Anchor
  147. ORDER BY r.start_time DESC, r.id DESC
  148. LIMIT 1
  149. """;
  150. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  151. {
  152. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
  153. });
  154. return rows.FirstOrDefault();
  155. }
  156. /// <summary>
  157. /// 锚点之后是否存在尚未终态的运行,取最早那条的开始时刻。
  158. /// <para><c>end_time IS NULL</c> 与 <c>status = 'RUNNING'</c> 在该表上实测完全等价(298/298 双向成立),
  159. /// 用前者是因为它对将来新增的非终态状态同样成立。</para>
  160. /// </summary>
  161. private async Task<DateTime?> QueryNewerInFlightAsync(long tenantId, string jobCode, DateTime anchor)
  162. {
  163. const string sql =
  164. """
  165. SELECT MIN(r.start_time) AS StartTime
  166. FROM mdp_transform_run_log r
  167. WHERE r.job_code = @JobCode
  168. AND r.tenant_id = @TenantId
  169. AND r.end_time IS NULL
  170. AND r.start_time > @Anchor
  171. """;
  172. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  173. {
  174. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
  175. });
  176. return rows.FirstOrDefault()?.StartTime;
  177. }
  178. /// <summary>
  179. /// 快照型 Authority 的发布不变量。
  180. ///
  181. /// <para>当前批次数 = 0 时返回 <see cref="S8SnapshotInvariant.Unconfirmed"/> 而非「健康的空」:
  182. /// 生产侧的发布语句在本轮零行时命中 0 行、在库里不留任何 batch marker,
  183. /// 于是「源侧本轮确实没有数据」与「发布压根没跑成」字面等同。
  184. /// 这一点必须保守 —— 猜成健康会让真实预警被整体判成已恢复。
  185. /// 消除它需要生产侧补一个发布标记,属下一批的工作。</para>
  186. ///
  187. /// <para>作用域只按租户。快照的真实作用域是 (租户, 工厂),但运行日志表没有 factory 列,
  188. /// 无法把两侧对齐;若某租户将来出现多工厂,本查询会看到多个当前批次并判 VIOLATED ——
  189. /// 是拦截而非放行,失败方向安全。</para>
  190. /// </summary>
  191. private async Task<(string Invariant, string? Detail)> ResolveSnapshotInvariantAsync(
  192. long tenantId, S8AuthoritySpec spec, RunRow? lastSuccess)
  193. {
  194. var lastSuccessBatchId = lastSuccess?.BatchId;
  195. if (!string.Equals(spec.AuthorityKind, S8AuthorityKind.PublishedSnapshot, StringComparison.Ordinal))
  196. return (S8SnapshotInvariant.NotApplicable, null);
  197. if (string.IsNullOrWhiteSpace(spec.SnapshotTable))
  198. return (S8SnapshotInvariant.Violated, "AuthorityKind 声明为 PUBLISHED_SNAPSHOT 但未声明 SnapshotTable");
  199. var table = RequireIdentifier(spec.SnapshotTable, nameof(spec.SnapshotTable));
  200. var batchCol = RequireIdentifier(spec.SnapshotBatchColumn, nameof(spec.SnapshotBatchColumn));
  201. var flagCol = RequireIdentifier(spec.SnapshotCurrentFlagColumn, nameof(spec.SnapshotCurrentFlagColumn));
  202. var batchSql =
  203. $"""
  204. SELECT COUNT(DISTINCT s.`{batchCol}`) AS CurrentBatches, MIN(s.`{batchCol}`) AS BatchId
  205. FROM `{table}` s
  206. WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1
  207. """;
  208. var batchRows = await _db.Ado.SqlQueryAsync<SnapshotRow>(
  209. batchSql, new SugarParameter[] { new("@TenantId", tenantId) });
  210. var batch = batchRows.FirstOrDefault();
  211. var currentBatches = batch?.CurrentBatches ?? 0;
  212. if (currentBatches == 0)
  213. {
  214. // 零个当前批次有两种截然不同的成因,DWD 表本身分不出来:
  215. // ① 源侧本轮确实没有数据 —— 发布语句命中 0 行,成功发布了一个空快照(合法)
  216. // ② 发布压根没跑成 —— 上一批次被退休了,新批次没被翻牌(故障)
  217. // 生产侧的发布证据是唯一能区分二者的东西,且它与 status='SUCCESS' 写在同一条 UPDATE 里,
  218. // 不存在「成功了但没有证据」的中间态。
  219. var publishedBatchId = lastSuccess?.PublishedBatchId;
  220. var publishedRows = lastSuccess?.PublishedCurrentRows;
  221. if (!string.IsNullOrWhiteSpace(publishedBatchId)
  222. && string.Equals(publishedBatchId, lastSuccessBatchId, StringComparison.Ordinal)
  223. && publishedRows == 0)
  224. {
  225. return (S8SnapshotInvariant.Satisfied, null);
  226. }
  227. // 没有发布证据(历史行,或生产侧尚未升级)→ 维持保守判定,与升级前行为完全一致。
  228. var why = string.IsNullOrWhiteSpace(publishedBatchId)
  229. ? "最近一次成功运行没有留下发布证据"
  230. : $"发布证据指向批次 {publishedBatchId}(当前行数 {publishedRows?.ToString() ?? "未知"})," +
  231. $"与最近一次成功运行的批次 {lastSuccessBatchId} 不一致";
  232. return (S8SnapshotInvariant.Unconfirmed,
  233. $"{table} 在租户 {tenantId} 下当前批次数为 0,且{why}" +
  234. "(合法空快照与发布失败在库里字面等同,不猜)");
  235. }
  236. if (currentBatches > 1)
  237. return (S8SnapshotInvariant.Violated,
  238. $"{table} 在租户 {tenantId} 下同时存在 {currentBatches} 个当前批次");
  239. if (!string.IsNullOrWhiteSpace(lastSuccessBatchId)
  240. && !string.Equals(batch!.BatchId, lastSuccessBatchId, StringComparison.Ordinal))
  241. {
  242. return (S8SnapshotInvariant.Violated,
  243. $"{table} 的当前批次 {batch.BatchId} 与最近一次成功运行的批次 {lastSuccessBatchId} 不一致");
  244. }
  245. // 必填列校验用绝对计数,空快照上天然真空成立。
  246. foreach (var column in spec.SnapshotRequiredColumns)
  247. {
  248. var col = RequireIdentifier(column, nameof(spec.SnapshotRequiredColumns));
  249. var nullSql =
  250. $"""
  251. SELECT COUNT(*) AS NullCount
  252. FROM `{table}` s
  253. WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1 AND s.`{col}` IS NULL
  254. """;
  255. var nullRows = await _db.Ado.SqlQueryAsync<SnapshotRow>(
  256. nullSql, new SugarParameter[] { new("@TenantId", tenantId) });
  257. var nullCount = nullRows.FirstOrDefault()?.NullCount ?? 0;
  258. if (nullCount > 0)
  259. return (S8SnapshotInvariant.Violated,
  260. $"{table} 当前快照中 {col} 为空的行数 = {nullCount}");
  261. }
  262. return (S8SnapshotInvariant.Satisfied, null);
  263. }
  264. private static string RequireIdentifier(string value, string field)
  265. {
  266. if (!IdentifierPattern.IsMatch(value))
  267. throw new InvalidOperationException($"AuthoritySpec.{field} 不是合法标识符:{value}");
  268. return value;
  269. }
  270. /// <summary>运行日志投影行。属性名与 SELECT 别名一一对应。</summary>
  271. private sealed class RunRow
  272. {
  273. public string? BatchId { get; set; }
  274. public string? Status { get; set; }
  275. public DateTime? StartTime { get; set; }
  276. public DateTime? EndTime { get; set; }
  277. /// <summary>生产侧写入的发布证据;历史行没有该字段,取出为 null。</summary>
  278. public string? PublishedBatchId { get; set; }
  279. /// <summary>发布后的当前行数。0 是合法值(成功发布了一个空快照)。</summary>
  280. public int? PublishedCurrentRows { get; set; }
  281. }
  282. /// <summary>快照校验投影行。</summary>
  283. private sealed class SnapshotRow
  284. {
  285. public int CurrentBatches { get; set; }
  286. public string? BatchId { get; set; }
  287. public int NullCount { get; set; }
  288. }
  289. }