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; /// /// 定时探活 mdp_source:DB SELECT 1 / API GET baseUrl(或 /health)。 /// [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(); var scopeFactory = scope.ServiceProvider.GetRequiredService(); using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(15) }; var sources = await db.Queryable().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); } }