From d4d216e93e492c3c7c2b3d481e966e22841b6825 Mon Sep 17 00:00:00 2001 From: Jesse Hallam Date: Thu, 30 Jul 2026 09:48:59 -0300 Subject: [PATCH] Add ClusterInterface.Shutdown to surface skipped cluster sends (#37753) * Add ClusterInterface.Shutdown to surface skipped cluster sends * Defer cluster interface shutdown so it runs on every exit path --- server/channels/app/busy_test.go | 2 ++ server/channels/app/channel_guards_test.go | 2 ++ server/channels/app/platform/busy_test.go | 2 ++ server/channels/app/platform/config_test.go | 1 + server/channels/app/platform/service.go | 10 ++++++++++ server/channels/testlib/cluster.go | 2 ++ server/einterfaces/cluster.go | 1 + server/einterfaces/mocks/ClusterInterface.go | 5 +++++ 8 files changed, 25 insertions(+) diff --git a/server/channels/app/busy_test.go b/server/channels/app/busy_test.go index feefc09652b..c6bf8de2e59 100644 --- a/server/channels/app/busy_test.go +++ b/server/channels/app/busy_test.go @@ -176,3 +176,5 @@ func (c *ClusterMock) WebConnCountForUser(userID string) (int, *model.AppError) func (c *ClusterMock) GetWSQueues(userID, connectionID string, seqNum int64) (map[string]*model.WSQueues, error) { return nil, nil } + +func (c *ClusterMock) Shutdown() {} diff --git a/server/channels/app/channel_guards_test.go b/server/channels/app/channel_guards_test.go index cee183000e9..13931b35d82 100644 --- a/server/channels/app/channel_guards_test.go +++ b/server/channels/app/channel_guards_test.go @@ -86,6 +86,8 @@ func (c *captureClusterMock) GetWSQueues(userID, connectionID string, seqNum int return nil, nil } +func (c *captureClusterMock) Shutdown() {} + func TestChannelGuardCacheBroadcastShape(t *testing.T) { mainHelper.Parallel(t) cluster := &captureClusterMock{} diff --git a/server/channels/app/platform/busy_test.go b/server/channels/app/platform/busy_test.go index ac757703dab..ef09e022a8e 100644 --- a/server/channels/app/platform/busy_test.go +++ b/server/channels/app/platform/busy_test.go @@ -177,3 +177,5 @@ func (c *ClusterMock) WebConnCountForUser(userID string) (int, *model.AppError) func (c *ClusterMock) GetWSQueues(userID, connectionID string, seqNum int64) (map[string]*model.WSQueues, error) { return nil, nil } + +func (c *ClusterMock) Shutdown() {} diff --git a/server/channels/app/platform/config_test.go b/server/channels/app/platform/config_test.go index 7aa584796a9..b3123fabf25 100644 --- a/server/channels/app/platform/config_test.go +++ b/server/channels/app/platform/config_test.go @@ -56,6 +56,7 @@ func TestConfigSave(t *testing.T) { mainHelper.Parallel(t) cm := &mocks.ClusterInterface{} cm.On("SendClusterMessage", mock.AnythingOfType("*model.ClusterMessage")).Return(nil) + cm.On("Shutdown").Return() th := SetupWithCluster(t, cm) t.Run("trigger a config changed event for the cluster", func(t *testing.T) { diff --git a/server/channels/app/platform/service.go b/server/channels/app/platform/service.go index 734ff01f2f6..54b71f372f8 100644 --- a/server/channels/app/platform/service.go +++ b/server/channels/app/platform/service.go @@ -573,6 +573,16 @@ func (ps *PlatformService) TotalWebsocketConnections() int { } func (ps *PlatformService) Shutdown() error { + // Deferred so it still runs even if a later step below returns early + // (e.g. cacheProvider.Close failing). Must run last: every other step + // below (notably HubStop) is a candidate to still attempt a cluster send + // after StopInterNodeCommunication has already run, so the cluster + // interface's own shutdown accounting needs everything below to have + // already happened, which defer guarantees regardless of the early return. + if ps.clusterIFace != nil { + defer ps.clusterIFace.Shutdown() + } + ps.HubStop() // Shutdown status processor. diff --git a/server/channels/testlib/cluster.go b/server/channels/testlib/cluster.go index daec0482681..3690f90fa70 100644 --- a/server/channels/testlib/cluster.go +++ b/server/channels/testlib/cluster.go @@ -116,3 +116,5 @@ func (c *FakeClusterInterface) WebConnCountForUser(userID string) (int, *model.A func (c *FakeClusterInterface) GetWSQueues(userID, connectionID string, seqNum int64) (map[string]*model.WSQueues, error) { return nil, nil } + +func (c *FakeClusterInterface) Shutdown() {} diff --git a/server/einterfaces/cluster.go b/server/einterfaces/cluster.go index 8ca113c8c97..dfff8ee71e2 100644 --- a/server/einterfaces/cluster.go +++ b/server/einterfaces/cluster.go @@ -24,6 +24,7 @@ type ClusterInterface interface { GetClusterInfos() ([]*model.ClusterInfo, error) SendClusterMessage(msg *model.ClusterMessage) SendClusterMessageToNode(nodeID string, msg *model.ClusterMessage) error + Shutdown() NotifyMsg(buf []byte) GetClusterStats(rctx request.CTX) ([]*model.ClusterStats, *model.AppError) GetLogs(rctx request.CTX, page, perPage int) ([]string, *model.AppError) diff --git a/server/einterfaces/mocks/ClusterInterface.go b/server/einterfaces/mocks/ClusterInterface.go index f8aa96d40a1..640ee7dc24d 100644 --- a/server/einterfaces/mocks/ClusterInterface.go +++ b/server/einterfaces/mocks/ClusterInterface.go @@ -361,6 +361,11 @@ func (_m *ClusterInterface) SendClusterMessageToNode(nodeID string, msg *model.C return r0 } +// Shutdown provides a mock function with no fields +func (_m *ClusterInterface) Shutdown() { + _m.Called() +} + // StartInterNodeCommunication provides a mock function with no fields func (_m *ClusterInterface) StartInterNodeCommunication() { _m.Called()