S8TimeoutRuleEvaluator.cs 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171
  1. using System.Text.Json;
  2. using Admin.NET.Plugin.AiDOP.Entity.S8;
  3. using Admin.NET.Plugin.AiDOP.Infrastructure.S8;
  4. using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess;
  5. using Microsoft.Extensions.Logging;
  6. namespace Admin.NET.Plugin.AiDOP.Service.S8.Rules;
  7. /// <summary>
  8. /// R2 TIMEOUT 类规则 evaluator MVP。
  9. /// params_json 约定(首版最小集合):
  10. /// { dueAtField, statusField, completedStates[], objectCodeField, objectIdField, graceMinutes, exceptionTypeCode }
  11. /// 判定:dueAt &lt;= now - graceMinutes 且 status 不在 completedStates 内 → HIT。
  12. /// 不做严重度阶梯、不做 SLA 升级、不做事件触发。
  13. ///
  14. /// S8-DATASET-PROVIDER-FOUNDATION-1:取数改由 <see cref="S8MonitoringDataGateway"/> 承担,
  15. /// 本类不再直接依赖 DataTable / 数据源 / SQL。判定算法与迁移前逐行一致。
  16. /// </summary>
  17. public class S8TimeoutRuleEvaluator : IS8RuleEvaluator, ITransient
  18. {
  19. public const string RuleTypeCode = "TIMEOUT";
  20. public string RuleType => RuleTypeCode;
  21. // S8-WATCH-EXPRESSION-COLUMN-CONTRACT-FIX-1:S8ConfigDraftService.BuildExpression 统一把结果列
  22. // 别名为以下 canonical 名(无论源表真实列名为何)。evaluator 优先按 canonical 读取,仅当结果集
  23. // 不含 canonical 列时才回退到 params_json 指定的真实列名(兼容历史未别名规则)。
  24. private const string CanonicalDueAtColumn = S8CanonicalColumns.DueAt;
  25. private const string CanonicalStatusColumn = S8CanonicalColumns.Status;
  26. private const string CanonicalSourceObjectIdColumn = S8CanonicalColumns.SourceObjectId;
  27. private const string CanonicalRelatedObjectCodeColumn = S8CanonicalColumns.RelatedObjectCode;
  28. private readonly S8MonitoringDataGateway _dataGateway;
  29. private readonly ILogger<S8TimeoutRuleEvaluator> _logger;
  30. public S8TimeoutRuleEvaluator(
  31. S8MonitoringDataGateway dataGateway,
  32. ILogger<S8TimeoutRuleEvaluator> logger)
  33. {
  34. _dataGateway = dataGateway;
  35. _logger = logger;
  36. }
  37. public async Task<List<S8RuleHit>> EvaluateAsync(
  38. long tenantId,
  39. long factoryId,
  40. AdoS8WatchRule rule,
  41. IReadOnlyList<AdoS8AlertRule> alertRules,
  42. CancellationToken cancellationToken = default)
  43. {
  44. // R5 evaluator 失败语义保护:所有"非命中判定"路径均改为抛出 S8RuleEvaluatorException,
  45. // 由 SchedulerService 标记 evaluate_failed 并跳过 recovery reconcile,避免对未确认未命中的 rule 误标 recovered_at。
  46. // LEGACY_SQL 仍要求 expression 存在;STANDARD_DATASET 不依赖 expression。
  47. var accessMode = S8DataAccessMode.Resolve(rule.DataAccessMode);
  48. var expressionRequired = accessMode == S8DataAccessMode.LegacySql;
  49. if ((expressionRequired && string.IsNullOrWhiteSpace(rule.Expression))
  50. || string.IsNullOrWhiteSpace(rule.ParamsJson))
  51. throw new S8RuleEvaluatorException("rule_not_configured", $"TIMEOUT 规则 {rule.RuleCode} 缺少 expression 或 params_json");
  52. S8TimeoutParams parameters;
  53. try { parameters = S8TimeoutParams.Parse(rule.ParamsJson!); }
  54. catch (Exception ex) { throw new S8RuleEvaluatorException("params_parse_failed", $"TIMEOUT 规则 {rule.RuleCode} params_json 解析失败:{ex.Message}", ex); }
  55. if (string.IsNullOrWhiteSpace(parameters.DueAtField)
  56. || string.IsNullOrWhiteSpace(parameters.StatusField)
  57. || string.IsNullOrWhiteSpace(parameters.ExceptionTypeCode))
  58. throw new S8RuleEvaluatorException("params_schema_invalid", $"TIMEOUT 规则 {rule.RuleCode} params 缺少必填字段 dueAtField/statusField/exceptionTypeCode");
  59. // S8-SQL-EVALUATOR-GUARD-P2-1:每次评估解析 timeout / maxRows(env 优先,回退代码默认)。
  60. var timeoutSeconds = S8EvaluatorGuard.ResolveCommandTimeoutSeconds(_logger);
  61. var maxRows = S8EvaluatorGuard.ResolveMaxRows(_logger);
  62. var data = await _dataGateway.LoadAsync(
  63. tenantId, factoryId, rule, RuleTypeCode, timeoutSeconds, maxRows, cancellationToken);
  64. return EvaluateRows(data.RowSet, parameters, rule, tenantId, factoryId, data.DataSourceId, DateTime.Now);
  65. }
  66. /// <summary>
  67. /// TIMEOUT 判定核心:只消费 canonical 行,不接触 SQL / DataTable / 数据源。
  68. /// 判定算法与迁移前逐行一致;internal 暴露供契约测试直接驱动。
  69. /// </summary>
  70. internal static List<S8RuleHit> EvaluateRows(
  71. S8MonitoringRowSet rowSet,
  72. S8TimeoutParams parameters,
  73. AdoS8WatchRule rule,
  74. long tenantId,
  75. long factoryId,
  76. long dataSourceId,
  77. DateTime detectedAt)
  78. {
  79. var hits = new List<S8RuleHit>();
  80. var threshold = detectedAt.AddMinutes(-parameters.GraceMinutes);
  81. var sourceObjectType = string.IsNullOrWhiteSpace(rule.SourceObjectType)
  82. ? rule.WatchObjectType
  83. : rule.SourceObjectType!;
  84. // 结果列名一次性解析(结果集列在整个结果集内稳定):canonical 优先,缺失回退 params 真实列名。
  85. var statusColumn = ResolveResultColumn(rowSet, CanonicalStatusColumn, parameters.StatusField);
  86. var dueAtColumn = ResolveResultColumn(rowSet, CanonicalDueAtColumn, parameters.DueAtField);
  87. var objectCodeColumn = ResolveResultColumn(rowSet, CanonicalRelatedObjectCodeColumn, parameters.ObjectCodeField);
  88. var objectIdColumn = ResolveResultColumn(rowSet, CanonicalSourceObjectIdColumn, parameters.ObjectIdField);
  89. foreach (var row in rowSet.Rows)
  90. {
  91. var status = row.GetString(statusColumn) ?? string.Empty;
  92. if (parameters.CompletedStates.Contains(status, StringComparer.OrdinalIgnoreCase))
  93. continue;
  94. var due = row.GetDateTime(dueAtColumn);
  95. if (due == null || due > threshold) continue;
  96. var relatedObjectCode = row.GetString(objectCodeColumn) ?? string.Empty;
  97. if (string.IsNullOrWhiteSpace(relatedObjectCode)) continue;
  98. var sourceObjectId = row.GetString(objectIdColumn) ?? relatedObjectCode;
  99. var dedupKey = BuildDedupKey(tenantId, factoryId, rule.RuleCode, sourceObjectType, sourceObjectId);
  100. hits.Add(new S8RuleHit
  101. {
  102. SourceRuleId = rule.Id,
  103. SourceRuleCode = rule.RuleCode,
  104. SourceObjectType = sourceObjectType,
  105. SourceObjectId = sourceObjectId,
  106. RelatedObjectCode = relatedObjectCode,
  107. ExceptionTypeCode = parameters.ExceptionTypeCode!,
  108. SceneCode = rule.SceneCode,
  109. Severity = S8SeverityCode.Normalize(rule.Severity),
  110. DedupKey = dedupKey,
  111. SourcePayload = BuildPayload(row, sourceObjectType, sourceObjectId, due.Value, status, parameters),
  112. DetectedAt = detectedAt,
  113. Title = $"[超时] {sourceObjectType} {sourceObjectId} 已超期至 {due.Value:yyyy-MM-dd HH:mm:ss}(状态 {status})",
  114. DataSourceId = dataSourceId,
  115. OccurrenceDeptId = row.GetLong(S8CanonicalColumns.OccurrenceDeptId),
  116. ResponsibleDeptId = row.GetLong(S8CanonicalColumns.ResponsibleDeptId)
  117. });
  118. }
  119. return hits;
  120. }
  121. /// <summary>构造 R2 dedup_key 稳定字符串:T{tenant}:F{factory}:R{ruleCode}:{sourceObjectType}:{sourceObjectId}。internal 暴露供测试。</summary>
  122. internal static string BuildDedupKey(long tenantId, long factoryId, string ruleCode, string sourceObjectType, string sourceObjectId) =>
  123. $"T{tenantId}:F{factoryId}:R{ruleCode}:{sourceObjectType}:{sourceObjectId}";
  124. private static string BuildPayload(S8MonitoringRow row, string sourceObjectType, string sourceObjectId, DateTime dueAt, string status, S8TimeoutParams parameters)
  125. {
  126. var payload = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  127. foreach (var kv in row.Values)
  128. payload[kv.Key] = kv.Value;
  129. payload["__ruleType"] = RuleTypeCode;
  130. payload["__sourceObjectType"] = sourceObjectType;
  131. payload["__sourceObjectId"] = sourceObjectId;
  132. payload["__dueAt"] = dueAt;
  133. payload["__status"] = status;
  134. payload["__graceMinutes"] = parameters.GraceMinutes;
  135. payload["__exceptionTypeCode"] = parameters.ExceptionTypeCode;
  136. return JsonSerializer.Serialize(payload);
  137. }
  138. /// <summary>
  139. /// 结果列名解析:BuildExpression 已把结果列统一别名为 canonical(due_at/status/source_object_id/
  140. /// related_object_code)。优先返回 canonical 列名;仅当结果集不含 canonical 列时,回退到 params_json
  141. /// 指定的真实列名(兼容历史未别名规则)。仅在 canonical 与 params 字段之间二选一,不新增无依据兜底字段。
  142. /// </summary>
  143. private static string ResolveResultColumn(S8MonitoringRowSet rowSet, string canonicalColumn, string? paramsColumn)
  144. {
  145. if (rowSet.HasColumn(canonicalColumn)) return canonicalColumn;
  146. return string.IsNullOrWhiteSpace(paramsColumn) ? canonicalColumn : paramsColumn!;
  147. }
  148. }