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() } }