| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475 |
- using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
- namespace Admin.NET.Plugin.AiDOP.Infrastructure;
- /// <summary>
- /// ETL 执行机实例注册表的写入侧。
- ///
- /// <para><b>「最多一台执行机」这条不变式由本类保证</b>,不要在控制器、前端或数据库约束里另做一份:
- /// MySQL 没有部分索引,无法用唯一索引表达「最多一行 <c>is_runner=1</c>」;
- /// 而前端保证等于没有保证(两个超管同时点两台就会破)。
- /// 故 <see cref="AssignRunnerAsync"/> 在**同一事务内**先清后置。</para>
- /// </summary>
- public interface IEtlInstanceStore
- {
- /// <summary>把指定实例指派为唯一执行机。同一事务内先把其余行置 0。</summary>
- /// <returns>false 表示该 instanceId 不存在(进程可能已退出并被清理)。</returns>
- Task<bool> AssignRunnerAsync(string instanceId, CancellationToken ct = default);
- /// <summary>取消所有指派。此后无人跑定时 ETL——这是允许的状态。</summary>
- Task ClearRunnerAsync(CancellationToken ct = default);
- Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default);
- }
- public sealed class EtlInstanceStore : IEtlInstanceStore, ITransient
- {
- private readonly ISqlSugarClient _db;
- public EtlInstanceStore(ISqlSugarClient db) => _db = db;
- public async Task<bool> 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<AdoEtlInstance>()
- .SetColumns(x => new AdoEtlInstance { IsRunner = false, UpdateTime = now })
- .Where(x => x.IsRunner && x.InstanceId != id)
- .ExecuteCommandAsync(ct);
- var n = await _db.Updateable<AdoEtlInstance>()
- .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<AdoEtlInstance>()
- .SetColumns(x => new AdoEtlInstance { IsRunner = false, UpdateTime = now })
- .Where(x => x.IsRunner)
- .ExecuteCommandAsync(ct);
- }
- public Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default) =>
- _db.Queryable<AdoEtlInstance>()
- .OrderBy(x => x.MachineName)
- .OrderBy(x => x.ProcessId)
- .ToListAsync(ct);
- }
|