refactor: 用 Eino 框架重构编排层 #133
@@ -10,14 +10,21 @@ CAMTALK_AI_TTS_API_KEY=sk-your-tts-key
|
|||||||
# JWT 认证
|
# JWT 认证
|
||||||
CAMTALK_AUTH_JWT_SECRET=your-jwt-secret-here
|
CAMTALK_AUTH_JWT_SECRET=your-jwt-secret-here
|
||||||
|
|
||||||
# PostgreSQL(storage.driver 为 postgres 时必填)
|
# 三级存储配置
|
||||||
|
# L2: Redis(热数据分布式会话层)
|
||||||
|
CAMTALK_STORAGE_REDIS_ENABLED=true
|
||||||
|
# 开发时填写远程服务器地址,部署时 docker-compose 会覆盖为容器内网地址
|
||||||
|
CAMTALK_REDIS_ADDR=your-remote-server:6379
|
||||||
|
CAMTALK_REDIS_PASSWORD=your-redis-password
|
||||||
|
|
||||||
|
# L3: PostgreSQL(冷数据持久化层)
|
||||||
|
CAMTALK_STORAGE_PERSISTENCE_ENABLED=true
|
||||||
|
CAMTALK_STORAGE_PERSISTENCE_DRIVER=postgres
|
||||||
POSTGRES_USER=camtalk
|
POSTGRES_USER=camtalk
|
||||||
POSTGRES_PASSWORD=your-postgres-password
|
POSTGRES_PASSWORD=your-postgres-password
|
||||||
CAMTALK_STORAGE_DSN=postgres://camtalk:your-postgres-password@postgres:5432/camtalk?sslmode=disable
|
# 开发时填写远程服务器地址,部署时 docker-compose 会覆盖为容器内网地址
|
||||||
|
CAMTALK_STORAGE_DSN=postgres://camtalk:your-postgres-password@your-remote-server:5432/camtalk?sslmode=disable
|
||||||
|
|
||||||
# 可选覆盖(默认值见 config.yaml)
|
# 可选覆盖(默认值见 config.yaml)
|
||||||
# CAMTALK_SERVER_PORT=8080
|
# CAMTALK_SERVER_PORT=8080
|
||||||
# CAMTALK_LOG_LEVEL=info
|
# CAMTALK_LOG_LEVEL=info
|
||||||
# CAMTALK_STORAGE_DRIVER=memory
|
|
||||||
# CAMTALK_REDIS_ADDR=localhost:6379
|
|
||||||
# CAMTALK_REDIS_PASSWORD=
|
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/gin-gonic/gin"
|
"github.com/gin-gonic/gin"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
|
||||||
"github.com/hhs/camtalk/internal/api"
|
"github.com/hhs/camtalk/internal/api"
|
||||||
"github.com/hhs/camtalk/internal/auth"
|
"github.com/hhs/camtalk/internal/auth"
|
||||||
@@ -47,7 +48,7 @@ func main() {
|
|||||||
"addr", cfg.Server.Addr(),
|
"addr", cfg.Server.Addr(),
|
||||||
)
|
)
|
||||||
|
|
||||||
// 初始化存储层(条件初始化 PostgreSQL)
|
// 初始化存储层(三级存储架构:L1 内存 → L2 Redis → L3 PostgreSQL)
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
@@ -55,12 +56,17 @@ func main() {
|
|||||||
var msgRepo store.MessageRepository
|
var msgRepo store.MessageRepository
|
||||||
var sessRepo store.SessionRepository
|
var sessRepo store.SessionRepository
|
||||||
|
|
||||||
if cfg.Storage.Driver == "postgres" {
|
// L3: PostgreSQL(冷数据持久化层)
|
||||||
if cfg.Storage.DSN == "" {
|
dsn := cfg.Storage.Persistence.DSN
|
||||||
logger.Log.Fatalw("storage.dsn is required when storage.driver is postgres",
|
if dsn == "" {
|
||||||
|
dsn = cfg.Storage.DSN // 兼容旧配置
|
||||||
|
}
|
||||||
|
if cfg.Storage.Persistence.Enabled && cfg.Storage.Persistence.Driver == "postgres" {
|
||||||
|
if dsn == "" {
|
||||||
|
logger.Log.Fatalw("storage.persistence.dsn is required when persistence is enabled",
|
||||||
"hint", "set CAMTALK_STORAGE_DSN environment variable")
|
"hint", "set CAMTALK_STORAGE_DSN environment variable")
|
||||||
}
|
}
|
||||||
pool, err := store.NewPostgresPool(ctx, cfg.Storage.DSN)
|
pool, err := store.NewPostgresPool(ctx, dsn)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Log.Fatalw("failed to connect to postgres", "error", err)
|
logger.Log.Fatalw("failed to connect to postgres", "error", err)
|
||||||
}
|
}
|
||||||
@@ -74,14 +80,56 @@ func main() {
|
|||||||
userRepo = store.NewPgUserRepository(pool)
|
userRepo = store.NewPgUserRepository(pool)
|
||||||
msgRepo = store.NewPgMessageRepository(pool)
|
msgRepo = store.NewPgMessageRepository(pool)
|
||||||
sessRepo = store.NewPgSessionRepository(pool)
|
sessRepo = store.NewPgSessionRepository(pool)
|
||||||
logger.Log.Infow("postgres storage initialized", "driver", cfg.Storage.Driver)
|
logger.Log.Infow("L3 PostgreSQL storage initialized", "driver", cfg.Storage.Persistence.Driver)
|
||||||
} else {
|
} else {
|
||||||
userRepo = store.NewMemUserRepository()
|
userRepo = store.NewMemUserRepository()
|
||||||
logger.Log.Info("using in-memory storage")
|
logger.Log.Info("using in-memory user storage")
|
||||||
}
|
}
|
||||||
|
|
||||||
// 初始化 Session Manager
|
// L2: Redis(热数据分布式会话层)
|
||||||
|
var redisMgr *session.RedisManager
|
||||||
|
if cfg.Storage.Redis.Enabled {
|
||||||
|
rdb := redis.NewClient(&redis.Options{
|
||||||
|
Addr: cfg.Redis.Addr,
|
||||||
|
Password: cfg.Redis.Password,
|
||||||
|
DB: cfg.Redis.DB,
|
||||||
|
})
|
||||||
|
// 验证 Redis 连接
|
||||||
|
if err := rdb.Ping(ctx).Err(); err != nil {
|
||||||
|
logger.Log.Fatalw("failed to connect to redis", "error", err)
|
||||||
|
}
|
||||||
|
redisMgr = session.NewRedisManager(
|
||||||
|
rdb,
|
||||||
|
time.Duration(cfg.Session.TTL)*time.Minute,
|
||||||
|
cfg.Session.MaxHistory,
|
||||||
|
)
|
||||||
|
logger.Log.Infow("L2 Redis storage initialized",
|
||||||
|
"addr", cfg.Redis.Addr,
|
||||||
|
"db", cfg.Redis.DB)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 初始化 Session Manager(三级存储)
|
||||||
var sessionMgr session.Manager
|
var sessionMgr session.Manager
|
||||||
|
if cfg.Storage.Redis.Enabled {
|
||||||
|
// L1 + L2 + L3 三级存储
|
||||||
|
var tieredOpts []session.TieredOption
|
||||||
|
if sessRepo != nil {
|
||||||
|
tieredOpts = append(tieredOpts, session.WithTieredSessionRepository(sessRepo))
|
||||||
|
}
|
||||||
|
if msgRepo != nil {
|
||||||
|
tieredOpts = append(tieredOpts, session.WithTieredMessageRepository(msgRepo))
|
||||||
|
}
|
||||||
|
tieredMgr := session.NewTieredManager(
|
||||||
|
time.Duration(cfg.Session.TTL)*time.Minute,
|
||||||
|
cfg.Session.MaxHistory,
|
||||||
|
redisMgr,
|
||||||
|
tieredOpts...,
|
||||||
|
)
|
||||||
|
sessionMgr = tieredMgr
|
||||||
|
defer tieredMgr.Stop()
|
||||||
|
logger.Log.Info("session manager initialized with L1+L2+L3 tiered storage")
|
||||||
|
} else {
|
||||||
|
// L1 + L3 两级存储(无 Redis)
|
||||||
var sessionOpts []session.Option
|
var sessionOpts []session.Option
|
||||||
if msgRepo != nil {
|
if msgRepo != nil {
|
||||||
sessionOpts = append(sessionOpts, session.WithMessageRepository(msgRepo))
|
sessionOpts = append(sessionOpts, session.WithMessageRepository(msgRepo))
|
||||||
@@ -89,12 +137,15 @@ func main() {
|
|||||||
if sessRepo != nil {
|
if sessRepo != nil {
|
||||||
sessionOpts = append(sessionOpts, session.WithSessionRepository(sessRepo))
|
sessionOpts = append(sessionOpts, session.WithSessionRepository(sessRepo))
|
||||||
}
|
}
|
||||||
sessionMgr = session.NewMemoryManager(
|
memMgr := session.NewMemoryManager(
|
||||||
time.Duration(cfg.Session.TTL)*time.Minute,
|
time.Duration(cfg.Session.TTL)*time.Minute,
|
||||||
cfg.Session.MaxHistory,
|
cfg.Session.MaxHistory,
|
||||||
sessionOpts...,
|
sessionOpts...,
|
||||||
)
|
)
|
||||||
defer sessionMgr.(*session.MemoryManager).Stop()
|
sessionMgr = memMgr
|
||||||
|
defer memMgr.Stop()
|
||||||
|
logger.Log.Info("session manager initialized with L1+L3 storage (Redis disabled)")
|
||||||
|
}
|
||||||
|
|
||||||
// 初始化 AI 服务
|
// 初始化 AI 服务
|
||||||
logger.Log.Infow("initializing AI services",
|
logger.Log.Infow("initializing AI services",
|
||||||
|
|||||||
@@ -42,7 +42,12 @@ ai:
|
|||||||
sample_rate: 24000 # 输出采样率
|
sample_rate: 24000 # 输出采样率
|
||||||
|
|
||||||
storage:
|
storage:
|
||||||
driver: memory # memory / redis / postgres
|
# 三级存储架构:L1 内存 → L2 Redis → L3 PostgreSQL
|
||||||
|
redis:
|
||||||
|
enabled: true # 是否启用 Redis(L2 热数据层)
|
||||||
|
persistence:
|
||||||
|
enabled: true # 是否启用持久化(L3 冷数据层)
|
||||||
|
driver: postgres # postgres
|
||||||
# dsn 通过环境变量 CAMTALK_STORAGE_DSN 设置
|
# dsn 通过环境变量 CAMTALK_STORAGE_DSN 设置
|
||||||
|
|
||||||
redis:
|
redis:
|
||||||
|
|||||||
@@ -91,6 +91,19 @@ type TTSConfig struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type StorageConfig struct {
|
type StorageConfig struct {
|
||||||
|
Redis RedisStorageConfig `mapstructure:"redis"`
|
||||||
|
Persistence PersistenceConfig `mapstructure:"persistence"`
|
||||||
|
// Deprecated: 使用 Redis 和 Persistence 替代
|
||||||
|
Driver string `mapstructure:"driver"`
|
||||||
|
DSN string `mapstructure:"dsn"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type RedisStorageConfig struct {
|
||||||
|
Enabled bool `mapstructure:"enabled"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type PersistenceConfig struct {
|
||||||
|
Enabled bool `mapstructure:"enabled"`
|
||||||
Driver string `mapstructure:"driver"`
|
Driver string `mapstructure:"driver"`
|
||||||
DSN string `mapstructure:"dsn"`
|
DSN string `mapstructure:"dsn"`
|
||||||
}
|
}
|
||||||
@@ -189,6 +202,9 @@ func setDefaults(v *viper.Viper) {
|
|||||||
|
|
||||||
// storage
|
// storage
|
||||||
v.SetDefault("storage.driver", "memory")
|
v.SetDefault("storage.driver", "memory")
|
||||||
|
v.SetDefault("storage.redis.enabled", false)
|
||||||
|
v.SetDefault("storage.persistence.enabled", false)
|
||||||
|
v.SetDefault("storage.persistence.driver", "postgres")
|
||||||
|
|
||||||
// redis
|
// redis
|
||||||
v.SetDefault("redis.addr", "localhost:6379")
|
v.SetDefault("redis.addr", "localhost:6379")
|
||||||
@@ -220,7 +236,12 @@ func bindEnvVars(v *viper.Viper) {
|
|||||||
|
|
||||||
// 数据库
|
// 数据库
|
||||||
v.BindEnv("storage.dsn", "CAMTALK_STORAGE_DSN")
|
v.BindEnv("storage.dsn", "CAMTALK_STORAGE_DSN")
|
||||||
|
v.BindEnv("storage.persistence.dsn", "CAMTALK_STORAGE_DSN")
|
||||||
|
v.BindEnv("storage.redis.enabled", "CAMTALK_STORAGE_REDIS_ENABLED")
|
||||||
|
v.BindEnv("storage.persistence.enabled", "CAMTALK_STORAGE_PERSISTENCE_ENABLED")
|
||||||
|
v.BindEnv("storage.persistence.driver", "CAMTALK_STORAGE_PERSISTENCE_DRIVER")
|
||||||
|
|
||||||
// Redis(密码可能包含特殊字符,通过环境变量设置更安全)
|
// Redis(密码可能包含特殊字符,通过环境变量设置更安全)
|
||||||
|
v.BindEnv("redis.addr", "CAMTALK_REDIS_ADDR")
|
||||||
v.BindEnv("redis.password", "CAMTALK_REDIS_PASSWORD")
|
v.BindEnv("redis.password", "CAMTALK_REDIS_PASSWORD")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -36,6 +36,11 @@ func NewRedisManager(rdb *redis.Client, ttl time.Duration, maxHistory int) *Redi
|
|||||||
return &RedisManager{rdb: rdb, ttl: ttl, maxHistory: maxHistory}
|
return &RedisManager{rdb: rdb, ttl: ttl, maxHistory: maxHistory}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Ping 检查 Redis 连接是否正常。
|
||||||
|
func (m *RedisManager) Ping(ctx context.Context) error {
|
||||||
|
return m.rdb.Ping(ctx).Err()
|
||||||
|
}
|
||||||
|
|
||||||
func metaKey(id string) string { return fmt.Sprintf("session:%s:meta", id) }
|
func metaKey(id string) string { return fmt.Sprintf("session:%s:meta", id) }
|
||||||
func histKey(id string) string { return fmt.Sprintf("session:%s:history", id) }
|
func histKey(id string) string { return fmt.Sprintf("session:%s:history", id) }
|
||||||
func userSessKey(id string) string { return fmt.Sprintf("user:%s:sessions", id) }
|
func userSessKey(id string) string { return fmt.Sprintf("user:%s:sessions", id) }
|
||||||
|
|||||||
361
backend/internal/session/tiered.go
Normal file
361
backend/internal/session/tiered.go
Normal file
@@ -0,0 +1,361 @@
|
|||||||
|
package session
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"sync/atomic"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/hhs/camtalk/internal/logger"
|
||||||
|
"github.com/hhs/camtalk/internal/models"
|
||||||
|
"github.com/hhs/camtalk/internal/store"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TieredManager 三级存储 SessionManager 实现。
|
||||||
|
//
|
||||||
|
// L1(内存)→ L2(Redis)→ L3(PostgreSQL)
|
||||||
|
//
|
||||||
|
// 读:L1 miss → L2 miss → L3,回填到 L1+L2
|
||||||
|
// 写:L1 → L2(同步)→ L3(异步)
|
||||||
|
// 降级:Redis 不可用时,回退到 L1+L3 模式
|
||||||
|
type TieredManager struct {
|
||||||
|
l1 *MemoryManager // L1: 内存缓存
|
||||||
|
l2 *RedisManager // L2: Redis(可选)
|
||||||
|
sessRepo store.SessionRepository // L3: PostgreSQL 会话持久化(可选)
|
||||||
|
msgRepo store.MessageRepository // L3: PostgreSQL 消息持久化(可选)
|
||||||
|
|
||||||
|
redisOK atomic.Bool // Redis 健康状态
|
||||||
|
stopCh chan struct{} // 停止信号
|
||||||
|
}
|
||||||
|
|
||||||
|
// TieredOption TieredManager 的函数式选项。
|
||||||
|
type TieredOption func(*TieredManager)
|
||||||
|
|
||||||
|
// WithTieredSessionRepository 注入 L3 会话持久化仓库。
|
||||||
|
func WithTieredSessionRepository(repo store.SessionRepository) TieredOption {
|
||||||
|
return func(m *TieredManager) {
|
||||||
|
m.sessRepo = repo
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// WithTieredMessageRepository 注入 L3 消息持久化仓库。
|
||||||
|
func WithTieredMessageRepository(repo store.MessageRepository) TieredOption {
|
||||||
|
return func(m *TieredManager) {
|
||||||
|
m.msgRepo = repo
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewTieredManager 创建三级存储 SessionManager。
|
||||||
|
// l2 为 nil 时降级为 L1+L3 模式。
|
||||||
|
func NewTieredManager(
|
||||||
|
ttl time.Duration,
|
||||||
|
maxHistory int,
|
||||||
|
l2 *RedisManager,
|
||||||
|
opts ...TieredOption,
|
||||||
|
) *TieredManager {
|
||||||
|
m := &TieredManager{
|
||||||
|
l2: l2,
|
||||||
|
stopCh: make(chan struct{}),
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, opt := range opts {
|
||||||
|
opt(m)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 初始化 L1(内存),注入 L3 仓库实现 Write-Through
|
||||||
|
var l1Opts []Option
|
||||||
|
if m.sessRepo != nil {
|
||||||
|
l1Opts = append(l1Opts, WithSessionRepository(m.sessRepo))
|
||||||
|
}
|
||||||
|
if m.msgRepo != nil {
|
||||||
|
l1Opts = append(l1Opts, WithMessageRepository(m.msgRepo))
|
||||||
|
}
|
||||||
|
m.l1 = NewMemoryManager(ttl, maxHistory, l1Opts...)
|
||||||
|
|
||||||
|
// 初始化 Redis 健康状态
|
||||||
|
if l2 != nil {
|
||||||
|
m.redisOK.Store(true)
|
||||||
|
go m.healthCheck()
|
||||||
|
}
|
||||||
|
|
||||||
|
return m
|
||||||
|
}
|
||||||
|
|
||||||
|
// healthCheck 定期检查 Redis 健康状态。
|
||||||
|
func (m *TieredManager) healthCheck() {
|
||||||
|
ticker := time.NewTicker(30 * time.Second)
|
||||||
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ticker.C:
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
err := m.l2.Ping(ctx)
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
wasOK := m.redisOK.Load()
|
||||||
|
isOK := err == nil
|
||||||
|
m.redisOK.Store(isOK)
|
||||||
|
|
||||||
|
if wasOK && !isOK {
|
||||||
|
logger.Log.Warn("Redis connection lost, degrading to L1+L3 mode")
|
||||||
|
} else if !wasOK && isOK {
|
||||||
|
logger.Log.Info("Redis connection restored, resuming L1+L2+L3 mode")
|
||||||
|
}
|
||||||
|
case <-m.stopCh:
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// isRedisOK 检查 Redis 是否可用。
|
||||||
|
func (m *TieredManager) isRedisOK() bool {
|
||||||
|
return m.l2 != nil && m.redisOK.Load()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Stop 停止 TieredManager(清理后台 goroutine)。
|
||||||
|
func (m *TieredManager) Stop() {
|
||||||
|
close(m.stopCh)
|
||||||
|
m.l1.Stop()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Create 创建新会话。
|
||||||
|
// 写入:L1 → L2(同步)→ L3(异步)
|
||||||
|
func (m *TieredManager) Create(ctx context.Context, userID string, config models.SessionConfig) (string, error) {
|
||||||
|
// L1: 内存
|
||||||
|
id, err := m.l1.Create(ctx, userID, config)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if _, err := m.l2.Create(ctx, userID, config); err != nil {
|
||||||
|
logger.Log.Warnw("Redis Create failed, continuing without L2",
|
||||||
|
"session", id, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// L3: PostgreSQL(异步,由 L1 的 Write-Through 处理)
|
||||||
|
|
||||||
|
return id, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Get 获取会话。
|
||||||
|
// 读取:L1 → L2(回填 L1)→ L3(回填 L1+L2)
|
||||||
|
func (m *TieredManager) Get(ctx context.Context, sessionID string) (*models.Session, error) {
|
||||||
|
// L1: 内存
|
||||||
|
sess, err := m.l1.Get(ctx, sessionID)
|
||||||
|
if err == nil {
|
||||||
|
return sess, nil
|
||||||
|
}
|
||||||
|
if err != ErrSessionNotFound {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis
|
||||||
|
if m.isRedisOK() {
|
||||||
|
sess, err = m.l2.Get(ctx, sessionID)
|
||||||
|
if err == nil {
|
||||||
|
// 回填 L1
|
||||||
|
history, _ := m.l2.GetHistory(ctx, sessionID, 0)
|
||||||
|
m.l1.LoadSession(sess, history)
|
||||||
|
return sess, nil
|
||||||
|
}
|
||||||
|
if err != ErrSessionNotFound {
|
||||||
|
logger.Log.Warnw("Redis Get failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// L3: PostgreSQL(由 L1 的 Cache-Aside 处理)
|
||||||
|
// L1.Get 已经实现了从 PostgreSQL 恢复的逻辑
|
||||||
|
return nil, ErrSessionNotFound
|
||||||
|
}
|
||||||
|
|
||||||
|
// UpdateConfig 更新会话配置。
|
||||||
|
// 写入:L1 → L2(同步)→ L3(异步)
|
||||||
|
func (m *TieredManager) UpdateConfig(ctx context.Context, sessionID string, patch models.SessionConfigPatch) error {
|
||||||
|
// L1: 内存
|
||||||
|
if err := m.l1.UpdateConfig(ctx, sessionID, patch); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.UpdateConfig(ctx, sessionID, patch); err != nil {
|
||||||
|
logger.Log.Warnw("Redis UpdateConfig failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// L3: PostgreSQL(异步,由 L1 的 Write-Through 处理)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// UpdateTitle 更新会话标题。
|
||||||
|
// 写入:L1 → L2(同步)→ L3(异步)
|
||||||
|
func (m *TieredManager) UpdateTitle(ctx context.Context, sessionID string, title string) error {
|
||||||
|
// L1: 内存
|
||||||
|
if err := m.l1.UpdateTitle(ctx, sessionID, title); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.UpdateTitle(ctx, sessionID, title); err != nil {
|
||||||
|
logger.Log.Warnw("Redis UpdateTitle failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// L3: PostgreSQL(异步,由 L1 的 Write-Through 处理)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ListByUser 查询用户的会话列表。
|
||||||
|
// 读取:L1 + L2 + L3 合并去重
|
||||||
|
func (m *TieredManager) ListByUser(ctx context.Context, userID string, page, size int) ([]ConversationSummary, int, error) {
|
||||||
|
// 优先使用 L1(已集成 L3 回退逻辑)
|
||||||
|
return m.l1.ListByUser(ctx, userID, page, size)
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetHistory 获取对话历史。
|
||||||
|
// 读取:L1 → L2(回填 L1)→ L3(回填 L1+L2)
|
||||||
|
func (m *TieredManager) GetHistory(ctx context.Context, sessionID string, limit int) ([]models.Message, error) {
|
||||||
|
// L1: 内存
|
||||||
|
msgs, err := m.l1.GetHistory(ctx, sessionID, limit)
|
||||||
|
if err == nil && len(msgs) > 0 {
|
||||||
|
return msgs, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis
|
||||||
|
if m.isRedisOK() {
|
||||||
|
msgs, err = m.l2.GetHistory(ctx, sessionID, limit)
|
||||||
|
if err == nil && len(msgs) > 0 {
|
||||||
|
// 回填 L1(通过 Get 触发)
|
||||||
|
m.l1.Get(ctx, sessionID)
|
||||||
|
return msgs, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// L3: PostgreSQL(由 L1 的 Cache-Aside 处理)
|
||||||
|
return nil, ErrSessionNotFound
|
||||||
|
}
|
||||||
|
|
||||||
|
// AppendMessage 追加消息。
|
||||||
|
// 写入:L1 → L2(同步)→ L3(异步)
|
||||||
|
func (m *TieredManager) AppendMessage(ctx context.Context, sessionID string, msg models.Message) error {
|
||||||
|
// L1: 内存(Write-Through 到 L3)
|
||||||
|
if err := m.l1.AppendMessage(ctx, sessionID, msg); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.AppendMessage(ctx, sessionID, msg); err != nil {
|
||||||
|
logger.Log.Warnw("Redis AppendMessage failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetActiveRequest 设置当前活跃请求。
|
||||||
|
func (m *TieredManager) SetActiveRequest(ctx context.Context, sessionID string, requestID string) error {
|
||||||
|
// L1: 内存
|
||||||
|
if err := m.l1.SetActiveRequest(ctx, sessionID, requestID); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.SetActiveRequest(ctx, sessionID, requestID); err != nil {
|
||||||
|
logger.Log.Warnw("Redis SetActiveRequest failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetActiveRequestID 获取当前活跃请求 ID。
|
||||||
|
func (m *TieredManager) GetActiveRequestID(ctx context.Context, sessionID string) (string, error) {
|
||||||
|
// L1: 内存
|
||||||
|
id, err := m.l1.GetActiveRequestID(ctx, sessionID)
|
||||||
|
if err == nil && id != "" {
|
||||||
|
return id, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis
|
||||||
|
if m.isRedisOK() {
|
||||||
|
id, err = m.l2.GetActiveRequestID(ctx, sessionID)
|
||||||
|
if err == nil && id != "" {
|
||||||
|
return id, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return "", nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ClearActiveRequest 清除当前活跃请求。
|
||||||
|
func (m *TieredManager) ClearActiveRequest(ctx context.Context, sessionID string) error {
|
||||||
|
// L1: 内存
|
||||||
|
if err := m.l1.ClearActiveRequest(ctx, sessionID); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.ClearActiveRequest(ctx, sessionID); err != nil {
|
||||||
|
logger.Log.Warnw("Redis ClearActiveRequest failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Touch 刷新会话活跃时间。
|
||||||
|
func (m *TieredManager) Touch(ctx context.Context, sessionID string) error {
|
||||||
|
// L1: 内存
|
||||||
|
if err := m.l1.Touch(ctx, sessionID); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.Touch(ctx, sessionID); err != nil {
|
||||||
|
logger.Log.Warnw("Redis Touch failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Destroy 销毁会话。
|
||||||
|
// 写入:L1 → L2 → L3
|
||||||
|
func (m *TieredManager) Destroy(ctx context.Context, sessionID string) error {
|
||||||
|
// L1: 内存(Write-Through 到 L3)
|
||||||
|
if err := m.l1.Destroy(ctx, sessionID); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// L2: Redis(同步)
|
||||||
|
if m.isRedisOK() {
|
||||||
|
if err := m.l2.Destroy(ctx, sessionID); err != nil {
|
||||||
|
logger.Log.Warnw("Redis Destroy failed",
|
||||||
|
"session", sessionID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ActiveCount 返回活跃会话数量。
|
||||||
|
func (m *TieredManager) ActiveCount() int {
|
||||||
|
return m.l1.ActiveCount()
|
||||||
|
}
|
||||||
@@ -20,11 +20,18 @@ services:
|
|||||||
env_file:
|
env_file:
|
||||||
- /opt/camtalk/.env
|
- /opt/camtalk/.env
|
||||||
environment:
|
environment:
|
||||||
- CAMTALK_STORAGE_DRIVER=${CAMTALK_STORAGE_DRIVER:-postgres}
|
# 三级存储配置
|
||||||
|
- CAMTALK_STORAGE_REDIS_ENABLED=${CAMTALK_STORAGE_REDIS_ENABLED:-true}
|
||||||
|
- CAMTALK_STORAGE_PERSISTENCE_ENABLED=${CAMTALK_STORAGE_PERSISTENCE_ENABLED:-true}
|
||||||
|
- CAMTALK_STORAGE_PERSISTENCE_DRIVER=${CAMTALK_STORAGE_PERSISTENCE_DRIVER:-postgres}
|
||||||
- CAMTALK_STORAGE_DSN=postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@postgres:5432/camtalk?sslmode=disable
|
- CAMTALK_STORAGE_DSN=postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@postgres:5432/camtalk?sslmode=disable
|
||||||
|
- CAMTALK_REDIS_ADDR=redis:6379
|
||||||
|
- CAMTALK_REDIS_PASSWORD=${CAMTALK_REDIS_PASSWORD:-}
|
||||||
depends_on:
|
depends_on:
|
||||||
postgres:
|
postgres:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
|
redis:
|
||||||
|
condition: service_healthy
|
||||||
networks:
|
networks:
|
||||||
- camtalk-net
|
- camtalk-net
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
@@ -48,8 +55,24 @@ services:
|
|||||||
retries: 10
|
retries: 10
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
|
|
||||||
|
redis:
|
||||||
|
image: docker.m.daocloud.io/library/redis:7-alpine
|
||||||
|
container_name: camtalk-redis
|
||||||
|
command: redis-server --appendonly yes --requirepass ${CAMTALK_REDIS_PASSWORD:-}
|
||||||
|
volumes:
|
||||||
|
- redisdata:/data
|
||||||
|
networks:
|
||||||
|
- camtalk-net
|
||||||
|
healthcheck:
|
||||||
|
test: ["CMD", "redis-cli", "-a", "${CAMTALK_REDIS_PASSWORD:-}", "ping"]
|
||||||
|
interval: 5s
|
||||||
|
timeout: 3s
|
||||||
|
retries: 10
|
||||||
|
restart: unless-stopped
|
||||||
|
|
||||||
volumes:
|
volumes:
|
||||||
pgdata:
|
pgdata:
|
||||||
|
redisdata:
|
||||||
|
|
||||||
networks:
|
networks:
|
||||||
camtalk-net:
|
camtalk-net:
|
||||||
|
|||||||
Reference in New Issue
Block a user