InventoryInboundLockGuard.cs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258
  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>重算链路等待同一把库存锁的上限。后台增量仍走非阻塞 <see cref="TryAcquireAsync"/>。</summary>
  36. public static readonly TimeSpan DefaultWaitTimeout = TimeSpan.FromMinutes(10);
  37. /// <summary>重算链路轮询 <c>GET_LOCK(key, 0)</c> 的间隔,避免长时间占用错误连接。</summary>
  38. public static readonly TimeSpan DefaultPollInterval = TimeSpan.FromSeconds(2);
  39. /// <summary>锁语句自身的命令超时(秒);只跑 GET_LOCK/RELEASE_LOCK,不需要长超时。</summary>
  40. private const int LockCommandTimeoutSeconds = 10;
  41. /// <summary>
  42. /// 进程内闸门,**按锁键分桶**。单实例内多个 HTTP 请求 / Job 并发时先在这里被挡掉,
  43. /// 不必每次都去 MySQL 抢咨询锁;同时保证同进程内不会有两条并发链路各自持有一条锁连接。
  44. /// <para>分桶是必须的:早期实现是全类共用一个 <c>SemaphoreSlim</c>,
  45. /// 于是不同业务域用不同 lockKey 调用本类时会互相阻塞——
  46. /// S5 库存全量可以持锁数分钟,期间任何其它 lockKey 都拿不到闸门,
  47. /// 表现为「锁没被占用却一直取不到」。同键行为与分桶前完全一致。</para>
  48. /// </summary>
  49. private static readonly System.Collections.Concurrent.ConcurrentDictionary<string, SemaphoreSlim> ProcessGates = new();
  50. private static SemaphoreSlim GateFor(string lockKey) =>
  51. ProcessGates.GetOrAdd(lockKey, _ => new SemaphoreSlim(1, 1));
  52. private readonly SqlSugarClient _lockDb;
  53. private readonly ILogger _logger;
  54. private readonly string _lockKey;
  55. private bool _processGateHeld;
  56. private bool _dbLockHeld;
  57. private bool _disposed;
  58. /// <summary>是否成功拿到互斥权(进程内闸门 + 跨实例咨询锁都拿到才为 true)。</summary>
  59. public bool Acquired { get; private set; }
  60. /// <summary>持锁会话的 MySQL CONNECTION_ID(),仅用于排障/断言,未取到锁时为 null。</summary>
  61. public long? OwnerConnectionId { get; private set; }
  62. /// <summary>未取到锁时的原因,便于日志区分「同进程忙」与「其它实例忙」。</summary>
  63. public string BusyReason { get; private set; }
  64. private InventoryInboundLockGuard(SqlSugarClient lockDb, ILogger logger, string lockKey)
  65. {
  66. _lockDb = lockDb;
  67. _logger = logger;
  68. _lockKey = lockKey;
  69. }
  70. /// <summary>
  71. /// 有界等待同一把库存锁。生产调用使用 10 分钟上限、每 2 秒非阻塞轮询。
  72. /// 超时抛 <see cref="TimeoutException"/>,消息固定为 <c>inventory materialization lock busy</c>。
  73. /// </summary>
  74. public static Task<InventoryInboundLockGuard> AcquireAsync(
  75. ISqlSugarClient db,
  76. ILogger logger,
  77. string lockKey = LockKey,
  78. CancellationToken cancellationToken = default)
  79. => AcquireAsync(db, logger, lockKey, DefaultWaitTimeout, DefaultPollInterval, cancellationToken);
  80. /// <summary>测试可传入更短的等待窗口;生产入口不暴露这两个参数。</summary>
  81. public static async Task<InventoryInboundLockGuard> AcquireAsync(
  82. ISqlSugarClient db,
  83. ILogger logger,
  84. string lockKey,
  85. TimeSpan waitTimeout,
  86. TimeSpan pollInterval,
  87. CancellationToken cancellationToken = default)
  88. {
  89. if (db == null) throw new ArgumentNullException(nameof(db));
  90. if (waitTimeout <= TimeSpan.Zero) throw new ArgumentOutOfRangeException(nameof(waitTimeout));
  91. if (pollInterval <= TimeSpan.Zero) throw new ArgumentOutOfRangeException(nameof(pollInterval));
  92. var deadline = DateTime.UtcNow + waitTimeout;
  93. while (true)
  94. {
  95. cancellationToken.ThrowIfCancellationRequested();
  96. var guard = await TryAcquireAsync(db, logger, lockKey, cancellationToken);
  97. if (guard.Acquired)
  98. return guard;
  99. await guard.DisposeAsync();
  100. if (DateTime.UtcNow >= deadline)
  101. throw new TimeoutException("inventory materialization lock busy");
  102. var remaining = deadline - DateTime.UtcNow;
  103. var delay = remaining < pollInterval ? remaining : pollInterval;
  104. if (delay > TimeSpan.Zero)
  105. await Task.Delay(delay, cancellationToken);
  106. }
  107. }
  108. /// <summary>
  109. /// 尝试取锁;无论成功与否都返回一个 guard,调用方以 <see cref="Acquired"/> 判定,
  110. /// 并始终以 <c>await using</c> 持有(未取到锁的 guard 释放时什么都不做)。
  111. /// 后台增量使用本方法,抢不到立即返回,不排队。
  112. /// </summary>
  113. public static async Task<InventoryInboundLockGuard> TryAcquireAsync(
  114. ISqlSugarClient db,
  115. ILogger logger,
  116. string lockKey = LockKey,
  117. CancellationToken cancellationToken = default)
  118. {
  119. if (db == null) throw new ArgumentNullException(nameof(db));
  120. var main = db.CurrentConnectionConfig;
  121. var lockDb = new SqlSugarClient(new ConnectionConfig
  122. {
  123. ConnectionString = main.ConnectionString,
  124. DbType = main.DbType,
  125. // 关键:专用连接不自动关闭,GET_LOCK 与 RELEASE_LOCK 之间连接不回池、会话不重置
  126. IsAutoCloseConnection = false,
  127. });
  128. var guard = new InventoryInboundLockGuard(lockDb, logger, lockKey);
  129. // 第一层:进程内闸门,非阻塞
  130. var gate = await GateFor(lockKey).WaitAsync(NoWaitMilliseconds, cancellationToken);
  131. if (!gate)
  132. {
  133. guard.BusyReason = "in-process busy";
  134. lockDb.Dispose();
  135. guard._disposed = true; // 什么都没持有,调用方的 await using 直接空转
  136. return guard;
  137. }
  138. guard._processGateHeld = true;
  139. try
  140. {
  141. lockDb.Ado.CommandTimeOut = LockCommandTimeoutSeconds;
  142. // 第二层:MySQL 咨询锁,非阻塞(timeout=0)
  143. var got = await ScalarAsync(lockDb, "SELECT GET_LOCK(@k, 0)",
  144. new List<SugarParameter> { new("@k", lockKey) });
  145. // GET_LOCK 返回 1=取到;0=超时未取到;NULL=出错/被中断
  146. if (got is null || Convert.ToInt64(got) != 1L)
  147. {
  148. guard.BusyReason = got is null ? "GET_LOCK returned NULL" : "cross-instance busy";
  149. await guard.DisposeAsync();
  150. return guard;
  151. }
  152. guard._dbLockHeld = true;
  153. // 同会话自证:持有者必须就是本连接。若不成立说明连接被换过(即本缺陷的回归),fail closed。
  154. var selfCheck = await lockDb.Ado.GetDataTableAsync(
  155. "SELECT IS_USED_LOCK(@k) AS owner_conn, CONNECTION_ID() AS my_conn",
  156. new List<SugarParameter> { new("@k", lockKey) });
  157. var owner = CellOrNull(selfCheck, 0);
  158. var me = CellOrNull(selfCheck, 1);
  159. if (owner is null || me is null || Convert.ToInt64(owner) != Convert.ToInt64(me))
  160. {
  161. logger?.LogError(
  162. "[InventoryInboundLock] 取锁后自证失败:owner={Owner} me={Me},判定跨实例互斥不可信,放弃本轮",
  163. owner, me);
  164. guard.BusyReason = "lock self-check failed";
  165. await guard.DisposeAsync();
  166. return guard;
  167. }
  168. guard.OwnerConnectionId = Convert.ToInt64(me);
  169. guard.Acquired = true;
  170. logger?.LogInformation(
  171. "[InventoryInboundLock] acquired key={Key} conn={Conn}", lockKey, guard.OwnerConnectionId);
  172. return guard;
  173. }
  174. catch (Exception ex)
  175. {
  176. logger?.LogError(ex, "[InventoryInboundLock] 取锁失败 key={Key}", lockKey);
  177. await guard.DisposeAsync();
  178. guard.BusyReason ??= "acquire failed: " + ex.Message;
  179. return guard;
  180. }
  181. }
  182. /// <summary>释放:先 RELEASE_LOCK(同一连接),再关闭专用连接兜底。</summary>
  183. public async ValueTask DisposeAsync()
  184. {
  185. if (_disposed) return;
  186. _disposed = true;
  187. Acquired = false;
  188. try
  189. {
  190. if (_dbLockHeld)
  191. {
  192. try
  193. {
  194. await _lockDb.Ado.ExecuteCommandAsync(
  195. "SELECT RELEASE_LOCK(@k)", new List<SugarParameter> { new("@k", _lockKey) });
  196. }
  197. catch (Exception ex)
  198. {
  199. // 连接 Dispose 时会话结束 / 连接归还池触发 reset,MySQL 会自动释放该会话的咨询锁,故这里只记录
  200. _logger?.LogWarning(ex, "[InventoryInboundLock] RELEASE_LOCK 失败,改由连接关闭兜底 key={Key}", _lockKey);
  201. }
  202. _dbLockHeld = false;
  203. }
  204. try { _lockDb.Dispose(); }
  205. catch (Exception ex) { _logger?.LogWarning(ex, "[InventoryInboundLock] 专用锁连接释放异常"); }
  206. }
  207. finally
  208. {
  209. if (_processGateHeld)
  210. {
  211. _processGateHeld = false;
  212. GateFor(_lockKey).Release();
  213. }
  214. }
  215. }
  216. private static async Task<object> ScalarAsync(SqlSugarClient db, string sql, List<SugarParameter> pars)
  217. {
  218. DataTable dt = await db.Ado.GetDataTableAsync(sql, pars);
  219. return CellOrNull(dt, 0);
  220. }
  221. private static object CellOrNull(DataTable dt, int columnIndex)
  222. {
  223. if (dt == null || dt.Rows.Count == 0 || dt.Columns.Count <= columnIndex) return null;
  224. var v = dt.Rows[0][columnIndex];
  225. return v == DBNull.Value ? null : v;
  226. }
  227. }