namespace Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh; public static class S1DashboardRebuildStatus { public const string Queued = "QUEUED"; public const string Running = "RUNNING"; public const string Success = "SUCCESS"; public const string Failed = "FAILED"; public const string Cancelled = "CANCELLED"; } public sealed class S1MdpAlreadyRunningException : Exception { public S1MdpAlreadyRunningException() : base("S1 数据重算正在执行,请勿重复提交") { } } public interface IS1MdpFullRunLease : IAsyncDisposable { Task HeartbeatAsync(CancellationToken ct = default); } public interface IS1MdpFullRunLock { Task TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default); } public sealed class S1MdpFullRunLock : IS1MdpFullRunLock, ITransient { public const string LegacyLockName = "S1_MDP_FULL"; public static readonly TimeSpan StaleAfter = TimeSpan.FromMinutes(30); private readonly ISqlSugarClient _db; public S1MdpFullRunLock(ISqlSugarClient db) => _db = db; public async Task TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default) { if (string.IsNullOrWhiteSpace(lockName)) throw new InvalidOperationException("S1 全量锁名不能为空"); var now = DateTime.Now; var staleBefore = now - StaleAfter; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO ado_s1_mdp_run_lock (lock_name, status) SELECT @LockName, 'FREE' FROM DUAL WHERE NOT EXISTS (SELECT 1 FROM ado_s1_mdp_run_lock WHERE lock_name=@LockName) """, new SugarParameter("@LockName", lockName)); var n = await _db.Ado.ExecuteCommandAsync( """ UPDATE ado_s1_mdp_run_lock SET holder_id=@HolderId, job_id=@JobId, status='HELD', acquired_at=@Now, heartbeat_at=@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", lockName), new SugarParameter("@StaleBefore", staleBefore)); if (n <= 0) return null; return new DbLease(_db, lockName, holderId); } private sealed class DbLease : IS1MdpFullRunLease { 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_s1_mdp_run_lock SET heartbeat_at=@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_s1_mdp_run_lock SET status='FREE', holder_id=NULL, job_id=NULL WHERE lock_name=@LockName AND holder_id=@HolderId """, new SugarParameter("@LockName", _lockName), new SugarParameter("@HolderId", _holderId)); } catch { // 释放失败留给租约过期 } } } } public sealed class InMemoryS1MdpFullRunLock : IS1MdpFullRunLock { private readonly Dictionary _holders = new(StringComparer.Ordinal); public Task TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default) { lock (this) { if (_holders.ContainsKey(lockName)) return Task.FromResult(null); _holders[lockName] = holderId; return Task.FromResult(new MemoryLease(this, lockName)); } } private sealed class MemoryLease : IS1MdpFullRunLease { private readonly InMemoryS1MdpFullRunLock _owner; private readonly string _lockName; public MemoryLease(InMemoryS1MdpFullRunLock 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; } } }