namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild; public interface IModuleRebuildLease : IAsyncDisposable { Task HeartbeatAsync(CancellationToken ct = default); } public interface IModuleRebuildLock { Task TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default); } public sealed class ModuleRebuildLock : IModuleRebuildLock, ITransient { public static readonly TimeSpan StaleAfter = TimeSpan.FromMinutes(30); private readonly ISqlSugarClient _db; public ModuleRebuildLock(ISqlSugarClient db) => _db = db; public async Task TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default) { scope = MdpRebuildScope.Create(scope.ModuleCode, scope.TenantId, scope.FactoryId); if (string.IsNullOrWhiteSpace(holderId)) throw new InvalidOperationException("重算锁持有者不能为空"); var now = DateTime.Now; var staleBefore = now - StaleAfter; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO ado_mdp_rebuild_lock (lock_name, module_code, tenant_id, factory_id, status, create_time, update_time) SELECT @LockName, @ModuleCode, @TenantId, @FactoryId, 'FREE', @Now, @Now FROM DUAL WHERE NOT EXISTS (SELECT 1 FROM ado_mdp_rebuild_lock WHERE lock_name=@LockName) """, new SugarParameter("@LockName", scope.LockName), new SugarParameter("@ModuleCode", scope.ModuleCode), new SugarParameter("@TenantId", scope.TenantId), new SugarParameter("@FactoryId", scope.FactoryId), new SugarParameter("@Now", now)); var n = await _db.Ado.ExecuteCommandAsync( """ UPDATE ado_mdp_rebuild_lock SET holder_id=@HolderId, job_id=@JobId, status='HELD', acquired_at=@Now, heartbeat_at=@Now, update_time=@Now WHERE lock_name=@LockName AND (status='FREE' OR heartbeat_at IS NULL OR heartbeat_at < @StaleBefore) """, new SugarParameter("@HolderId", holderId), new SugarParameter("@JobId", (object?)jobId ?? DBNull.Value), new SugarParameter("@Now", now), new SugarParameter("@LockName", scope.LockName), new SugarParameter("@StaleBefore", staleBefore)); if (n <= 0) return null; return new DbLease(_db, scope.LockName, holderId); } private sealed class DbLease : IModuleRebuildLease { private readonly ISqlSugarClient _db; private readonly string _lockName; private readonly string _holderId; private int _released; public DbLease(ISqlSugarClient db, string lockName, string holderId) { _db = db; _lockName = lockName; _holderId = holderId; } public async Task HeartbeatAsync(CancellationToken ct = default) { using var db = _db.CopyNew(); await db.Ado.ExecuteCommandAsync( """ UPDATE ado_mdp_rebuild_lock SET heartbeat_at=@Now, update_time=@Now WHERE lock_name=@LockName AND holder_id=@HolderId AND status='HELD' """, new SugarParameter("@Now", DateTime.Now), new SugarParameter("@LockName", _lockName), new SugarParameter("@HolderId", _holderId)); } public async ValueTask DisposeAsync() { if (Interlocked.Exchange(ref _released, 1) == 1) return; try { using var db = _db.CopyNew(); await db.Ado.ExecuteCommandAsync( """ UPDATE ado_mdp_rebuild_lock SET status='FREE', holder_id=NULL, job_id=NULL, update_time=@Now WHERE lock_name=@LockName AND holder_id=@HolderId """, new SugarParameter("@Now", DateTime.Now), new SugarParameter("@LockName", _lockName), new SugarParameter("@HolderId", _holderId)); } catch { // 释放失败留给租约过期 } } } } public sealed class InMemoryModuleRebuildLock : IModuleRebuildLock { private readonly Dictionary _holders = new(StringComparer.Ordinal); public Task TryAcquireAsync(MdpRebuildScope scope, string holderId, long? jobId = null, CancellationToken ct = default) { var lockName = MdpRebuildScope.Create(scope.ModuleCode, scope.TenantId, scope.FactoryId).LockName; lock (this) { if (_holders.ContainsKey(lockName)) return Task.FromResult(null); _holders[lockName] = holderId; return Task.FromResult(new MemoryLease(this, lockName)); } } private sealed class MemoryLease : IModuleRebuildLease { private readonly InMemoryModuleRebuildLock _owner; private readonly string _lockName; public MemoryLease(InMemoryModuleRebuildLock owner, string lockName) { _owner = owner; _lockName = lockName; } public Task HeartbeatAsync(CancellationToken ct = default) => Task.CompletedTask; public ValueTask DisposeAsync() { lock (_owner) _owner._holders.Remove(_lockName); return ValueTask.CompletedTask; } } }