Websocket: split welcome message (#27967)
This commit is contained in:
parent
625fe4d4b9
commit
a0dccea9a3
1 changed files with 25 additions and 7 deletions
|
|
@ -52,10 +52,11 @@ func (h *SocketHub) ServeWebsocket(w http.ResponseWriter, r *http.Request) {
|
|||
InsecureSkipVerify: true,
|
||||
}
|
||||
|
||||
// https://github.com/nhooyr/websocket/issues/218
|
||||
// Safari deflate message compression is broken, enable for others
|
||||
// see: https://github.com/gorilla/websocket/issues/731
|
||||
ua := strings.ToLower(r.Header.Get("User-Agent"))
|
||||
if strings.Contains(ua, "safari") && !strings.Contains(ua, "chrome") && !strings.Contains(ua, "android") {
|
||||
acceptOptions.CompressionMode = websocket.CompressionDisabled
|
||||
if strings.Contains(ua, "chrome") || strings.Contains(ua, "firefox") {
|
||||
acceptOptions.CompressionMode = websocket.CompressionContextTakeover
|
||||
}
|
||||
|
||||
conn, err := websocket.Accept(w, r, acceptOptions)
|
||||
|
|
@ -88,6 +89,7 @@ func (h *SocketHub) subscribe(ctx context.Context, conn *websocket.Conn) error {
|
|||
select {
|
||||
case msg := <-s.send:
|
||||
if err := writeTimeout(ctx, socketWriteTimeout, conn, msg); err != nil {
|
||||
log.TRACE.Printf("ws write error: %v, len: %.1fkB, msg: %s", err, float64(len(msg))/1024, msg[:min(len(msg), 20)])
|
||||
return err
|
||||
}
|
||||
case <-ctx.Done():
|
||||
|
|
@ -112,6 +114,7 @@ func (h *SocketHub) deleteSubscriber(s *socketSubscriber) {
|
|||
|
||||
func (h *SocketHub) welcome(subscriber *socketSubscriber, params []util.Param) {
|
||||
msg := make(map[string]json.RawMessage, len(params))
|
||||
sharders := make(map[string]util.Sharder)
|
||||
|
||||
for _, p := range params {
|
||||
k := p.Key
|
||||
|
|
@ -119,13 +122,28 @@ func (h *SocketHub) welcome(subscriber *socketSubscriber, params []util.Param) {
|
|||
k = "loadpoints." + p.UniqueID()
|
||||
}
|
||||
|
||||
msg[k] = json.RawMessage(socketEncode(p.Val))
|
||||
// Sharder values are split into shards and sent as a separate message
|
||||
if sharder, ok := (p.Val).(util.Sharder); ok {
|
||||
sharders[k] = sharder
|
||||
} else {
|
||||
msg[k] = json.RawMessage(socketEncode(p.Val))
|
||||
}
|
||||
}
|
||||
|
||||
b, _ := json.Marshal(msg)
|
||||
// send compact state first, then potentially large sharded data
|
||||
if b, err := json.Marshal(msg); err == nil {
|
||||
subscriber.send <- b
|
||||
}
|
||||
|
||||
// should not block
|
||||
subscriber.send <- b
|
||||
for k, sharder := range sharders {
|
||||
for key, val := range sharder.AllShards() {
|
||||
if b, err := json.Marshal(map[string]json.RawMessage{
|
||||
k + "." + key: json.RawMessage(socketEncode(val)),
|
||||
}); err == nil {
|
||||
subscriber.send <- b
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (h *SocketHub) broadcast(p util.Param) {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue