fix a data race in the platform service (#21242)
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
5e69c6b02f
Коммит
10655a2eb3
@@ -181,7 +181,7 @@ func setupTestHelper(dbStore store.Store, enterprise bool, includeCacheLayer boo
|
|||||||
th.Service.SetLicense(nil)
|
th.Service.SetLicense(nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
th.Service.HubStart(th.Suite)
|
th.Service.Start(th.Suite)
|
||||||
|
|
||||||
return th
|
return th
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -243,32 +243,6 @@ func New(sc ServiceConfig, options ...Option) (*PlatformService, error) {
|
|||||||
|
|
||||||
ps.Busy = NewBusy(ps.clusterIFace)
|
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
|
// Enable developer settings if this is a "dev" build
|
||||||
if model.BuildNumber == "dev" {
|
if model.BuildNumber == "dev" {
|
||||||
ps.UpdateConfig(func(cfg *model.Config) { *cfg.ServiceSettings.EnableDeveloper = true })
|
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
|
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 {
|
func (ps *PlatformService) ShutdownMetrics() error {
|
||||||
if ps.metrics != nil {
|
if ps.metrics != nil {
|
||||||
return ps.metrics.stopMetricsServer()
|
return ps.metrics.stopMetricsServer()
|
||||||
|
|||||||
@@ -94,8 +94,8 @@ func newWebHub(ps *PlatformService) *Hub {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// HubStart starts all the hubs.
|
// hubStart starts all the hubs.
|
||||||
func (ps *PlatformService) HubStart(suite SuiteIFace) {
|
func (ps *PlatformService) hubStart(suite SuiteIFace) {
|
||||||
// Total number of hubs is twice the number of CPUs.
|
// Total number of hubs is twice the number of CPUs.
|
||||||
numberOfHubs := runtime.NumCPU() * 2
|
numberOfHubs := runtime.NumCPU() * 2
|
||||||
ps.logger.Info("Starting websocket hubs", mlog.Int("number_of_hubs", numberOfHubs))
|
ps.logger.Info("Starting websocket hubs", mlog.Int("number_of_hubs", numberOfHubs))
|
||||||
|
|||||||
@@ -68,7 +68,7 @@ func TestHubStopWithMultipleConnections(t *testing.T) {
|
|||||||
})
|
})
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
th.Service.HubStart(th.Suite)
|
th.Service.Start(th.Suite)
|
||||||
wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
||||||
wc2 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
wc2 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
||||||
wc3 := 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)
|
require.NoError(t, err)
|
||||||
|
|
||||||
th.Service.HubStart(th.Suite)
|
th.Service.Start(th.Suite)
|
||||||
wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
||||||
defer wc1.Close()
|
defer wc1.Close()
|
||||||
|
|
||||||
@@ -479,7 +479,7 @@ func TestHubIsRegistered(t *testing.T) {
|
|||||||
s := httptest.NewServer(dummyWebsocketHandler(t))
|
s := httptest.NewServer(dummyWebsocketHandler(t))
|
||||||
defer s.Close()
|
defer s.Close()
|
||||||
|
|
||||||
th.Service.HubStart(th.Suite)
|
th.Service.Start(th.Suite)
|
||||||
wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
wc1 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
||||||
wc2 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
wc2 := registerDummyWebConn(t, th, s.Listener.Addr(), session)
|
||||||
wc3 := 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()
|
th := Setup(b).InitBasic()
|
||||||
defer th.TearDown()
|
defer th.TearDown()
|
||||||
|
|
||||||
th.Service.HubStart(th.Suite)
|
th.Service.Start(th.Suite)
|
||||||
|
|
||||||
b.ResetTimer()
|
b.ResetTimer()
|
||||||
for i := 0; i < b.N; i++ {
|
for i := 0; i < b.N; i++ {
|
||||||
|
|||||||
@@ -305,8 +305,8 @@ func NewServer(options ...Option) (*Server, error) {
|
|||||||
|
|
||||||
// It is important to initialize the hub only after the global logger is set
|
// It is important to initialize the hub only after the global logger is set
|
||||||
// to avoid race conditions while logging from inside the hub.
|
// to avoid race conditions while logging from inside the hub.
|
||||||
// Step 5: Hub depends on s.Channels() (step 8)
|
// Step 5: Start hub in platform which the hub depends on s.Channels() (step 4)
|
||||||
s.platform.HubStart(New(ServerConnector(s.Channels())))
|
s.platform.Start(New(ServerConnector(s.Channels())))
|
||||||
|
|
||||||
// -------------------------------------------------------------------------
|
// -------------------------------------------------------------------------
|
||||||
// Everything below this is not order sensitive and safe to be moved around.
|
// Everything below this is not order sensitive and safe to be moved around.
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user