S8MdpAuthorityHealthResolver.cs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295
  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?.BatchId);
  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. FROM mdp_transform_run_log r
  115. WHERE r.job_code = @JobCode
  116. AND r.tenant_id = @TenantId
  117. AND r.status = 'SUCCESS'
  118. AND r.end_time IS NOT NULL
  119. AND r.start_time >= @LookbackFrom
  120. ORDER BY r.start_time DESC, r.id DESC
  121. LIMIT 1
  122. """;
  123. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  124. {
  125. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@LookbackFrom", lookbackFrom)
  126. });
  127. return rows.FirstOrDefault();
  128. }
  129. /// <summary>
  130. /// 锚点之后是否出现终态非成功的运行。
  131. /// <para>用 <c>status &lt;&gt; 'SUCCESS'</c> 而非 <c>IN ('FAILED', ...)</c>:
  132. /// 将来若新增未知的终态状态,它会被归为「失败」而不是被静默当成健康 —— 失败方向必须保守。</para>
  133. /// </summary>
  134. private async Task<RunRow?> QueryNewerTerminalFailureAsync(long tenantId, string jobCode, DateTime anchor)
  135. {
  136. const string sql =
  137. """
  138. SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime
  139. FROM mdp_transform_run_log r
  140. WHERE r.job_code = @JobCode
  141. AND r.tenant_id = @TenantId
  142. AND r.end_time IS NOT NULL
  143. AND r.status <> 'SUCCESS'
  144. AND r.start_time > @Anchor
  145. ORDER BY r.start_time DESC, r.id DESC
  146. LIMIT 1
  147. """;
  148. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  149. {
  150. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
  151. });
  152. return rows.FirstOrDefault();
  153. }
  154. /// <summary>
  155. /// 锚点之后是否存在尚未终态的运行,取最早那条的开始时刻。
  156. /// <para><c>end_time IS NULL</c> 与 <c>status = 'RUNNING'</c> 在该表上实测完全等价(298/298 双向成立),
  157. /// 用前者是因为它对将来新增的非终态状态同样成立。</para>
  158. /// </summary>
  159. private async Task<DateTime?> QueryNewerInFlightAsync(long tenantId, string jobCode, DateTime anchor)
  160. {
  161. const string sql =
  162. """
  163. SELECT MIN(r.start_time) AS StartTime
  164. FROM mdp_transform_run_log r
  165. WHERE r.job_code = @JobCode
  166. AND r.tenant_id = @TenantId
  167. AND r.end_time IS NULL
  168. AND r.start_time > @Anchor
  169. """;
  170. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  171. {
  172. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
  173. });
  174. return rows.FirstOrDefault()?.StartTime;
  175. }
  176. /// <summary>
  177. /// 快照型 Authority 的发布不变量。
  178. ///
  179. /// <para>当前批次数 = 0 时返回 <see cref="S8SnapshotInvariant.Unconfirmed"/> 而非「健康的空」:
  180. /// 生产侧的发布语句在本轮零行时命中 0 行、在库里不留任何 batch marker,
  181. /// 于是「源侧本轮确实没有数据」与「发布压根没跑成」字面等同。
  182. /// 这一点必须保守 —— 猜成健康会让真实预警被整体判成已恢复。
  183. /// 消除它需要生产侧补一个发布标记,属下一批的工作。</para>
  184. ///
  185. /// <para>作用域只按租户。快照的真实作用域是 (租户, 工厂),但运行日志表没有 factory 列,
  186. /// 无法把两侧对齐;若某租户将来出现多工厂,本查询会看到多个当前批次并判 VIOLATED ——
  187. /// 是拦截而非放行,失败方向安全。</para>
  188. /// </summary>
  189. private async Task<(string Invariant, string? Detail)> ResolveSnapshotInvariantAsync(
  190. long tenantId, S8AuthoritySpec spec, string? lastSuccessBatchId)
  191. {
  192. if (!string.Equals(spec.AuthorityKind, S8AuthorityKind.PublishedSnapshot, StringComparison.Ordinal))
  193. return (S8SnapshotInvariant.NotApplicable, null);
  194. if (string.IsNullOrWhiteSpace(spec.SnapshotTable))
  195. return (S8SnapshotInvariant.Violated, "AuthorityKind 声明为 PUBLISHED_SNAPSHOT 但未声明 SnapshotTable");
  196. var table = RequireIdentifier(spec.SnapshotTable, nameof(spec.SnapshotTable));
  197. var batchCol = RequireIdentifier(spec.SnapshotBatchColumn, nameof(spec.SnapshotBatchColumn));
  198. var flagCol = RequireIdentifier(spec.SnapshotCurrentFlagColumn, nameof(spec.SnapshotCurrentFlagColumn));
  199. var batchSql =
  200. $"""
  201. SELECT COUNT(DISTINCT s.`{batchCol}`) AS CurrentBatches, MIN(s.`{batchCol}`) AS BatchId
  202. FROM `{table}` s
  203. WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1
  204. """;
  205. var batchRows = await _db.Ado.SqlQueryAsync<SnapshotRow>(
  206. batchSql, new SugarParameter[] { new("@TenantId", tenantId) });
  207. var batch = batchRows.FirstOrDefault();
  208. var currentBatches = batch?.CurrentBatches ?? 0;
  209. if (currentBatches == 0)
  210. return (S8SnapshotInvariant.Unconfirmed,
  211. $"{table} 在租户 {tenantId} 下当前批次数为 0(合法空快照与发布失败在库里无法区分)");
  212. if (currentBatches > 1)
  213. return (S8SnapshotInvariant.Violated,
  214. $"{table} 在租户 {tenantId} 下同时存在 {currentBatches} 个当前批次");
  215. if (!string.IsNullOrWhiteSpace(lastSuccessBatchId)
  216. && !string.Equals(batch!.BatchId, lastSuccessBatchId, StringComparison.Ordinal))
  217. {
  218. return (S8SnapshotInvariant.Violated,
  219. $"{table} 的当前批次 {batch.BatchId} 与最近一次成功运行的批次 {lastSuccessBatchId} 不一致");
  220. }
  221. // 必填列校验用绝对计数,空快照上天然真空成立。
  222. foreach (var column in spec.SnapshotRequiredColumns)
  223. {
  224. var col = RequireIdentifier(column, nameof(spec.SnapshotRequiredColumns));
  225. var nullSql =
  226. $"""
  227. SELECT COUNT(*) AS NullCount
  228. FROM `{table}` s
  229. WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1 AND s.`{col}` IS NULL
  230. """;
  231. var nullRows = await _db.Ado.SqlQueryAsync<SnapshotRow>(
  232. nullSql, new SugarParameter[] { new("@TenantId", tenantId) });
  233. var nullCount = nullRows.FirstOrDefault()?.NullCount ?? 0;
  234. if (nullCount > 0)
  235. return (S8SnapshotInvariant.Violated,
  236. $"{table} 当前快照中 {col} 为空的行数 = {nullCount}");
  237. }
  238. return (S8SnapshotInvariant.Satisfied, null);
  239. }
  240. private static string RequireIdentifier(string value, string field)
  241. {
  242. if (!IdentifierPattern.IsMatch(value))
  243. throw new InvalidOperationException($"AuthoritySpec.{field} 不是合法标识符:{value}");
  244. return value;
  245. }
  246. /// <summary>运行日志投影行。属性名与 SELECT 别名一一对应。</summary>
  247. private sealed class RunRow
  248. {
  249. public string? BatchId { get; set; }
  250. public string? Status { get; set; }
  251. public DateTime? StartTime { get; set; }
  252. public DateTime? EndTime { get; set; }
  253. }
  254. /// <summary>快照校验投影行。</summary>
  255. private sealed class SnapshotRow
  256. {
  257. public int CurrentBatches { get; set; }
  258. public string? BatchId { get; set; }
  259. public int NullCount { get; set; }
  260. }
  261. }