fix: 第二轮修复 — 数据库启动检查、会话持久化、URL路由、设备排序等
1. DevTools 启动前检查数据库状态,失败时自动尝试启动 2. ai-core 添加数据库断线重连机制 (30秒间隔) 3. Dashboard 添加数据库状态卡片 (启动/停止/重启) 4. Gateway 会话空闲超时管理 (30分钟标记空闲) 5. 会话/消息 PostgreSQL 持久化 (SessionStore + REST API) 6. 前端服务端会话持久化 + URL hash 路由 + 侧边栏管理 7. 管理员回到主对话按钮 8. IoT 设备卡片固定排序 9. 更新相关文档
This commit is contained in:
@@ -0,0 +1,270 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
_ "github.com/lib/pq"
|
||||
)
|
||||
|
||||
// Session 会话模型
|
||||
type Session struct {
|
||||
ID string `json:"id"`
|
||||
UserID string `json:"user_id"`
|
||||
Title string `json:"title"`
|
||||
IsMain bool `json:"is_main"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at"`
|
||||
}
|
||||
|
||||
// Message 消息模型
|
||||
type Message struct {
|
||||
ID int `json:"id"`
|
||||
SessionID string `json:"session_id"`
|
||||
Role string `json:"role"`
|
||||
Content string `json:"content"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
}
|
||||
|
||||
// SessionStore 会话持久化存储
|
||||
type SessionStore struct {
|
||||
db *sql.DB
|
||||
}
|
||||
|
||||
// NewSessionStore 初始化数据库连接并自动建表
|
||||
// 如果连接失败,返回 nil 和错误(调用方可以选择降级为仅内存模式)
|
||||
func NewSessionStore(databaseURL string) (*SessionStore, error) {
|
||||
db, err := sql.Open("postgres", databaseURL)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("无法打开数据库连接: %w", err)
|
||||
}
|
||||
|
||||
// 配置连接池
|
||||
db.SetMaxOpenConns(25)
|
||||
db.SetMaxIdleConns(5)
|
||||
db.SetConnMaxLifetime(5 * time.Minute)
|
||||
|
||||
// 验证连接
|
||||
if err := db.Ping(); err != nil {
|
||||
db.Close()
|
||||
return nil, fmt.Errorf("数据库连接验证失败: %w", err)
|
||||
}
|
||||
|
||||
store := &SessionStore{db: db}
|
||||
|
||||
// 自动建表
|
||||
if err := store.migrate(); err != nil {
|
||||
db.Close()
|
||||
return nil, fmt.Errorf("数据库迁移失败: %w", err)
|
||||
}
|
||||
|
||||
log.Println("[SessionStore] PostgreSQL 持久化存储已初始化")
|
||||
return store, nil
|
||||
}
|
||||
|
||||
// migrate 自动创建表结构
|
||||
func (s *SessionStore) migrate() error {
|
||||
queries := []string{
|
||||
`CREATE TABLE IF NOT EXISTS sessions (
|
||||
id VARCHAR(64) PRIMARY KEY,
|
||||
user_id VARCHAR(128) NOT NULL,
|
||||
title VARCHAR(256) DEFAULT '',
|
||||
is_main BOOLEAN DEFAULT FALSE,
|
||||
created_at TIMESTAMP DEFAULT NOW(),
|
||||
updated_at TIMESTAMP DEFAULT NOW()
|
||||
)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_sessions_user_id ON sessions(user_id)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_sessions_updated_at ON sessions(updated_at DESC)`,
|
||||
|
||||
`CREATE TABLE IF NOT EXISTS messages (
|
||||
id SERIAL PRIMARY KEY,
|
||||
session_id VARCHAR(64) REFERENCES sessions(id) ON DELETE CASCADE,
|
||||
role VARCHAR(16) NOT NULL,
|
||||
content TEXT NOT NULL,
|
||||
created_at TIMESTAMP DEFAULT NOW()
|
||||
)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_messages_session_id ON messages(session_id)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_messages_created_at ON messages(session_id, created_at)`,
|
||||
}
|
||||
|
||||
for _, q := range queries {
|
||||
if _, err := s.db.Exec(q); err != nil {
|
||||
return fmt.Errorf("迁移SQL执行失败: %w\nSQL: %s", err, q)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateSession 创建新会话
|
||||
func (s *SessionStore) CreateSession(userID, sessionID, title string, isMain bool) error {
|
||||
now := time.Now()
|
||||
_, err := s.db.Exec(
|
||||
`INSERT INTO sessions (id, user_id, title, is_main, created_at, updated_at)
|
||||
VALUES ($1, $2, $3, $4, $5, $5)
|
||||
ON CONFLICT (id) DO UPDATE SET updated_at = $5`,
|
||||
sessionID, userID, title, isMain, now,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("创建会话失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetUserSessions 获取用户的所有会话(按 updated_at DESC 排序)
|
||||
func (s *SessionStore) GetUserSessions(userID string) ([]Session, error) {
|
||||
rows, err := s.db.Query(
|
||||
`SELECT id, user_id, title, is_main, created_at, updated_at
|
||||
FROM sessions WHERE user_id = $1
|
||||
ORDER BY updated_at DESC`,
|
||||
userID,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("查询用户会话失败: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var sessions []Session
|
||||
for rows.Next() {
|
||||
var sess Session
|
||||
if err := rows.Scan(&sess.ID, &sess.UserID, &sess.Title, &sess.IsMain, &sess.CreatedAt, &sess.UpdatedAt); err != nil {
|
||||
return nil, fmt.Errorf("扫描会话行失败: %w", err)
|
||||
}
|
||||
sessions = append(sessions, sess)
|
||||
}
|
||||
|
||||
if sessions == nil {
|
||||
sessions = []Session{}
|
||||
}
|
||||
return sessions, rows.Err()
|
||||
}
|
||||
|
||||
// GetSession 获取单个会话
|
||||
func (s *SessionStore) GetSession(sessionID string) (*Session, error) {
|
||||
var sess Session
|
||||
err := s.db.QueryRow(
|
||||
`SELECT id, user_id, title, is_main, created_at, updated_at
|
||||
FROM sessions WHERE id = $1`,
|
||||
sessionID,
|
||||
).Scan(&sess.ID, &sess.UserID, &sess.Title, &sess.IsMain, &sess.CreatedAt, &sess.UpdatedAt)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, nil
|
||||
}
|
||||
return nil, fmt.Errorf("查询会话失败: %w", err)
|
||||
}
|
||||
return &sess, nil
|
||||
}
|
||||
|
||||
// UpdateSessionTitle 更新会话标题
|
||||
func (s *SessionStore) UpdateSessionTitle(sessionID, title string) error {
|
||||
_, err := s.db.Exec(
|
||||
`UPDATE sessions SET title = $1, updated_at = NOW() WHERE id = $2`,
|
||||
title, sessionID,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("更新会话标题失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// UpdateSessionTime 更新会话的 updated_at 时间戳
|
||||
func (s *SessionStore) UpdateSessionTime(sessionID string) error {
|
||||
_, err := s.db.Exec(
|
||||
`UPDATE sessions SET updated_at = NOW() WHERE id = $1`,
|
||||
sessionID,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("更新会话时间失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteSession 删除会话(级联删除消息,但不删除记忆)
|
||||
func (s *SessionStore) DeleteSession(sessionID string) error {
|
||||
_, err := s.db.Exec(`DELETE FROM sessions WHERE id = $1`, sessionID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("删除会话失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteAllUserSessions 删除用户的所有会话(但不删除记忆)
|
||||
func (s *SessionStore) DeleteAllUserSessions(userID string) error {
|
||||
_, err := s.db.Exec(`DELETE FROM sessions WHERE user_id = $1`, userID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("删除用户所有会话失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AddMessage 添加一条消息到会话
|
||||
func (s *SessionStore) AddMessage(sessionID, role, content string) error {
|
||||
_, err := s.db.Exec(
|
||||
`INSERT INTO messages (session_id, role, content) VALUES ($1, $2, $3)`,
|
||||
sessionID, role, content,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("添加消息失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetMessages 获取会话的消息列表(按时间正序)
|
||||
func (s *SessionStore) GetMessages(sessionID string, limit int) ([]Message, error) {
|
||||
if limit <= 0 {
|
||||
limit = 50
|
||||
}
|
||||
|
||||
rows, err := s.db.Query(
|
||||
`SELECT id, session_id, role, content, created_at
|
||||
FROM messages WHERE session_id = $1
|
||||
ORDER BY created_at ASC
|
||||
LIMIT $2`,
|
||||
sessionID, limit,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("查询消息失败: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var messages []Message
|
||||
for rows.Next() {
|
||||
var msg Message
|
||||
if err := rows.Scan(&msg.ID, &msg.SessionID, &msg.Role, &msg.Content, &msg.CreatedAt); err != nil {
|
||||
return nil, fmt.Errorf("扫描消息行失败: %w", err)
|
||||
}
|
||||
messages = append(messages, msg)
|
||||
}
|
||||
|
||||
if messages == nil {
|
||||
messages = []Message{}
|
||||
}
|
||||
return messages, rows.Err()
|
||||
}
|
||||
|
||||
// ClearSessionMessages 清空会话的所有消息但不删除会话本身
|
||||
func (s *SessionStore) ClearSessionMessages(sessionID string) error {
|
||||
_, err := s.db.Exec(`DELETE FROM messages WHERE session_id = $1`, sessionID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("清空会话消息失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close 关闭数据库连接
|
||||
func (s *SessionStore) Close() error {
|
||||
if s.db != nil {
|
||||
return s.db.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// IsAvailable 检查存储是否可用(数据库连接正常)
|
||||
func (s *SessionStore) IsAvailable() bool {
|
||||
if s.db == nil {
|
||||
return false
|
||||
}
|
||||
return s.db.Ping() == nil
|
||||
}
|
||||
Reference in New Issue
Block a user