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