From 1b210080c61879c575546fbf56bf9917fd02e6e7 Mon Sep 17 00:00:00 2001 From: Nia Weiss Date: Tue, 29 Sep 2020 13:54:16 -0400 Subject: [PATCH] fix race condition when resetting translation Previously, we never waited for translation sync goroutines to stop. That issue should be mostly harmless in the normal path. Additionally, this waits for the translation sync to shut down when stopping the server. --- boltdb/translate_test.go | 25 +++++++++++++++++++++++++ holder.go | 16 +++++++++++++--- server.go | 5 +++++ 3 files changed, 43 insertions(+), 3 deletions(-) diff --git a/boltdb/translate_test.go b/boltdb/translate_test.go index 8c7a23c18..8c3915f00 100644 --- a/boltdb/translate_test.go +++ b/boltdb/translate_test.go @@ -20,6 +20,7 @@ import ( "io/ioutil" "os" "reflect" + "strconv" "testing" "time" @@ -216,6 +217,30 @@ func TestTranslateStore_TranslateIDs(t *testing.T) { } } +func TestTranslateStore_MaxID(t *testing.T) { + s := MustOpenNewTranslateStore() + defer MustCloseTranslateStore(s) + + // Generate a bunch of keys. + var lastk uint64 + for i := 0; i < 1026; i++ { + k, err := s.TranslateKey(strconv.Itoa(i), true) + if err != nil { + t.Fatalf("translating %d: %v", i, err) + } + lastk = k + } + + // Verify the max ID. + max, err := s.MaxID() + if err != nil { + t.Fatalf("checking max ID: %v", err) + } + if max != lastk { + t.Fatalf("last key is %d but max is %d", lastk, max) + } +} + func TestTranslateStore_EntryReader(t *testing.T) { t.Run("OK", func(t *testing.T) { s := MustOpenNewTranslateStore() diff --git a/holder.go b/holder.go index 1b8700020..5d8691a92 100644 --- a/holder.go +++ b/holder.go @@ -1275,6 +1275,8 @@ type holderSyncer struct { // Translation sync handling. readers []TranslateEntryReader + syncers errgroup.Group + // Stats Stats stats.StatsClient @@ -1569,6 +1571,7 @@ func (s *holderSyncer) stopTranslationSync() error { return rd.Close() }) } + g.Go(s.syncers.Wait) return g.Wait() } @@ -1658,7 +1661,11 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error { } s.readers = append(s.readers, rd) - go func() { defer rd.Close(); s.readIndexTranslateReader(rd) }() + s.syncers.Go(func() error { + defer rd.Close() + s.readIndexTranslateReader(rd) + return nil + }) } return nil @@ -1697,8 +1704,11 @@ func (s *holderSyncer) initializeFieldTranslateReplication() error { } s.readers = append(s.readers, rd) - go func() { defer rd.Close(); s.readFieldTranslateReader(rd) }() - + s.syncers.Go(func() error { + defer rd.Close() + s.readFieldTranslateReader(rd) + return nil + }) return nil } diff --git a/server.go b/server.go index a2ceadc01..a93df180a 100644 --- a/server.go +++ b/server.go @@ -534,10 +534,12 @@ func (s *Server) Close() error { s.wg.Wait() var errh error + var errhs error var errc error if s.cluster != nil { errc = s.cluster.close() } + errhs = s.syncer.stopTranslationSync() if s.holder != nil { errh = s.holder.Close() } @@ -553,6 +555,9 @@ func (s *Server) Close() error { if errh != nil { return errors.Wrap(errh, "closing holder") } + if errhs != nil { + return errors.Wrap(errhs, "terminating holder translation sync") + } if errc != nil { return errors.Wrap(errc, "closing cluster") }