MdpSourceHealthCheckJob.cs 9.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  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 illegal = string.Equals(src.SourceType, "DB", StringComparison.OrdinalIgnoreCase)
  77. ? MdpSourceScopeFactory.IllegalExtraKeys(src.DbType, src.DbExtraParams)
  78. : Array.Empty<string>();
  79. var prefix = illegal.Count == 0
  80. ? ""
  81. : $"db_extra_params 含非法关键字 {string.Join(",", illegal)},已忽略。";
  82. var msg = prefix + ex.Message;
  83. src.HealthMsg = msg.Length > 480 ? msg[..480] : msg;
  84. fail++;
  85. }
  86. src.LastHealthCheck = now;
  87. src.UpdateTime = now;
  88. await db.Updateable(src)
  89. .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
  90. .ExecuteCommandAsync(stoppingToken);
  91. }
  92. if (sources.Count > 0)
  93. _logger.LogInformation("[MdpSourceHealthCheckJob] total={Total} ok={Ok} fail={Fail}", sources.Count, ok, fail);
  94. await ExpireInboundSnapshotsAsync(scope, stoppingToken);
  95. await CheckInboundSilenceAsync(db, sources, stoppingToken);
  96. }
  97. /// <summary>顺带把 OPEN 且已过期的快照置 EXPIRED;未 commit 不产生删除语义。</summary>
  98. private async Task ExpireInboundSnapshotsAsync(IServiceScope scope, CancellationToken ct)
  99. {
  100. try
  101. {
  102. var snaps = scope.ServiceProvider.GetRequiredService<MdpInboundSnapshotService>();
  103. await snaps.ExpireStaleOpenAsync(ct);
  104. }
  105. catch (Exception ex)
  106. {
  107. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound snapshot expire failed");
  108. }
  109. }
  110. /// <summary>
  111. /// 入站实体沉默:inbound_enabled=1 且配置了 silence_alert_hours,
  112. /// 最近 COMMITTED 超过阈值且当天命中 silence_calendar(空=每天)则写回源健康告警。
  113. /// 从未推送过的实体以 update_time/create_time 为基线。
  114. /// </summary>
  115. private async Task CheckInboundSilenceAsync(
  116. ISqlSugarClient db, List<MdpSource> sources, CancellationToken ct)
  117. {
  118. List<MdpEntity> entities;
  119. try
  120. {
  121. entities = await db.Queryable<MdpEntity>()
  122. .Where(e => e.InboundEnabled == 1 && e.SilenceAlertHours != null)
  123. .ToListAsync(ct);
  124. }
  125. catch (Exception ex)
  126. {
  127. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound silence query entities failed");
  128. return;
  129. }
  130. if (entities.Count == 0)
  131. return;
  132. var today = DateTime.Now;
  133. var isoDow = today.DayOfWeek == DayOfWeek.Sunday ? 7 : (int)today.DayOfWeek;
  134. var alertsBySource = new Dictionary<long, List<string>>();
  135. foreach (var entity in entities)
  136. {
  137. ct.ThrowIfCancellationRequested();
  138. if (!HitsSilenceCalendar(entity.SilenceCalendar, isoDow))
  139. continue;
  140. DateTime? lastCommitted = null;
  141. try
  142. {
  143. var last = await db.Queryable<MdpInboundRequest>()
  144. .Where(r => r.Status == "COMMITTED")
  145. .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entity.EntityCode.ToUpperInvariant()))
  146. .OrderBy(r => r.CreateTime, OrderByType.Desc)
  147. .FirstAsync(ct);
  148. lastCommitted = last?.CreateTime;
  149. }
  150. catch (Exception ex)
  151. {
  152. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound silence query requests failed entity={Entity}", entity.EntityCode);
  153. continue;
  154. }
  155. var baseline = lastCommitted
  156. ?? (entity.UpdateTime != default ? entity.UpdateTime : entity.CreateTime);
  157. var hours = entity.SilenceAlertHours.GetValueOrDefault();
  158. if (hours <= 0)
  159. continue;
  160. if ((today - baseline).TotalHours <= hours)
  161. continue;
  162. var lastText = lastCommitted.HasValue
  163. ? lastCommitted.Value.ToString("yyyy-MM-dd HH:mm:ss")
  164. : "never";
  165. var msg = $"INBOUND silence: {entity.EntityCode} last COMMITTED {lastText} (threshold {hours}h)";
  166. _logger.LogWarning("[MdpSourceHealthCheckJob] {Message}", msg);
  167. if (!alertsBySource.TryGetValue(entity.SourceId, out var list))
  168. {
  169. list = [];
  170. alertsBySource[entity.SourceId] = list;
  171. }
  172. list.Add(msg);
  173. }
  174. var sourceById = sources.ToDictionary(s => s.Id);
  175. foreach (var (sourceId, messages) in alertsBySource)
  176. {
  177. if (!sourceById.TryGetValue(sourceId, out var src))
  178. {
  179. src = await db.Queryable<MdpSource>().Where(s => s.Id == sourceId).FirstAsync(ct);
  180. if (src == null)
  181. continue;
  182. }
  183. var now = DateTime.Now;
  184. src.HealthStatus = 0;
  185. src.HealthMsg = Truncate(string.Join("; ", messages), 480);
  186. src.LastHealthCheck = now;
  187. src.UpdateTime = now;
  188. await db.Updateable(src)
  189. .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
  190. .ExecuteCommandAsync(ct);
  191. }
  192. }
  193. private static bool HitsSilenceCalendar(string calendar, int isoDow)
  194. {
  195. if (string.IsNullOrWhiteSpace(calendar))
  196. return true;
  197. var days = calendar.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
  198. return days.Any(d => int.TryParse(d, out var n) && n == isoDow);
  199. }
  200. private static string Truncate(string s, int max) =>
  201. string.IsNullOrEmpty(s) || s.Length <= max ? s : s[..max];
  202. }