99 lines
2.6 KiB
Go
99 lines
2.6 KiB
Go
// 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 向该客户端推送一条消息(满队列时丢弃,防止慢客户端拖垮全局)
|
||
func (c *Client) push(msgType string, data any) {
|
||
payload, _ := json.Marshal(map[string]any{"type": msgType, "data": data})
|
||
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(8 * 1024)
|
||
// 心跳: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
|
||
}
|
||
}
|
||
}
|
||
}
|