S8TimeoutRuleEvaluator.cs 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196
  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 Admin.NET.Plugin.AiDOP.Service.S8.Rules.Definitions;
  6. using Microsoft.Extensions.Logging;
  7. namespace Admin.NET.Plugin.AiDOP.Service.S8.Rules;
  8. /// <summary>
  9. /// TIMEOUT 类规则 evaluator。
  10. /// 判定:dueAt &lt;= now - graceMinutes 且 status 不在 completedStates 内 → HIT。
  11. /// 不做严重度阶梯、不做 SLA 升级、不做事件触发。
  12. ///
  13. /// S8-DATASET-PROVIDER-FOUNDATION-1:取数由 <see cref="S8MonitoringDataGateway"/> 承担,
  14. /// 本类不直接依赖 DataTable / 数据源 / SQL。
  15. ///
  16. /// <para><b>S8-RULE-GOVERNANCE-BATCH1:判定语义改由代码定义供给。</b>
  17. /// 此前 <c>dueAtField</c> / <c>statusField</c> / <c>completedStates</c> / <c>objectIdField</c> /
  18. /// <c>exceptionTypeCode</c> 全部从 <c>params_json</c> 解析 —— 而 <c>params_json</c> 是配置页可整块
  19. /// 覆写的自由文本。于是「这条规则判什么」实际由业务用户决定,改一个字符串就换一套业务含义。
  20. /// 现在这些一律取自 <see cref="S8RuleDefinition"/>(<see cref="IS8RuleCatalog"/> 供给),
  21. /// params_json 只剩 graceMinutes 等真正的运行策略。</para>
  22. ///
  23. /// <para>连带修复:params_json 为 NULL 不再抛 <c>rule_not_configured</c> —— 定义在代码里,
  24. /// 一条没有运行参数的规则完全可以按默认值跑。</para>
  25. /// </summary>
  26. public class S8TimeoutRuleEvaluator : IS8RuleEvaluator, ITransient
  27. {
  28. public const string RuleTypeCode = "TIMEOUT";
  29. public string RuleType => RuleTypeCode;
  30. private readonly S8MonitoringDataGateway _dataGateway;
  31. private readonly IS8RuleCatalog _ruleCatalog;
  32. private readonly ILogger<S8TimeoutRuleEvaluator> _logger;
  33. public S8TimeoutRuleEvaluator(
  34. S8MonitoringDataGateway dataGateway,
  35. IS8RuleCatalog ruleCatalog,
  36. ILogger<S8TimeoutRuleEvaluator> logger)
  37. {
  38. _dataGateway = dataGateway;
  39. _ruleCatalog = ruleCatalog;
  40. _logger = logger;
  41. }
  42. public async Task<List<S8RuleHit>> EvaluateAsync(
  43. long tenantId,
  44. AdoS8WatchRule rule,
  45. CancellationToken cancellationToken = default)
  46. {
  47. // R5 evaluator 失败语义保护:所有"非命中判定"路径均改为抛出 S8RuleEvaluatorException,
  48. // 由 SchedulerService 标记 evaluate_failed 并跳过 recovery reconcile,避免对未确认未命中的 rule 误标 recovered_at。
  49. //
  50. // S8-RULE-GOVERNANCE-BATCH1:第一道门是「这条规则在当前代码版本里有没有定义」。
  51. // 没有定义 = 它的判定语义无处可取,绝不允许猜(历史实现会去 params_json 里猜,
  52. // 于是任何人建一条规则、写一段 JSON 就能让调度器替他跑)。
  53. var effective = ResolveEffective(rule);
  54. // S8-SQL-EVALUATOR-GUARD-P2-1:每次评估解析 maxRows(env 优先,回退代码默认)。
  55. // S8-STANDARD-DATASET-HARD-CUTOVER-1:不再解析 commandTimeout —— S8 已不执行 SQL,
  56. // 超时归 Provider 自己的取数实现负责。
  57. var maxRows = S8EvaluatorGuard.ResolveMaxRows(_logger);
  58. var data = await _dataGateway.LoadAsync(
  59. tenantId, rule, RuleTypeCode, maxRows, cancellationToken);
  60. return EvaluateRows(data.RowSet, effective, tenantId, DateTime.Now);
  61. }
  62. /// <summary>
  63. /// 解析生效规则:定义必须存在、类型必须匹配、TIMEOUT 语义必须齐备。
  64. /// 三者任一不满足都是**发布态问题**,不是数据问题,故 fail-fast 而非静默返空命中。
  65. /// </summary>
  66. private S8EffectiveRule ResolveEffective(AdoS8WatchRule rule)
  67. {
  68. var definition = _ruleCatalog.TryGet(rule.RuleCode)
  69. ?? throw new S8RuleEvaluatorException(
  70. S8RuleCatalog.ReasonNotFound,
  71. $"规则 {rule.RuleCode} 在当前版本中没有代码定义,不予执行");
  72. if (!string.Equals(definition.RuleType, RuleTypeCode, StringComparison.Ordinal))
  73. throw new S8RuleEvaluatorException(
  74. "rule_type_mismatch",
  75. $"规则 {rule.RuleCode} 的代码定义类型为 {definition.RuleType},不能由 TIMEOUT evaluator 执行");
  76. if (definition.Timeout == null)
  77. throw new S8RuleEvaluatorException(
  78. "rule_semantics_missing",
  79. $"规则 {rule.RuleCode} 的代码定义缺少 TIMEOUT 判定语义");
  80. return S8EffectiveRule.Resolve(rule, definition);
  81. }
  82. /// <summary>
  83. /// TIMEOUT 判定核心:只消费 canonical 行,不接触 SQL / DataTable / 数据源。
  84. /// 判定算法与迁移前逐行一致;internal 暴露供契约测试直接驱动。
  85. /// </summary>
  86. internal static List<S8RuleHit> EvaluateRows(
  87. S8MonitoringRowSet rowSet,
  88. S8EffectiveRule effective,
  89. long tenantId,
  90. DateTime detectedAt)
  91. {
  92. var definition = effective.Definition;
  93. var semantics = definition.Timeout!;
  94. var parameters = effective.Parameters;
  95. var rule = effective.Row;
  96. var hits = new List<S8RuleHit>();
  97. var threshold = detectedAt.AddMinutes(-parameters.GraceMinutes);
  98. // 源对象类型来自定义。此前是 rule.SourceObjectType ?? rule.WatchObjectType 两列回落 ——
  99. // 两列都是用户可写的,而该值直接进 dedup_key,改动即造成历史异常断代。
  100. var sourceObjectType = definition.SourceObjectType;
  101. // 判定列来自定义。Provider 按 canonical 名产出行,因此这些值就是 canonical 名;
  102. // 若结果集确实没有该列,读取自然得到 null → 该行不命中,不再退回"某个 JSON 里写的列名"。
  103. var statusColumn = semantics.StatusColumn;
  104. var dueAtColumn = semantics.DueAtColumn;
  105. var objectCodeColumn = semantics.RelatedObjectCodeColumn;
  106. var objectIdColumn = semantics.SourceObjectIdColumn;
  107. foreach (var row in rowSet.Rows)
  108. {
  109. var status = row.GetString(statusColumn) ?? string.Empty;
  110. if (semantics.CompletedStates.Contains(status, StringComparer.OrdinalIgnoreCase))
  111. continue;
  112. var due = row.GetDateTime(dueAtColumn);
  113. if (due == null || due > threshold) continue;
  114. var relatedObjectCode = row.GetString(objectCodeColumn) ?? string.Empty;
  115. if (string.IsNullOrWhiteSpace(relatedObjectCode)) continue;
  116. var sourceObjectId = row.GetString(objectIdColumn) ?? relatedObjectCode;
  117. var dedupKey = BuildDedupKey(tenantId, definition.RuleCode, sourceObjectType, sourceObjectId);
  118. hits.Add(new S8RuleHit
  119. {
  120. SourceRuleId = rule.Id,
  121. SourceRuleCode = definition.RuleCode,
  122. SourceObjectType = sourceObjectType,
  123. SourceObjectId = sourceObjectId,
  124. RelatedObjectCode = relatedObjectCode,
  125. ExceptionTypeCode = definition.ExceptionTypeCode,
  126. SceneCode = definition.SceneCode,
  127. Severity = parameters.Severity,
  128. DedupKey = dedupKey,
  129. SourcePayload = BuildPayload(row, sourceObjectType, sourceObjectId, due.Value, status, definition, parameters),
  130. DetectedAt = detectedAt,
  131. Title = $"[超时] {sourceObjectType} {sourceObjectId} 已超期至 {due.Value:yyyy-MM-dd HH:mm:ss}(状态 {status})",
  132. OccurrenceDeptId = row.GetLong(S8CanonicalColumns.OccurrenceDeptId) ?? parameters.DefaultOccurrenceDeptId,
  133. ResponsibleDeptId = row.GetLong(S8CanonicalColumns.ResponsibleDeptId) ?? parameters.DefaultResponsibleDeptId
  134. });
  135. }
  136. return hits;
  137. }
  138. /// <summary>
  139. /// 构造 dedup_key:<c>T{tenant}:R{ruleCode}:{sourceObjectType}:{sourceObjectId}</c>。
  140. ///
  141. /// <para><b>S8-TENANT-ONLY-BATCH5:去掉了 <c>:F{factory}</c> 段。</b>
  142. /// 租户隔离由首段 <c>T{tenant}</c> 与各查询的 <c>tenant_id</c> 谓词双重保证;
  143. /// 工厂段既不提供隔离,又让同一个业务对象在不同工厂元数据下被当成两个异常。</para>
  144. ///
  145. /// <para>真库证据(迁移前取证):<c>dwd_supplier_delivery</c> 全表 39854 行 / 4 租户 / 97 个快照,
  146. /// <c>(tenant_id, stat_date, po_no, po_line)</c> 重复组为 0,跨工厂碰撞组为 0。
  147. /// 若将来某个规则的 SourceObjectId 在租户内不唯一,唯一性应由该规则的
  148. /// <b>SourceObjectId contract</b> 负责(例如把工厂编号并进对象标识),
  149. /// 而不是把工厂重新变回隔离维度。</para>
  150. ///
  151. /// internal 暴露供测试。
  152. /// </summary>
  153. internal static string BuildDedupKey(long tenantId, string ruleCode, string sourceObjectType, string sourceObjectId) =>
  154. $"T{tenantId}:R{ruleCode}:{sourceObjectType}:{sourceObjectId}";
  155. private static string BuildPayload(
  156. S8MonitoringRow row, string sourceObjectType, string sourceObjectId, DateTime dueAt, string status,
  157. Definitions.S8RuleDefinition definition, S8RuleRuntimeParameters parameters)
  158. {
  159. var payload = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  160. foreach (var kv in row.Values)
  161. payload[kv.Key] = kv.Value;
  162. payload["__ruleType"] = RuleTypeCode;
  163. payload["__sourceObjectType"] = sourceObjectType;
  164. payload["__sourceObjectId"] = sourceObjectId;
  165. payload["__dueAt"] = dueAt;
  166. payload["__status"] = status;
  167. payload["__graceMinutes"] = parameters.GraceMinutes;
  168. payload["__exceptionTypeCode"] = definition.ExceptionTypeCode;
  169. return JsonSerializer.Serialize(payload);
  170. }
  171. }