EtlInstanceStore.cs 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115
  1. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  2. namespace Admin.NET.Plugin.AiDOP.Infrastructure;
  3. /// <summary>
  4. /// ETL 执行机指派与实例注册表的写入侧。
  5. ///
  6. /// <para><b>「最多一台执行机」由 <c>ado_etl_runner_designation</c> 的结构保证</b>:
  7. /// 该表主键固定为 <see cref="AdoEtlRunnerDesignation.SingletonId"/>,放不下第二行,
  8. /// 因此不再需要旧模型里「同一事务内先清后置」那套写入侧约束。
  9. /// 并发指派会被主键与条件更新挡住,而不是依赖调用方守规矩。</para>
  10. ///
  11. /// <para><b>指派的对象是槽位,不是实例</b>:实例标识含进程ID与启动时间,
  12. /// 指派挂上去每次重启都会丢(见 <c>AdoEtlRunnerDesignation</c> 注释)。</para>
  13. /// </summary>
  14. public interface IEtlInstanceStore
  15. {
  16. /// <summary>
  17. /// 把某个部署槽位指派为唯一执行机。槽位无需当前有存活实例——
  18. /// 指派可以先于部署,也可以在目标机器重启期间保持有效,这正是本模型的要点。
  19. /// </summary>
  20. Task AssignRunnerSlotAsync(string slotCode, string assignedBy, CancellationToken ct = default);
  21. /// <summary>取消指派。此后无人跑定时 ETL——这是允许的状态。</summary>
  22. Task ClearRunnerAsync(string clearedBy = null, CancellationToken ct = default);
  23. /// <summary>当前被指派的槽位;无指派返回 null。</summary>
  24. Task<string> GetDesignatedSlotAsync(CancellationToken ct = default);
  25. /// <summary>指派记录全文,供页面展示「谁在什么时候指派的」。无指派返回 null。</summary>
  26. Task<AdoEtlRunnerDesignation> GetDesignationAsync(CancellationToken ct = default);
  27. Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default);
  28. }
  29. public sealed class EtlInstanceStore : IEtlInstanceStore, ITransient
  30. {
  31. private readonly ISqlSugarClient _db;
  32. public EtlInstanceStore(ISqlSugarClient db) => _db = db;
  33. public async Task AssignRunnerSlotAsync(string slotCode, string assignedBy, CancellationToken ct = default)
  34. {
  35. var slot = string.IsNullOrWhiteSpace(slotCode) ? null : slotCode.Trim().ToLowerInvariant();
  36. await UpsertDesignationAsync(slot, assignedBy, ct);
  37. }
  38. public async Task ClearRunnerAsync(string clearedBy = null, CancellationToken ct = default) =>
  39. await UpsertDesignationAsync(null, clearedBy, ct);
  40. public async Task<string> GetDesignatedSlotAsync(CancellationToken ct = default)
  41. {
  42. var row = await GetDesignationAsync(ct);
  43. return string.IsNullOrWhiteSpace(row?.SlotCode) ? null : row.SlotCode;
  44. }
  45. public Task<AdoEtlRunnerDesignation> GetDesignationAsync(CancellationToken ct = default) =>
  46. _db.Queryable<AdoEtlRunnerDesignation>()
  47. .Where(x => x.Id == AdoEtlRunnerDesignation.SingletonId)
  48. .FirstAsync(ct);
  49. public Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default) =>
  50. _db.Queryable<AdoEtlInstance>()
  51. .OrderBy(x => x.MachineName)
  52. .OrderBy(x => x.ProcessId)
  53. .ToListAsync(ct);
  54. /// <summary>
  55. /// 写唯一那一行。先 UPDATE 再按影响行数决定是否 INSERT:
  56. /// 两个超管同时指派时,后到的 INSERT 会撞主键而不是插出第二行,
  57. /// 撞了就退回 UPDATE 重写一次,最终仍是单行单值。
  58. /// </summary>
  59. private async Task UpsertDesignationAsync(string slot, string operatorName, CancellationToken ct)
  60. {
  61. var now = DateTime.Now;
  62. var by = string.IsNullOrWhiteSpace(operatorName) ? null : operatorName.Trim();
  63. // 三元运算符不放进表达式树:SqlSugar 要把 SetColumns 的初始化器翻成 SQL,
  64. // 取消指派时 assigned_at 应回到 NULL,这个判断在 C# 侧先算完更稳。
  65. DateTime? assignedAt = slot == null ? null : now;
  66. var updated = await UpdateDesignationAsync(slot, by, assignedAt, now, ct);
  67. if (updated > 0) return;
  68. try
  69. {
  70. await _db.Insertable(new AdoEtlRunnerDesignation
  71. {
  72. Id = AdoEtlRunnerDesignation.SingletonId,
  73. SlotCode = slot,
  74. AssignedBy = by,
  75. AssignedAt = assignedAt,
  76. CreateTime = now,
  77. UpdateTime = now
  78. }).ExecuteCommandAsync(ct);
  79. }
  80. catch
  81. {
  82. // 并发下另一方已插入,改写即可;再失败就让异常上抛给控制器。
  83. await UpdateDesignationAsync(slot, by, assignedAt, now, ct);
  84. }
  85. }
  86. private Task<int> UpdateDesignationAsync(string slot, string by, DateTime? assignedAt, DateTime now, CancellationToken ct) =>
  87. _db.Updateable<AdoEtlRunnerDesignation>()
  88. .SetColumns(x => new AdoEtlRunnerDesignation
  89. {
  90. SlotCode = slot,
  91. AssignedBy = by,
  92. AssignedAt = assignedAt,
  93. UpdateTime = now
  94. })
  95. .Where(x => x.Id == AdoEtlRunnerDesignation.SingletonId)
  96. .ExecuteCommandAsync(ct);
  97. }