| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154 |
- 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<IS1MdpFullRunLease?> 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<IS1MdpFullRunLease?> 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<string, string> _holders = new(StringComparer.Ordinal);
- public Task<IS1MdpFullRunLease?> TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default)
- {
- lock (this)
- {
- if (_holders.ContainsKey(lockName))
- return Task.FromResult<IS1MdpFullRunLease?>(null);
- _holders[lockName] = holderId;
- return Task.FromResult<IS1MdpFullRunLease?>(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;
- }
- }
- }
|