using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace Admin.NET.Plugin.AiDOP.Infrastructure;
///
/// ETL 执行机实例注册器:把本进程登记进 ado_etl_instance 并定期心跳,
/// 同时把本实例的指派位发布到 供 读取。
///
/// 只读指派、从不自行指派:本注册器插入新行时 is_runner 恒为 0。
/// 指派是人工动作(超管在 MDP 运行监控页操作,或直接改库)。
/// 禁止在这里实现「抢锁自荐」或「心跳超时自动接管」——
/// 那会让开发机在某次重启后意外变成执行机,正是本次治理要消除的情形。
///
/// 本服务不受 管辖,必须常开。
/// 它是闸门的数据来源,被闸门关掉就成了死锁:没人心跳 → 没人发布指派 → 永远不是执行机。
///
public sealed class EtlInstanceRegistrar : BackgroundService
{
/// 心跳周期。 按其 3 倍取值。
public static readonly TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(20);
/// 启动延时:等宿主配置与建表脚本就绪,避免首拍必然失败刷一条 warn。
private static readonly TimeSpan StartupDelay = TimeSpan.FromSeconds(5);
/// 下线实例行的保留期。超期且未被指派的行会被清理,防止表随重启无限增长。
private static readonly TimeSpan DeadInstanceRetention = TimeSpan.FromDays(7);
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger _logger;
private int _consecutiveFailures;
private int _guardrailsLogged;
public EtlInstanceRegistrar(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
{
_scopeFactory = scopeFactory;
_logger = loggerFactory.CreateLogger(nameof(EtlInstanceRegistrar));
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
try { await Task.Delay(StartupDelay, stoppingToken); }
catch (OperationCanceledException) { return; }
_logger.LogInformation(
"[EtlInstanceRegistrar] 本实例 instanceId={InstanceId} machine={Machine} pid={Pid} version={Version}",
AidopInstanceIdentity.InstanceId, AidopInstanceIdentity.MachineName,
AidopInstanceIdentity.ProcessId, AidopInstanceIdentity.AppVersion);
var nextCleanupAt = DateTimeOffset.UtcNow;
while (!stoppingToken.IsCancellationRequested)
{
try
{
var assigned = await HeartbeatAsync(stoppingToken);
AidopRunnerState.Publish(assigned);
if (Interlocked.Exchange(ref _guardrailsLogged, 1) == 0)
await LogMySqlGuardrailsAsync(assigned, stoppingToken);
if (_consecutiveFailures > 0)
{
_logger.LogInformation(
"[EtlInstanceRegistrar] 心跳恢复,连续失败 {Count} 次后已重新发布指派 isRunner={IsRunner}",
_consecutiveFailures, assigned);
_consecutiveFailures = 0;
}
if (DateTimeOffset.UtcNow >= nextCleanupAt)
{
await CleanupDeadAsync(stoppingToken);
nextCleanupAt = DateTimeOffset.UtcNow + TimeSpan.FromHours(6);
}
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
break;
}
catch (Exception ex)
{
// 查库失败必须打回未知,让闸门回退到配置值:
// 指派可能已被运维改掉,继续沿用旧值会让两台实例同时自认执行机。
AidopRunnerState.Invalidate();
// 库不可用时本循环每 20 秒一次,不做节流会刷屏淹没真正的错误。
_consecutiveFailures++;
if (_consecutiveFailures == 1 || _consecutiveFailures % 30 == 0)
_logger.LogWarning(ex,
"[EtlInstanceRegistrar] 心跳失败(连续第 {Count} 次),指派状态已置未知,闸门回退到配置值",
_consecutiveFailures);
}
try { await Task.Delay(HeartbeatInterval, stoppingToken); }
catch (OperationCanceledException) { break; }
}
}
///
/// 心跳一拍:更新本实例行(不存在则插入),回读指派位。
/// 返回本实例是否被指派为执行机。
///
private async Task HeartbeatAsync(CancellationToken ct)
{
using var scope = _scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
var instanceId = AidopInstanceIdentity.InstanceId;
var now = DateTime.Now;
// 刻意不更新 is_runner:指派归人工,注册器只报到。
var affected = await db.Updateable()
.SetColumns(x => new AdoEtlInstance
{
LastHeartbeatAt = now,
AppVersion = AidopInstanceIdentity.AppVersion,
UpdateTime = now
})
.Where(x => x.InstanceId == instanceId)
.ExecuteCommandAsync(ct);
if (affected <= 0)
{
await db.Insertable(new AdoEtlInstance
{
InstanceId = instanceId,
MachineName = AidopInstanceIdentity.MachineName,
ProcessId = AidopInstanceIdentity.ProcessId,
AppVersion = AidopInstanceIdentity.AppVersion,
IsRunner = false,
StartedAt = AidopInstanceIdentity.StartedAt,
LastHeartbeatAt = now,
CreateTime = now,
UpdateTime = now
}).ExecuteCommandAsync(ct);
}
return await db.Queryable()
.Where(x => x.InstanceId == instanceId)
.Select(x => x.IsRunner)
.FirstAsync(ct);
}
///
/// 首次心跳成功后打一条现场参数,全进程只打一次。
///
/// 为什么要打:2026-09-24 的 1040 事故里,165 的
/// innodb_buffer_pool_size 一直是默认 128MB、max_connections 是默认 151,
/// 而没有任何人能从日志看出这一点——只能事后 SSH 上去查。
/// 把它连同执行机身份打进启动日志,下次排障第一眼就能看到。
///
/// 不打连接串、不打任何口令。失败只 warn,不影响心跳。
///
private async Task LogMySqlGuardrailsAsync(bool assigned, CancellationToken ct)
{
try
{
using var scope = _scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
var row = await db.Ado.SqlQuerySingleAsync(
"SELECT @@max_connections AS MaxConnections, "
+ "@@innodb_buffer_pool_size AS BufferPoolBytes, "
+ "DATABASE() AS DbName");
_logger.LogInformation(
"[EtlInstanceRegistrar] 现场参数 db={Db} maxConnections={MaxConn} bufferPoolMB={PoolMb} "
+ "isRunner={IsRunner} instanceId={InstanceId}",
row?.DbName, row?.MaxConnections, row?.BufferPoolBytes / 1024 / 1024,
assigned, AidopInstanceIdentity.InstanceId);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "[EtlInstanceRegistrar] 读取现场 MySQL 参数失败,跳过");
}
}
private sealed class MySqlGuardrails
{
public long MaxConnections { get; set; }
public long BufferPoolBytes { get; set; }
public string DbName { get; set; }
}
///
/// 清理早已下线的实例行。被指派的那行即使心跳超期也保留——
/// 运维需要看到「指派的机器挂了」,而不是记录凭空消失。
///
private async Task CleanupDeadAsync(CancellationToken ct)
{
try
{
using var scope = _scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
var cutoff = DateTime.Now - DeadInstanceRetention;
var removed = await db.Deleteable()
.Where(x => !x.IsRunner && x.LastHeartbeatAt < cutoff)
.ExecuteCommandAsync(ct);
if (removed > 0)
_logger.LogInformation("[EtlInstanceRegistrar] 已清理下线实例 {Count} 行", removed);
}
catch (Exception ex)
{
// 清理失败不影响心跳与指派,吞掉即可
_logger.LogWarning(ex, "[EtlInstanceRegistrar] 清理下线实例失败");
}
}
}