package ws import ( "time" "github.com/gorilla/websocket" "go.uber.org/zap" ) const ( writeWait = 10 * time.Second pongWait = 60 * time.Second pingPeriod = pongWait * 9 / 10 maxMessageSize = 8 * 1024 ) type Client struct { hub *Hub conn *websocket.Conn send chan OutgoingMessage room string clientID string } func NewClient(hub *Hub, conn *websocket.Conn, room, clientID string) *Client { return &Client{ hub: hub, conn: conn, send: make(chan OutgoingMessage, 64), room: room, clientID: clientID, } } func (c *Client) ReadPump() { defer func() { remaining := c.hub.Unregister(c) _ = c.conn.Close() c.hub.Broadcast(c.room, presenceMessage(c.room, c.clientID, "left", remaining)) }() c.conn.SetReadLimit(maxMessageSize) _ = c.conn.SetReadDeadline(time.Now().Add(pongWait)) c.conn.SetPongHandler(func(string) error { return c.conn.SetReadDeadline(time.Now().Add(pongWait)) }) for { var incoming IncomingMessage if err := c.conn.ReadJSON(&incoming); err != nil { if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) { zap.L().Warn("WebSocket读取失败", zap.String("client_id", c.clientID), zap.Error(err)) } return } if incoming.Type != EventMessage || incoming.Data == "" { c.send <- OutgoingMessage{Type: EventError, Data: "仅支持非空 message 事件", Timestamp: time.Now()} continue } c.hub.Broadcast(c.room, OutgoingMessage{ Type: EventMessage, Room: c.room, ClientID: c.clientID, Data: incoming.Data, Timestamp: time.Now(), }) } } 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 { _ = c.conn.WriteMessage(websocket.CloseMessage, []byte{}) return } if err := c.conn.WriteJSON(message); err != nil { return } case <-ticker.C: _ = c.conn.SetWriteDeadline(time.Now().Add(writeWait)) if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil { return } } } } func presenceMessage(room, clientID, action string, online int) OutgoingMessage { return OutgoingMessage{ Type: EventPresence, Room: room, ClientID: clientID, Data: action, Online: online, Timestamp: time.Now(), } }