ModuleRebuildLock.cs 5.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143
  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. private readonly ISqlSugarClient _db;
  14. public ModuleRebuildLock(ISqlSugarClient db) => _db = db;
  15. public async Task<IModuleRebuildLease?> TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default)
  16. {
  17. scope = MdpRebuildScope.Create(scope.ModuleCode, scope.TenantId, scope.FactoryId);
  18. if (string.IsNullOrWhiteSpace(holderId))
  19. throw new InvalidOperationException("重算锁持有者不能为空");
  20. var now = DateTime.Now;
  21. var staleBefore = now - StaleAfter;
  22. await _db.Ado.ExecuteCommandAsync(
  23. """
  24. INSERT INTO ado_mdp_rebuild_lock
  25. (lock_name, module_code, tenant_id, factory_id, status, create_time, update_time)
  26. SELECT @LockName, @ModuleCode, @TenantId, @FactoryId, 'FREE', @Now, @Now
  27. FROM DUAL
  28. WHERE NOT EXISTS (SELECT 1 FROM ado_mdp_rebuild_lock WHERE lock_name=@LockName)
  29. """,
  30. new SugarParameter("@LockName", scope.LockName),
  31. new SugarParameter("@ModuleCode", scope.ModuleCode),
  32. new SugarParameter("@TenantId", scope.TenantId),
  33. new SugarParameter("@FactoryId", scope.FactoryId),
  34. new SugarParameter("@Now", now));
  35. var n = await _db.Ado.ExecuteCommandAsync(
  36. """
  37. UPDATE ado_mdp_rebuild_lock
  38. SET holder_id=@HolderId, job_id=@JobId, status='HELD',
  39. acquired_at=@Now, heartbeat_at=@Now, update_time=@Now
  40. WHERE lock_name=@LockName
  41. AND (status='FREE' OR heartbeat_at IS NULL OR heartbeat_at < @StaleBefore)
  42. """,
  43. new SugarParameter("@HolderId", holderId),
  44. new SugarParameter("@JobId", (object?)jobId ?? DBNull.Value),
  45. new SugarParameter("@Now", now),
  46. new SugarParameter("@LockName", scope.LockName),
  47. new SugarParameter("@StaleBefore", staleBefore));
  48. if (n <= 0)
  49. return null;
  50. return new DbLease(_db, scope.LockName, holderId);
  51. }
  52. private sealed class DbLease : IModuleRebuildLease
  53. {
  54. private readonly ISqlSugarClient _db;
  55. private readonly string _lockName;
  56. private readonly string _holderId;
  57. private int _released;
  58. public DbLease(ISqlSugarClient db, string lockName, string holderId)
  59. {
  60. _db = db;
  61. _lockName = lockName;
  62. _holderId = holderId;
  63. }
  64. public async Task HeartbeatAsync(CancellationToken ct = default)
  65. {
  66. using var db = _db.CopyNew();
  67. await db.Ado.ExecuteCommandAsync(
  68. """
  69. UPDATE ado_mdp_rebuild_lock
  70. SET heartbeat_at=@Now, update_time=@Now
  71. WHERE lock_name=@LockName AND holder_id=@HolderId AND status='HELD'
  72. """,
  73. new SugarParameter("@Now", DateTime.Now),
  74. new SugarParameter("@LockName", _lockName),
  75. new SugarParameter("@HolderId", _holderId));
  76. }
  77. public async ValueTask DisposeAsync()
  78. {
  79. if (Interlocked.Exchange(ref _released, 1) == 1)
  80. return;
  81. try
  82. {
  83. using var db = _db.CopyNew();
  84. await db.Ado.ExecuteCommandAsync(
  85. """
  86. UPDATE ado_mdp_rebuild_lock
  87. SET status='FREE', holder_id=NULL, job_id=NULL, update_time=@Now
  88. WHERE lock_name=@LockName AND holder_id=@HolderId
  89. """,
  90. new SugarParameter("@Now", DateTime.Now),
  91. new SugarParameter("@LockName", _lockName),
  92. new SugarParameter("@HolderId", _holderId));
  93. }
  94. catch
  95. {
  96. // 释放失败留给租约过期
  97. }
  98. }
  99. }
  100. }
  101. public sealed class InMemoryModuleRebuildLock : IModuleRebuildLock
  102. {
  103. private readonly Dictionary<string, string> _holders = new(StringComparer.Ordinal);
  104. public Task<IModuleRebuildLease?> TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default)
  105. {
  106. var lockName = MdpRebuildScope.Create(scope.ModuleCode, scope.TenantId, scope.FactoryId).LockName;
  107. lock (this)
  108. {
  109. if (_holders.ContainsKey(lockName))
  110. return Task.FromResult<IModuleRebuildLease?>(null);
  111. _holders[lockName] = holderId;
  112. return Task.FromResult<IModuleRebuildLease?>(new MemoryLease(this, lockName));
  113. }
  114. }
  115. private sealed class MemoryLease : IModuleRebuildLease
  116. {
  117. private readonly InMemoryModuleRebuildLock _owner;
  118. private readonly string _lockName;
  119. public MemoryLease(InMemoryModuleRebuildLock owner, string lockName)
  120. {
  121. _owner = owner;
  122. _lockName = lockName;
  123. }
  124. public Task HeartbeatAsync(CancellationToken ct = default) => Task.CompletedTask;
  125. public ValueTask DisposeAsync()
  126. {
  127. lock (_owner) _owner._holders.Remove(_lockName);
  128. return ValueTask.CompletedTask;
  129. }
  130. }
  131. }