| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153 |
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
- public interface IModuleRebuildLease : IAsyncDisposable
- {
- Task HeartbeatAsync(CancellationToken ct = default);
- }
- public interface IModuleRebuildLock
- {
- Task<IModuleRebuildLease?> 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);
- /// <summary>
- /// QUEUED 滞留多久后判定为「已无消费可能」并收口为 FAILED。
- ///
- /// <para><b>为什么是 24 小时而不是更短</b>:<c>GlobalMaxParallelScopes = 1</c> 且单 scope
- /// 跑 10~20 分钟,夜间全量一次扇出 28 个 scope 时,队尾那个**合法**等待可达 7 小时左右。
- /// 阈值取到日级才能确保「超时」只捕获真正无人消费的行(指派丢失、Worker 全停),
- /// 不会误杀正常排队。改小前请先核算当前 scope 数 × 单次耗时。</para>
- /// </summary>
- public static readonly TimeSpan QueuedStaleAfter = TimeSpan.FromHours(24);
- private readonly ISqlSugarClient _db;
- public ModuleRebuildLock(ISqlSugarClient db) => _db = db;
- public async Task<IModuleRebuildLease?> 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<string, string> _holders = new(StringComparer.Ordinal);
- public Task<IModuleRebuildLease?> 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<IModuleRebuildLease?>(null);
- _holders[lockName] = holderId;
- return Task.FromResult<IModuleRebuildLease?>(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;
- }
- }
- }
|