From 10655a2eb3c5bfeeaeb5314bd95a70b31cd341c2 Mon Sep 17 00:00:00 2001 From: Ibrahim Serdar Acikgoz Date: Thu, 6 Oct 2022 14:14:03 +0300 Subject: [PATCH] fix a data race in the platform service (#21242) --- app/platform/helper_test.go | 2 +- app/platform/service.go | 57 ++++++++++++++++++++---------------- app/platform/web_hub.go | 4 +-- app/platform/web_hub_test.go | 8 ++--- app/server.go | 4 +-- 5 files changed, 40 insertions(+), 35 deletions(-) diff --git a/app/platform/helper_test.go b/app/platform/helper_test.go index 20d92653ea..17e40627a3 100644 --- a/app/platform/helper_test.go +++ b/app/platform/helper_test.go @@ -181,7 +181,7 @@ func setupTestHelper(dbStore store.Store, enterprise bool, includeCacheLayer boo th.Service.SetLicense(nil) } - th.Service.HubStart(th.Suite) + th.Service.Start(th.Suite) return th } diff --git a/app/platform/service.go b/app/platform/service.go index 69cba002e0..261297b2fc 100644 --- a/app/platform/service.go +++ b/app/platform/service.go @@ -243,32 +243,6 @@ func New(sc ServiceConfig, options ...Option) (*PlatformService, error) { ps.Busy = NewBusy(ps.clusterIFace) - ps.configListenerId = ps.AddConfigListener(func(_, _ *model.Config) { - ps.regenerateClientConfig() - - message := model.NewWebSocketEvent(model.WebsocketEventConfigChanged, "", "", "", nil, "") - - message.Add("config", ps.ClientConfigWithComputed()) - ps.Go(func() { - ps.Publish(message) - }) - - if err = ps.ReconfigureLogger(); err != nil { - mlog.Error("Error re-configuring logging after config change", mlog.Err(err)) - return - } - }) - ps.licenseListenerId = ps.AddLicenseListener(func(oldLicense, newLicense *model.License) { - ps.regenerateClientConfig() - - message := model.NewWebSocketEvent(model.WebsocketEventLicenseChanged, "", "", "", nil, "") - message.Add("license", ps.GetSanitizedClientLicense()) - ps.Go(func() { - ps.Publish(message) - }) - - }) - // Enable developer settings if this is a "dev" build if model.BuildNumber == "dev" { ps.UpdateConfig(func(cfg *model.Config) { *cfg.ServiceSettings.EnableDeveloper = true }) @@ -296,6 +270,37 @@ func New(sc ServiceConfig, options ...Option) (*PlatformService, error) { return ps, nil } +func (ps *PlatformService) Start(suite SuiteIFace) error { + ps.hubStart(suite) + + ps.configListenerId = ps.AddConfigListener(func(_, _ *model.Config) { + ps.regenerateClientConfig() + + message := model.NewWebSocketEvent(model.WebsocketEventConfigChanged, "", "", "", nil, "") + + message.Add("config", ps.ClientConfigWithComputed()) + ps.Go(func() { + ps.Publish(message) + }) + + if err := ps.ReconfigureLogger(); err != nil { + mlog.Error("Error re-configuring logging after config change", mlog.Err(err)) + return + } + }) + ps.licenseListenerId = ps.AddLicenseListener(func(oldLicense, newLicense *model.License) { + ps.regenerateClientConfig() + + message := model.NewWebSocketEvent(model.WebsocketEventLicenseChanged, "", "", "", nil, "") + message.Add("license", ps.GetSanitizedClientLicense()) + ps.Go(func() { + ps.Publish(message) + }) + + }) + return nil +} + func (ps *PlatformService) ShutdownMetrics() error { if ps.metrics != nil { return ps.metrics.stopMetricsServer() diff --git a/app/platform/web_hub.go b/app/platform/web_hub.go index 6352be02f8..3a0d5fd6bb 100644 --- a/app/platform/web_hub.go +++ b/app/platform/web_hub.go @@ -94,8 +94,8 @@ func newWebHub(ps *PlatformService) *Hub { } } -// HubStart starts all the hubs. -func (ps *PlatformService) HubStart(suite SuiteIFace) { +// hubStart starts all the hubs. +func (ps *PlatformService) hubStart(suite SuiteIFace) { // Total number of hubs is twice the number of CPUs. numberOfHubs := runtime.NumCPU() * 2 ps.logger.Info("Starting websocket hubs", mlog.Int("number_of_hubs", numberOfHubs)) diff --git a/app/platform/web_hub_test.go b/app/platform/web_hub_test.go index a3e1c8d0ae..f73836e1ad 100644 --- a/app/platform/web_hub_test.go +++ b/app/platform/web_hub_test.go @@ -68,7 +68,7 @@ func TestHubStopWithMultipleConnections(t *testing.T) { }) require.NoError(t, err) - th.Service.HubStart(th.Suite) + th.Service.Start(th.Suite) wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session) wc2 := registerDummyWebConn(t, th, s.Listener.Addr(), session) wc3 := registerDummyWebConn(t, th, s.Listener.Addr(), session) @@ -91,7 +91,7 @@ func TestHubStopRaceCondition(t *testing.T) { }) require.NoError(t, err) - th.Service.HubStart(th.Suite) + th.Service.Start(th.Suite) wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session) defer wc1.Close() @@ -479,7 +479,7 @@ func TestHubIsRegistered(t *testing.T) { s := httptest.NewServer(dummyWebsocketHandler(t)) defer s.Close() - th.Service.HubStart(th.Suite) + th.Service.Start(th.Suite) wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session) wc2 := registerDummyWebConn(t, th, s.Listener.Addr(), session) wc3 := registerDummyWebConn(t, th, s.Listener.Addr(), session) @@ -552,7 +552,7 @@ func BenchmarkGetHubForUserId(b *testing.B) { th := Setup(b).InitBasic() defer th.TearDown() - th.Service.HubStart(th.Suite) + th.Service.Start(th.Suite) b.ResetTimer() for i := 0; i < b.N; i++ { diff --git a/app/server.go b/app/server.go index cddce3da40..591b5751b1 100644 --- a/app/server.go +++ b/app/server.go @@ -305,8 +305,8 @@ func NewServer(options ...Option) (*Server, error) { // It is important to initialize the hub only after the global logger is set // to avoid race conditions while logging from inside the hub. - // Step 5: Hub depends on s.Channels() (step 8) - s.platform.HubStart(New(ServerConnector(s.Channels()))) + // Step 5: Start hub in platform which the hub depends on s.Channels() (step 4) + s.platform.Start(New(ServerConnector(s.Channels()))) // ------------------------------------------------------------------------- // Everything below this is not order sensitive and safe to be moved around.