| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329 |
- 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;
- /// <summary>
- /// <see cref="IS8AuthorityHealthResolver"/> 的数据中台适配实现:
- /// 从 <c>mdp_transform_run_log</c> 采集生产运行事实,必要时再校验快照表的发布不变量,
- /// 然后交给纯函数 <see cref="S8AuthorityHealthEvaluator"/> 判定。
- ///
- /// <para>本类是「临时适配」:中台一旦提供数据集级 freshness 接口,
- /// 换掉本类即可,端口与调用方不动。</para>
- /// </summary>
- public class S8MdpAuthorityHealthResolver : IS8AuthorityHealthResolver, ITransient
- {
- /// <summary>合法标识符(表名 / 列名)。声明值虽由开发者书写,仍不直接拼进 SQL 而先做白名单校验。</summary>
- 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<S8MdpAuthorityHealthResolver> _logger;
- public S8MdpAuthorityHealthResolver(
- ISqlSugarClient db,
- IS8DatasetCatalog datasetCatalog,
- ILogger<S8MdpAuthorityHealthResolver> logger)
- {
- _db = db;
- _datasetCatalog = datasetCatalog;
- _logger = logger;
- }
- public async Task<S8AuthorityHealthResult> 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}"
- };
- }
- }
- /// <summary>
- /// 最近一次<b>终态成功</b>的生产运行。
- ///
- /// <para><b>三处刻意为之,改动前先读完</b>:</para>
- /// <list type="number">
- /// <item>租户条件是严格相等,<b>绝不能</b>写成 <c>(tenant_id = @TenantId OR tenant_id = 0)</c>。
- /// 仓内 <c>MdpMonitorService.BuildMdpRunLogTenantWhere</c> 就是后者,本类不得复用:
- /// 该表确实存在 <c>tenant_id = 0</c> 的平台行(实测 S1 有 111 条、
- /// 且 <c>S5_PURCHASE_RECEIPT_MDP_SYNC</c> 有 325 条平台级 SUCCESS),
- /// 一条这样的行会让所有租户同时判成健康。</item>
- /// <item>必须显式过滤终态。该表有 298 条永不回收的 RUNNING 行(S1 占 42 条,最老 570 小时),
- /// 裸 <c>ORDER BY start_time DESC LIMIT 1</c> 会取到孤儿行,其 <c>end_time</c> 为 NULL,
- /// 后续新鲜度运算全部失效。</item>
- /// <item>按 <c>start_time</c> 排序(走 <c>idx_job_start</c> 的反向索引扫描,无 filesort),
- /// 但新鲜度用取到那行的 <c>end_time</c> 计算 —— 那才是 Authority 真正变成当前态的时刻。</item>
- /// </list>
- /// </summary>
- private async Task<RunRow?> 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<RunRow>(sql, new SugarParameter[]
- {
- new("@JobCode", jobCode), new("@TenantId", tenantId), new("@LookbackFrom", lookbackFrom)
- });
- return rows.FirstOrDefault();
- }
- /// <summary>
- /// 锚点之后是否出现终态非成功的运行。
- /// <para>用 <c>status <> 'SUCCESS'</c> 而非 <c>IN ('FAILED', ...)</c>:
- /// 将来若新增未知的终态状态,它会被归为「失败」而不是被静默当成健康 —— 失败方向必须保守。</para>
- /// </summary>
- private async Task<RunRow?> 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<RunRow>(sql, new SugarParameter[]
- {
- new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
- });
- return rows.FirstOrDefault();
- }
- /// <summary>
- /// 锚点之后是否存在尚未终态的运行,取最早那条的开始时刻。
- /// <para><c>end_time IS NULL</c> 与 <c>status = 'RUNNING'</c> 在该表上实测完全等价(298/298 双向成立),
- /// 用前者是因为它对将来新增的非终态状态同样成立。</para>
- /// </summary>
- private async Task<DateTime?> 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<RunRow>(sql, new SugarParameter[]
- {
- new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
- });
- return rows.FirstOrDefault()?.StartTime;
- }
- /// <summary>
- /// 快照型 Authority 的发布不变量。
- ///
- /// <para>当前批次数 = 0 时返回 <see cref="S8SnapshotInvariant.Unconfirmed"/> 而非「健康的空」:
- /// 生产侧的发布语句在本轮零行时命中 0 行、在库里不留任何 batch marker,
- /// 于是「源侧本轮确实没有数据」与「发布压根没跑成」字面等同。
- /// 这一点必须保守 —— 猜成健康会让真实预警被整体判成已恢复。
- /// 消除它需要生产侧补一个发布标记,属下一批的工作。</para>
- ///
- /// <para>作用域只按租户。快照的真实作用域是 (租户, 工厂),但运行日志表没有 factory 列,
- /// 无法把两侧对齐;若某租户将来出现多工厂,本查询会看到多个当前批次并判 VIOLATED ——
- /// 是拦截而非放行,失败方向安全。</para>
- /// </summary>
- 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<SnapshotRow>(
- 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<SnapshotRow>(
- 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;
- }
- /// <summary>运行日志投影行。属性名与 SELECT 别名一一对应。</summary>
- private sealed class RunRow
- {
- public string? BatchId { get; set; }
- public string? Status { get; set; }
- public DateTime? StartTime { get; set; }
- public DateTime? EndTime { get; set; }
- /// <summary>生产侧写入的发布证据;历史行没有该字段,取出为 null。</summary>
- public string? PublishedBatchId { get; set; }
- /// <summary>发布后的当前行数。0 是合法值(成功发布了一个空快照)。</summary>
- public int? PublishedCurrentRows { get; set; }
- }
- /// <summary>快照校验投影行。</summary>
- private sealed class SnapshotRow
- {
- public int CurrentBatches { get; set; }
- public string? BatchId { get; set; }
- public int NullCount { get; set; }
- }
- }
|