| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205 |
- using System.Data;
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
- /// <summary>
- /// S5 库存冷链(InventoryInbound)互斥闸门:进程内 + 跨实例双层。
- /// <para>
- /// <b>为什么不能直接用共享 <c>ISqlSugarClient</c> 跑 GET_LOCK</b>:主库 SqlSugar 配置
- /// <c>IsAutoCloseConnection = true</c>(Admin.NET.Core/SqlSugar/SqlSugarSetup.cs:128),
- /// 每条命令执行完连接立即归还 MySqlConnector 连接池;归还时驱动默认发送
- /// <c>COM_RESET_CONNECTION</c>,而 MySQL 明确规定该命令会释放 <c>GET_LOCK()</c> 取得的咨询锁。
- /// 于是「取锁 → 归还连接 → 锁当场没了」,`IS_USED_LOCK` 恒为 NULL,跨实例互斥形同虚设;
- /// 且后续 RELEASE_LOCK 多半跑在另一条连接上(非持有者,返回 0,什么都没释放)。
- /// </para>
- /// <para>
- /// <b>本类的不变量</b>:取锁的 connection == 任务全程保持打开的 connection == 释放锁的 connection。
- /// 做法是另开一个 <c>IsAutoCloseConnection = false</c> 的专用 <see cref="SqlSugarClient"/>
- /// (与 <c>KpiSqlReadOnlyExecutor</c> 的专用隔离连接同款),它在 Dispose 前不会把连接还回池子。
- /// </para>
- /// <para>
- /// <b>非阻塞</b>:进程内闸门 <c>WaitAsync(0)</c>、DB 锁 <c>GET_LOCK(key, 0)</c>,
- /// 第二个并发请求立即得到「未取到」,由调用方返回 <c>Skipped=true / "lock busy"</c>,绝不排队等待。
- /// </para>
- /// <para>
- /// <b>释放</b>:<see cref="DisposeAsync"/> 先 RELEASE_LOCK,再 Dispose 专用连接。
- /// 即使 RELEASE_LOCK 因故失败,连接关闭 / 归还池时的 reset 也会兜底释放该会话持有的咨询锁。
- /// 调用方必须用 <c>await using</c>,异常路径才会走到释放。
- /// </para>
- /// </summary>
- public sealed class InventoryInboundLockGuard : IAsyncDisposable
- {
- /// <summary>跨实例咨询锁名(MySQL GET_LOCK 键)。</summary>
- public const string LockKey = "aidop:s5:inventory-inbound";
- /// <summary>取锁等待毫秒数:0 = 立即返回,不排队。</summary>
- private const int NoWaitMilliseconds = 0;
- /// <summary>锁语句自身的命令超时(秒);只跑 GET_LOCK/RELEASE_LOCK,不需要长超时。</summary>
- private const int LockCommandTimeoutSeconds = 10;
- /// <summary>
- /// 进程内闸门。单实例内多个 HTTP 请求 / Job 并发时先在这里被挡掉,
- /// 不必每次都去 MySQL 抢咨询锁;同时保证同进程内不会有两条并发链路各自持有一条锁连接。
- /// </summary>
- 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;
- /// <summary>是否成功拿到互斥权(进程内闸门 + 跨实例咨询锁都拿到才为 true)。</summary>
- public bool Acquired { get; private set; }
- /// <summary>持锁会话的 MySQL CONNECTION_ID(),仅用于排障/断言,未取到锁时为 null。</summary>
- public long? OwnerConnectionId { get; private set; }
- /// <summary>未取到锁时的原因,便于日志区分「同进程忙」与「其它实例忙」。</summary>
- public string BusyReason { get; private set; }
- private InventoryInboundLockGuard(SqlSugarClient lockDb, ILogger logger, string lockKey)
- {
- _lockDb = lockDb;
- _logger = logger;
- _lockKey = lockKey;
- }
- /// <summary>
- /// 尝试取锁;无论成功与否都返回一个 guard,调用方以 <see cref="Acquired"/> 判定,
- /// 并始终以 <c>await using</c> 持有(未取到锁的 guard 释放时什么都不做)。
- /// </summary>
- /// <param name="db">主库客户端,仅用于读取连接配置,不在其上取锁。</param>
- /// <param name="logger">日志。</param>
- /// <param name="lockKey">咨询锁键;默认 <see cref="LockKey"/>,测试可传独立键避免打扰生产链路。</param>
- /// <param name="cancellationToken">取消令牌。</param>
- public static async Task<InventoryInboundLockGuard> 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<SugarParameter> { 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<SugarParameter> { 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;
- }
- }
- /// <summary>释放:先 RELEASE_LOCK(同一连接),再关闭专用连接兜底。</summary>
- 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<SugarParameter> { 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<object> ScalarAsync(SqlSugarClient db, string sql, List<SugarParameter> 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;
- }
- }
|