From 5d43d414f7f14185e029bc2ea1adf0762a7a558d Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 27 Jun 2018 16:59:35 -0500 Subject: [PATCH 1/2] Fix a few data races --- http/translator_test.go | 17 +++++++++++------ translate.go | 4 +++- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/http/translator_test.go b/http/translator_test.go index d317944e4..fdfc93f3b 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -5,6 +5,7 @@ import ( "io" "io/ioutil" gohttp "net/http" + "sync/atomic" "testing" "time" @@ -37,9 +38,10 @@ func TestTranslateStore_Reader(t *testing.T) { return 0, nil } } - var closeInvoked bool + closeInvoked := atomic.Value{} + closeInvoked.Store(false) mrc.CloseFunc = func() error { - closeInvoked = true + closeInvoked.Store(true) return nil } @@ -85,7 +87,7 @@ func TestTranslateStore_Reader(t *testing.T) { t.Fatal(err) } - if !closeInvoked { + if !closeInvoked.Load().(bool) { t.Fatal("expected server close") } }) @@ -100,9 +102,12 @@ func TestTranslateStore_Reader(t *testing.T) { <-done return 0, io.EOF } - var closeInvoked bool + + closeInvoked := atomic.Value{} + closeInvoked.Store(false) + mrc.CloseFunc = func() error { - closeInvoked = true + closeInvoked.Store(true) return nil } @@ -127,7 +132,7 @@ func TestTranslateStore_Reader(t *testing.T) { // Cancel the context and check if server is closed. cancel() time.Sleep(100 * time.Millisecond) - if !closeInvoked { + if !closeInvoked.Load().(bool) { t.Fatal("expected server-side close") } }) diff --git a/translate.go b/translate.go index 398055fdb..7c716640b 100644 --- a/translate.go +++ b/translate.go @@ -306,11 +306,13 @@ func (s *TranslateFile) replicate(ctx context.Context) error { } else if err != nil { return err } - + s.mu.Lock() // Write to local store. if err := s.appendEntry(&entry); err != nil { + s.mu.Unlock() return err } + s.mu.Unlock() } } From 25be5c0f2fe0ca18605d3ac1dd587773c5d57e48 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Thu, 28 Jun 2018 10:38:37 -0500 Subject: [PATCH 2/2] Use channel to notify on server close instead of atomic.Value. Ensure CloseFunc() only called once. --- http/translator_test.go | 46 +++++++++++++++++++++++------------------ mock/mock.go | 9 +++++++- 2 files changed, 34 insertions(+), 21 deletions(-) diff --git a/http/translator_test.go b/http/translator_test.go index fdfc93f3b..1eea725a2 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -5,7 +5,6 @@ import ( "io" "io/ioutil" gohttp "net/http" - "sync/atomic" "testing" "time" @@ -16,6 +15,17 @@ import ( "github.com/pilosa/pilosa/test" ) +func newMockReadCloser() *mock.ReadCloser { + return &mock.ReadCloser{ + ReadFunc: func(p []byte) (int, error) { + return 0, io.EOF + }, + CloseFunc: func() error { + return nil + }, + } +} + func TestTranslateStore_Reader(t *testing.T) { // Ensure client can connect and stream the translate store data. t.Run("OK", func(t *testing.T) { @@ -38,10 +48,9 @@ func TestTranslateStore_Reader(t *testing.T) { return 0, nil } } - closeInvoked := atomic.Value{} - closeInvoked.Store(false) + closeInvoked := make(chan struct{}) mrc.CloseFunc = func() error { - closeInvoked.Store(true) + close(closeInvoked) return nil } @@ -57,15 +66,7 @@ func TestTranslateStore_Reader(t *testing.T) { } return &mrc, nil } - mrc2 := mock.ReadCloser{ - ReadFunc: func(p []byte) (int, error) { - return 0, io.EOF - }, - CloseFunc: func() error { - return nil - }, - } - return &mrc2, nil + return newMockReadCloser(), nil } opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore)) @@ -87,8 +88,11 @@ func TestTranslateStore_Reader(t *testing.T) { t.Fatal(err) } - if !closeInvoked.Load().(bool) { + select { + case <-time.NewTimer(time.Millisecond * 100).C: t.Fatal("expected server close") + case <-closeInvoked: + return } }) @@ -103,15 +107,15 @@ func TestTranslateStore_Reader(t *testing.T) { return 0, io.EOF } - closeInvoked := atomic.Value{} - closeInvoked.Store(false) + closeInvoked := make(chan struct{}) mrc.CloseFunc = func() error { - closeInvoked.Store(true) + close(closeInvoked) return nil } var translateStore mock.TranslateStore + translateStore.ReaderFunc = func(ctx context.Context, off int64) (io.ReadCloser, error) { return &mrc, nil } @@ -131,9 +135,11 @@ func TestTranslateStore_Reader(t *testing.T) { // Cancel the context and check if server is closed. cancel() - time.Sleep(100 * time.Millisecond) - if !closeInvoked.Load().(bool) { - t.Fatal("expected server-side close") + select { + case <-time.NewTimer(time.Millisecond * 100).C: + t.Fatal("expected server close") + case <-closeInvoked: + return } }) }) diff --git a/mock/mock.go b/mock/mock.go index 46469c2f9..97ebf8641 100644 --- a/mock/mock.go +++ b/mock/mock.go @@ -1,8 +1,11 @@ package mock +import "sync" + type ReadCloser struct { ReadFunc func(p []byte) (int, error) CloseFunc func() error + once sync.Once } func (rc *ReadCloser) Read(p []byte) (int, error) { @@ -10,5 +13,9 @@ func (rc *ReadCloser) Read(p []byte) (int, error) { } func (rc *ReadCloser) Close() error { - return rc.CloseFunc() + var err error = nil + rc.once.Do(func() { + err = rc.CloseFunc() + }) + return err }