184 lines
6.9 KiB
Go
184 lines
6.9 KiB
Go
// Package utils 提供统一的 Redis 分布式锁(防并发互斥)。
|
||
//
|
||
// 收敛全仓「SET key token EX ttl NX + 归属校验释放」范式
|
||
// (此前 ai-agent / shop-user-trade / model-gateway 三处逐字复制)。
|
||
// 只暴露一个入口:
|
||
// - WithLock:抢锁 + 自动续期 + 临界区 fn + 返回时保证释放。
|
||
//
|
||
// 正确性要点:
|
||
// - value 存 uuid token 标识持有者(不存 true),key 标识资源;
|
||
// - 释放 / 续期均按 token 比对归属:锁 TTL 过期易主后,旧持有者的释放不会误删新持有者的锁,
|
||
// 续期也不会把锁无限续到他人头上;
|
||
// - 归属比对经 go-redis WATCH/MULTI/EXEC 原子事务完成(通过 gogf 官方 escape hatch
|
||
// GetAdapter().Client() 取底层客户端),不手写 Lua 脚本。
|
||
//
|
||
// 兼容红线:本包只用 gf v2.9.5 已有 API(Set+SetOption{NX,EX}、GetAdapter().Client())。
|
||
package utils
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"time"
|
||
|
||
"github.com/gogf/gf/v2/database/gredis"
|
||
"github.com/gogf/gf/v2/frame/g"
|
||
"github.com/google/uuid"
|
||
goredis "github.com/redis/go-redis/v9"
|
||
)
|
||
|
||
// Lock 分布式锁(防并发互斥),委托 common/redislock.WithLock(4 次尝试,间隔 500ms)。
|
||
// 释放 / 续期按 token 归属比对(WATCH/MULTI/EXEC 原子事务,不写 Lua):
|
||
// 锁 TTL 过期易主后,旧持有者的释放不会误删新持有者的锁。
|
||
// fn 返回 err 时 success=false、err 上抛。
|
||
func Lock(ctx context.Context, key string, expireSeconds int64, fn func(ctx context.Context) error) (success bool, err error) {
|
||
_, err = WithLock(ctx, key, expireSeconds, fn, 4)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
return true, nil
|
||
}
|
||
|
||
// WithLock 函数式锁:抢锁(可选重试次数),成功后自动续期(默认超过一半 TTL 就续期),
|
||
// 执行 fn,返回时停止续期并保证释放(WATCH/MULTI 按 token 归属释放,防误删他人锁)。
|
||
//
|
||
// retryTimes 传了循环次数:最多尝试那么多次,仍抢不到返回锁占用错误;
|
||
// retryTimes 未传:无限等待,一直重试直到抢到锁或 ctx 取消。
|
||
// 两次尝试间隔固定 500ms。
|
||
//
|
||
// 返回值:
|
||
// - (true, nil):抢到锁并执行成功;
|
||
// - (true, fn 的 err):抢到锁、fn 执行失败(业务错误由 err 单独表达);
|
||
// - (false, err):重试次数耗尽(锁占用)、Redis 故障或 ctx 被取消,均以 err 表达。
|
||
//
|
||
// 续期与释放均用 context.WithoutCancel(ctx):即使 fn 中途 ctx 被取消,锁仍能正常续期并释放,
|
||
// 不会残留到 TTL 造成下一个持有者等待。
|
||
func WithLock(ctx context.Context, key string, ttlSeconds int64, fn func(ctx context.Context) error, retryTimes ...int) (bool, error) {
|
||
maxRetries := -1
|
||
if len(retryTimes) > 0 {
|
||
maxRetries = retryTimes[0]
|
||
}
|
||
l := &lock{key: key, token: uuid.NewString(), ttl: ttlSeconds}
|
||
|
||
for attempt := 0; ; attempt++ {
|
||
if maxRetries >= 0 && attempt >= maxRetries {
|
||
return false, errors.New("redis lock busy")
|
||
}
|
||
ok, err := l.acquire(ctx)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
if ok {
|
||
break
|
||
}
|
||
select {
|
||
case <-ctx.Done():
|
||
return false, ctx.Err()
|
||
case <-time.After(500 * time.Millisecond):
|
||
}
|
||
}
|
||
|
||
stop := make(chan struct{})
|
||
go l.renewLoop(context.WithoutCancel(ctx), stop)
|
||
defer func() {
|
||
close(stop)
|
||
_ = l.release(context.WithoutCancel(ctx))
|
||
}()
|
||
|
||
if err := fn(ctx); err != nil {
|
||
return true, err
|
||
}
|
||
return true, nil
|
||
}
|
||
|
||
// lock 单次锁会话:token 标识持有者身份,释放 / 续期按 token 原子比对。
|
||
// 不导出——生命周期全部由 WithLock 编排。
|
||
type lock struct {
|
||
key string
|
||
token string
|
||
ttl int64 // 秒
|
||
}
|
||
|
||
// acquire 抢锁;返回 true 表示抢到(SET NX 成功)。
|
||
// 注意:gogf 新版 SetNX(ctx,key,value) 不接收 TTL 参数,直接 SETNX 会永不过期(崩溃后死锁);
|
||
// 故改用 Set + SetOption{NX,TTLOption{EX}},原子地执行 `SET key token EX <ttl> NX`。
|
||
func (l *lock) acquire(ctx context.Context) (bool, error) {
|
||
r, err := g.Redis().Set(ctx, l.key, l.token, gredis.SetOption{
|
||
TTLOption: gredis.TTLOption{EX: &l.ttl},
|
||
NX: true,
|
||
})
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
// SET NX 失败时 Redis 返回空回复(nil),gogf 转为值 nil 的 gvar;成功时返回 "OK"
|
||
return !r.IsNil(), nil
|
||
}
|
||
|
||
// release 释放锁:WATCH key → GET 比对 token → MULTI/DEL/EXEC,仅当仍为本锁 token 时删除。
|
||
// 非持有者调用是 no-op:锁过期易主后,旧持有者的释放不会删掉新持有者的锁。
|
||
func (l *lock) release(ctx context.Context) error {
|
||
return l.compareAndWrite(ctx, false)
|
||
}
|
||
|
||
// renew 续期:WATCH key → GET 比对 token → MULTI/EXPIRE/EXEC,仅当仍为本锁 token 时重置 TTL。
|
||
// 锁已易主时 no-op,避免续到他人锁上(把新持有者的锁无限延长)。
|
||
func (l *lock) renew(ctx context.Context) error {
|
||
return l.compareAndWrite(ctx, true)
|
||
}
|
||
|
||
// compareAndWrite 归属比对 + 原子写(WATCH/MULTI/EXEC,无 Lua):
|
||
// - 锁不存在(已过期 / 被删):no-op;
|
||
// - key 值 ≠ 本锁 token(已易主):no-op;
|
||
// - key 值 == 本锁 token:renew=true 时重置 TTL,否则删除。
|
||
//
|
||
// WATCH 保证:比对与写之间若 key 被其他客户端改动,EXEC 会被服务端中止(返回空),写不生效。
|
||
// 经 gogf 官方 escape hatch(GetAdapter().Client())取底层 go-redis 客户端执行原子事务;
|
||
// 各服务 main.go 空导入 contrib/nosql/redis/v2 后 Adapter 为 go-redis 实现,
|
||
// Client() 返回 redis.UniversalClient(单节点 / 哨兵 / 集群均实现 Watch)。
|
||
func (l *lock) compareAndWrite(ctx context.Context, renew bool) error {
|
||
universal, ok := g.Redis().GetAdapter().Client().(goredis.UniversalClient)
|
||
if !ok {
|
||
return errors.New("redis 底层客户端非 UniversalClient,无法执行 WATCH 原子事务")
|
||
}
|
||
return universal.Watch(ctx, func(tx *goredis.Tx) error {
|
||
val, err := tx.Get(ctx, l.key).Result()
|
||
if errors.Is(err, goredis.Nil) {
|
||
return nil // 锁已过期 / 被删
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if val != l.token {
|
||
return nil // 锁已易主,不动他人锁
|
||
}
|
||
_, err = tx.TxPipelined(ctx, func(pipe goredis.Pipeliner) error {
|
||
if renew {
|
||
pipe.Expire(ctx, l.key, time.Duration(l.ttl)*time.Second)
|
||
} else {
|
||
pipe.Del(ctx, l.key)
|
||
}
|
||
return nil
|
||
})
|
||
return err
|
||
}, l.key)
|
||
}
|
||
|
||
// renewLoop 自动续期看门狗:每 ttl/2 续一次(「超过一半就续期」)。
|
||
// 续期失败仅记日志,不打断临界区(WATCH 事务在锁易主时会因 token 不匹配安全地 no-op)。
|
||
// WithLock 内部启动,临界区结束即停止。
|
||
func (l *lock) renewLoop(ctx context.Context, stop <-chan struct{}) {
|
||
ticker := time.NewTicker(time.Duration(l.ttl) * time.Second / 2)
|
||
defer ticker.Stop()
|
||
for {
|
||
select {
|
||
case <-stop:
|
||
return
|
||
case <-ctx.Done():
|
||
return
|
||
case <-ticker.C:
|
||
if err := l.renew(context.WithoutCancel(ctx)); err != nil {
|
||
g.Log().Warningf(ctx, "redislock 续期失败: key=%s err=%v", l.key, err)
|
||
}
|
||
}
|
||
}
|
||
}
|