MdpSourceHealthCheckJob.cs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Inbound;
  3. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  4. using Furion.Schedule;
  5. using Microsoft.Extensions.DependencyInjection;
  6. using Microsoft.Extensions.Logging;
  7. namespace Admin.NET.Plugin.AiDOP.Job;
  8. /// <summary>
  9. /// 定时探活 mdp_source:DB SELECT 1 / API GET baseUrl(或 /health)。
  10. /// </summary>
  11. [JobDetail("job_mdp_source_health", Description = "MDP 数据源健康检查", GroupName = "default", Concurrent = false)]
  12. [PeriodSeconds(300, TriggerId = "trigger_mdp_source_health", Description = "每 5 分钟探活数据源", RunOnStart = false)]
  13. public class MdpSourceHealthCheckJob : IJob
  14. {
  15. private readonly IServiceScopeFactory _scopeFactory;
  16. private readonly ILogger _logger;
  17. public MdpSourceHealthCheckJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  18. {
  19. _scopeFactory = scopeFactory;
  20. _logger = loggerFactory.CreateLogger(nameof(MdpSourceHealthCheckJob));
  21. }
  22. public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
  23. {
  24. if (!AidopJobGate.ShouldRun(nameof(MdpSourceHealthCheckJob), _logger)) return;
  25. using var scope = _scopeFactory.CreateScope();
  26. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  27. var scopeFactory = scope.ServiceProvider.GetRequiredService<MdpSourceScopeFactory>();
  28. using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(15) };
  29. var sources = await db.Queryable<MdpSource>().Where(x => x.Status == 1).ToListAsync(stoppingToken);
  30. var ok = 0;
  31. var fail = 0;
  32. foreach (var src in sources)
  33. {
  34. stoppingToken.ThrowIfCancellationRequested();
  35. var now = DateTime.Now;
  36. try
  37. {
  38. if (string.Equals(src.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  39. {
  40. var illegal = MdpSourceScopeFactory.IllegalExtraKeys(src.DbType, src.DbExtraParams);
  41. var illegalNote = illegal.Count == 0
  42. ? null
  43. : $"db_extra_params 含非法关键字 {string.Join(",", illegal)},已忽略";
  44. var scopeDb = await scopeFactory.GetScopeAsync(src.SourceCode, stoppingToken);
  45. await scopeDb.Ado.GetIntAsync("SELECT 1");
  46. src.HealthStatus = 1;
  47. src.HealthMsg = illegalNote == null ? "OK" : $"OK({illegalNote})";
  48. ok++;
  49. }
  50. else if (string.Equals(src.SourceType, "API", StringComparison.OrdinalIgnoreCase))
  51. {
  52. if (string.IsNullOrWhiteSpace(src.ApiBaseUrl))
  53. throw new InvalidOperationException("api_base_url 为空");
  54. var url = src.ApiBaseUrl.TrimEnd('/') + "/";
  55. using var resp = await http.GetAsync(url, stoppingToken);
  56. src.HealthStatus = (int)resp.StatusCode is >= 200 and < 500 ? 1 : 0;
  57. src.HealthMsg = $"HTTP {(int)resp.StatusCode}";
  58. if (src.HealthStatus == 1) ok++; else fail++;
  59. }
  60. else if (string.Equals(src.SourceType, "API_INBOUND", StringComparison.OrdinalIgnoreCase))
  61. {
  62. src.HealthStatus = 1;
  63. src.HealthMsg = "OK";
  64. ok++;
  65. }
  66. else
  67. {
  68. src.HealthStatus = 0;
  69. src.HealthMsg = $"未知 source_type={src.SourceType}";
  70. fail++;
  71. }
  72. }
  73. catch (Exception ex)
  74. {
  75. src.HealthStatus = 0;
  76. var code = ClassifyHealthFailure(ex);
  77. var illegal = string.Equals(src.SourceType, "DB", StringComparison.OrdinalIgnoreCase)
  78. ? MdpSourceScopeFactory.IllegalExtraKeys(src.DbType, src.DbExtraParams)
  79. : Array.Empty<string>();
  80. var prefix = illegal.Count == 0
  81. ? ""
  82. : $"db_extra_params 含非法关键字 {string.Join(",", illegal)},已忽略。";
  83. src.HealthMsg = prefix + code + " " + SafeHealthMessage(ex.Message);
  84. if (src.HealthMsg.Length > 480)
  85. src.HealthMsg = src.HealthMsg[..480];
  86. _logger.LogWarning(
  87. "[MdpSourceHealthCheckJob] source={SourceCode} code={Code}",
  88. src.SourceCode, code);
  89. fail++;
  90. }
  91. src.LastHealthCheck = now;
  92. src.UpdateTime = now;
  93. await db.Updateable(src)
  94. .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
  95. .ExecuteCommandAsync(stoppingToken);
  96. }
  97. if (sources.Count > 0)
  98. _logger.LogInformation("[MdpSourceHealthCheckJob] total={Total} ok={Ok} fail={Fail}", sources.Count, ok, fail);
  99. await ExpireInboundSnapshotsAsync(scope, stoppingToken);
  100. await CheckInboundSilenceAsync(db, sources, stoppingToken);
  101. }
  102. /// <summary>顺带把 OPEN 且已过期的快照置 EXPIRED;未 commit 不产生删除语义。</summary>
  103. private async Task ExpireInboundSnapshotsAsync(IServiceScope scope, CancellationToken ct)
  104. {
  105. try
  106. {
  107. var snaps = scope.ServiceProvider.GetRequiredService<MdpInboundSnapshotService>();
  108. await snaps.ExpireStaleOpenAsync(ct);
  109. }
  110. catch (Exception ex)
  111. {
  112. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound snapshot expire failed");
  113. }
  114. }
  115. /// <summary>
  116. /// 入站实体沉默:inbound_enabled=1 且配置了 silence_alert_hours,
  117. /// 最近 COMMITTED 超过阈值且当天命中 silence_calendar(空=每天)则写回源健康告警。
  118. /// 从未推送过的实体以 update_time/create_time 为基线。
  119. /// </summary>
  120. private async Task CheckInboundSilenceAsync(
  121. ISqlSugarClient db, List<MdpSource> sources, CancellationToken ct)
  122. {
  123. List<MdpEntity> entities;
  124. try
  125. {
  126. entities = await db.Queryable<MdpEntity>()
  127. .Where(e => e.InboundEnabled == 1 && e.SilenceAlertHours != null)
  128. .ToListAsync(ct);
  129. }
  130. catch (Exception ex)
  131. {
  132. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound silence query entities failed");
  133. return;
  134. }
  135. if (entities.Count == 0)
  136. return;
  137. var today = DateTime.Now;
  138. var isoDow = today.DayOfWeek == DayOfWeek.Sunday ? 7 : (int)today.DayOfWeek;
  139. var alertsBySource = new Dictionary<long, List<string>>();
  140. foreach (var entity in entities)
  141. {
  142. ct.ThrowIfCancellationRequested();
  143. if (!HitsSilenceCalendar(entity.SilenceCalendar, isoDow))
  144. continue;
  145. DateTime? lastCommitted = null;
  146. try
  147. {
  148. var last = await db.Queryable<MdpInboundRequest>()
  149. .Where(r => r.Status == "COMMITTED")
  150. .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entity.EntityCode.ToUpperInvariant()))
  151. .OrderBy(r => r.CreateTime, OrderByType.Desc)
  152. .FirstAsync(ct);
  153. lastCommitted = last?.CreateTime;
  154. }
  155. catch (Exception ex)
  156. {
  157. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound silence query requests failed entity={Entity}", entity.EntityCode);
  158. continue;
  159. }
  160. var baseline = lastCommitted
  161. ?? (entity.UpdateTime != default ? entity.UpdateTime : entity.CreateTime);
  162. var hours = entity.SilenceAlertHours.GetValueOrDefault();
  163. if (hours <= 0)
  164. continue;
  165. if ((today - baseline).TotalHours <= hours)
  166. continue;
  167. var lastText = lastCommitted.HasValue
  168. ? lastCommitted.Value.ToString("yyyy-MM-dd HH:mm:ss")
  169. : "never";
  170. var msg = $"INBOUND silence: {entity.EntityCode} last COMMITTED {lastText} (threshold {hours}h)";
  171. _logger.LogWarning("[MdpSourceHealthCheckJob] {Message}", msg);
  172. if (!alertsBySource.TryGetValue(entity.SourceId, out var list))
  173. {
  174. list = [];
  175. alertsBySource[entity.SourceId] = list;
  176. }
  177. list.Add(msg);
  178. }
  179. var sourceById = sources.ToDictionary(s => s.Id);
  180. foreach (var (sourceId, messages) in alertsBySource)
  181. {
  182. if (!sourceById.TryGetValue(sourceId, out var src))
  183. {
  184. src = await db.Queryable<MdpSource>().Where(s => s.Id == sourceId).FirstAsync(ct);
  185. if (src == null)
  186. continue;
  187. }
  188. var now = DateTime.Now;
  189. src.HealthStatus = 0;
  190. src.HealthMsg = Truncate(string.Join("; ", messages), 480);
  191. src.LastHealthCheck = now;
  192. src.UpdateTime = now;
  193. await db.Updateable(src)
  194. .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
  195. .ExecuteCommandAsync(ct);
  196. }
  197. }
  198. private static bool HitsSilenceCalendar(string calendar, int isoDow)
  199. {
  200. if (string.IsNullOrWhiteSpace(calendar))
  201. return true;
  202. var days = calendar.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
  203. return days.Any(d => int.TryParse(d, out var n) && n == isoDow);
  204. }
  205. private static string Truncate(string s, int max) =>
  206. string.IsNullOrEmpty(s) || s.Length <= max ? s : s[..max];
  207. internal static string ClassifyHealthFailure(Exception ex)
  208. {
  209. if (ex is MdpSourceSecretException secret)
  210. return secret.Code;
  211. for (var current = ex; current != null; current = current.InnerException)
  212. {
  213. var number = current.GetType().GetProperty("Number")?.GetValue(current);
  214. if (number is int n && (n == 1045 || n == 18456))
  215. return "AUTH_FAILED";
  216. var msg = current.Message ?? "";
  217. if (msg.Contains("登录失败", StringComparison.Ordinal)
  218. || msg.Contains("Login failed", StringComparison.OrdinalIgnoreCase)
  219. || msg.Contains("Access denied", StringComparison.OrdinalIgnoreCase))
  220. return "AUTH_FAILED";
  221. if (msg.Contains("timeout", StringComparison.OrdinalIgnoreCase)
  222. || msg.Contains("Unable to connect", StringComparison.OrdinalIgnoreCase)
  223. || msg.Contains("network-related", StringComparison.OrdinalIgnoreCase)
  224. || msg.Contains("连接", StringComparison.Ordinal))
  225. return "NETWORK_FAILED";
  226. }
  227. return "FAILED";
  228. }
  229. internal static string SafeHealthMessage(string? message)
  230. {
  231. var text = message ?? "";
  232. if (text.Contains("Password=", StringComparison.OrdinalIgnoreCase)
  233. || text.Contains("Pwd=", StringComparison.OrdinalIgnoreCase)
  234. || text.Contains("ConnectionString", StringComparison.OrdinalIgnoreCase)
  235. || text.Contains("Server=", StringComparison.OrdinalIgnoreCase))
  236. return "连接失败,详情见运行日志";
  237. return text.Length > 400 ? text[..400] : text;
  238. }
  239. }