MM-25516: Changed to byte slice instead of string for cluster messages (#17998)

* MM-25516: Changed to byte slice instead of string for cluster messages

https://mattermost.atlassian.net/browse/MM-25116

Testing:
Manually tested.
Load-tested with Cluster Controller.

I looked into changing the serialization method to use msgpack,
but the ClusterMessage struct was mainly used for only 3 fields
which didn't lead to much of a CPU time improvement, whereas
actually led to more allocations using msgpack. Hence, I chose
to remain with JSON.

```
name              old time/op    new time/op    delta
ClusterMarshal-8    3.51µs ± 1%    3.10µs ± 2%  -11.59%  (p=0.000 n=9+10)

name              old alloc/op   new alloc/op   delta
ClusterMarshal-8      776B ± 0%     1000B ± 0%  +28.87%  (p=0.000 n=10+10)

name              old allocs/op  new allocs/op  delta
ClusterMarshal-8      12.0 ± 0%      13.0 ± 0%   +8.33%  (p=0.000 n=10+10)
```

```release-note
Changed the field type of Data in model.ClusterMessage to []byte from string.
```

* Trigger CI
```release-note
NONE
```
Этот коммит содержится в:
Agniva De Sarker
2021-07-26 13:41:20 +05:30
коммит произвёл GitHub
родитель 23800326a0
Коммит 7be61af24f
28 изменённых файлов: 112 добавлений и 130 удалений

Просмотреть файл

@@ -4,6 +4,7 @@
package app
import (
"encoding/json"
"sync"
"sync/atomic"
"time"
@@ -109,11 +110,12 @@ func (b *Busy) notifyServerBusyChange(sbs *model.ServerBusyState) {
if b.cluster == nil {
return
}
buf, _ := json.Marshal(sbs)
msg := &model.ClusterMessage{
Event: model.ClusterEventBusyStateChanged,
SendType: model.ClusterSendReliable,
WaitForAllToSend: true,
Data: sbs.ToJson(),
Data: buf,
}
b.cluster.SendClusterMessage(msg)
}

Просмотреть файл

@@ -4,7 +4,7 @@
package app
import (
"strings"
"bytes"
"testing"
"time"
@@ -116,7 +116,7 @@ type ClusterMock struct {
}
func (c *ClusterMock) SendClusterMessage(msg *model.ClusterMessage) {
sbs := model.ServerBusyStateFromJson(strings.NewReader(msg.Data))
sbs := model.ServerBusyStateFromJson(bytes.NewReader(msg.Data))
c.Busy.ClusterEventChanged(sbs)
}

Просмотреть файл

@@ -4,7 +4,7 @@
package app
import (
"strings"
"bytes"
"github.com/mattermost/mattermost-server/v6/model"
"github.com/mattermost/mattermost-server/v6/plugin"
@@ -12,11 +12,11 @@ import (
)
func (s *Server) clusterInstallPluginHandler(msg *model.ClusterMessage) {
s.installPluginFromData(model.PluginEventDataFromJson(strings.NewReader(msg.Data)))
s.installPluginFromData(model.PluginEventDataFromJson(bytes.NewReader(msg.Data)))
}
func (s *Server) clusterRemovePluginHandler(msg *model.ClusterMessage) {
s.removePluginFromData(model.PluginEventDataFromJson(strings.NewReader(msg.Data)))
s.removePluginFromData(model.PluginEventDataFromJson(bytes.NewReader(msg.Data)))
}
func (s *Server) clusterPluginEventHandler(msg *model.ClusterMessage) {
@@ -44,7 +44,7 @@ func (s *Server) clusterPluginEventHandler(msg *model.ClusterMessage) {
hooks.OnPluginClusterEvent(&plugin.Context{}, model.PluginClusterEvent{
Id: eventID,
Data: []byte(msg.Data),
Data: msg.Data,
})
}
@@ -69,7 +69,7 @@ func (s *Server) registerClusterHandlers() {
}
func (s *Server) clusterPublishHandler(msg *model.ClusterMessage) {
event := model.WebSocketEventFromJson(strings.NewReader(msg.Data))
event := model.WebSocketEventFromJson(bytes.NewReader(msg.Data))
if event == nil {
return
}
@@ -77,7 +77,7 @@ func (s *Server) clusterPublishHandler(msg *model.ClusterMessage) {
}
func (s *Server) clusterUpdateStatusHandler(msg *model.ClusterMessage) {
status := model.StatusFromJson(strings.NewReader(msg.Data))
status := model.StatusFromJson(bytes.NewReader(msg.Data))
s.statusCache.Set(status.UserId, status)
}
@@ -86,7 +86,7 @@ func (s *Server) clusterInvalidateAllCachesHandler(msg *model.ClusterMessage) {
}
func (s *Server) clusterInvalidateCacheForChannelMembersNotifyPropHandler(msg *model.ClusterMessage) {
s.invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(msg.Data)
s.invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(string(msg.Data))
}
func (s *Server) clusterInvalidateCacheForChannelByNameHandler(msg *model.ClusterMessage) {
@@ -94,11 +94,11 @@ func (s *Server) clusterInvalidateCacheForChannelByNameHandler(msg *model.Cluste
}
func (s *Server) clusterInvalidateCacheForUserHandler(msg *model.ClusterMessage) {
s.invalidateCacheForUserSkipClusterSend(msg.Data)
s.invalidateCacheForUserSkipClusterSend(string(msg.Data))
}
func (s *Server) clusterInvalidateCacheForUserTeamsHandler(msg *model.ClusterMessage) {
s.invalidateWebConnSessionCacheForUser(msg.Data)
s.invalidateWebConnSessionCacheForUser(string(msg.Data))
}
func (s *Server) clearSessionCacheForUserSkipClusterSend(userID string) {
@@ -112,7 +112,7 @@ func (s *Server) clearSessionCacheForAllUsersSkipClusterSend() {
}
func (s *Server) clusterClearSessionCacheForUserHandler(msg *model.ClusterMessage) {
s.clearSessionCacheForUserSkipClusterSend(msg.Data)
s.clearSessionCacheForUserSkipClusterSend(string(msg.Data))
}
func (s *Server) clusterClearSessionCacheForAllUsersHandler(msg *model.ClusterMessage) {
@@ -120,7 +120,7 @@ func (s *Server) clusterClearSessionCacheForAllUsersHandler(msg *model.ClusterMe
}
func (s *Server) clusterBusyStateChgHandler(msg *model.ClusterMessage) {
s.serverBusyStateChanged(model.ServerBusyStateFromJson(strings.NewReader(msg.Data)))
s.serverBusyStateChanged(model.ServerBusyStateFromJson(bytes.NewReader(msg.Data)))
}
func (s *Server) invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(channelID string) {

Просмотреть файл

@@ -1139,7 +1139,7 @@ func (api *PluginAPI) PublishPluginClusterEvent(ev model.PluginClusterEvent,
"PluginID": api.id,
"EventID": ev.Id,
},
Data: string(ev.Data),
Data: ev.Data,
}
// If TargetId is empty we broadcast to all other cluster nodes.

Просмотреть файл

@@ -34,7 +34,7 @@ func (p *Plugin) ServeHTTP(_ *plugin.Context, w http.ResponseWriter, r *http.Req
}
req := model.WebSocketRequestFromJson(bytes.NewReader(msg))
resp := model.NewWebSocketResponse("OK", req.Seq, map[string]interface{}{"action": req.Action, "value": req.Data["value"]})
if err = ws.WriteMessage(mt, []byte(resp.ToJson())); err != nil {
if err = ws.WriteMessage(mt, resp.ToJson()); err != nil {
break
}
}

Просмотреть файл

@@ -4,16 +4,19 @@
package app
import (
"encoding/json"
"github.com/mattermost/mattermost-server/v6/model"
)
func (s *Server) notifyClusterPluginEvent(event model.ClusterEvent, data model.PluginEventData) {
buf, _ := json.Marshal(data)
if s.Cluster != nil {
s.Cluster.SendClusterMessage(&model.ClusterMessage{
Event: event,
SendType: model.ClusterSendReliable,
WaitForAllToSend: true,
Data: data.ToJson(),
Data: buf,
})
}
}

Просмотреть файл

@@ -4,7 +4,7 @@
package slashcommands
import (
"strings"
"bytes"
"testing"
"github.com/stretchr/testify/assert"
@@ -50,7 +50,7 @@ func TestShareProviderDoCommand(t *testing.T) {
require.Equal(t, "##### "+args.T("api.command_share.channel_shared"), response.Text)
channelConvertedMessages := testCluster.SelectMessages(func(msg *model.ClusterMessage) bool {
event := model.WebSocketEventFromJson(strings.NewReader(msg.Data))
event := model.WebSocketEventFromJson(bytes.NewReader(msg.Data))
return event != nil && event.EventType() == model.WebsocketEventChannelConverted
})
assert.Len(t, channelConvertedMessages, 1)
@@ -85,7 +85,7 @@ func TestShareProviderDoCommand(t *testing.T) {
require.Equal(t, "##### "+args.T("api.command_share.shared_channel_unavailable"), response.Text)
channelConvertedMessages := testCluster.SelectMessages(func(msg *model.ClusterMessage) bool {
event := model.WebSocketEventFromJson(strings.NewReader(msg.Data))
event := model.WebSocketEventFromJson(bytes.NewReader(msg.Data))
return event != nil && event.EventType() == model.WebsocketEventChannelConverted
})
require.Len(t, channelConvertedMessages, 1)

Просмотреть файл

@@ -24,7 +24,7 @@ func (a *App) AddStatusCache(status *model.Status) {
msg := &model.ClusterMessage{
Event: model.ClusterEventUpdateStatus,
SendType: model.ClusterSendBestEffort,
Data: status.ToClusterJson(),
Data: []byte(status.ToClusterJson()),
}
a.Cluster().SendClusterMessage(msg)
}

Просмотреть файл

@@ -241,7 +241,7 @@ func (a *App) invalidateCacheForChannelMembersNotifyProps(channelID string) {
msg := &model.ClusterMessage{
Event: model.ClusterEventInvalidateCacheForChannelMembersNotifyProps,
SendType: model.ClusterSendBestEffort,
Data: channelID,
Data: []byte(channelID),
}
a.Cluster().SendClusterMessage(msg)
}
@@ -266,7 +266,7 @@ func (a *App) invalidateCacheForUserTeams(userID string) {
msg := &model.ClusterMessage{
Event: model.ClusterEventInvalidateCacheForUserTeams,
SendType: model.ClusterSendBestEffort,
Data: userID,
Data: []byte(userID),
}
a.Cluster().SendClusterMessage(msg)
}