From 09eb6fb4d713441b8d16cdd4a4ef854d69f67b6f Mon Sep 17 00:00:00 2001 From: Claudio Costa Date: Tue, 22 Feb 2022 13:43:38 +0100 Subject: [PATCH] Allow sending intra cluster WS events reliably (#19594) --- app/web_hub.go | 3 ++- app/web_hub_test.go | 22 ++++++++++++++++++++++ model/websocket_message.go | 3 +++ 3 files changed, 27 insertions(+), 1 deletion(-) diff --git a/app/web_hub.go b/app/web_hub.go index 1364447dd0..bbaf4fc5e6 100644 --- a/app/web_hub.go +++ b/app/web_hub.go @@ -178,7 +178,8 @@ func (s *Server) Publish(message *model.WebSocketEvent) { message.EventType() == model.WebsocketEventPostEdited || message.EventType() == model.WebsocketEventDirectAdded || message.EventType() == model.WebsocketEventGroupAdded || - message.EventType() == model.WebsocketEventAddedToTeam { + message.EventType() == model.WebsocketEventAddedToTeam || + message.GetBroadcast().ReliableClusterSend { cm.SendType = model.ClusterSendReliable } diff --git a/app/web_hub_test.go b/app/web_hub_test.go index 630a10ce33..e7ef4f2b87 100644 --- a/app/web_hub_test.go +++ b/app/web_hub_test.go @@ -19,6 +19,7 @@ import ( "github.com/mattermost/mattermost-server/v6/model" "github.com/mattermost/mattermost-server/v6/shared/i18n" "github.com/mattermost/mattermost-server/v6/store/storetest/mocks" + "github.com/mattermost/mattermost-server/v6/testlib" ) func dummyWebsocketHandler(t *testing.T) http.HandlerFunc { @@ -336,6 +337,27 @@ func TestHubConnIndexInactive(t *testing.T) { assert.Len(t, connIndex.All(), 2) } +func TestReliableWebSocketSend(t *testing.T) { + th := Setup(t) + defer th.TearDown() + + testCluster := &testlib.FakeClusterInterface{} + th.Server.Cluster = testCluster + + ev := model.NewWebSocketEvent("test_reliable_event", "", "", "", nil) + ev = ev.SetBroadcast(&model.WebsocketBroadcast{}) + th.App.Publish(ev) + ev = ev.SetBroadcast(&model.WebsocketBroadcast{ + ReliableClusterSend: true, + }) + th.App.Publish(ev) + + messages := testCluster.GetMessages() + require.Len(t, messages, 2) + require.Equal(t, model.ClusterSendBestEffort, messages[0].SendType) + require.Equal(t, model.ClusterSendReliable, messages[1].SendType) +} + func TestHubIsRegistered(t *testing.T) { th := Setup(t).InitBasic() defer th.TearDown() diff --git a/model/websocket_message.go b/model/websocket_message.go index 53d8b4c464..38e42bb2e9 100644 --- a/model/websocket_message.go +++ b/model/websocket_message.go @@ -90,6 +90,9 @@ type WebsocketBroadcast struct { TeamId string `json:"team_id"` // broadcast only occurs for users in this team ContainsSanitizedData bool `json:"-"` ContainsSensitiveData bool `json:"-"` + // ReliableClusterSend indicates whether or not the message should + // be sent through the cluster using the reliable, TCP backed channel. + ReliableClusterSend bool `json:"-"` } func (wb *WebsocketBroadcast) copy() *WebsocketBroadcast {