using System.Text.Json; using Admin.NET.Plugin.AiDOP.Entity.S8; using Admin.NET.Plugin.AiDOP.Infrastructure.S8; using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.Service.S8.Rules; /// /// R2 TIMEOUT 类规则 evaluator MVP。 /// params_json 约定(首版最小集合): /// { dueAtField, statusField, completedStates[], objectCodeField, objectIdField, graceMinutes, exceptionTypeCode } /// 判定:dueAt <= now - graceMinutes 且 status 不在 completedStates 内 → HIT。 /// 不做严重度阶梯、不做 SLA 升级、不做事件触发。 /// /// S8-DATASET-PROVIDER-FOUNDATION-1:取数改由 承担, /// 本类不再直接依赖 DataTable / 数据源 / SQL。判定算法与迁移前逐行一致。 /// public class S8TimeoutRuleEvaluator : IS8RuleEvaluator, ITransient { public const string RuleTypeCode = "TIMEOUT"; public string RuleType => RuleTypeCode; // S8-WATCH-EXPRESSION-COLUMN-CONTRACT-FIX-1:S8ConfigDraftService.BuildExpression 统一把结果列 // 别名为以下 canonical 名(无论源表真实列名为何)。evaluator 优先按 canonical 读取,仅当结果集 // 不含 canonical 列时才回退到 params_json 指定的真实列名(兼容历史未别名规则)。 private const string CanonicalDueAtColumn = S8CanonicalColumns.DueAt; private const string CanonicalStatusColumn = S8CanonicalColumns.Status; private const string CanonicalSourceObjectIdColumn = S8CanonicalColumns.SourceObjectId; private const string CanonicalRelatedObjectCodeColumn = S8CanonicalColumns.RelatedObjectCode; private readonly S8MonitoringDataGateway _dataGateway; private readonly ILogger _logger; public S8TimeoutRuleEvaluator( S8MonitoringDataGateway dataGateway, ILogger logger) { _dataGateway = dataGateway; _logger = logger; } public async Task> EvaluateAsync( long tenantId, long factoryId, AdoS8WatchRule rule, IReadOnlyList alertRules, CancellationToken cancellationToken = default) { // R5 evaluator 失败语义保护:所有"非命中判定"路径均改为抛出 S8RuleEvaluatorException, // 由 SchedulerService 标记 evaluate_failed 并跳过 recovery reconcile,避免对未确认未命中的 rule 误标 recovered_at。 // LEGACY_SQL 仍要求 expression 存在;STANDARD_DATASET 不依赖 expression。 var accessMode = S8DataAccessMode.Resolve(rule.DataAccessMode); var expressionRequired = accessMode == S8DataAccessMode.LegacySql; if ((expressionRequired && string.IsNullOrWhiteSpace(rule.Expression)) || string.IsNullOrWhiteSpace(rule.ParamsJson)) throw new S8RuleEvaluatorException("rule_not_configured", $"TIMEOUT 规则 {rule.RuleCode} 缺少 expression 或 params_json"); S8TimeoutParams parameters; try { parameters = S8TimeoutParams.Parse(rule.ParamsJson!); } catch (Exception ex) { throw new S8RuleEvaluatorException("params_parse_failed", $"TIMEOUT 规则 {rule.RuleCode} params_json 解析失败:{ex.Message}", ex); } if (string.IsNullOrWhiteSpace(parameters.DueAtField) || string.IsNullOrWhiteSpace(parameters.StatusField) || string.IsNullOrWhiteSpace(parameters.ExceptionTypeCode)) throw new S8RuleEvaluatorException("params_schema_invalid", $"TIMEOUT 规则 {rule.RuleCode} params 缺少必填字段 dueAtField/statusField/exceptionTypeCode"); // S8-SQL-EVALUATOR-GUARD-P2-1:每次评估解析 timeout / maxRows(env 优先,回退代码默认)。 var timeoutSeconds = S8EvaluatorGuard.ResolveCommandTimeoutSeconds(_logger); var maxRows = S8EvaluatorGuard.ResolveMaxRows(_logger); var data = await _dataGateway.LoadAsync( tenantId, factoryId, rule, RuleTypeCode, timeoutSeconds, maxRows, cancellationToken); return EvaluateRows(data.RowSet, parameters, rule, tenantId, factoryId, data.DataSourceId, DateTime.Now); } /// /// TIMEOUT 判定核心:只消费 canonical 行,不接触 SQL / DataTable / 数据源。 /// 判定算法与迁移前逐行一致;internal 暴露供契约测试直接驱动。 /// internal static List EvaluateRows( S8MonitoringRowSet rowSet, S8TimeoutParams parameters, AdoS8WatchRule rule, long tenantId, long factoryId, long dataSourceId, DateTime detectedAt) { var hits = new List(); var threshold = detectedAt.AddMinutes(-parameters.GraceMinutes); var sourceObjectType = string.IsNullOrWhiteSpace(rule.SourceObjectType) ? rule.WatchObjectType : rule.SourceObjectType!; // 结果列名一次性解析(结果集列在整个结果集内稳定):canonical 优先,缺失回退 params 真实列名。 var statusColumn = ResolveResultColumn(rowSet, CanonicalStatusColumn, parameters.StatusField); var dueAtColumn = ResolveResultColumn(rowSet, CanonicalDueAtColumn, parameters.DueAtField); var objectCodeColumn = ResolveResultColumn(rowSet, CanonicalRelatedObjectCodeColumn, parameters.ObjectCodeField); var objectIdColumn = ResolveResultColumn(rowSet, CanonicalSourceObjectIdColumn, parameters.ObjectIdField); foreach (var row in rowSet.Rows) { var status = row.GetString(statusColumn) ?? string.Empty; if (parameters.CompletedStates.Contains(status, StringComparer.OrdinalIgnoreCase)) continue; var due = row.GetDateTime(dueAtColumn); if (due == null || due > threshold) continue; var relatedObjectCode = row.GetString(objectCodeColumn) ?? string.Empty; if (string.IsNullOrWhiteSpace(relatedObjectCode)) continue; var sourceObjectId = row.GetString(objectIdColumn) ?? relatedObjectCode; var dedupKey = BuildDedupKey(tenantId, factoryId, rule.RuleCode, sourceObjectType, sourceObjectId); hits.Add(new S8RuleHit { SourceRuleId = rule.Id, SourceRuleCode = rule.RuleCode, SourceObjectType = sourceObjectType, SourceObjectId = sourceObjectId, RelatedObjectCode = relatedObjectCode, ExceptionTypeCode = parameters.ExceptionTypeCode!, SceneCode = rule.SceneCode, Severity = S8SeverityCode.Normalize(rule.Severity), DedupKey = dedupKey, SourcePayload = BuildPayload(row, sourceObjectType, sourceObjectId, due.Value, status, parameters), DetectedAt = detectedAt, Title = $"[超时] {sourceObjectType} {sourceObjectId} 已超期至 {due.Value:yyyy-MM-dd HH:mm:ss}(状态 {status})", DataSourceId = dataSourceId, OccurrenceDeptId = row.GetLong(S8CanonicalColumns.OccurrenceDeptId), ResponsibleDeptId = row.GetLong(S8CanonicalColumns.ResponsibleDeptId) }); } return hits; } /// 构造 R2 dedup_key 稳定字符串:T{tenant}:F{factory}:R{ruleCode}:{sourceObjectType}:{sourceObjectId}。internal 暴露供测试。 internal static string BuildDedupKey(long tenantId, long factoryId, string ruleCode, string sourceObjectType, string sourceObjectId) => $"T{tenantId}:F{factoryId}:R{ruleCode}:{sourceObjectType}:{sourceObjectId}"; private static string BuildPayload(S8MonitoringRow row, string sourceObjectType, string sourceObjectId, DateTime dueAt, string status, S8TimeoutParams parameters) { var payload = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var kv in row.Values) payload[kv.Key] = kv.Value; payload["__ruleType"] = RuleTypeCode; payload["__sourceObjectType"] = sourceObjectType; payload["__sourceObjectId"] = sourceObjectId; payload["__dueAt"] = dueAt; payload["__status"] = status; payload["__graceMinutes"] = parameters.GraceMinutes; payload["__exceptionTypeCode"] = parameters.ExceptionTypeCode; return JsonSerializer.Serialize(payload); } /// /// 结果列名解析:BuildExpression 已把结果列统一别名为 canonical(due_at/status/source_object_id/ /// related_object_code)。优先返回 canonical 列名;仅当结果集不含 canonical 列时,回退到 params_json /// 指定的真实列名(兼容历史未别名规则)。仅在 canonical 与 params 字段之间二选一,不新增无依据兜底字段。 /// private static string ResolveResultColumn(S8MonitoringRowSet rowSet, string canonicalColumn, string? paramsColumn) { if (rowSet.HasColumn(canonicalColumn)) return canonicalColumn; return string.IsNullOrWhiteSpace(paramsColumn) ? canonicalColumn : paramsColumn!; } }