From 6b23925bd764fe2a7552fa3d739ba2625ebd1c9b Mon Sep 17 00:00:00 2001 From: reesporte Date: Wed, 23 Feb 2022 14:29:15 -0600 Subject: [PATCH] improve Server WaitGroup concurrent usage Add a lock to the Server WaitGroup so that if the Server WaitGroup is already waiting, we won't concurrently add to it and cause a data race. Also, when adding to the Server WaitGroup, check that the server is not closing already, since that means we really shouldn't be doing more work. --- api.go | 14 +++++++---- server.go | 53 ++++++++++++++++++++++++++++++++++------- server_internal_test.go | 31 ++++++++++++++++++++++++ 3 files changed, 85 insertions(+), 13 deletions(-) diff --git a/api.go b/api.go index 45892dce1..783b22e0e 100644 --- a/api.go +++ b/api.go @@ -1049,9 +1049,17 @@ func (api *API) requestUsageOfNodes() { // Calculates disk usage from scratch if cache has expired for each index and stores the results in the usage cache func (api *API) calculateUsage() { + // don't need to calculateUsage if we're about to close! + if api.isClosing() { + return + } + api.usageCache.muCalculate.Lock() defer api.usageCache.muCalculate.Unlock() - api.server.wg.Add(1) + if ok := api.server.addToWaitGroup(1); !ok { + // the server is closing, so just stop! + return + } defer api.server.wg.Done() api.usageCache.muAssign.Lock() @@ -1065,10 +1073,6 @@ func (api *API) calculateUsage() { if err != nil { api.server.logger.Infof("couldn't get index usage details: %s", err) } - if api.isClosing() { - return - } - totalSize := nodeMetadataBytes for _, s := range indexDetails { totalSize += s.Total diff --git a/server.go b/server.go index 531e56737..b58b1b714 100644 --- a/server.go +++ b/server.go @@ -44,6 +44,7 @@ var _ broadcaster = &Server{} type Server struct { // nolint: maligned // Close management. wg sync.WaitGroup + muWG sync.Mutex closing chan struct{} // Internal @@ -99,6 +100,26 @@ func (s *Server) Holder() *Holder { return s.holder } +// addToWaitGroup adds to the server WaitGroup but makes sure the server isn't +// closing, and that the WaitGroup is not already waiting before it adds +func (s *Server) addToWaitGroup(delta int) bool { + select { + case <-s.closing: + return false + default: + s.muWG.Lock() + defer s.muWG.Unlock() + select { + case <-s.closing: + // if we're closing after having gotten the lock, stop!! + return false + default: + s.wg.Add(delta) + return true + } + } +} + // ServerOption is a functional option type for pilosa.Server type ServerOption func(s *Server) error @@ -590,7 +611,10 @@ func (s *Server) Open() error { // Start background process listening for translation // sync resets. - s.wg.Add(1) + if ok := s.addToWaitGroup(1); !ok { + return fmt.Errorf("closing server while opening server is NOT allowed") + } + go func() { defer s.wg.Done(); s.monitorResetTranslationSync() }() go func() { _ = s.translationSyncer.Reset() }() @@ -617,7 +641,9 @@ func (s *Server) Open() error { return errors.Wrap(err, "setting nodeState") } - s.wg.Add(3) + if ok := s.addToWaitGroup(3); !ok { + return fmt.Errorf("closing server while opening server is NOT allowed") + } go func() { defer s.wg.Done(); s.monitorAntiEntropy() }() go func() { defer s.wg.Done(); s.monitorRuntime() }() go func() { defer s.wg.Done(); s.monitorDiagnostics() }() @@ -631,14 +657,18 @@ func (s *Server) Open() error { return toSend }() - s.wg.Add(1) + if ok := s.addToWaitGroup(1); !ok { + return fmt.Errorf("closing server while opening server is NOT allowed") + } go func() { defer s.wg.Done() ctx, cancel := context.WithCancel(context.Background()) defer cancel() - - s.wg.Add(1) + if ok := s.addToWaitGroup(1); !ok { + // the server is closing, stop!! + return + } go func() { defer s.wg.Done() defer cancel() @@ -716,11 +746,15 @@ func (s *Server) Close() error { case <-s.closing: return nil default: - errE := s.executor.Close() - + // get the muWG lock so that noone adds to the WaitGroup while it Waits + s.muWG.Lock() + defer s.muWG.Unlock() // Notify goroutines to stop. close(s.closing) s.wg.Wait() + + errE := s.executor.Close() + var errh, errd error var errhs error var errc error @@ -776,8 +810,11 @@ func (s *Server) monitorResetTranslationSync() { case <-s.closing: return case <-s.resetTranslationSyncCh: + if ok := s.addToWaitGroup(1); !ok { + // the server is closing!!! stop!! + return + } s.logger.Infof("holder translation sync beginning") - s.wg.Add(1) go func() { // Obtaining this lock ensures that there is only // one instance of resetTranslationSync() running diff --git a/server_internal_test.go b/server_internal_test.go index da6d57578..b2e5d1116 100644 --- a/server_internal_test.go +++ b/server_internal_test.go @@ -35,3 +35,34 @@ func TestMonitorAntiEntropyZero(t *testing.T) { t.Fatalf("monitorAntiEntropy should have returned immediately with duration 0") } } + +func TestAddToWaitGroup(t *testing.T) { + // if this test times out / panics we have a problem, otherwise we're fine + td := t.TempDir() + cfg := &storage.Config{FsyncEnabled: false, Backend: storage.DefaultBackend} + s, err := NewServer(OptServerDataDir(td), OptServerStorageConfig(cfg)) + if err != nil { + t.Fatalf("making new server: %v", err) + } + + oks := make(chan bool, 10) + for i := 0; i < 10; i++ { + go func() { + oks <- s.addToWaitGroup(1) + time.Sleep(10 * time.Millisecond) + defer s.wg.Done() + }() + } + + for i := 0; i < 10; i++ { + ok := <-oks + if !ok { + t.Fatalf("unexpected close during WaitGroup add") + } + } + + s.Close() + if ok := s.addToWaitGroup(1); ok { + t.Fatalf("shouldn't be able to add while server is closing") + } +}