71f0a1abdb
- 所有Go模块路径从 github.com/yourname/cyrene-ai 迁移到 git.yeij.top/AskaEth/Cyrene - 5个Go Dockerfile添加 GOPROXY=https://goproxy.cn,direct 解决国内构建问题 - ai-core go.mod 添加 pkg/plugins replace 指令 - Caddyfile 简化为 http:// 通配 + handle 保留 /api 前缀 - ethend Dockerfile 适配 (npm install + 仅 COPY package.json) - ethend 新增 RUNNING_IN_DOCKER 环境变量,健康检查改用Docker服务名 - ethend 数据库状态检查支持Docker hostname (postgres/redis/qdrant/minio) - process-manager 新增 CONTAINER_SVC_MAP + Docker模式自动检测 - 统一 docker-compose.dev.db.yml 卷名 (pg_data/redis_data/qdrant_data/minio_data) - docker-compose.yml ethend服务挂载docker.sock + 端口变量化 - 清理 .env 统一后的残留文件与提示信息 Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
152 lines
3.3 KiB
Go
152 lines
3.3 KiB
Go
package ws
|
|
|
|
import (
|
|
"encoding/json"
|
|
"git.yeij.top/AskaEth/Cyrene/pkg/logger"
|
|
"time"
|
|
|
|
"github.com/gorilla/websocket"
|
|
)
|
|
|
|
const (
|
|
// 写入超时
|
|
writeWait = 10 * time.Second
|
|
|
|
// 读取pong超时
|
|
pongWait = 60 * time.Second
|
|
|
|
// pong发送后等待下一次ping的间隔
|
|
pingPeriod = (pongWait * 9) / 10
|
|
|
|
// 最大消息大小
|
|
maxMessageSize = 2 * 1024 * 1024 // 2MB 支持语音消息
|
|
)
|
|
|
|
// Client WebSocket客户端
|
|
type Client struct {
|
|
Hub *Hub
|
|
Conn *websocket.Conn
|
|
Send chan []byte
|
|
UserID string
|
|
SessionID string
|
|
ClientID string
|
|
DeviceName string
|
|
UserAgent string
|
|
}
|
|
|
|
// NewClient 创建WebSocket客户端
|
|
func NewClient(hub *Hub, conn *websocket.Conn, userID, sessionID, clientID, deviceName, userAgent string) *Client {
|
|
return &Client{
|
|
Hub: hub,
|
|
Conn: conn,
|
|
Send: make(chan []byte, 256),
|
|
UserID: userID,
|
|
SessionID: sessionID,
|
|
ClientID: clientID,
|
|
DeviceName: deviceName,
|
|
UserAgent: userAgent,
|
|
}
|
|
}
|
|
|
|
// 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) {
|
|
logger.Printf("[WS] 读取错误: user=%s err=%v", c.UserID, err)
|
|
}
|
|
break
|
|
}
|
|
|
|
// 解析消息
|
|
var msg ClientMessage
|
|
if err := json.Unmarshal(rawMessage, &msg); err != nil {
|
|
logger.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 {
|
|
logger.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
|
|
}
|
|
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
logger.Printf("[WS] 发送消息时连接已关闭: type=%s user=%s", msg.Type, c.UserID)
|
|
}
|
|
}()
|
|
|
|
select {
|
|
case c.Send <- data:
|
|
return nil
|
|
default:
|
|
logger.Printf("[WS] 发送通道已满,丢弃消息: type=%s user=%s", msg.Type, c.UserID)
|
|
return nil
|
|
}
|
|
}
|