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 NULL 与 status = '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; }
}
}