// 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 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) } } } }