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);
}
}