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