MdpSourceHealthCheckJob.cs 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212
  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 scopeDb = await scopeFactory.GetScopeAsync(src.SourceCode, stoppingToken);
  41. await scopeDb.Ado.GetIntAsync("SELECT 1");
  42. src.HealthStatus = 1;
  43. src.HealthMsg = "OK";
  44. ok++;
  45. }
  46. else if (string.Equals(src.SourceType, "API", StringComparison.OrdinalIgnoreCase))
  47. {
  48. if (string.IsNullOrWhiteSpace(src.ApiBaseUrl))
  49. throw new InvalidOperationException("api_base_url 为空");
  50. var url = src.ApiBaseUrl.TrimEnd('/') + "/";
  51. using var resp = await http.GetAsync(url, stoppingToken);
  52. src.HealthStatus = (int)resp.StatusCode is >= 200 and < 500 ? 1 : 0;
  53. src.HealthMsg = $"HTTP {(int)resp.StatusCode}";
  54. if (src.HealthStatus == 1) ok++; else fail++;
  55. }
  56. else if (string.Equals(src.SourceType, "API_INBOUND", StringComparison.OrdinalIgnoreCase))
  57. {
  58. src.HealthStatus = 1;
  59. src.HealthMsg = "OK";
  60. ok++;
  61. }
  62. else
  63. {
  64. src.HealthStatus = 0;
  65. src.HealthMsg = $"未知 source_type={src.SourceType}";
  66. fail++;
  67. }
  68. }
  69. catch (Exception ex)
  70. {
  71. src.HealthStatus = 0;
  72. src.HealthMsg = ex.Message.Length > 480 ? ex.Message[..480] : ex.Message;
  73. fail++;
  74. }
  75. src.LastHealthCheck = now;
  76. src.UpdateTime = now;
  77. await db.Updateable(src)
  78. .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
  79. .ExecuteCommandAsync(stoppingToken);
  80. }
  81. if (sources.Count > 0)
  82. _logger.LogInformation("[MdpSourceHealthCheckJob] total={Total} ok={Ok} fail={Fail}", sources.Count, ok, fail);
  83. await ExpireInboundSnapshotsAsync(scope, stoppingToken);
  84. await CheckInboundSilenceAsync(db, sources, stoppingToken);
  85. }
  86. /// <summary>顺带把 OPEN 且已过期的快照置 EXPIRED;未 commit 不产生删除语义。</summary>
  87. private async Task ExpireInboundSnapshotsAsync(IServiceScope scope, CancellationToken ct)
  88. {
  89. try
  90. {
  91. var snaps = scope.ServiceProvider.GetRequiredService<MdpInboundSnapshotService>();
  92. await snaps.ExpireStaleOpenAsync(ct);
  93. }
  94. catch (Exception ex)
  95. {
  96. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound snapshot expire failed");
  97. }
  98. }
  99. /// <summary>
  100. /// 入站实体沉默:inbound_enabled=1 且配置了 silence_alert_hours,
  101. /// 最近 COMMITTED 超过阈值且当天命中 silence_calendar(空=每天)则写回源健康告警。
  102. /// 从未推送过的实体以 update_time/create_time 为基线。
  103. /// </summary>
  104. private async Task CheckInboundSilenceAsync(
  105. ISqlSugarClient db, List<MdpSource> sources, CancellationToken ct)
  106. {
  107. List<MdpEntity> entities;
  108. try
  109. {
  110. entities = await db.Queryable<MdpEntity>()
  111. .Where(e => e.InboundEnabled == 1 && e.SilenceAlertHours != null)
  112. .ToListAsync(ct);
  113. }
  114. catch (Exception ex)
  115. {
  116. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound silence query entities failed");
  117. return;
  118. }
  119. if (entities.Count == 0)
  120. return;
  121. var today = DateTime.Now;
  122. var isoDow = today.DayOfWeek == DayOfWeek.Sunday ? 7 : (int)today.DayOfWeek;
  123. var alertsBySource = new Dictionary<long, List<string>>();
  124. foreach (var entity in entities)
  125. {
  126. ct.ThrowIfCancellationRequested();
  127. if (!HitsSilenceCalendar(entity.SilenceCalendar, isoDow))
  128. continue;
  129. DateTime? lastCommitted = null;
  130. try
  131. {
  132. var last = await db.Queryable<MdpInboundRequest>()
  133. .Where(r => r.Status == "COMMITTED")
  134. .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entity.EntityCode.ToUpperInvariant()))
  135. .OrderBy(r => r.CreateTime, OrderByType.Desc)
  136. .FirstAsync(ct);
  137. lastCommitted = last?.CreateTime;
  138. }
  139. catch (Exception ex)
  140. {
  141. _logger.LogWarning(ex, "[MdpSourceHealthCheckJob] inbound silence query requests failed entity={Entity}", entity.EntityCode);
  142. continue;
  143. }
  144. var baseline = lastCommitted
  145. ?? (entity.UpdateTime != default ? entity.UpdateTime : entity.CreateTime);
  146. var hours = entity.SilenceAlertHours.GetValueOrDefault();
  147. if (hours <= 0)
  148. continue;
  149. if ((today - baseline).TotalHours <= hours)
  150. continue;
  151. var lastText = lastCommitted.HasValue
  152. ? lastCommitted.Value.ToString("yyyy-MM-dd HH:mm:ss")
  153. : "never";
  154. var msg = $"INBOUND silence: {entity.EntityCode} last COMMITTED {lastText} (threshold {hours}h)";
  155. _logger.LogWarning("[MdpSourceHealthCheckJob] {Message}", msg);
  156. if (!alertsBySource.TryGetValue(entity.SourceId, out var list))
  157. {
  158. list = [];
  159. alertsBySource[entity.SourceId] = list;
  160. }
  161. list.Add(msg);
  162. }
  163. var sourceById = sources.ToDictionary(s => s.Id);
  164. foreach (var (sourceId, messages) in alertsBySource)
  165. {
  166. if (!sourceById.TryGetValue(sourceId, out var src))
  167. {
  168. src = await db.Queryable<MdpSource>().Where(s => s.Id == sourceId).FirstAsync(ct);
  169. if (src == null)
  170. continue;
  171. }
  172. var now = DateTime.Now;
  173. src.HealthStatus = 0;
  174. src.HealthMsg = Truncate(string.Join("; ", messages), 480);
  175. src.LastHealthCheck = now;
  176. src.UpdateTime = now;
  177. await db.Updateable(src)
  178. .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
  179. .ExecuteCommandAsync(ct);
  180. }
  181. }
  182. private static bool HitsSilenceCalendar(string calendar, int isoDow)
  183. {
  184. if (string.IsNullOrWhiteSpace(calendar))
  185. return true;
  186. var days = calendar.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
  187. return days.Any(d => int.TryParse(d, out var n) && n == isoDow);
  188. }
  189. private static string Truncate(string s, int max) =>
  190. string.IsNullOrEmpty(s) || s.Length <= max ? s : s[..max];
  191. }