using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
namespace Admin.NET.Plugin.AiDOP.Infrastructure;
///
/// ETL 执行机实例注册表的写入侧。
///
/// 「最多一台执行机」这条不变式由本类保证,不要在控制器、前端或数据库约束里另做一份:
/// MySQL 没有部分索引,无法用唯一索引表达「最多一行 is_runner=1」;
/// 而前端保证等于没有保证(两个超管同时点两台就会破)。
/// 故 在**同一事务内**先清后置。
///
public interface IEtlInstanceStore
{
/// 把指定实例指派为唯一执行机。同一事务内先把其余行置 0。
/// false 表示该 instanceId 不存在(进程可能已退出并被清理)。
Task AssignRunnerAsync(string instanceId, CancellationToken ct = default);
/// 取消所有指派。此后无人跑定时 ETL——这是允许的状态。
Task ClearRunnerAsync(CancellationToken ct = default);
Task> ListAsync(CancellationToken ct = default);
}
public sealed class EtlInstanceStore : IEtlInstanceStore, ITransient
{
private readonly ISqlSugarClient _db;
public EtlInstanceStore(ISqlSugarClient db) => _db = db;
public async Task AssignRunnerAsync(string instanceId, CancellationToken ct = default)
{
if (string.IsNullOrWhiteSpace(instanceId))
return false;
var id = instanceId.Trim();
var assigned = false;
// 事务边界不可省:先清后置若被拆成两个自动提交语句,
// 中间那一瞬会出现「零台执行机」或(并发指派时)「两台执行机」。
await _db.Ado.UseTranAsync(async () =>
{
var now = DateTime.Now;
await _db.Updateable()
.SetColumns(x => new AdoEtlInstance { IsRunner = false, UpdateTime = now })
.Where(x => x.IsRunner && x.InstanceId != id)
.ExecuteCommandAsync(ct);
var n = await _db.Updateable()
.SetColumns(x => new AdoEtlInstance { IsRunner = true, UpdateTime = now })
.Where(x => x.InstanceId == id)
.ExecuteCommandAsync(ct);
assigned = n > 0;
});
return assigned;
}
public async Task ClearRunnerAsync(CancellationToken ct = default)
{
var now = DateTime.Now;
await _db.Updateable()
.SetColumns(x => new AdoEtlInstance { IsRunner = false, UpdateTime = now })
.Where(x => x.IsRunner)
.ExecuteCommandAsync(ct);
}
public Task> ListAsync(CancellationToken ct = default) =>
_db.Queryable()
.OrderBy(x => x.MachineName)
.OrderBy(x => x.ProcessId)
.ToListAsync(ct);
}