MM-29979: make websocket writes zero-alloc (#16115)
This reverts commit 31dac80e66.
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
fb9c2faa1a
Коммит
4758408699
@@ -4,6 +4,8 @@
|
|||||||
package app
|
package app
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
@@ -169,6 +171,11 @@ func (wc *WebConn) writePump() {
|
|||||||
wc.WebSocket.Close()
|
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 {
|
for {
|
||||||
select {
|
select {
|
||||||
case msg, ok := <-wc.send:
|
case msg, ok := <-wc.send:
|
||||||
@@ -201,20 +208,25 @@ func (wc *WebConn) writePump() {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
var msgBytes []byte
|
buf.Reset()
|
||||||
|
var err error
|
||||||
if evtOk {
|
if evtOk {
|
||||||
cpyEvt := evt.SetSequence(wc.Sequence)
|
cpyEvt := evt.SetSequence(wc.Sequence)
|
||||||
msgBytes = []byte(cpyEvt.ToJson())
|
err = cpyEvt.Encode(enc)
|
||||||
wc.Sequence++
|
wc.Sequence++
|
||||||
} else {
|
} 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 {
|
if len(wc.send) >= sendFullWarn {
|
||||||
logData := []mlog.Field{
|
logData := []mlog.Field{
|
||||||
mlog.String("user_id", wc.UserId),
|
mlog.String("user_id", wc.UserId),
|
||||||
mlog.String("type", msg.EventType()),
|
mlog.String("type", msg.EventType()),
|
||||||
mlog.Int("size", len(msgBytes)),
|
mlog.Int("size", buf.Len()),
|
||||||
}
|
}
|
||||||
if evtOk {
|
if evtOk {
|
||||||
logData = append(logData, mlog.String("channel_id", evt.GetBroadcast().ChannelId))
|
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))
|
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)
|
wc.logSocketErr("websocket.send", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -203,6 +203,22 @@ func (ev *WebSocketEvent) ToJson() string {
|
|||||||
return string(b)
|
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 {
|
func WebSocketEventFromJson(data io.Reader) *WebSocketEvent {
|
||||||
var ev WebSocketEvent
|
var ev WebSocketEvent
|
||||||
var o webSocketEventJSON
|
var o webSocketEventJSON
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user