26a61cb57c
## 🐛 Bug 修复 - 修复前端对话无响应:消除 ChatContainer 中的双重 WebSocket 连接,优化 sendMessage 失败提示 - 修复 Memory-Service 数据库迁移失败:ai-core 和 memory-service 均添加 ALTER TABLE ADD COLUMN IF NOT EXISTS 模式演化 - 修复语音/STT 不可用:添加 MediaRecorder API 降级方案,修复 whisper-cli 输出文件名错误 - 修复仪表盘数据库按钮失效:补充按钮 ID 属性,重写 controlDB() 控制逻辑 ## 🎨 UI 修复 - 修正用户消息头像位置:从 flex-row-reverse 改为 justify-end - 移除空聊天列表的 emoji 占位图标 ## ✨ 新功能 - devtools 新增 STT 处理日志面板(环形缓冲区 + WebSocket 广播 + 可视化表格) - 新增 ADMIN_NICKNAME 环境变量,支持自定义管理员昵称 ## 🔧 改进 - 注册流程增加昵称必填字段(前后端同步) ## 🏗️ 架构重构 - 重构自主思考逻辑:从定时器轮询改为事件驱动(对话后触发 + 静默检测),优化提示词使其更自然人性化 - 实现主-子会话架构:新增 4 种子会话类型(general/memory/iot/knowledge),意图分析 → 并行分发 → 结果合成流程 ## 📄 新增文档 - docs/architecture/main-session-sub-session-design.md — 子会话架构设计文档
141 lines
3.0 KiB
Go
141 lines
3.0 KiB
Go
package ws
|
|
|
|
import (
|
|
"encoding/json"
|
|
"log"
|
|
"time"
|
|
|
|
"github.com/gorilla/websocket"
|
|
)
|
|
|
|
const (
|
|
// 写入超时
|
|
writeWait = 10 * time.Second
|
|
|
|
// 读取pong超时
|
|
pongWait = 60 * time.Second
|
|
|
|
// pong发送后等待下一次ping的间隔
|
|
pingPeriod = (pongWait * 9) / 10
|
|
|
|
// 最大消息大小
|
|
maxMessageSize = 65536
|
|
)
|
|
|
|
// Client WebSocket客户端
|
|
type Client struct {
|
|
Hub *Hub
|
|
Conn *websocket.Conn
|
|
Send chan []byte
|
|
UserID string
|
|
SessionID string
|
|
}
|
|
|
|
// NewClient 创建WebSocket客户端
|
|
func NewClient(hub *Hub, conn *websocket.Conn, userID, sessionID string) *Client {
|
|
return &Client{
|
|
Hub: hub,
|
|
Conn: conn,
|
|
Send: make(chan []byte, 256),
|
|
UserID: userID,
|
|
SessionID: sessionID,
|
|
}
|
|
}
|
|
|
|
// ReadPump 读取协程 —— 从WebSocket连接读取消息
|
|
func (c *Client) ReadPump(onMessage func(client *Client, msg ClientMessage)) {
|
|
defer func() {
|
|
c.Hub.unregister <- c
|
|
c.Conn.Close()
|
|
}()
|
|
|
|
c.Conn.SetReadLimit(maxMessageSize)
|
|
c.Conn.SetReadDeadline(time.Now().Add(pongWait))
|
|
c.Conn.SetPongHandler(func(string) error {
|
|
c.Conn.SetReadDeadline(time.Now().Add(pongWait))
|
|
return nil
|
|
})
|
|
|
|
for {
|
|
_, rawMessage, err := c.Conn.ReadMessage()
|
|
if err != nil {
|
|
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseNormalClosure, websocket.CloseAbnormalClosure) {
|
|
log.Printf("[WS] 读取错误: user=%s err=%v", c.UserID, err)
|
|
}
|
|
break
|
|
}
|
|
|
|
// 解析消息
|
|
var msg ClientMessage
|
|
if err := json.Unmarshal(rawMessage, &msg); err != nil {
|
|
log.Printf("[WS] 消息解析失败: user=%s err=%v", c.UserID, err)
|
|
continue
|
|
}
|
|
|
|
// 处理ping
|
|
if msg.Type == "ping" {
|
|
pongMsg := ServerMessage{
|
|
Type: "pong",
|
|
Timestamp: time.Now().UnixMilli(),
|
|
}
|
|
data, _ := json.Marshal(pongMsg)
|
|
c.Send <- data
|
|
continue
|
|
}
|
|
|
|
// 调用消息处理器
|
|
if onMessage != nil {
|
|
onMessage(c, msg)
|
|
}
|
|
}
|
|
}
|
|
|
|
// WritePump 写入协程 —— 向WebSocket连接写入消息
|
|
func (c *Client) WritePump() {
|
|
ticker := time.NewTicker(pingPeriod)
|
|
defer func() {
|
|
ticker.Stop()
|
|
c.Conn.Close()
|
|
}()
|
|
|
|
for {
|
|
select {
|
|
case message, ok := <-c.Send:
|
|
c.Conn.SetWriteDeadline(time.Now().Add(writeWait))
|
|
if !ok {
|
|
// Hub关闭了通道
|
|
c.Conn.WriteMessage(websocket.CloseMessage, []byte{})
|
|
return
|
|
}
|
|
|
|
if err := c.Conn.WriteMessage(websocket.TextMessage, message); err != nil {
|
|
log.Printf("[WS] 写入错误: user=%s err=%v", c.UserID, err)
|
|
return
|
|
}
|
|
|
|
case <-ticker.C:
|
|
c.Conn.SetWriteDeadline(time.Now().Add(writeWait))
|
|
if err := c.Conn.WriteMessage(websocket.PingMessage, nil); err != nil {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// SendMessage 向客户端发送消息
|
|
func (c *Client) SendMessage(msg ServerMessage) error {
|
|
data, err := json.Marshal(msg)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
select {
|
|
case c.Send <- data:
|
|
return nil
|
|
default:
|
|
// 通道满:记录警告并返回错误(避免静默丢弃)
|
|
log.Printf("[WS] 发送通道已满,丢弃消息: type=%s user=%s", msg.Type, c.UserID)
|
|
return nil
|
|
}
|
|
}
|