diff --git a/app/web_conn.go b/app/web_conn.go index 6082751a4d..363e7c310e 100644 --- a/app/web_conn.go +++ b/app/web_conn.go @@ -4,6 +4,8 @@ package app import ( + "bytes" + "encoding/json" "fmt" "sync" "sync/atomic" @@ -169,6 +171,11 @@ func (wc *WebConn) writePump() { wc.WebSocket.Close() }() + var buf bytes.Buffer + // 2k is seen to be a good heuristic under which 98.5% of message sizes remain. + buf.Grow(1024 * 2) + enc := json.NewEncoder(&buf) + for { select { case msg, ok := <-wc.send: @@ -201,20 +208,25 @@ func (wc *WebConn) writePump() { continue } - var msgBytes []byte + buf.Reset() + var err error if evtOk { cpyEvt := evt.SetSequence(wc.Sequence) - msgBytes = []byte(cpyEvt.ToJson()) + err = cpyEvt.Encode(enc) wc.Sequence++ } else { - msgBytes = []byte(msg.ToJson()) + err = enc.Encode(msg) + } + if err != nil { + mlog.Warn("Error in encoding websocket message", mlog.Err(err)) + continue } if len(wc.send) >= sendFullWarn { logData := []mlog.Field{ mlog.String("user_id", wc.UserId), mlog.String("type", msg.EventType()), - mlog.Int("size", len(msgBytes)), + mlog.Int("size", buf.Len()), } if evtOk { logData = append(logData, mlog.String("channel_id", evt.GetBroadcast().ChannelId)) @@ -224,7 +236,7 @@ func (wc *WebConn) writePump() { } wc.WebSocket.SetWriteDeadline(time.Now().Add(writeWaitTime)) - if err := wc.WebSocket.WriteMessage(websocket.TextMessage, msgBytes); err != nil { + if err := wc.WebSocket.WriteMessage(websocket.TextMessage, buf.Bytes()); err != nil { wc.logSocketErr("websocket.send", err) return } diff --git a/model/websocket_message.go b/model/websocket_message.go index 8885ed4e5d..a4f92f80b2 100644 --- a/model/websocket_message.go +++ b/model/websocket_message.go @@ -203,6 +203,22 @@ func (ev *WebSocketEvent) ToJson() string { return string(b) } +// Encode encodes the event to the given encoder. +func (ev *WebSocketEvent) Encode(enc *json.Encoder) error { + if ev.precomputedJSON != nil { + return enc.Encode(json.RawMessage( + fmt.Sprintf(`{"event": %s, "data": %s, "broadcast": %s, "seq": %d}`, ev.precomputedJSON.Event, ev.precomputedJSON.Data, ev.precomputedJSON.Broadcast, ev.Sequence), + )) + } + + return enc.Encode(webSocketEventJSON{ + ev.Event, + ev.Data, + ev.Broadcast, + ev.Sequence, + }) +} + func WebSocketEventFromJson(data io.Reader) *WebSocketEvent { var ev WebSocketEvent var o webSocketEventJSON