Files
Blizzard f926f6fd41 feat(obs): exec 执行轨迹 Redis 回放 —— 连晚/重连不再丢轨迹事件
执行轨迹原本只走瞬时 core NATS(sundynix.exec.<id>),SSE 连晚或刷新重连就丢掉
已发生的节点点亮/工具调用/推理过程事件。本提交把它做成与 token 流同构的可回放流:

- store: Redis Stream 函数加 channel 维度(ChannelToken="stream" / ChannelExec="exec"),
  同一套 XADD/XREAD/TTL 复用;key 按 channel 分命名空间互不串扰。
- gateway: 提交即启 startExecRecorder 后台订阅轨迹落 Redis(与 SSE 是否在线无关,
  12min 兜底含 HITL 审批等待);StreamExec 改为优先 Redis 回放 + Last-Event-ID 断点续传,
  Redis 降级回退 live NATS(streamExecLive)。

单测 streamKey channel 隔离;live:任务 done 后再连 /exec,仍从 Redis 完整回放
全程轨迹(含推理过程),Redis XLEN 对账一致。

至此「可靠性细节」两项(优雅停机 + 轨迹回放)补齐。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-26 14:03:31 +08:00

157 lines
4.9 KiB
Go
Raw Permalink 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 }
// 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()
}
}