| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283 |
- using Admin.NET.Plugin.AiDOP.DataPlatform;
- using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
- using Furion.Schedule;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.Job;
- /// <summary>
- /// 定时探活 mdp_source:DB SELECT 1 / API GET baseUrl(或 /health)。
- /// </summary>
- [JobDetail("job_mdp_source_health", Description = "MDP 数据源健康检查", GroupName = "default", Concurrent = false)]
- [PeriodSeconds(300, TriggerId = "trigger_mdp_source_health", Description = "每 5 分钟探活数据源", RunOnStart = false)]
- public class MdpSourceHealthCheckJob : IJob
- {
- private readonly IServiceScopeFactory _scopeFactory;
- private readonly ILogger _logger;
- public MdpSourceHealthCheckJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
- {
- _scopeFactory = scopeFactory;
- _logger = loggerFactory.CreateLogger(nameof(MdpSourceHealthCheckJob));
- }
- public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
- {
- using var scope = _scopeFactory.CreateScope();
- var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
- var scopeFactory = scope.ServiceProvider.GetRequiredService<MdpSourceScopeFactory>();
- using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(15) };
- var sources = await db.Queryable<MdpSource>().Where(x => x.Status == 1).ToListAsync(stoppingToken);
- var ok = 0;
- var fail = 0;
- foreach (var src in sources)
- {
- stoppingToken.ThrowIfCancellationRequested();
- var now = DateTime.Now;
- try
- {
- if (string.Equals(src.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
- {
- var scopeDb = await scopeFactory.GetScopeAsync(src.SourceCode, stoppingToken);
- await scopeDb.Ado.GetIntAsync("SELECT 1");
- src.HealthStatus = 1;
- src.HealthMsg = "OK";
- ok++;
- }
- else if (string.Equals(src.SourceType, "API", StringComparison.OrdinalIgnoreCase))
- {
- if (string.IsNullOrWhiteSpace(src.ApiBaseUrl))
- throw new InvalidOperationException("api_base_url 为空");
- var url = src.ApiBaseUrl.TrimEnd('/') + "/";
- using var resp = await http.GetAsync(url, stoppingToken);
- src.HealthStatus = (int)resp.StatusCode is >= 200 and < 500 ? 1 : 0;
- src.HealthMsg = $"HTTP {(int)resp.StatusCode}";
- if (src.HealthStatus == 1) ok++; else fail++;
- }
- else
- {
- src.HealthStatus = 0;
- src.HealthMsg = $"未知 source_type={src.SourceType}";
- fail++;
- }
- }
- catch (Exception ex)
- {
- src.HealthStatus = 0;
- src.HealthMsg = ex.Message.Length > 480 ? ex.Message[..480] : ex.Message;
- fail++;
- }
- src.LastHealthCheck = now;
- src.UpdateTime = now;
- await db.Updateable(src)
- .UpdateColumns(x => new { x.HealthStatus, x.HealthMsg, x.LastHealthCheck, x.UpdateTime })
- .ExecuteCommandAsync(stoppingToken);
- }
- if (sources.Count > 0)
- _logger.LogInformation("[MdpSourceHealthCheckJob] total={Total} ok={Ok} fail={Fail}", sources.Count, ok, fail);
- }
- }
|