diff --git a/enterprise/replicasync/replicasync.go b/enterprise/replicasync/replicasync.go index 508672cd4d..3dc3b3414d 100644 --- a/enterprise/replicasync/replicasync.go +++ b/enterprise/replicasync/replicasync.go @@ -276,6 +276,8 @@ func (m *Manager) syncReplicas(ctx context.Context) error { return xerrors.Errorf("ping database: %w", err) } + m.mutex.Lock() + defer m.mutex.Unlock() replica, err := m.db.UpdateReplica(ctx, database.UpdateReplicaParams{ ID: m.self.ID, UpdatedAt: database.Now(), @@ -291,8 +293,6 @@ func (m *Manager) syncReplicas(ctx context.Context) error { if err != nil { return xerrors.Errorf("update replica: %w", err) } - m.mutex.Lock() - defer m.mutex.Unlock() if m.self.Error != replica.Error { // Publish an update occurred! err = m.pubsub.Publish(PubsubEvent, []byte(m.self.ID.String())) diff --git a/enterprise/replicasync/replicasync_test.go b/enterprise/replicasync/replicasync_test.go index 7538b48a38..61296a61c5 100644 --- a/enterprise/replicasync/replicasync_test.go +++ b/enterprise/replicasync/replicasync_test.go @@ -206,7 +206,12 @@ func TestReplica(t *testing.T) { _ = server.Close() }) done := false + + var m sync.Mutex server.SetCallback(func() { + m.Lock() + defer m.Unlock() + if len(server.All()) != count { return }