using System.Data; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse; /// /// S5 库存冷链(InventoryInbound)互斥闸门:进程内 + 跨实例双层。 /// /// 为什么不能直接用共享 ISqlSugarClient 跑 GET_LOCK:主库 SqlSugar 配置 /// IsAutoCloseConnection = true(Admin.NET.Core/SqlSugar/SqlSugarSetup.cs:128), /// 每条命令执行完连接立即归还 MySqlConnector 连接池;归还时驱动默认发送 /// COM_RESET_CONNECTION,而 MySQL 明确规定该命令会释放 GET_LOCK() 取得的咨询锁。 /// 于是「取锁 → 归还连接 → 锁当场没了」,`IS_USED_LOCK` 恒为 NULL,跨实例互斥形同虚设; /// 且后续 RELEASE_LOCK 多半跑在另一条连接上(非持有者,返回 0,什么都没释放)。 /// /// /// 本类的不变量:取锁的 connection == 任务全程保持打开的 connection == 释放锁的 connection。 /// 做法是另开一个 IsAutoCloseConnection = false 的专用 /// (与 KpiSqlReadOnlyExecutor 的专用隔离连接同款),它在 Dispose 前不会把连接还回池子。 /// /// /// 非阻塞:进程内闸门 WaitAsync(0)、DB 锁 GET_LOCK(key, 0), /// 第二个并发请求立即得到「未取到」,由调用方返回 Skipped=true / "lock busy",绝不排队等待。 /// /// /// 释放 先 RELEASE_LOCK,再 Dispose 专用连接。 /// 即使 RELEASE_LOCK 因故失败,连接关闭 / 归还池时的 reset 也会兜底释放该会话持有的咨询锁。 /// 调用方必须用 await using,异常路径才会走到释放。 /// /// public sealed class InventoryInboundLockGuard : IAsyncDisposable { /// 跨实例咨询锁名(MySQL GET_LOCK 键)。 public const string LockKey = "aidop:s5:inventory-inbound"; /// 取锁等待毫秒数:0 = 立即返回,不排队。 private const int NoWaitMilliseconds = 0; /// 锁语句自身的命令超时(秒);只跑 GET_LOCK/RELEASE_LOCK,不需要长超时。 private const int LockCommandTimeoutSeconds = 10; /// /// 进程内闸门。单实例内多个 HTTP 请求 / Job 并发时先在这里被挡掉, /// 不必每次都去 MySQL 抢咨询锁;同时保证同进程内不会有两条并发链路各自持有一条锁连接。 /// private static readonly SemaphoreSlim ProcessGate = new(1, 1); private readonly SqlSugarClient _lockDb; private readonly ILogger _logger; private readonly string _lockKey; private bool _processGateHeld; private bool _dbLockHeld; private bool _disposed; /// 是否成功拿到互斥权(进程内闸门 + 跨实例咨询锁都拿到才为 true)。 public bool Acquired { get; private set; } /// 持锁会话的 MySQL CONNECTION_ID(),仅用于排障/断言,未取到锁时为 null。 public long? OwnerConnectionId { get; private set; } /// 未取到锁时的原因,便于日志区分「同进程忙」与「其它实例忙」。 public string BusyReason { get; private set; } private InventoryInboundLockGuard(SqlSugarClient lockDb, ILogger logger, string lockKey) { _lockDb = lockDb; _logger = logger; _lockKey = lockKey; } /// /// 尝试取锁;无论成功与否都返回一个 guard,调用方以 判定, /// 并始终以 await using 持有(未取到锁的 guard 释放时什么都不做)。 /// /// 主库客户端,仅用于读取连接配置,不在其上取锁。 /// 日志。 /// 咨询锁键;默认 ,测试可传独立键避免打扰生产链路。 /// 取消令牌。 public static async Task TryAcquireAsync( ISqlSugarClient db, ILogger logger, string lockKey = LockKey, CancellationToken cancellationToken = default) { if (db == null) throw new ArgumentNullException(nameof(db)); var main = db.CurrentConnectionConfig; var lockDb = new SqlSugarClient(new ConnectionConfig { ConnectionString = main.ConnectionString, DbType = main.DbType, // 关键:专用连接不自动关闭,GET_LOCK 与 RELEASE_LOCK 之间连接不回池、会话不重置 IsAutoCloseConnection = false, }); var guard = new InventoryInboundLockGuard(lockDb, logger, lockKey); // 第一层:进程内闸门,非阻塞 var gate = await ProcessGate.WaitAsync(NoWaitMilliseconds, cancellationToken); if (!gate) { guard.BusyReason = "in-process busy"; lockDb.Dispose(); guard._disposed = true; // 什么都没持有,调用方的 await using 直接空转 return guard; } guard._processGateHeld = true; try { lockDb.Ado.CommandTimeOut = LockCommandTimeoutSeconds; // 第二层:MySQL 咨询锁,非阻塞(timeout=0) var got = await ScalarAsync(lockDb, "SELECT GET_LOCK(@k, 0)", new List { new("@k", lockKey) }); // GET_LOCK 返回 1=取到;0=超时未取到;NULL=出错/被中断 if (got is null || Convert.ToInt64(got) != 1L) { guard.BusyReason = got is null ? "GET_LOCK returned NULL" : "cross-instance busy"; await guard.DisposeAsync(); return guard; } guard._dbLockHeld = true; // 同会话自证:持有者必须就是本连接。若不成立说明连接被换过(即本缺陷的回归),fail closed。 var selfCheck = await lockDb.Ado.GetDataTableAsync( "SELECT IS_USED_LOCK(@k) AS owner_conn, CONNECTION_ID() AS my_conn", new List { new("@k", lockKey) }); var owner = CellOrNull(selfCheck, 0); var me = CellOrNull(selfCheck, 1); if (owner is null || me is null || Convert.ToInt64(owner) != Convert.ToInt64(me)) { logger?.LogError( "[InventoryInboundLock] 取锁后自证失败:owner={Owner} me={Me},判定跨实例互斥不可信,放弃本轮", owner, me); guard.BusyReason = "lock self-check failed"; await guard.DisposeAsync(); return guard; } guard.OwnerConnectionId = Convert.ToInt64(me); guard.Acquired = true; logger?.LogInformation( "[InventoryInboundLock] acquired key={Key} conn={Conn}", lockKey, guard.OwnerConnectionId); return guard; } catch (Exception ex) { logger?.LogError(ex, "[InventoryInboundLock] 取锁失败 key={Key}", lockKey); await guard.DisposeAsync(); guard.BusyReason ??= "acquire failed: " + ex.Message; return guard; } } /// 释放:先 RELEASE_LOCK(同一连接),再关闭专用连接兜底。 public async ValueTask DisposeAsync() { if (_disposed) return; _disposed = true; Acquired = false; try { if (_dbLockHeld) { try { await _lockDb.Ado.ExecuteCommandAsync( "SELECT RELEASE_LOCK(@k)", new List { new("@k", _lockKey) }); } catch (Exception ex) { // 连接 Dispose 时会话结束 / 连接归还池触发 reset,MySQL 会自动释放该会话的咨询锁,故这里只记录 _logger?.LogWarning(ex, "[InventoryInboundLock] RELEASE_LOCK 失败,改由连接关闭兜底 key={Key}", _lockKey); } _dbLockHeld = false; } try { _lockDb.Dispose(); } catch (Exception ex) { _logger?.LogWarning(ex, "[InventoryInboundLock] 专用锁连接释放异常"); } } finally { if (_processGateHeld) { _processGateHeld = false; ProcessGate.Release(); } } } private static async Task ScalarAsync(SqlSugarClient db, string sql, List pars) { DataTable dt = await db.Ado.GetDataTableAsync(sql, pars); return CellOrNull(dt, 0); } private static object CellOrNull(DataTable dt, int columnIndex) { if (dt == null || dt.Rows.Count == 0 || dt.Columns.Count <= columnIndex) return null; var v = dt.Rows[0][columnIndex]; return v == DBNull.Value ? null : v; } }