| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115 |
- using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
- namespace Admin.NET.Plugin.AiDOP.Infrastructure;
- /// <summary>
- /// ETL 执行机指派与实例注册表的写入侧。
- ///
- /// <para><b>「最多一台执行机」由 <c>ado_etl_runner_designation</c> 的结构保证</b>:
- /// 该表主键固定为 <see cref="AdoEtlRunnerDesignation.SingletonId"/>,放不下第二行,
- /// 因此不再需要旧模型里「同一事务内先清后置」那套写入侧约束。
- /// 并发指派会被主键与条件更新挡住,而不是依赖调用方守规矩。</para>
- ///
- /// <para><b>指派的对象是槽位,不是实例</b>:实例标识含进程ID与启动时间,
- /// 指派挂上去每次重启都会丢(见 <c>AdoEtlRunnerDesignation</c> 注释)。</para>
- /// </summary>
- public interface IEtlInstanceStore
- {
- /// <summary>
- /// 把某个部署槽位指派为唯一执行机。槽位无需当前有存活实例——
- /// 指派可以先于部署,也可以在目标机器重启期间保持有效,这正是本模型的要点。
- /// </summary>
- Task AssignRunnerSlotAsync(string slotCode, string assignedBy, CancellationToken ct = default);
- /// <summary>取消指派。此后无人跑定时 ETL——这是允许的状态。</summary>
- Task ClearRunnerAsync(string clearedBy = null, CancellationToken ct = default);
- /// <summary>当前被指派的槽位;无指派返回 null。</summary>
- Task<string> GetDesignatedSlotAsync(CancellationToken ct = default);
- /// <summary>指派记录全文,供页面展示「谁在什么时候指派的」。无指派返回 null。</summary>
- Task<AdoEtlRunnerDesignation> GetDesignationAsync(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 AssignRunnerSlotAsync(string slotCode, string assignedBy, CancellationToken ct = default)
- {
- var slot = string.IsNullOrWhiteSpace(slotCode) ? null : slotCode.Trim().ToLowerInvariant();
- await UpsertDesignationAsync(slot, assignedBy, ct);
- }
- public async Task ClearRunnerAsync(string clearedBy = null, CancellationToken ct = default) =>
- await UpsertDesignationAsync(null, clearedBy, ct);
- public async Task<string> GetDesignatedSlotAsync(CancellationToken ct = default)
- {
- var row = await GetDesignationAsync(ct);
- return string.IsNullOrWhiteSpace(row?.SlotCode) ? null : row.SlotCode;
- }
- public Task<AdoEtlRunnerDesignation> GetDesignationAsync(CancellationToken ct = default) =>
- _db.Queryable<AdoEtlRunnerDesignation>()
- .Where(x => x.Id == AdoEtlRunnerDesignation.SingletonId)
- .FirstAsync(ct);
- public Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default) =>
- _db.Queryable<AdoEtlInstance>()
- .OrderBy(x => x.MachineName)
- .OrderBy(x => x.ProcessId)
- .ToListAsync(ct);
- /// <summary>
- /// 写唯一那一行。先 UPDATE 再按影响行数决定是否 INSERT:
- /// 两个超管同时指派时,后到的 INSERT 会撞主键而不是插出第二行,
- /// 撞了就退回 UPDATE 重写一次,最终仍是单行单值。
- /// </summary>
- private async Task UpsertDesignationAsync(string slot, string operatorName, CancellationToken ct)
- {
- var now = DateTime.Now;
- var by = string.IsNullOrWhiteSpace(operatorName) ? null : operatorName.Trim();
- // 三元运算符不放进表达式树:SqlSugar 要把 SetColumns 的初始化器翻成 SQL,
- // 取消指派时 assigned_at 应回到 NULL,这个判断在 C# 侧先算完更稳。
- DateTime? assignedAt = slot == null ? null : now;
- var updated = await UpdateDesignationAsync(slot, by, assignedAt, now, ct);
- if (updated > 0) return;
- try
- {
- await _db.Insertable(new AdoEtlRunnerDesignation
- {
- Id = AdoEtlRunnerDesignation.SingletonId,
- SlotCode = slot,
- AssignedBy = by,
- AssignedAt = assignedAt,
- CreateTime = now,
- UpdateTime = now
- }).ExecuteCommandAsync(ct);
- }
- catch
- {
- // 并发下另一方已插入,改写即可;再失败就让异常上抛给控制器。
- await UpdateDesignationAsync(slot, by, assignedAt, now, ct);
- }
- }
- private Task<int> UpdateDesignationAsync(string slot, string by, DateTime? assignedAt, DateTime now, CancellationToken ct) =>
- _db.Updateable<AdoEtlRunnerDesignation>()
- .SetColumns(x => new AdoEtlRunnerDesignation
- {
- SlotCode = slot,
- AssignedBy = by,
- AssignedAt = assignedAt,
- UpdateTime = now
- })
- .Where(x => x.Id == AdoEtlRunnerDesignation.SingletonId)
- .ExecuteCommandAsync(ct);
- }
|