// Package room 实现 WebSocket 联机对战:房间管理、斗地主/象棋对局流程、AI 补位与托管 package room import ( "encoding/json" "time" "github.com/gorilla/websocket" ) // Client 一个 WebSocket 连接对应的客户端 type Client struct { hub *Hub // 所属房间管理器 conn *websocket.Conn // 底层连接 send chan []byte // 待发送消息队列(写泵消费) userID int // 登录用户ID name string // 昵称 avatar string // 头像 } // wsMessage 客户端与服务端通用的消息信封 type wsMessage struct { Type string `json:"type"` // 消息类型 Data json.RawMessage `json:"data"` // 消息负载(各类型自定义) } // push 向该客户端推送一条消息(满队列时丢弃,防止慢客户端拖垮全局) // 注意:data 若是 json.RawMessage([]byte),必须走 wsMessage 结构体嵌入, // 否则 encoding/json 会把 []byte 编成 base64 字符串,饥荒世界快照会损坏。 func (c *Client) push(msgType string, data any) { var dataRaw json.RawMessage switch v := data.(type) { case nil: dataRaw = json.RawMessage("null") case json.RawMessage: if len(v) == 0 { dataRaw = json.RawMessage("null") } else { dataRaw = v } default: b, err := json.Marshal(v) if err != nil { return } dataRaw = b } payload, err := json.Marshal(wsMessage{Type: msgType, Data: dataRaw}) if err != nil { return } select { case c.send <- payload: default: } } // pushError 推送错误提示 func (c *Client) pushError(msg string) { c.push("error", map[string]string{"msg": msg}) } // readPump 读泵:持续读取客户端消息并分发处理,连接断开时做离线清理 func (c *Client) readPump() { defer func() { c.hub.onDisconnect(c) c.conn.Close() }() c.conn.SetReadLimit(4 * 1024 * 1024) // 4MB:饥荒全量世界快照(tiles+ents)远超默认 8KB // 心跳:60秒收不到任何数据判定断线(前端每30秒发 ping) c.conn.SetReadDeadline(time.Now().Add(60 * time.Second)) c.conn.SetPongHandler(func(string) error { c.conn.SetReadDeadline(time.Now().Add(60 * time.Second)) return nil }) for { _, raw, err := c.conn.ReadMessage() if err != nil { return } c.conn.SetReadDeadline(time.Now().Add(60 * time.Second)) var msg wsMessage if err := json.Unmarshal(raw, &msg); err != nil { c.pushError("消息格式错误") continue } // 客户端应用层心跳直接回应 if msg.Type == "ping" { c.push("pong", nil) continue } c.hub.dispatch(c, msg.Type, msg.Data) } } // writePump 写泵:消费发送队列写入连接,定期发协议层 ping 保活 func (c *Client) writePump() { ticker := time.NewTicker(25 * time.Second) defer func() { ticker.Stop() c.conn.Close() }() for { select { case payload, ok := <-c.send: if !ok { return } c.conn.SetWriteDeadline(time.Now().Add(10 * time.Second)) if err := c.conn.WriteMessage(websocket.TextMessage, payload); err != nil { return } case <-ticker.C: c.conn.SetWriteDeadline(time.Now().Add(10 * time.Second)) if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil { return } } } }