S8ActiveFlowWatchService.cs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281
  1. using Admin.NET.Plugin.AiDOP.Entity.S8;
  2. using Admin.NET.Plugin.ApprovalFlow;
  3. using Microsoft.Extensions.Logging;
  4. using Microsoft.Extensions.Options;
  5. using SqlSugar;
  6. namespace Admin.NET.Plugin.AiDOP.Service.S8;
  7. /// <summary>
  8. /// 扫描疑似卡死的 ActiveFlow 异常,并输出可检索告警。
  9. /// 当前只做发现,不做自动恢复。
  10. /// </summary>
  11. public class S8ActiveFlowWatchService : ITransient
  12. {
  13. public const string AlertChannel = "s8-active-flow-stuck";
  14. public const string AlertChannelOrphan = "s8-orphan-flow-instance";
  15. public static readonly string[] S8FlowBizTypes = new[] { "EXCEPTION_ESCALATION", "EXCEPTION_CLOSURE" };
  16. private readonly SqlSugarRepository<AdoS8Exception> _exceptionRep;
  17. private readonly SqlSugarRepository<AdoS8NotificationLog> _notificationLogRep;
  18. private readonly SqlSugarRepository<ApprovalFlowInstance> _flowInstanceRep;
  19. private readonly S8NotificationService _notificationService;
  20. private readonly S8ActiveFlowWatchOptions _options;
  21. private readonly ILogger<S8ActiveFlowWatchService> _logger;
  22. public S8ActiveFlowWatchService(
  23. SqlSugarRepository<AdoS8Exception> exceptionRep,
  24. SqlSugarRepository<AdoS8NotificationLog> notificationLogRep,
  25. SqlSugarRepository<ApprovalFlowInstance> flowInstanceRep,
  26. S8NotificationService notificationService,
  27. IOptions<S8ActiveFlowWatchOptions> options,
  28. ILogger<S8ActiveFlowWatchService> logger)
  29. {
  30. _exceptionRep = exceptionRep;
  31. _notificationLogRep = notificationLogRep;
  32. _flowInstanceRep = flowInstanceRep;
  33. _notificationService = notificationService;
  34. _options = options.Value;
  35. _logger = logger;
  36. }
  37. /// <summary>
  38. /// 列出仍有 running 审批实例的**租户**。
  39. ///
  40. /// <para><b>S8-TENANT-ONLY-BATCH6:删掉了 <c>AND e.factory_id &gt; 0</c>。</b>
  41. /// 这一条不是清理,是止血:本批之后新建异常的 <c>factory_id</c> 恒为 0,
  42. /// 该谓词会让卡死扫描**完全看不到本批之后产生的任何异常**——审批卡死不再有人告警,
  43. /// 且失效方式是静默的(job 照跑、日志照打、扫描集为空)。</para>
  44. ///
  45. /// <para>返回租户而不是 (租户, 工厂) 对:<see cref="ScanAsync"/> 的实际查询本来就只用 tenant_id,
  46. /// factory 段从来只是把同一批异常重复扫若干遍的分片键。</para>
  47. /// </summary>
  48. public async Task<List<long>> ListActiveTenantsAsync()
  49. {
  50. return await _exceptionRep.Context.Ado.SqlQueryAsync<long>(
  51. """
  52. SELECT DISTINCT e.tenant_id
  53. FROM ado_s8_exception e
  54. INNER JOIN SysTenant t ON t.Id = e.tenant_id AND t.Status = 1
  55. WHERE e.is_deleted = 0
  56. AND e.active_flow_instance_id IS NOT NULL
  57. AND e.tenant_id > 0
  58. ORDER BY e.tenant_id
  59. """);
  60. }
  61. public async Task<int> ScanAsync(long tenantId, CancellationToken cancellationToken = default)
  62. {
  63. if (!_options.Enabled)
  64. return 0;
  65. var now = DateTime.Now;
  66. var thresholdHours = _options.ThresholdHours > 0 ? _options.ThresholdHours : 4;
  67. var cooldownMinutes = _options.AlertCooldownMinutes > 0 ? _options.AlertCooldownMinutes : 60;
  68. var batchSize = _options.BatchSize > 0 ? _options.BatchSize : 100;
  69. var staleBefore = now.AddHours(-thresholdHours);
  70. var alertedAfter = now.AddMinutes(-cooldownMinutes);
  71. var staleExceptions = await _exceptionRep.AsQueryable()
  72. .Where(x => !x.IsDeleted && x.ActiveFlowInstanceId.HasValue &&
  73. x.TenantId == tenantId &&
  74. SqlFunc.IsNull(x.UpdatedAt, x.CreatedAt) < staleBefore)
  75. .OrderBy(x => x.UpdatedAt, OrderByType.Asc)
  76. .OrderBy(x => x.Id, OrderByType.Asc)
  77. .Take(batchSize)
  78. .ToListAsync();
  79. if (staleExceptions.Count == 0)
  80. return 0;
  81. var exceptionIds = staleExceptions.Select(x => x.Id).ToHashSet();
  82. var recentAlertLogs = await _notificationLogRep.AsQueryable()
  83. .Where(x => x.TenantId == tenantId &&
  84. x.Channel == AlertChannel && x.CreatedAt >= alertedAfter && x.ExceptionId != null)
  85. .ToListAsync();
  86. var alertedIds = recentAlertLogs
  87. .Where(x => x.ExceptionId.HasValue && exceptionIds.Contains(x.ExceptionId.Value))
  88. .Select(x => x.ExceptionId!.Value)
  89. .ToHashSet();
  90. var alertCount = 0;
  91. foreach (var item in staleExceptions.Where(x => !alertedIds.Contains(x.Id)))
  92. {
  93. cancellationToken.ThrowIfCancellationRequested();
  94. var lastTouchedAt = item.UpdatedAt ?? item.CreatedAt;
  95. var payload = new
  96. {
  97. type = "S8_ACTIVE_FLOW_STUCK",
  98. message = "S8 异常存在进行中的审批实例,但已超过阈值未更新,请人工检查审批流状态与回调链路",
  99. exceptionId = item.Id,
  100. exceptionCode = item.ExceptionCode,
  101. title = item.Title,
  102. status = item.Status,
  103. activeFlowInstanceId = item.ActiveFlowInstanceId,
  104. activeFlowBizType = item.ActiveFlowBizType,
  105. tenantId = item.TenantId,
  106. lastTouchedAt,
  107. thresholdHours,
  108. detectedAt = now
  109. };
  110. try
  111. {
  112. // factory 段写 0:告警日志的归属边界与异常本身一致,只由 tenant 决定。
  113. await _notificationService.SendAsync(item.TenantId, Infrastructure.S8ConfigScope.GlobalFactoryId, item.Id, AlertChannel, payload);
  114. _logger.LogWarning(
  115. "S8 ActiveFlow 疑似卡死: ExceptionId={ExceptionId}, Code={ExceptionCode}, Status={Status}, FlowInstanceId={FlowInstanceId}, LastTouchedAt={LastTouchedAt}, ThresholdHours={ThresholdHours}",
  116. item.Id,
  117. item.ExceptionCode,
  118. item.Status,
  119. item.ActiveFlowInstanceId,
  120. lastTouchedAt,
  121. thresholdHours);
  122. alertCount++;
  123. }
  124. catch (Exception ex)
  125. {
  126. _logger.LogError(
  127. ex,
  128. "S8 ActiveFlow 卡死告警写入失败: ExceptionId={ExceptionId}, Code={ExceptionCode}, FlowInstanceId={FlowInstanceId}",
  129. item.Id,
  130. item.ExceptionCode,
  131. item.ActiveFlowInstanceId);
  132. }
  133. }
  134. if (alertCount > 0)
  135. {
  136. _logger.LogInformation(
  137. "S8 ActiveFlow 卡死扫描完成,本轮新增告警 {AlertCount} 条,候选 {CandidateCount} 条",
  138. alertCount,
  139. staleExceptions.Count);
  140. }
  141. var orphanAlertCount = await ScanOrphanFlowInstancesAsync(staleBefore, alertedAfter, batchSize, now, thresholdHours, cancellationToken);
  142. return alertCount + orphanAlertCount;
  143. }
  144. /// <summary>
  145. /// N-2:扫描孤立的 S8 审批实例(FlowInstance 存在但无任何 AdoS8Exception 引用),疑似 OnFlowStarted 失败或回调链路丢失。
  146. /// </summary>
  147. private async Task<int> ScanOrphanFlowInstancesAsync(
  148. DateTime staleBefore,
  149. DateTime alertedAfter,
  150. int batchSize,
  151. DateTime now,
  152. int thresholdHours,
  153. CancellationToken cancellationToken)
  154. {
  155. // 本地副本:SqlSugar 表达式树不允许直接引用静态/私有字段,必须先拷到 lambda 闭包变量。
  156. var bizTypes = S8FlowBizTypes;
  157. var candidates = await _flowInstanceRep.AsQueryable()
  158. .Where(fi => bizTypes.Contains(fi.BizType)
  159. && fi.Status == FlowInstanceStatusEnum.Running
  160. && fi.StartTime < staleBefore)
  161. .OrderBy(fi => fi.StartTime, OrderByType.Asc)
  162. .Take(batchSize)
  163. .ToListAsync();
  164. if (candidates.Count == 0)
  165. return 0;
  166. var bizIds = candidates.Select(fi => fi.BizId).ToHashSet();
  167. var instanceIds = candidates.Select(fi => fi.Id).ToHashSet();
  168. var linked = await _exceptionRep.AsQueryable()
  169. .Where(e => !e.IsDeleted
  170. && bizIds.Contains(e.Id)
  171. && e.ActiveFlowInstanceId.HasValue
  172. && instanceIds.Contains(e.ActiveFlowInstanceId!.Value))
  173. .Select(e => new { e.Id, FlowId = e.ActiveFlowInstanceId!.Value })
  174. .ToListAsync();
  175. var linkedPairs = linked.Select(x => (x.Id, x.FlowId)).ToHashSet();
  176. var orphans = candidates.Where(fi => !linkedPairs.Contains((fi.BizId, fi.Id))).ToList();
  177. if (orphans.Count == 0)
  178. return 0;
  179. var recentPayloads = await _notificationLogRep.AsQueryable()
  180. .Where(x => x.Channel == AlertChannelOrphan && x.CreatedAt >= alertedAfter)
  181. .Select(x => x.Payload)
  182. .ToListAsync();
  183. var alertedInstanceIds = new HashSet<long>();
  184. foreach (var p in recentPayloads)
  185. {
  186. var id = TryExtractFlowInstanceId(p);
  187. if (id.HasValue) alertedInstanceIds.Add(id.Value);
  188. }
  189. var alertCount = 0;
  190. foreach (var fi in orphans.Where(x => !alertedInstanceIds.Contains(x.Id)))
  191. {
  192. cancellationToken.ThrowIfCancellationRequested();
  193. var payload = new
  194. {
  195. type = "S8_ORPHAN_FLOW_INSTANCE",
  196. message = "S8 审批实例存在但未被任何业务异常单据引用,疑似 OnFlowStarted 失败或回调链路丢失",
  197. flowInstanceId = fi.Id,
  198. bizType = fi.BizType,
  199. bizId = fi.BizId,
  200. bizNo = fi.BizNo,
  201. flowStatus = fi.Status.ToString(),
  202. startTime = fi.StartTime,
  203. thresholdHours,
  204. detectedAt = now
  205. };
  206. try
  207. {
  208. await _notificationService.SendAsync(0, 0, null, AlertChannelOrphan, payload);
  209. _logger.LogWarning(
  210. "S8 孤立审批实例: FlowInstanceId={InstanceId}, BizType={BizType}, BizId={BizId}, StartTime={StartTime}",
  211. fi.Id, fi.BizType, fi.BizId, fi.StartTime);
  212. alertCount++;
  213. }
  214. catch (Exception ex)
  215. {
  216. _logger.LogError(ex,
  217. "S8 孤立审批实例告警写入失败: FlowInstanceId={InstanceId}, BizType={BizType}, BizId={BizId}",
  218. fi.Id, fi.BizType, fi.BizId);
  219. }
  220. }
  221. if (alertCount > 0)
  222. {
  223. _logger.LogInformation(
  224. "S8 孤立审批实例扫描完成,本轮新增告警 {AlertCount} 条,候选 {CandidateCount} 条",
  225. alertCount, orphans.Count);
  226. }
  227. return alertCount;
  228. }
  229. /// <summary>
  230. /// 从历史告警 payload JSON 中提取 flowInstanceId 用于冷却去重。失败返回 null。
  231. /// </summary>
  232. internal static long? TryExtractFlowInstanceId(string? payload)
  233. {
  234. if (string.IsNullOrWhiteSpace(payload))
  235. return null;
  236. try
  237. {
  238. using var doc = System.Text.Json.JsonDocument.Parse(payload);
  239. if (doc.RootElement.TryGetProperty("flowInstanceId", out var prop)
  240. && prop.ValueKind == System.Text.Json.JsonValueKind.Number
  241. && prop.TryGetInt64(out var id))
  242. {
  243. return id;
  244. }
  245. }
  246. catch
  247. {
  248. // payload 非法 JSON 时忽略,不影响后续告警。
  249. }
  250. return null;
  251. }
  252. }