diff --git a/internal/session/fallback_store.go b/internal/session/fallback_store.go index cb24da0..391e3f2 100644 --- a/internal/session/fallback_store.go +++ b/internal/session/fallback_store.go @@ -58,9 +58,9 @@ func (fs *FallbackStore) healthLoop() { ok := fs.ping() if ok && !prev { - // Redis 恢复:先切到 Redis(新会话立即走 Redis),再迁移内存存量 - fs.healthy.Store(true) + // Redis 恢复:先迁移(内部先清旧数据)再切 healthy,避免迁移期间新请求创建 Redis 会话被误删 fs.migrateSessions() + fs.healthy.Store(true) fs.degraded.Store(false) } else if !ok && prev { // Redis 掉线,标记降级(快速恢复由 recoveryLoop 负责) @@ -74,7 +74,6 @@ func (fs *FallbackStore) healthLoop() { } // recoveryLoop 启动后持续监听退化状态,一旦降级立即启动快速恢复轮询(5s 间隔) -// 退出条件:Redis 恢复成功或 connectOnce 竞争到恢复权限后执行迁移 func (fs *FallbackStore) recoveryLoop() { ticker := time.NewTicker(recoveryInterval) defer ticker.Stop() @@ -92,9 +91,9 @@ func (fs *FallbackStore) recoveryLoop() { continue // 仍未恢复 } - // Redis 已恢复 - fs.healthy.Store(true) + // Redis 已恢复:先迁后切,避免旧会话残留 + 新请求竞争 fs.migrateSessions() + fs.healthy.Store(true) fs.degraded.Store(false) fs.recovering.Store(false) } @@ -113,7 +112,22 @@ func (fs *FallbackStore) migrateSessions() { return } + // 收集迁移涉及的用户 + userSet := make(map[uint]bool, len(sessions)) + for _, s := range sessions { + userSet[s.UserID] = true + } + ctx := context.Background() + + // 先清理 Redis 中同用户的旧会话(降级期间用户重新登录,旧会话已失效) + // 避免 Redis 恢复后出现新旧两条会话并存、设备列表重复 + for uid := range userSet { + if err := fs.redis.DeleteByUID(uid); err != nil { + log.Printf("[FallbackStore] 迁移前清理旧会话失败 uid=%d err=%v", uid, err) + } + } + pipe := fs.client.Pipeline() for _, s := range sessions { @@ -141,7 +155,7 @@ func (fs *FallbackStore) migrateSessions() { _ = fs.memory.Delete(s.ID) } - log.Printf("[FallbackStore] Redis 已恢复,静默迁移 %d 条会话(内存已清理)", len(sessions)) + log.Printf("[FallbackStore] Redis 已恢复,静默迁移 %d 条会话(旧会话已清理,内存已释放)", len(sessions)) } // ping 检查 Redis 连接状态(带超时)