S8WatchSchedulerService.cs 42 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000
  1. using Admin.NET.Plugin.AiDOP.Entity.S8;
  2. using Admin.NET.Plugin.AiDOP.Infrastructure.S8;
  3. using Admin.NET.Plugin.AiDOP.Service.S8.Rules;
  4. using Microsoft.Extensions.Logging;
  5. using SqlSugar;
  6. using System.Data;
  7. using System.Globalization;
  8. using System.Text.Json;
  9. namespace Admin.NET.Plugin.AiDOP.Service.S8;
  10. /// <summary>
  11. /// 监视规则轮询调度服务(首轮存根)。
  12. /// 后续接入 Admin.NET 定时任务机制后,由调度器周期调用 <see cref="RunOnceAsync"/>,
  13. /// 按各规则的 PollIntervalSeconds 逐条评估并生成异常记录。
  14. /// </summary>
  15. public class S8WatchSchedulerService : ITransient
  16. {
  17. private readonly SqlSugarRepository<AdoS8WatchRule> _ruleRep;
  18. private readonly SqlSugarRepository<AdoS8AlertRule> _alertRuleRep;
  19. private readonly SqlSugarRepository<AdoS8DataSource> _dataSourceRep;
  20. private readonly SqlSugarRepository<AdoS8Exception> _exceptionRep;
  21. private readonly SqlSugarRepository<AdoS8ExceptionType> _exceptionTypeRep;
  22. private readonly S8NotificationService _notificationService;
  23. private readonly S8ManualReportService _manualReportService;
  24. private readonly S8TimeoutRuleEvaluator _timeoutEvaluator;
  25. private readonly S8ShortageRuleEvaluator _shortageEvaluator;
  26. private readonly S8OutOfRangeRuleEvaluator _outOfRangeEvaluator;
  27. private readonly ILogger<S8WatchSchedulerService> _logger;
  28. private const string DefaultTriggerType = "VALUE_DEVIATION";
  29. private const string SqlDataSourceType = "SQL";
  30. // G01-05 未闭环状态集合:复用自 S8ExceptionService 当前 pendingStatuses 事实口径
  31. // (见 S8ExceptionService.GetPagedAsync 中 pendingStatuses 的定义,两处必须保持一致)。
  32. // 这不是“自定义未闭环集合”;若现有口径调整,两处需同步修改。
  33. private static readonly string[] UnclosedExceptionStatuses =
  34. { "NEW", "ASSIGNED", "IN_PROGRESS", "PENDING_VERIFICATION" };
  35. public S8WatchSchedulerService(
  36. SqlSugarRepository<AdoS8WatchRule> ruleRep,
  37. SqlSugarRepository<AdoS8AlertRule> alertRuleRep,
  38. SqlSugarRepository<AdoS8DataSource> dataSourceRep,
  39. SqlSugarRepository<AdoS8Exception> exceptionRep,
  40. SqlSugarRepository<AdoS8ExceptionType> exceptionTypeRep,
  41. S8NotificationService notificationService,
  42. S8ManualReportService manualReportService,
  43. S8TimeoutRuleEvaluator timeoutEvaluator,
  44. S8ShortageRuleEvaluator shortageEvaluator,
  45. S8OutOfRangeRuleEvaluator outOfRangeEvaluator,
  46. ILogger<S8WatchSchedulerService> logger)
  47. {
  48. _ruleRep = ruleRep;
  49. _alertRuleRep = alertRuleRep;
  50. _dataSourceRep = dataSourceRep;
  51. _exceptionRep = exceptionRep;
  52. _exceptionTypeRep = exceptionTypeRep;
  53. _notificationService = notificationService;
  54. _manualReportService = manualReportService;
  55. _timeoutEvaluator = timeoutEvaluator;
  56. _shortageEvaluator = shortageEvaluator;
  57. _outOfRangeEvaluator = outOfRangeEvaluator;
  58. _logger = logger;
  59. }
  60. public async Task<List<S8WatchExecutionRule>> LoadExecutionRulesAsync(long tenantId, long factoryId)
  61. {
  62. // R3 OUT_OF_RANGE 重写后,新三类(OUT_OF_RANGE/TIMEOUT/SHORTAGE)改走 ProcessRulesByTypeAsync。
  63. // 旧 AlertRule 兼容主链此处只装载 RuleType 为空/未分类的历史规则,避免新旧双跑导致重复建单。
  64. var watchRules = await _ruleRep.AsQueryable()
  65. .Where(x => x.TenantId == tenantId
  66. && x.FactoryId == factoryId
  67. && x.Enabled
  68. && x.SceneCode == S8SceneCode.S2S6Production
  69. && (x.RuleType == null || x.RuleType == ""))
  70. .ToListAsync();
  71. var deviceRules = watchRules
  72. .Where(x => IsDeviceWatchObjectType(x.WatchObjectType))
  73. .ToList();
  74. if (deviceRules.Count == 0) return new();
  75. var dataSourceIds = deviceRules
  76. .Select(x => x.DataSourceId)
  77. .Distinct()
  78. .ToList();
  79. var dataSources = await _dataSourceRep.AsQueryable()
  80. .Where(x => x.TenantId == tenantId
  81. && x.FactoryId == factoryId
  82. && x.Enabled
  83. && dataSourceIds.Contains(x.Id))
  84. .ToListAsync();
  85. var dataSourceMap = dataSources.ToDictionary(x => x.Id);
  86. if (dataSourceMap.Count == 0) return new();
  87. var alertRules = (await _alertRuleRep.AsQueryable()
  88. .Where(x => x.TenantId == tenantId
  89. && x.FactoryId == factoryId
  90. && x.SceneCode == S8SceneCode.S2S6Production)
  91. .ToListAsync())
  92. .Where(IsSupportedAlertRule)
  93. .ToList();
  94. // G-01 首版 AlertRule 冲突口径(C 收口):
  95. // 当前场景存在多条可运行 AlertRule 时,视为“当前规则配置冲突”并跳过该规则,
  96. // 不按“首条”继续运行,也不扩大为“整场景停摆”。
  97. // 当前模型下所有 device watchRule 共享同场景 AlertRule,故冲突态下所有 device 规则均跳过,
  98. // 但此处按“逐规则跳过”的语义实现,避免被误读为“整场景 return empty 停摆”。
  99. var alertRule = alertRules.Count == 1 ? alertRules[0] : null;
  100. var executionRules = new List<S8WatchExecutionRule>();
  101. foreach (var watchRule in deviceRules.OrderBy(x => x.Id))
  102. {
  103. // 配置冲突:当前规则跳过(不停摆其他规则)。
  104. if (alertRule == null)
  105. continue;
  106. if (!dataSourceMap.TryGetValue(watchRule.DataSourceId, out var dataSource))
  107. continue;
  108. if (!IsSupportedSqlDataSource(dataSource))
  109. continue;
  110. executionRules.Add(new S8WatchExecutionRule
  111. {
  112. WatchRuleId = watchRule.Id,
  113. WatchRuleCode = watchRule.RuleCode,
  114. SceneCode = watchRule.SceneCode,
  115. TriggerType = DefaultTriggerType,
  116. WatchObjectType = watchRule.WatchObjectType.Trim(),
  117. DataSourceId = dataSource.Id,
  118. DataSourceCode = dataSource.DataSourceCode,
  119. DataSourceType = dataSource.Type,
  120. DataSourceConnection = dataSource.Endpoint?.Trim() ?? string.Empty,
  121. QueryExpression = watchRule.Expression?.Trim() ?? string.Empty,
  122. PollIntervalSeconds = watchRule.PollIntervalSeconds,
  123. AlertRuleId = alertRule.Id,
  124. AlertRuleCode = alertRule.RuleCode,
  125. TriggerCondition = alertRule.TriggerCondition!.Trim(),
  126. ThresholdValue = alertRule.ThresholdVal!.Trim(),
  127. Severity = alertRule.Severity
  128. });
  129. }
  130. return executionRules;
  131. }
  132. public async Task<List<S8WatchDeviceQueryResult>> QueryDeviceRowsAsync(long tenantId, long factoryId)
  133. {
  134. var executionRules = await LoadExecutionRulesAsync(tenantId, factoryId);
  135. var results = new List<S8WatchDeviceQueryResult>();
  136. foreach (var rule in executionRules)
  137. results.Add(await QueryDeviceRowsAsync(rule));
  138. return results;
  139. }
  140. public async Task<S8WatchDeviceQueryResult> QueryDeviceRowsAsync(S8WatchExecutionRule rule)
  141. {
  142. if (!string.Equals(rule.DataSourceType, SqlDataSourceType, StringComparison.OrdinalIgnoreCase))
  143. return S8WatchDeviceQueryResult.Fail(rule, "数据源类型不是 SQL,已跳过");
  144. if (string.IsNullOrWhiteSpace(rule.QueryExpression))
  145. return S8WatchDeviceQueryResult.Fail(rule, "查询表达式为空,已跳过");
  146. try
  147. {
  148. using var db = CreateSqlQueryScope(rule.DataSourceConnection);
  149. var table = await db.Ado.GetDataTableAsync(rule.QueryExpression);
  150. if (!HasRequiredColumns(table))
  151. return S8WatchDeviceQueryResult.Fail(rule, "查询结果缺少 required columns: related_object_code/current_value");
  152. var rows = table.Rows.Cast<DataRow>()
  153. .Select(MapDeviceRow)
  154. .Where(x => !string.IsNullOrWhiteSpace(x.RelatedObjectCode))
  155. .ToList();
  156. return S8WatchDeviceQueryResult.Ok(rule, rows);
  157. }
  158. catch (Exception ex)
  159. {
  160. return S8WatchDeviceQueryResult.Fail(rule, $"查询执行失败: {ex.Message}");
  161. }
  162. }
  163. /// <summary>
  164. /// G01-04:基于设备级结果行集做首版 VALUE_DEVIATION 单阈值判定,
  165. /// 产出命中结果对象列表,供 G01-05 去重与 G01-06 建单消费。
  166. /// 本方法不做去重、不做建单、不做严重度重算、不做时间线。
  167. /// </summary>
  168. public async Task<List<S8WatchHitResult>> EvaluateHitsAsync(long tenantId, long factoryId)
  169. {
  170. var executionRules = await LoadExecutionRulesAsync(tenantId, factoryId);
  171. var ruleMap = executionRules.ToDictionary(x => x.WatchRuleId);
  172. var queryResults = new List<S8WatchDeviceQueryResult>();
  173. foreach (var rule in executionRules)
  174. queryResults.Add(await QueryDeviceRowsAsync(rule));
  175. var hits = new List<S8WatchHitResult>();
  176. foreach (var queryResult in queryResults)
  177. {
  178. // G01-03 查询失败:跳过,不进入判定。
  179. if (!queryResult.Success) continue;
  180. if (!ruleMap.TryGetValue(queryResult.WatchRuleId, out var rule)) continue;
  181. // 判定参数缺失:跳过当前规则。
  182. if (string.IsNullOrWhiteSpace(rule.TriggerCondition)
  183. || string.IsNullOrWhiteSpace(rule.ThresholdValue))
  184. continue;
  185. // 比较符非法:跳过当前规则。
  186. var op = TryParseTriggerCondition(rule.TriggerCondition);
  187. if (op == null) continue;
  188. // ThresholdValue 非法:跳过当前规则。
  189. if (!TryParseDecimal(rule.ThresholdValue, out var threshold)) continue;
  190. foreach (var row in queryResult.Rows)
  191. {
  192. // CurrentValue 非法(null / 无法解析数值):跳过当前行,不进入判定。
  193. if (row.CurrentValue == null) continue;
  194. // 未命中:不进入后续链路。
  195. if (!EvaluateHit(row.CurrentValue.Value, op, threshold)) continue;
  196. hits.Add(new S8WatchHitResult
  197. {
  198. SourceRuleId = rule.WatchRuleId,
  199. SourceRuleCode = rule.WatchRuleCode,
  200. AlertRuleId = rule.AlertRuleId,
  201. DataSourceId = rule.DataSourceId,
  202. RelatedObjectCode = row.RelatedObjectCode,
  203. CurrentValue = row.CurrentValue.Value,
  204. ThresholdValue = threshold,
  205. TriggerCondition = op,
  206. Severity = rule.Severity,
  207. OccurrenceDeptId = row.OccurrenceDeptId,
  208. ResponsibleDeptId = row.ResponsibleDeptId,
  209. SourcePayload = row.SourcePayload
  210. });
  211. }
  212. }
  213. return hits;
  214. }
  215. /// <summary>
  216. /// G01-05:未闭环异常去重最小实现。
  217. /// 消费 G01-04 产出的 <see cref="S8WatchHitResult"/>,按 (SourceRuleId + RelatedObjectCode)
  218. /// 在未闭环状态集合内判重,只回答“是否允许建单”。
  219. /// 首版明确不做:原单刷新 / 时间线追加 / payload 更新 / 次数累计 / 严重度重算 / 状态修复。
  220. /// </summary>
  221. public async Task<List<S8WatchDedupResult>> EvaluateDedupAsync(long tenantId, long factoryId)
  222. {
  223. var hits = await EvaluateHitsAsync(tenantId, factoryId);
  224. var results = new List<S8WatchDedupResult>(hits.Count);
  225. foreach (var hit in hits)
  226. {
  227. // 防御性分支:正常情况下 SourceRuleId 与 RelatedObjectCode 已由上游
  228. // (G01-02 规则装配 + G01-03 查询结果列校验)保证;此处仅作兜底,
  229. // 不是首版正常路径。
  230. if (hit.SourceRuleId <= 0 || string.IsNullOrWhiteSpace(hit.RelatedObjectCode))
  231. {
  232. results.Add(new S8WatchDedupResult
  233. {
  234. Hit = hit,
  235. CanCreate = false,
  236. MatchedExceptionId = null,
  237. Reason = "missing_dedup_key"
  238. });
  239. continue;
  240. }
  241. long matchedId;
  242. try
  243. {
  244. matchedId = await FindPendingExceptionIdAsync(
  245. tenantId, factoryId, hit.SourceRuleId, hit.RelatedObjectCode);
  246. }
  247. catch
  248. {
  249. // 去重查询失败:首版偏保守,宁可阻止也不重复建单,不扩展为补偿。
  250. results.Add(new S8WatchDedupResult
  251. {
  252. Hit = hit,
  253. CanCreate = false,
  254. MatchedExceptionId = null,
  255. Reason = "query_failed"
  256. });
  257. continue;
  258. }
  259. if (matchedId > 0)
  260. {
  261. // 命中已有未闭环异常 → 阻止建单。
  262. results.Add(new S8WatchDedupResult
  263. {
  264. Hit = hit,
  265. CanCreate = false,
  266. MatchedExceptionId = matchedId,
  267. Reason = "duplicate_pending"
  268. });
  269. }
  270. else
  271. {
  272. // 未命中 → 允许建单,交 G01-06 消费。
  273. results.Add(new S8WatchDedupResult
  274. {
  275. Hit = hit,
  276. CanCreate = true,
  277. MatchedExceptionId = null,
  278. Reason = "no_pending"
  279. });
  280. }
  281. }
  282. return results;
  283. }
  284. // 取任意一条匹配的未闭环异常 Id 作为“是否存在重复单”的拦截依据。
  285. // 首版只需要“存在性”,不关心“最早 / 最新”;不在 G01-05 处理排序语义。
  286. private async Task<long> FindPendingExceptionIdAsync(
  287. long tenantId, long factoryId, long sourceRuleId, string relatedObjectCode)
  288. {
  289. // SqlSugar 表达式翻译要求 Contains 数组变量必须可访问;此处将类级常量
  290. // 承接到方法内局部变量,仅为表达式翻译服务,值与 UnclosedExceptionStatuses 一致。
  291. var statuses = UnclosedExceptionStatuses;
  292. var ids = await _exceptionRep.AsQueryable()
  293. .Where(x => x.TenantId == tenantId
  294. && x.FactoryId == factoryId
  295. && !x.IsDeleted
  296. && x.SourceRuleId == sourceRuleId
  297. && x.RelatedObjectCode == relatedObjectCode
  298. && statuses.Contains(x.Status))
  299. .Select(x => x.Id)
  300. .Take(1)
  301. .ToListAsync();
  302. return ids.Count > 0 ? ids[0] : 0L;
  303. }
  304. /// <summary>
  305. /// G01-06:自动建单入口。消费 G01-05 去重结果,对 CanCreate==true 的命中
  306. /// 复用 S8ManualReportService.CreateFromWatchAsync(同一主链的自动建单分支)落成标准 AdoS8Exception。
  307. /// CanCreate==false 直接跳过;创建失败返回最小失败结果,不补偿、不重试、不对账。
  308. /// </summary>
  309. public async Task<List<S8WatchCreationResult>> CreateExceptionsAsync(long tenantId, long factoryId)
  310. {
  311. var dedupResults = await EvaluateDedupAsync(tenantId, factoryId);
  312. var results = new List<S8WatchCreationResult>(dedupResults.Count);
  313. foreach (var dedup in dedupResults)
  314. {
  315. if (!dedup.CanCreate)
  316. {
  317. results.Add(new S8WatchCreationResult
  318. {
  319. DedupResult = dedup,
  320. Created = false,
  321. Skipped = true,
  322. CreatedExceptionId = null,
  323. Reason = dedup.Reason,
  324. ErrorMessage = null
  325. });
  326. continue;
  327. }
  328. try
  329. {
  330. var entity = await _manualReportService.CreateFromWatchAsync(dedup.Hit);
  331. results.Add(new S8WatchCreationResult
  332. {
  333. DedupResult = dedup,
  334. Created = true,
  335. Skipped = false,
  336. CreatedExceptionId = entity.Id,
  337. Reason = "auto_created",
  338. ErrorMessage = null
  339. });
  340. }
  341. catch (Exception ex)
  342. {
  343. results.Add(new S8WatchCreationResult
  344. {
  345. DedupResult = dedup,
  346. Created = false,
  347. Skipped = false,
  348. CreatedExceptionId = null,
  349. Reason = "create_failed",
  350. ErrorMessage = ex.Message
  351. });
  352. }
  353. }
  354. // R3-OUT_OF_RANGE-REWRITE-1:三类正式 evaluator 路径。
  355. // 旧 AlertRule 兼容主链(上方 dedupResults 循环)已在 LoadExecutionRulesAsync 处过滤掉
  356. // OUT_OF_RANGE/TIMEOUT/SHORTAGE,只保留未分类历史规则;故新旧不双跑。
  357. results.AddRange(await ProcessTimeoutRulesAsync(tenantId, factoryId));
  358. results.AddRange(await ProcessShortageRulesAsync(tenantId, factoryId));
  359. results.AddRange(await ProcessOutOfRangeRulesAsync(tenantId, factoryId));
  360. return results;
  361. }
  362. /// <summary>
  363. /// R2 TIMEOUT 类规则主链:薄包装,复用 <see cref="ProcessRulesByTypeAsync"/>。
  364. /// </summary>
  365. public Task<List<S8WatchCreationResult>> ProcessTimeoutRulesAsync(long tenantId, long factoryId) =>
  366. ProcessRulesByTypeAsync(tenantId, factoryId, _timeoutEvaluator, S8TimeoutRuleEvaluator.RuleTypeCode);
  367. /// <summary>
  368. /// R3 SHORTAGE 类规则主链:薄包装,复用 <see cref="ProcessRulesByTypeAsync"/>。
  369. /// </summary>
  370. public Task<List<S8WatchCreationResult>> ProcessShortageRulesAsync(long tenantId, long factoryId) =>
  371. ProcessRulesByTypeAsync(tenantId, factoryId, _shortageEvaluator, S8ShortageRuleEvaluator.RuleTypeCode);
  372. /// <summary>
  373. /// R3-OUT_OF_RANGE-REWRITE-1:OUT_OF_RANGE 类规则主链。
  374. /// 复用 <see cref="ProcessRulesByTypeAsync"/>,并对历史 dedup_key=NULL 的旧记录做 compat fallback:
  375. /// (source_rule_id=rule.Id AND related_object_code=hit AND status!=CLOSED AND dedup_key IS NULL AND is_deleted=0)
  376. /// 命中则 backfill 6 列,避免重复建单。
  377. /// </summary>
  378. public Task<List<S8WatchCreationResult>> ProcessOutOfRangeRulesAsync(long tenantId, long factoryId) =>
  379. ProcessRulesByTypeAsync(tenantId, factoryId, _outOfRangeEvaluator, S8OutOfRangeRuleEvaluator.RuleTypeCode);
  380. /// <summary>
  381. /// R2/R3 通用规则主链:装载 enabled WatchRule.RuleType=ruleType → evaluator → dedup_key 去重 → 建单/刷新。
  382. /// dedup 命中:UPDATE last_detected_at + source_payload,不重复建单;
  383. /// dedup 未命中:校验 ExceptionTypeCode 是否在 baseline(tenant=0/factory=0 全局或本租户工厂),缺则跳过;
  384. /// 通过 → S8ManualReportService.CreateFromHitAsync 落标准 AdoS8Exception,新列全部回填。
  385. /// 不做 SLA 升级、不做事件触发、不做 RecoveredAt。
  386. /// </summary>
  387. private async Task<List<S8WatchCreationResult>> ProcessRulesByTypeAsync(
  388. long tenantId, long factoryId, IS8RuleEvaluator evaluator, string ruleType)
  389. {
  390. var results = new List<S8WatchCreationResult>();
  391. var rules = await _ruleRep.AsQueryable()
  392. .Where(x => x.TenantId == tenantId
  393. && x.FactoryId == factoryId
  394. && x.Enabled
  395. && x.RuleType == ruleType)
  396. .ToListAsync();
  397. if (rules.Count == 0) return results;
  398. var alertRules = (await _alertRuleRep.AsQueryable()
  399. .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId)
  400. .ToListAsync()).AsReadOnly();
  401. foreach (var rule in rules.OrderBy(x => x.Id))
  402. {
  403. List<S8RuleHit> hits;
  404. try
  405. {
  406. hits = await evaluator.EvaluateAsync(tenantId, factoryId, rule, alertRules);
  407. }
  408. catch (Exception ex)
  409. {
  410. results.Add(BuildSkipResult(rule, "evaluate_failed", ex.Message));
  411. continue;
  412. }
  413. // R5-RECOVERY-MINIMAL-1:evaluator 成功执行后调和 recovered_at。
  414. // 仅在 evaluator 成功(throw 不会到这里)时才能判定"未命中"。仅写 recovered_at + updated_at,
  415. // 绝不动 status / assignee / verifier / source_payload / last_detected_at。
  416. try
  417. {
  418. await ReconcileRecoveriesForRuleAsync(tenantId, factoryId, rule, ruleType, hits);
  419. }
  420. catch (Exception ex)
  421. {
  422. _logger.LogWarning(ex, "recovery_reconcile_failed ruleCode={RuleCode} ruleType={RuleType}", rule.RuleCode, ruleType);
  423. }
  424. foreach (var hit in hits)
  425. {
  426. if (string.IsNullOrWhiteSpace(hit.DedupKey))
  427. {
  428. results.Add(BuildSkipResult(rule, "missing_dedup_key", null, hit));
  429. continue;
  430. }
  431. long matchedId;
  432. try
  433. {
  434. matchedId = await FindOpenExceptionByDedupKeyAsync(tenantId, factoryId, hit.DedupKey);
  435. }
  436. catch (Exception ex)
  437. {
  438. results.Add(BuildSkipResult(rule, "query_failed", ex.Message, hit));
  439. continue;
  440. }
  441. if (matchedId > 0)
  442. {
  443. try
  444. {
  445. await RefreshDetectionAsync(matchedId, hit);
  446. results.Add(BuildSkippedDuplicate(rule, hit, matchedId));
  447. }
  448. catch (Exception ex)
  449. {
  450. results.Add(BuildSkipResult(rule, "refresh_failed", ex.Message, hit));
  451. }
  452. continue;
  453. }
  454. // R3-OUT_OF_RANGE-REWRITE-1 compat fallback:dedup_key 未命中且 ruleType=OUT_OF_RANGE 时,
  455. // 尝试匹配历史 dedup_key=NULL 的旧 AlertRule 主链记录(如 id=34),命中则 backfill 6 列,避免重复建单。
  456. if (string.Equals(ruleType, S8OutOfRangeRuleEvaluator.RuleTypeCode, StringComparison.OrdinalIgnoreCase))
  457. {
  458. long compatId;
  459. try
  460. {
  461. compatId = await FindLegacyOutOfRangeExceptionAsync(tenantId, factoryId, rule.Id, hit.RelatedObjectCode);
  462. }
  463. catch (Exception ex)
  464. {
  465. results.Add(BuildSkipResult(rule, "query_failed", ex.Message, hit));
  466. continue;
  467. }
  468. if (compatId > 0)
  469. {
  470. try
  471. {
  472. await BackfillLegacyExceptionAsync(compatId, hit);
  473. results.Add(BuildSkippedDuplicate(rule, hit, compatId));
  474. }
  475. catch (Exception ex)
  476. {
  477. results.Add(BuildSkipResult(rule, "refresh_failed", ex.Message, hit));
  478. }
  479. continue;
  480. }
  481. }
  482. bool typeExists;
  483. try
  484. {
  485. typeExists = await _exceptionTypeRep.AsQueryable()
  486. .Where(t => t.TypeCode == hit.ExceptionTypeCode
  487. && (t.TenantId == 0 || t.TenantId == tenantId)
  488. && (t.FactoryId == 0 || t.FactoryId == factoryId)
  489. && t.Enabled)
  490. .AnyAsync();
  491. }
  492. catch (Exception ex)
  493. {
  494. results.Add(BuildSkipResult(rule, "query_failed", ex.Message, hit));
  495. continue;
  496. }
  497. if (!typeExists)
  498. {
  499. results.Add(BuildSkipResult(rule, "exception_type_missing", null, hit));
  500. continue;
  501. }
  502. try
  503. {
  504. var entity = await _manualReportService.CreateFromHitAsync(hit);
  505. results.Add(BuildCreatedResult(rule, hit, entity.Id));
  506. }
  507. catch (Exception ex)
  508. {
  509. results.Add(BuildSkipResult(rule, "create_failed", ex.Message, hit));
  510. }
  511. }
  512. }
  513. return results;
  514. }
  515. private async Task<long> FindOpenExceptionByDedupKeyAsync(long tenantId, long factoryId, string dedupKey)
  516. {
  517. var ids = await _exceptionRep.AsQueryable()
  518. .Where(x => x.TenantId == tenantId
  519. && x.FactoryId == factoryId
  520. && !x.IsDeleted
  521. && x.Status != "CLOSED"
  522. && x.DedupKey == dedupKey)
  523. .Select(x => x.Id)
  524. .Take(1)
  525. .ToListAsync();
  526. return ids.Count > 0 ? ids[0] : 0L;
  527. }
  528. /// <summary>
  529. /// R3 OUT_OF_RANGE compat fallback 查找:用 (source_rule_id, related_object_code, status!=CLOSED,
  530. /// dedup_key IS NULL, is_deleted=0) 严格条件定位旧 AlertRule 主链留下的历史记录。
  531. /// </summary>
  532. private async Task<long> FindLegacyOutOfRangeExceptionAsync(long tenantId, long factoryId, long sourceRuleId, string relatedObjectCode)
  533. {
  534. if (string.IsNullOrWhiteSpace(relatedObjectCode)) return 0L;
  535. var ids = await _exceptionRep.AsQueryable()
  536. .Where(x => x.TenantId == tenantId
  537. && x.FactoryId == factoryId
  538. && !x.IsDeleted
  539. && x.Status != "CLOSED"
  540. && x.DedupKey == null
  541. && x.SourceRuleId == sourceRuleId
  542. && x.RelatedObjectCode == relatedObjectCode)
  543. .Select(x => x.Id)
  544. .Take(1)
  545. .ToListAsync();
  546. return ids.Count > 0 ? ids[0] : 0L;
  547. }
  548. /// <summary>
  549. /// R3 OUT_OF_RANGE compat fallback backfill:把历史记录的 R1 新 6 列(dedup_key/source_rule_code/
  550. /// source_object_type/source_object_id/source_payload/last_detected_at)写入,并刷新 updated_at。
  551. /// </summary>
  552. private async Task BackfillLegacyExceptionAsync(long exceptionId, S8RuleHit hit)
  553. {
  554. await _exceptionRep.Context.Updateable<AdoS8Exception>()
  555. .SetColumns(x => new AdoS8Exception
  556. {
  557. DedupKey = hit.DedupKey,
  558. SourceRuleCode = hit.SourceRuleCode,
  559. SourceObjectType = hit.SourceObjectType,
  560. SourceObjectId = hit.SourceObjectId,
  561. SourcePayload = hit.SourcePayload,
  562. LastDetectedAt = hit.DetectedAt,
  563. UpdatedAt = DateTime.Now
  564. })
  565. .Where(x => x.Id == exceptionId)
  566. .ExecuteCommandAsync();
  567. }
  568. /// <summary>
  569. /// R5 恢复时间最小闭环:对当前 rule 下未关闭、有 dedup_key、recovered_at 仍为 NULL 的异常,
  570. /// 凡不在本轮 hits.dedup_key 集合内的,写入 recovered_at = now、updated_at = now。
  571. /// 仅写这 2 列;不动 status / assignee / verifier / source_payload / last_detected_at;
  572. /// recovered_at 一旦写入,本轮不做复发清空。
  573. /// </summary>
  574. private async Task ReconcileRecoveriesForRuleAsync(
  575. long tenantId, long factoryId, AdoS8WatchRule rule, string ruleType, List<S8RuleHit> hits)
  576. {
  577. var hitDedupKeys = hits
  578. .Where(h => !string.IsNullOrWhiteSpace(h.DedupKey))
  579. .Select(h => h.DedupKey)
  580. .ToHashSet(StringComparer.Ordinal);
  581. var candidates = await _exceptionRep.AsQueryable()
  582. .Where(x => x.TenantId == tenantId
  583. && x.FactoryId == factoryId
  584. && !x.IsDeleted
  585. && x.Status != "CLOSED"
  586. && x.SourceRuleCode == rule.RuleCode
  587. && x.DedupKey != null
  588. && x.RecoveredAt == null)
  589. .Select(x => new { x.Id, x.DedupKey })
  590. .ToListAsync();
  591. if (candidates.Count == 0) return;
  592. var now = DateTime.Now;
  593. var recoveredIds = new List<long>();
  594. foreach (var c in candidates)
  595. {
  596. if (hitDedupKeys.Contains(c.DedupKey!)) continue;
  597. await _exceptionRep.Context.Updateable<AdoS8Exception>()
  598. .SetColumns(x => new AdoS8Exception
  599. {
  600. RecoveredAt = now,
  601. UpdatedAt = now
  602. })
  603. .Where(x => x.Id == c.Id)
  604. .ExecuteCommandAsync();
  605. recoveredIds.Add(c.Id);
  606. }
  607. if (recoveredIds.Count > 0)
  608. {
  609. _logger.LogInformation(
  610. "rule_recovered ruleCode={RuleCode} ruleType={RuleType} recoveredCount={Count} recoveredIds={Ids}",
  611. rule.RuleCode, ruleType, recoveredIds.Count, string.Join(",", recoveredIds));
  612. }
  613. }
  614. private async Task RefreshDetectionAsync(long exceptionId, S8RuleHit hit)
  615. {
  616. await _exceptionRep.Context.Updateable<AdoS8Exception>()
  617. .SetColumns(x => new AdoS8Exception
  618. {
  619. LastDetectedAt = hit.DetectedAt,
  620. SourcePayload = hit.SourcePayload,
  621. UpdatedAt = DateTime.Now
  622. })
  623. .Where(x => x.Id == exceptionId)
  624. .ExecuteCommandAsync();
  625. }
  626. private static S8WatchCreationResult BuildCreatedResult(AdoS8WatchRule rule, S8RuleHit hit, long exceptionId) =>
  627. new()
  628. {
  629. DedupResult = new S8WatchDedupResult
  630. {
  631. Hit = ToWatchHit(rule, hit),
  632. CanCreate = true,
  633. MatchedExceptionId = null,
  634. Reason = "no_pending"
  635. },
  636. Created = true,
  637. Skipped = false,
  638. CreatedExceptionId = exceptionId,
  639. Reason = "auto_created",
  640. ErrorMessage = null
  641. };
  642. private static S8WatchCreationResult BuildSkippedDuplicate(AdoS8WatchRule rule, S8RuleHit hit, long matchedId) =>
  643. new()
  644. {
  645. DedupResult = new S8WatchDedupResult
  646. {
  647. Hit = ToWatchHit(rule, hit),
  648. CanCreate = false,
  649. MatchedExceptionId = matchedId,
  650. Reason = "duplicate_pending"
  651. },
  652. Created = false,
  653. Skipped = true,
  654. CreatedExceptionId = null,
  655. Reason = "duplicate_pending",
  656. ErrorMessage = null
  657. };
  658. private static S8WatchCreationResult BuildSkipResult(AdoS8WatchRule rule, string reason, string? error, S8RuleHit? hit = null) =>
  659. new()
  660. {
  661. DedupResult = new S8WatchDedupResult
  662. {
  663. Hit = hit != null ? ToWatchHit(rule, hit) : new S8WatchHitResult { SourceRuleId = rule.Id, SourceRuleCode = rule.RuleCode },
  664. CanCreate = false,
  665. MatchedExceptionId = null,
  666. Reason = reason
  667. },
  668. Created = false,
  669. Skipped = true,
  670. CreatedExceptionId = null,
  671. Reason = reason,
  672. ErrorMessage = error
  673. };
  674. private static S8WatchHitResult ToWatchHit(AdoS8WatchRule rule, S8RuleHit hit) => new()
  675. {
  676. SourceRuleId = hit.SourceRuleId == 0 ? rule.Id : hit.SourceRuleId,
  677. SourceRuleCode = string.IsNullOrEmpty(hit.SourceRuleCode) ? rule.RuleCode : hit.SourceRuleCode,
  678. DataSourceId = hit.DataSourceId,
  679. RelatedObjectCode = hit.RelatedObjectCode,
  680. Severity = hit.Severity,
  681. OccurrenceDeptId = hit.OccurrenceDeptId,
  682. ResponsibleDeptId = hit.ResponsibleDeptId,
  683. SourcePayload = hit.SourcePayload
  684. };
  685. /// <summary>
  686. /// 单次轮询入口。当前仅完成规则读取与组装,返回可执行规则数量,不做实际数据采集。
  687. /// </summary>
  688. public async Task<int> RunOnceAsync()
  689. {
  690. var executionRules = await LoadExecutionRulesAsync(1, 1);
  691. return executionRules.Count;
  692. }
  693. // G01-04 首版最小比较符集合:>, >=, <, <=。
  694. // 允许首尾空格;非此集合的一律视为“比较符非法”,由调用方跳过。
  695. private static string? TryParseTriggerCondition(string raw)
  696. {
  697. var normalized = raw.Trim();
  698. return normalized switch
  699. {
  700. ">" => ">",
  701. ">=" => ">=",
  702. "<" => "<",
  703. "<=" => "<=",
  704. _ => null
  705. };
  706. }
  707. private static bool TryParseDecimal(string raw, out decimal value) =>
  708. decimal.TryParse(raw.Trim(), NumberStyles.Any, CultureInfo.InvariantCulture, out value);
  709. private static bool EvaluateHit(decimal current, string op, decimal threshold) => op switch
  710. {
  711. ">" => current > threshold,
  712. ">=" => current >= threshold,
  713. "<" => current < threshold,
  714. "<=" => current <= threshold,
  715. _ => false
  716. };
  717. private static bool IsSupportedAlertRule(AdoS8AlertRule alertRule) =>
  718. !string.IsNullOrWhiteSpace(alertRule.TriggerCondition)
  719. && !string.IsNullOrWhiteSpace(alertRule.ThresholdVal);
  720. private static bool IsSupportedSqlDataSource(AdoS8DataSource dataSource) =>
  721. dataSource.Enabled
  722. && string.Equals(dataSource.Type?.Trim(), SqlDataSourceType, StringComparison.OrdinalIgnoreCase)
  723. && !string.IsNullOrWhiteSpace(dataSource.Endpoint);
  724. private static bool IsDeviceWatchObjectType(string? watchObjectType)
  725. {
  726. if (string.IsNullOrWhiteSpace(watchObjectType)) return false;
  727. var normalized = watchObjectType.Trim().ToUpperInvariant();
  728. return normalized is "DEVICE" or "EQUIPMENT" || watchObjectType.Trim() == "设备";
  729. }
  730. private SqlSugarScope CreateSqlQueryScope(string connectionString)
  731. {
  732. var dbType = _ruleRep.Context.CurrentConnectionConfig.DbType;
  733. return new SqlSugarScope(new ConnectionConfig
  734. {
  735. ConfigId = $"s8-watch-sql-{Guid.NewGuid():N}",
  736. DbType = dbType,
  737. ConnectionString = connectionString,
  738. InitKeyType = InitKeyType.Attribute,
  739. IsAutoCloseConnection = true
  740. });
  741. }
  742. private static bool HasRequiredColumns(DataTable table) =>
  743. TryGetColumnName(table.Columns, "related_object_code") != null
  744. && TryGetColumnName(table.Columns, "current_value") != null;
  745. private static S8WatchDeviceRow MapDeviceRow(DataRow row)
  746. {
  747. var columns = row.Table.Columns;
  748. var relatedObjectCodeColumn = TryGetColumnName(columns, "related_object_code");
  749. var currentValueColumn = TryGetColumnName(columns, "current_value");
  750. var occurrenceDeptIdColumn = TryGetColumnName(columns, "occurrence_dept_id");
  751. var responsibleDeptIdColumn = TryGetColumnName(columns, "responsible_dept_id");
  752. return new S8WatchDeviceRow
  753. {
  754. RelatedObjectCode = ReadString(row, relatedObjectCodeColumn),
  755. CurrentValue = ReadDecimal(row, currentValueColumn),
  756. OccurrenceDeptId = ReadLong(row, occurrenceDeptIdColumn),
  757. ResponsibleDeptId = ReadLong(row, responsibleDeptIdColumn),
  758. SourcePayload = BuildSourcePayload(row)
  759. };
  760. }
  761. private static string? TryGetColumnName(DataColumnCollection columns, string expectedName)
  762. {
  763. var normalizedExpected = NormalizeColumnName(expectedName);
  764. foreach (DataColumn column in columns)
  765. {
  766. if (NormalizeColumnName(column.ColumnName) == normalizedExpected)
  767. return column.ColumnName;
  768. }
  769. return null;
  770. }
  771. private static string NormalizeColumnName(string columnName) =>
  772. columnName.Replace("_", string.Empty, StringComparison.Ordinal).Trim().ToUpperInvariant();
  773. private static string ReadString(DataRow row, string? columnName)
  774. {
  775. if (string.IsNullOrWhiteSpace(columnName)) return string.Empty;
  776. var value = row[columnName];
  777. return value == DBNull.Value ? string.Empty : Convert.ToString(value)?.Trim() ?? string.Empty;
  778. }
  779. private static long? ReadLong(DataRow row, string? columnName)
  780. {
  781. if (string.IsNullOrWhiteSpace(columnName)) return null;
  782. var value = row[columnName];
  783. if (value == DBNull.Value) return null;
  784. return long.TryParse(Convert.ToString(value, CultureInfo.InvariantCulture), out var result) ? result : null;
  785. }
  786. private static decimal? ReadDecimal(DataRow row, string? columnName)
  787. {
  788. if (string.IsNullOrWhiteSpace(columnName)) return null;
  789. var value = row[columnName];
  790. if (value == DBNull.Value) return null;
  791. return decimal.TryParse(Convert.ToString(value, CultureInfo.InvariantCulture), NumberStyles.Any, CultureInfo.InvariantCulture, out var result)
  792. ? result
  793. : null;
  794. }
  795. private static string BuildSourcePayload(DataRow row)
  796. {
  797. var payload = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  798. foreach (DataColumn column in row.Table.Columns)
  799. {
  800. var value = row[column];
  801. payload[column.ColumnName] = value == DBNull.Value ? null : value;
  802. }
  803. return JsonSerializer.Serialize(payload);
  804. }
  805. }
  806. public sealed class S8WatchExecutionRule
  807. {
  808. public long WatchRuleId { get; set; }
  809. public string WatchRuleCode { get; set; } = string.Empty;
  810. public string SceneCode { get; set; } = string.Empty;
  811. public string TriggerType { get; set; } = string.Empty;
  812. public string WatchObjectType { get; set; } = string.Empty;
  813. public long DataSourceId { get; set; }
  814. public string DataSourceCode { get; set; } = string.Empty;
  815. public string DataSourceType { get; set; } = string.Empty;
  816. public string DataSourceConnection { get; set; } = string.Empty;
  817. public string QueryExpression { get; set; } = string.Empty;
  818. public int PollIntervalSeconds { get; set; }
  819. public long AlertRuleId { get; set; }
  820. public string AlertRuleCode { get; set; } = string.Empty;
  821. public string TriggerCondition { get; set; } = string.Empty;
  822. public string ThresholdValue { get; set; } = string.Empty;
  823. public string Severity { get; set; } = string.Empty;
  824. }
  825. public sealed class S8WatchDeviceQueryResult
  826. {
  827. public long WatchRuleId { get; set; }
  828. public string WatchRuleCode { get; set; } = string.Empty;
  829. public bool Success { get; set; }
  830. public string? FailureReason { get; set; }
  831. public List<S8WatchDeviceRow> Rows { get; set; } = new();
  832. public static S8WatchDeviceQueryResult Ok(S8WatchExecutionRule rule, List<S8WatchDeviceRow> rows) =>
  833. new()
  834. {
  835. WatchRuleId = rule.WatchRuleId,
  836. WatchRuleCode = rule.WatchRuleCode,
  837. Success = true,
  838. Rows = rows
  839. };
  840. public static S8WatchDeviceQueryResult Fail(S8WatchExecutionRule rule, string reason) =>
  841. new()
  842. {
  843. WatchRuleId = rule.WatchRuleId,
  844. WatchRuleCode = rule.WatchRuleCode,
  845. Success = false,
  846. FailureReason = reason
  847. };
  848. }
  849. public sealed class S8WatchDeviceRow
  850. {
  851. public string RelatedObjectCode { get; set; } = string.Empty;
  852. public decimal? CurrentValue { get; set; }
  853. public long? OccurrenceDeptId { get; set; }
  854. public long? ResponsibleDeptId { get; set; }
  855. public string SourcePayload { get; set; } = string.Empty;
  856. }
  857. /// <summary>
  858. /// G01-05 去重结果对象。仅服务 G01-06 建单前拦截,由 CanCreate 单决策位决定是否建单。
  859. /// 只服务首版唯一场景 S2S6_PRODUCTION + 唯一 trigger_type VALUE_DEVIATION + 设备对象。
  860. /// 不预留多 trigger_type / 平台化去重扩展结构。
  861. /// Reason 值域:no_pending / duplicate_pending / missing_dedup_key / query_failed。
  862. /// </summary>
  863. public sealed class S8WatchDedupResult
  864. {
  865. public S8WatchHitResult Hit { get; set; } = new();
  866. public bool CanCreate { get; set; }
  867. public long? MatchedExceptionId { get; set; }
  868. public string Reason { get; set; } = string.Empty;
  869. }
  870. /// <summary>
  871. /// G01-06 建单结果对象。仅服务 G-01 首版主线验收,由 Created / Skipped 两位决定结局。
  872. /// 只服务首版唯一场景 S2S6_PRODUCTION + 唯一 trigger_type VALUE_DEVIATION + 设备对象。
  873. /// 不预留多 trigger_type / 平台化工单扩展结构。
  874. /// Reason 值域:auto_created / create_failed / 透传自 DedupResult.Reason。
  875. /// </summary>
  876. public sealed class S8WatchCreationResult
  877. {
  878. public S8WatchDedupResult DedupResult { get; set; } = new();
  879. public bool Created { get; set; }
  880. public bool Skipped { get; set; }
  881. public long? CreatedExceptionId { get; set; }
  882. public string Reason { get; set; } = string.Empty;
  883. public string? ErrorMessage { get; set; }
  884. }
  885. /// <summary>
  886. /// G01-04 命中结果对象。承载 G01-05 去重与 G01-06 建单所需最小追溯字段,
  887. /// 仅服务首版唯一场景 S2S6_PRODUCTION + 唯一 trigger_type VALUE_DEVIATION + 设备对象。
  888. /// 不预留多 trigger_type / 多场景 / 平台化扩展结构。
  889. /// </summary>
  890. public sealed class S8WatchHitResult
  891. {
  892. public long SourceRuleId { get; set; }
  893. public string SourceRuleCode { get; set; } = string.Empty;
  894. public long AlertRuleId { get; set; }
  895. public long DataSourceId { get; set; }
  896. public string RelatedObjectCode { get; set; } = string.Empty;
  897. public decimal CurrentValue { get; set; }
  898. public decimal ThresholdValue { get; set; }
  899. public string TriggerCondition { get; set; } = string.Empty;
  900. public string Severity { get; set; } = string.Empty;
  901. public long? OccurrenceDeptId { get; set; }
  902. public long? ResponsibleDeptId { get; set; }
  903. public string SourcePayload { get; set; } = string.Empty;
  904. }