package store import ( "context" "log" "time" "github.com/redis/go-redis/v9" ) // Redis 持有 CacheDB 连接(Session / Rate Limit)。 // rdb 为 nil 表示降级模式(连接失败时放行,不限流)。 type Redis struct { rdb *redis.Client } // OpenRedis 连接 CacheDB。连接失败不 fatal:返回降级实例(限流放行)。 func OpenRedis(addr string) *Redis { rdb := redis.NewClient(&redis.Options{Addr: addr}) ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() if err := rdb.Ping(ctx).Err(); err != nil { log.Printf("[store] redis 不可用,降级运行(不限流): %v", err) _ = rdb.Close() return &Redis{} } log.Println("[store] redis connected") return &Redis{rdb: rdb} } // Enabled 报告是否处于真实限流模式。 func (r *Redis) Enabled() bool { return r.rdb != nil } // Ping 活性探测:发一次 PING,验证 Redis 此刻真的可达(非仅启动时连过)。 func (r *Redis) Ping(ctx context.Context) bool { if r.rdb == nil { return false } return r.rdb.Ping(ctx).Err() == nil } // Allow 滑动窗口计数限流:在 window 内对 key 累加,超过 limit 即拒绝。 // 降级模式(rdb==nil)始终放行。 func (r *Redis) Allow(ctx context.Context, key string, limit int64, window time.Duration) (bool, error) { if r.rdb == nil { return true, nil } rk := "sundynix:ratelimit:" + key n, err := r.rdb.Incr(ctx, rk).Result() if err != nil { return true, err // 限流后端故障时放行,不阻断业务 } if n == 1 { _ = r.rdb.Expire(ctx, rk, window).Err() // 首次计数设置窗口过期 } return n <= limit, nil } // ---- Token 流持久化(Redis Stream:可回放的追加日志,根治 SSE 连晚/重连丢 token)---- // ---- Token 用量日预算(成本护栏:按用户按天累计 token,供提交前门控)---- // usageDailyKey 是某用户某天的累计 token key(day 形如 20260626,便于到期自然滚动)。 func usageDailyKey(userID, day string) string { return "sundynix:usage:" + userID + ":" + day } // AddUsage 给某用户当天累计 token,并返回累计后总量(首次设 48h 过期自动清理)。降级返回 (0,nil)。 func (r *Redis) AddUsage(ctx context.Context, userID, day string, tokens int) (int64, error) { if r.rdb == nil || userID == "" { return 0, nil } rk := usageDailyKey(userID, day) n, err := r.rdb.IncrBy(ctx, rk, int64(tokens)).Result() if err != nil { return 0, err } if n == int64(tokens) { // 首次累计:设过期 _ = r.rdb.Expire(ctx, rk, 48*time.Hour).Err() } return n, nil } // GetUsage 取某用户当天已用 token(不存在/降级返回 0)。 func (r *Redis) GetUsage(ctx context.Context, userID, day string) int64 { if r.rdb == nil || userID == "" { return 0 } n, err := r.rdb.Get(ctx, usageDailyKey(userID, day)).Int64() if err != nil { return 0 } return n } // ---- 流持久化(Redis Stream:可回放的追加日志,根治 SSE 连晚/重连丢事件)---- // 按 channel 分流:token 走 "stream",执行轨迹走 "exec"——同一套 XADD/XREAD/TTL 逻辑复用。 const streamTTL = 10 * time.Minute // Channel 是回放流的种类(决定 Redis key 命名空间,互不串扰)。 const ( ChannelToken = "stream" // Token 流(与历史 key 兼容) ChannelExec = "exec" // 执行轨迹流 ) func streamKey(channel, taskID string) string { return "sundynix:" + channel + ":" + taskID } // StreamEntry 是回放流里的一条记录(ID 用于 SSE 的 Last-Event-ID 断点续传)。 type StreamEntry struct { ID string Kind string // token / exec / done Data string } // StreamAppend 把一条记录追加到任务某 channel 的 Redis Stream(带 TTL 自动清理)。 func (r *Redis) StreamAppend(ctx context.Context, channel, taskID, kind, data string) error { if r.rdb == nil { return nil } k := streamKey(channel, taskID) if err := r.rdb.XAdd(ctx, &redis.XAddArgs{Stream: k, Values: map[string]any{"kind": kind, "data": data}}).Err(); err != nil { return err } return r.rdb.Expire(ctx, k, streamTTL).Err() } // StreamRead 从 lastID 之后阻塞读取某 channel 的新条目(XREAD BLOCK)。lastID="0" 表示从头回放。 // 阻塞超时无新数据时返回空切片 + 原 lastID(调用方据此继续轮询)。 func (r *Redis) StreamRead(ctx context.Context, channel, taskID, lastID string, block time.Duration) ([]StreamEntry, string, error) { if r.rdb == nil { return nil, lastID, nil } res, err := r.rdb.XRead(ctx, &redis.XReadArgs{ Streams: []string{streamKey(channel, taskID), lastID}, Block: block, Count: 256, }).Result() if err == redis.Nil { return nil, lastID, nil // 阻塞超时、无新条目 } if err != nil { return nil, lastID, err } var out []StreamEntry nl := lastID for _, st := range res { for _, m := range st.Messages { out = append(out, StreamEntry{ID: m.ID, Kind: asString(m.Values["kind"]), Data: asString(m.Values["data"])}) nl = m.ID } } return out, nl, nil } func asString(v any) string { if s, ok := v.(string); ok { return s } return "" } // ---- 微信登录 ticket ---- // ticket 是一次性、短命、高频轮询的临时态,用 Redis 存(不落 PG)。 // Redis 降级时回退进程内内存 map —— 本地开发无 Redis 也能登录;生产多实例下内存回退 // 会因 PC 轮询与微信回调可能落在不同实例而失效,所以生产**必须**有 Redis(服务状态页会亮)。 // WxTicketSet 写入 ticket 状态(JSON 值),带 TTL。 func (r *Redis) WxTicketSet(ctx context.Context, ticket, value string, ttl time.Duration) error { if r.rdb == nil { memTicketSet(ticket, value, ttl) return nil } return r.rdb.Set(ctx, "wxlogin:"+ticket, value, ttl).Err() } // WxTicketGet 读取 ticket 状态;不存在或过期返回空串。 func (r *Redis) WxTicketGet(ctx context.Context, ticket string) string { if r.rdb == nil { return memTicketGet(ticket) } v, err := r.rdb.Get(ctx, "wxlogin:"+ticket).Result() if err != nil { return "" } return v } // WxTokenGet/Set 缓存微信 access_token(跨实例共享,避免重复拉取互相失效)。 // 无 Redis 时返回空 → 调用方每次现拉(单实例开发可接受)。 func (r *Redis) WxTokenGet(ctx context.Context, appID string) string { if r.rdb == nil { return "" } v, err := r.rdb.Get(ctx, "wxtoken:"+appID).Result() if err != nil { return "" } return v } func (r *Redis) WxTokenSet(ctx context.Context, appID, token string, ttl time.Duration) { if r.rdb == nil { return } _ = r.rdb.Set(ctx, "wxtoken:"+appID, token, ttl).Err() } // Close 释放底层连接。 func (r *Redis) Close() { if r.rdb != nil { _ = r.rdb.Close() } }