From 6b23925bd764fe2a7552fa3d739ba2625ebd1c9b Mon Sep 17 00:00:00 2001 From: reesporte Date: Wed, 23 Feb 2022 14:29:15 -0600 Subject: [PATCH 1/2] 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") + } +} From 45633e23a5c5ddce39e4a9ebbbe625adcfc6bb93 Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 25 Feb 2022 10:53:45 -0600 Subject: [PATCH 2/2] only set bits after the holder is completely setup This should help prevent a data race. SetBit can, in some cases, cause an asynchronous task to run which tries to update the stats counter. But if that task runs while we are modifying the stats counter itself, we have a data race. --- stats/stats_test.go | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/stats/stats_test.go b/stats/stats_test.go index 81b62a5f2..4515636c9 100644 --- a/stats/stats_test.go +++ b/stats/stats_test.go @@ -81,11 +81,6 @@ func TestStatsCount_TopN(t *testing.T) { defer c.Close() hldr := test.Holder{Holder: c.GetNode(0).Server.Holder()} - hldr.SetBit("d", "f", 0, 0) - hldr.SetBit("d", "f", 0, 1) - hldr.SetBit("d", "f", 0, pilosa.ShardWidth) - hldr.SetBit("d", "f", 0, pilosa.ShardWidth+2) - // Execute query. called := false hldr.Holder.Stats = &MockStats{ @@ -101,6 +96,12 @@ func TestStatsCount_TopN(t *testing.T) { called = true }, } + + hldr.SetBit("d", "f", 0, 0) + hldr.SetBit("d", "f", 0, 1) + hldr.SetBit("d", "f", 0, pilosa.ShardWidth) + hldr.SetBit("d", "f", 0, pilosa.ShardWidth+2) + if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: "d", Query: `TopN(field=f, n=2)`}); err != nil { t.Fatal(err) }