Files
sundynix-agentix/sundynix-gateway/internal/store/redis.go
T
Blizzard 373167b705 feat(admin): 补全平台任务观测 + 空间管理 + 基建加 MinIO/实时探针
对照已实现系统功能补 admin 缺失模块(feat/site):

- 全平台任务/运行观测:GET /admin/tasks(跨租户查 sundynix_task,join 租户名/提交人邮箱/评测,
  按状态/租户筛 + 状态计数;含 HITL 待审批=筛 waiting)+ TasksPage(状态卡+筛+表)。
- 空间(Space)管理:GET /admin/spaces(跨租户列 Space + 成员数子查询 + kind/归档态)+ SpacesPage。
- 基建观测:infra 加 MinIO(126 对象存储);postgres/redis/minio 从启动标志升级为实时 ping
  (blob.Ping/Postgres.Ping/Redis.Ping),能反映中途掉线。StatusPage 加 minio 标签、6 个依赖。
- 模型健康/熔断:确认 DashboardPage 已有「运行时链路态」渲染 m.health,无需重做。

验证:admin tsc + vitest 41 过、gateway build/vet/test 过;两新页浏览器实测渲染+优雅错误处理;
两新端点起临时 gateway 打真 PG 实测——tasks(counts+3路join)、spaces(成员数子查询)均返真数据。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-18 17:06:25 +08:00

165 lines
5.1 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 keyday 形如 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 ""
}
// Close 释放底层连接。
func (r *Redis) Close() {
if r.rdb != nil {
_ = r.rdb.Close()
}
}