InventoryInboundLockGuard.cs 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212
  1. using System.Data;
  2. using Microsoft.Extensions.Logging;
  3. namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  4. /// <summary>
  5. /// S5 库存冷链(InventoryInbound)互斥闸门:进程内 + 跨实例双层。
  6. /// <para>
  7. /// <b>为什么不能直接用共享 <c>ISqlSugarClient</c> 跑 GET_LOCK</b>:主库 SqlSugar 配置
  8. /// <c>IsAutoCloseConnection = true</c>(Admin.NET.Core/SqlSugar/SqlSugarSetup.cs:128),
  9. /// 每条命令执行完连接立即归还 MySqlConnector 连接池;归还时驱动默认发送
  10. /// <c>COM_RESET_CONNECTION</c>,而 MySQL 明确规定该命令会释放 <c>GET_LOCK()</c> 取得的咨询锁。
  11. /// 于是「取锁 → 归还连接 → 锁当场没了」,`IS_USED_LOCK` 恒为 NULL,跨实例互斥形同虚设;
  12. /// 且后续 RELEASE_LOCK 多半跑在另一条连接上(非持有者,返回 0,什么都没释放)。
  13. /// </para>
  14. /// <para>
  15. /// <b>本类的不变量</b>:取锁的 connection == 任务全程保持打开的 connection == 释放锁的 connection。
  16. /// 做法是另开一个 <c>IsAutoCloseConnection = false</c> 的专用 <see cref="SqlSugarClient"/>
  17. /// (与 <c>KpiSqlReadOnlyExecutor</c> 的专用隔离连接同款),它在 Dispose 前不会把连接还回池子。
  18. /// </para>
  19. /// <para>
  20. /// <b>非阻塞</b>:进程内闸门 <c>WaitAsync(0)</c>、DB 锁 <c>GET_LOCK(key, 0)</c>,
  21. /// 第二个并发请求立即得到「未取到」,由调用方返回 <c>Skipped=true / "lock busy"</c>,绝不排队等待。
  22. /// </para>
  23. /// <para>
  24. /// <b>释放</b>:<see cref="DisposeAsync"/> 先 RELEASE_LOCK,再 Dispose 专用连接。
  25. /// 即使 RELEASE_LOCK 因故失败,连接关闭 / 归还池时的 reset 也会兜底释放该会话持有的咨询锁。
  26. /// 调用方必须用 <c>await using</c>,异常路径才会走到释放。
  27. /// </para>
  28. /// </summary>
  29. public sealed class InventoryInboundLockGuard : IAsyncDisposable
  30. {
  31. /// <summary>跨实例咨询锁名(MySQL GET_LOCK 键)。</summary>
  32. public const string LockKey = "aidop:s5:inventory-inbound";
  33. /// <summary>取锁等待毫秒数:0 = 立即返回,不排队。</summary>
  34. private const int NoWaitMilliseconds = 0;
  35. /// <summary>锁语句自身的命令超时(秒);只跑 GET_LOCK/RELEASE_LOCK,不需要长超时。</summary>
  36. private const int LockCommandTimeoutSeconds = 10;
  37. /// <summary>
  38. /// 进程内闸门,**按锁键分桶**。单实例内多个 HTTP 请求 / Job 并发时先在这里被挡掉,
  39. /// 不必每次都去 MySQL 抢咨询锁;同时保证同进程内不会有两条并发链路各自持有一条锁连接。
  40. /// <para>分桶是必须的:早期实现是全类共用一个 <c>SemaphoreSlim</c>,
  41. /// 于是不同业务域用不同 lockKey 调用本类时会互相阻塞——
  42. /// S5 库存全量可以持锁数分钟,期间任何其它 lockKey 都拿不到闸门,
  43. /// 表现为「锁没被占用却一直取不到」。同键行为与分桶前完全一致。</para>
  44. /// </summary>
  45. private static readonly System.Collections.Concurrent.ConcurrentDictionary<string, SemaphoreSlim> ProcessGates = new();
  46. private static SemaphoreSlim GateFor(string lockKey) =>
  47. ProcessGates.GetOrAdd(lockKey, _ => new SemaphoreSlim(1, 1));
  48. private readonly SqlSugarClient _lockDb;
  49. private readonly ILogger _logger;
  50. private readonly string _lockKey;
  51. private bool _processGateHeld;
  52. private bool _dbLockHeld;
  53. private bool _disposed;
  54. /// <summary>是否成功拿到互斥权(进程内闸门 + 跨实例咨询锁都拿到才为 true)。</summary>
  55. public bool Acquired { get; private set; }
  56. /// <summary>持锁会话的 MySQL CONNECTION_ID(),仅用于排障/断言,未取到锁时为 null。</summary>
  57. public long? OwnerConnectionId { get; private set; }
  58. /// <summary>未取到锁时的原因,便于日志区分「同进程忙」与「其它实例忙」。</summary>
  59. public string BusyReason { get; private set; }
  60. private InventoryInboundLockGuard(SqlSugarClient lockDb, ILogger logger, string lockKey)
  61. {
  62. _lockDb = lockDb;
  63. _logger = logger;
  64. _lockKey = lockKey;
  65. }
  66. /// <summary>
  67. /// 尝试取锁;无论成功与否都返回一个 guard,调用方以 <see cref="Acquired"/> 判定,
  68. /// 并始终以 <c>await using</c> 持有(未取到锁的 guard 释放时什么都不做)。
  69. /// </summary>
  70. /// <param name="db">主库客户端,仅用于读取连接配置,不在其上取锁。</param>
  71. /// <param name="logger">日志。</param>
  72. /// <param name="lockKey">咨询锁键;默认 <see cref="LockKey"/>,测试可传独立键避免打扰生产链路。</param>
  73. /// <param name="cancellationToken">取消令牌。</param>
  74. public static async Task<InventoryInboundLockGuard> TryAcquireAsync(
  75. ISqlSugarClient db,
  76. ILogger logger,
  77. string lockKey = LockKey,
  78. CancellationToken cancellationToken = default)
  79. {
  80. if (db == null) throw new ArgumentNullException(nameof(db));
  81. var main = db.CurrentConnectionConfig;
  82. var lockDb = new SqlSugarClient(new ConnectionConfig
  83. {
  84. ConnectionString = main.ConnectionString,
  85. DbType = main.DbType,
  86. // 关键:专用连接不自动关闭,GET_LOCK 与 RELEASE_LOCK 之间连接不回池、会话不重置
  87. IsAutoCloseConnection = false,
  88. });
  89. var guard = new InventoryInboundLockGuard(lockDb, logger, lockKey);
  90. // 第一层:进程内闸门,非阻塞
  91. var gate = await GateFor(lockKey).WaitAsync(NoWaitMilliseconds, cancellationToken);
  92. if (!gate)
  93. {
  94. guard.BusyReason = "in-process busy";
  95. lockDb.Dispose();
  96. guard._disposed = true; // 什么都没持有,调用方的 await using 直接空转
  97. return guard;
  98. }
  99. guard._processGateHeld = true;
  100. try
  101. {
  102. lockDb.Ado.CommandTimeOut = LockCommandTimeoutSeconds;
  103. // 第二层:MySQL 咨询锁,非阻塞(timeout=0)
  104. var got = await ScalarAsync(lockDb, "SELECT GET_LOCK(@k, 0)",
  105. new List<SugarParameter> { new("@k", lockKey) });
  106. // GET_LOCK 返回 1=取到;0=超时未取到;NULL=出错/被中断
  107. if (got is null || Convert.ToInt64(got) != 1L)
  108. {
  109. guard.BusyReason = got is null ? "GET_LOCK returned NULL" : "cross-instance busy";
  110. await guard.DisposeAsync();
  111. return guard;
  112. }
  113. guard._dbLockHeld = true;
  114. // 同会话自证:持有者必须就是本连接。若不成立说明连接被换过(即本缺陷的回归),fail closed。
  115. var selfCheck = await lockDb.Ado.GetDataTableAsync(
  116. "SELECT IS_USED_LOCK(@k) AS owner_conn, CONNECTION_ID() AS my_conn",
  117. new List<SugarParameter> { new("@k", lockKey) });
  118. var owner = CellOrNull(selfCheck, 0);
  119. var me = CellOrNull(selfCheck, 1);
  120. if (owner is null || me is null || Convert.ToInt64(owner) != Convert.ToInt64(me))
  121. {
  122. logger?.LogError(
  123. "[InventoryInboundLock] 取锁后自证失败:owner={Owner} me={Me},判定跨实例互斥不可信,放弃本轮",
  124. owner, me);
  125. guard.BusyReason = "lock self-check failed";
  126. await guard.DisposeAsync();
  127. return guard;
  128. }
  129. guard.OwnerConnectionId = Convert.ToInt64(me);
  130. guard.Acquired = true;
  131. logger?.LogInformation(
  132. "[InventoryInboundLock] acquired key={Key} conn={Conn}", lockKey, guard.OwnerConnectionId);
  133. return guard;
  134. }
  135. catch (Exception ex)
  136. {
  137. logger?.LogError(ex, "[InventoryInboundLock] 取锁失败 key={Key}", lockKey);
  138. await guard.DisposeAsync();
  139. guard.BusyReason ??= "acquire failed: " + ex.Message;
  140. return guard;
  141. }
  142. }
  143. /// <summary>释放:先 RELEASE_LOCK(同一连接),再关闭专用连接兜底。</summary>
  144. public async ValueTask DisposeAsync()
  145. {
  146. if (_disposed) return;
  147. _disposed = true;
  148. Acquired = false;
  149. try
  150. {
  151. if (_dbLockHeld)
  152. {
  153. try
  154. {
  155. await _lockDb.Ado.ExecuteCommandAsync(
  156. "SELECT RELEASE_LOCK(@k)", new List<SugarParameter> { new("@k", _lockKey) });
  157. }
  158. catch (Exception ex)
  159. {
  160. // 连接 Dispose 时会话结束 / 连接归还池触发 reset,MySQL 会自动释放该会话的咨询锁,故这里只记录
  161. _logger?.LogWarning(ex, "[InventoryInboundLock] RELEASE_LOCK 失败,改由连接关闭兜底 key={Key}", _lockKey);
  162. }
  163. _dbLockHeld = false;
  164. }
  165. try { _lockDb.Dispose(); }
  166. catch (Exception ex) { _logger?.LogWarning(ex, "[InventoryInboundLock] 专用锁连接释放异常"); }
  167. }
  168. finally
  169. {
  170. if (_processGateHeld)
  171. {
  172. _processGateHeld = false;
  173. GateFor(_lockKey).Release();
  174. }
  175. }
  176. }
  177. private static async Task<object> ScalarAsync(SqlSugarClient db, string sql, List<SugarParameter> pars)
  178. {
  179. DataTable dt = await db.Ado.GetDataTableAsync(sql, pars);
  180. return CellOrNull(dt, 0);
  181. }
  182. private static object CellOrNull(DataTable dt, int columnIndex)
  183. {
  184. if (dt == null || dt.Rows.Count == 0 || dt.Columns.Count <= columnIndex) return null;
  185. var v = dt.Rows[0][columnIndex];
  186. return v == DBNull.Value ? null : v;
  187. }
  188. }