ModuleRebuildLock.cs 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153
  1. namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  2. public interface IModuleRebuildLease : IAsyncDisposable
  3. {
  4. Task HeartbeatAsync(CancellationToken ct = default);
  5. }
  6. public interface IModuleRebuildLock
  7. {
  8. Task<IModuleRebuildLease?> TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default);
  9. }
  10. public sealed class ModuleRebuildLock : IModuleRebuildLock, ITransient
  11. {
  12. public static readonly TimeSpan StaleAfter = TimeSpan.FromMinutes(30);
  13. /// <summary>
  14. /// QUEUED 滞留多久后判定为「已无消费可能」并收口为 FAILED。
  15. ///
  16. /// <para><b>为什么是 24 小时而不是更短</b>:<c>GlobalMaxParallelScopes = 1</c> 且单 scope
  17. /// 跑 10~20 分钟,夜间全量一次扇出 28 个 scope 时,队尾那个**合法**等待可达 7 小时左右。
  18. /// 阈值取到日级才能确保「超时」只捕获真正无人消费的行(指派丢失、Worker 全停),
  19. /// 不会误杀正常排队。改小前请先核算当前 scope 数 × 单次耗时。</para>
  20. /// </summary>
  21. public static readonly TimeSpan QueuedStaleAfter = TimeSpan.FromHours(24);
  22. private readonly ISqlSugarClient _db;
  23. public ModuleRebuildLock(ISqlSugarClient db) => _db = db;
  24. public async Task<IModuleRebuildLease?> TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default)
  25. {
  26. scope = MdpRebuildScope.Create(scope.ModuleCode, scope.TenantId, scope.FactoryId);
  27. if (string.IsNullOrWhiteSpace(holderId))
  28. throw new InvalidOperationException("重算锁持有者不能为空");
  29. var now = DateTime.Now;
  30. var staleBefore = now - StaleAfter;
  31. await _db.Ado.ExecuteCommandAsync(
  32. """
  33. INSERT INTO ado_mdp_rebuild_lock
  34. (lock_name, module_code, tenant_id, factory_id, status, create_time, update_time)
  35. SELECT @LockName, @ModuleCode, @TenantId, @FactoryId, 'FREE', @Now, @Now
  36. FROM DUAL
  37. WHERE NOT EXISTS (SELECT 1 FROM ado_mdp_rebuild_lock WHERE lock_name=@LockName)
  38. """,
  39. new SugarParameter("@LockName", scope.LockName),
  40. new SugarParameter("@ModuleCode", scope.ModuleCode),
  41. new SugarParameter("@TenantId", scope.TenantId),
  42. new SugarParameter("@FactoryId", scope.FactoryId),
  43. new SugarParameter("@Now", now));
  44. var n = await _db.Ado.ExecuteCommandAsync(
  45. """
  46. UPDATE ado_mdp_rebuild_lock
  47. SET holder_id=@HolderId, job_id=@JobId, status='HELD',
  48. acquired_at=@Now, heartbeat_at=@Now, update_time=@Now
  49. WHERE lock_name=@LockName
  50. AND (status='FREE' OR heartbeat_at IS NULL OR heartbeat_at < @StaleBefore)
  51. """,
  52. new SugarParameter("@HolderId", holderId),
  53. new SugarParameter("@JobId", (object?)jobId ?? DBNull.Value),
  54. new SugarParameter("@Now", now),
  55. new SugarParameter("@LockName", scope.LockName),
  56. new SugarParameter("@StaleBefore", staleBefore));
  57. if (n <= 0)
  58. return null;
  59. return new DbLease(_db, scope.LockName, holderId);
  60. }
  61. private sealed class DbLease : IModuleRebuildLease
  62. {
  63. private readonly ISqlSugarClient _db;
  64. private readonly string _lockName;
  65. private readonly string _holderId;
  66. private int _released;
  67. public DbLease(ISqlSugarClient db, string lockName, string holderId)
  68. {
  69. _db = db;
  70. _lockName = lockName;
  71. _holderId = holderId;
  72. }
  73. public async Task HeartbeatAsync(CancellationToken ct = default)
  74. {
  75. using var db = _db.CopyNew();
  76. await db.Ado.ExecuteCommandAsync(
  77. """
  78. UPDATE ado_mdp_rebuild_lock
  79. SET heartbeat_at=@Now, update_time=@Now
  80. WHERE lock_name=@LockName AND holder_id=@HolderId AND status='HELD'
  81. """,
  82. new SugarParameter("@Now", DateTime.Now),
  83. new SugarParameter("@LockName", _lockName),
  84. new SugarParameter("@HolderId", _holderId));
  85. }
  86. public async ValueTask DisposeAsync()
  87. {
  88. if (Interlocked.Exchange(ref _released, 1) == 1)
  89. return;
  90. try
  91. {
  92. using var db = _db.CopyNew();
  93. await db.Ado.ExecuteCommandAsync(
  94. """
  95. UPDATE ado_mdp_rebuild_lock
  96. SET status='FREE', holder_id=NULL, job_id=NULL, update_time=@Now
  97. WHERE lock_name=@LockName AND holder_id=@HolderId
  98. """,
  99. new SugarParameter("@Now", DateTime.Now),
  100. new SugarParameter("@LockName", _lockName),
  101. new SugarParameter("@HolderId", _holderId));
  102. }
  103. catch
  104. {
  105. // 释放失败留给租约过期
  106. }
  107. }
  108. }
  109. }
  110. public sealed class InMemoryModuleRebuildLock : IModuleRebuildLock
  111. {
  112. private readonly Dictionary<string, string> _holders = new(StringComparer.Ordinal);
  113. public Task<IModuleRebuildLease?> TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default)
  114. {
  115. var lockName = MdpRebuildScope.Create(scope.ModuleCode, scope.TenantId, scope.FactoryId).LockName;
  116. lock (this)
  117. {
  118. if (_holders.ContainsKey(lockName))
  119. return Task.FromResult<IModuleRebuildLease?>(null);
  120. _holders[lockName] = holderId;
  121. return Task.FromResult<IModuleRebuildLease?>(new MemoryLease(this, lockName));
  122. }
  123. }
  124. private sealed class MemoryLease : IModuleRebuildLease
  125. {
  126. private readonly InMemoryModuleRebuildLock _owner;
  127. private readonly string _lockName;
  128. public MemoryLease(InMemoryModuleRebuildLock owner, string lockName)
  129. {
  130. _owner = owner;
  131. _lockName = lockName;
  132. }
  133. public Task HeartbeatAsync(CancellationToken ct = default) => Task.CompletedTask;
  134. public ValueTask DisposeAsync()
  135. {
  136. lock (_owner) _owner._holders.Remove(_lockName);
  137. return ValueTask.CompletedTask;
  138. }
  139. }
  140. }