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