Merge pull request #1931 from jaffee/1919-cluster-race

address race condition by getting cluster nodes with lock
This commit is contained in:
Matthew Jaffee 2019-04-06 11:11:49 -05:00 committed by GitHub
commit 2c4f4d7d63
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
5 changed files with 270 additions and 13 deletions

View file

@ -296,18 +296,18 @@ func (c *cluster) unprotectedIsCoordinator() bool {
// nodes with its version of Cluster.Status.
func (c *cluster) setCoordinator(n *Node) error {
c.mu.Lock()
defer c.mu.Unlock()
// Verify that the new Coordinator value matches
// this node.
if c.Node.ID != n.ID {
c.mu.Unlock()
return fmt.Errorf("coordinator node does not match this node")
}
// Update IsCoordinator on all nodes (locally).
_ = c.unprotectedUpdateCoordinator(n)
c.mu.Unlock()
// Send the update coordinator message to all nodes.
err := c.broadcaster.SendSync(
err := c.unprotectedSendSync(
&UpdateCoordinatorMessage{
New: n,
})
@ -316,7 +316,25 @@ func (c *cluster) setCoordinator(n *Node) error {
}
// Broadcast cluster status.
return c.broadcaster.SendSync(c.status())
return c.unprotectedSendSync(c.unprotectedStatus())
}
// unprotectedSendSync is used in place of c.broadcaster.SendSync (which is
// Server.SendSync) because Server.SendSync needs to obtain a cluster lock to
// get the list of nodes. TODO: the reference loop from
// Server->cluster->broadcaster(Server) will likely continue to cause confusion
// and should be refactored.
func (c *cluster) unprotectedSendSync(m Message) error {
var eg errgroup.Group
for _, node := range c.nodes {
node := node
// Don't send to myself.
if node.ID == c.Node.ID {
continue
}
eg.Go(func() error { return c.broadcaster.SendTo(node, m) })
}
return eg.Wait()
}
// updateCoordinator updates this nodes Coordinator value as well as
@ -533,12 +551,6 @@ func (c *cluster) determineClusterState() (clusterState string) {
return ClusterStateStarting
}
func (c *cluster) status() *ClusterStatus {
c.mu.RLock()
defer c.mu.RUnlock()
return c.unprotectedStatus()
}
// unprotectedStatus returns the the cluster's status including what nodes it contains, its ID, and current state.
func (c *cluster) unprotectedStatus() *ClusterStatus {
return &ClusterStatus{
@ -1083,7 +1095,7 @@ func (c *cluster) unprotectedSetStateAndBroadcast(state string) error {
}
// Broadcast cluster status changes to the cluster.
status := c.unprotectedStatus()
return c.broadcaster.SendSync(status) // TODO fix c.Status
return c.unprotectedSendSync(status) // TODO fix c.Status
}

View file

@ -589,7 +589,8 @@ func (s *Server) SendSync(m Message) error {
return fmt.Errorf("marshaling message: %v", err)
}
msg = append([]byte{getMessageType(m)}, msg...)
for _, node := range s.cluster.nodes {
for _, node := range s.cluster.Nodes() {
node := node
// Don't forward the message to ourselves.
if s.uri == node.URI {

View file

@ -309,6 +309,189 @@ func TestClusterResize_AddNode(t *testing.T) {
})
}
// Ensure that adding a node correctly resizes the cluster.
func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
t.Run("WithIndex", func(t *testing.T) {
// Configure node0
m0 := test.MustRunCluster(t, 1)[0]
defer m0.Close()
seed := m0.GossipAddress()
// Create a client for each node.
client0 := m0.Client()
// Create indexes and fields on one node.
if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
} else if err := client0.CreateField(context.Background(), "i", "f"); err != nil {
t.Fatal(err)
}
errc := make(chan error)
go func() {
_, err := m0.API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{})
errc <- err
}()
// Configure node1
m1 := test.NewCommandNode(false)
m1.Config.Gossip.Port = "0"
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.API.State())
} else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node1 cluster state: %s", m1.API.State())
}
if err := <-errc; err != nil {
t.Fatalf("error from index creation: %v", err)
}
})
t.Run("ContinuousShards", func(t *testing.T) {
// Configure node0
m0 := test.MustRunCluster(t, 1)[0]
defer m0.Close()
seed := m0.GossipAddress()
// Create a client for each node.
client0 := m0.Client()
// Create indexes and fields on one node.
if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
} else if err := client0.CreateField(context.Background(), "i", "f"); err != nil {
t.Fatal(err)
}
// Write data on first node.
if _, err := m0.Query("i", "", `
Set(1, f=1)
Set(1300000, f=1)
`); err != nil {
t.Fatal(err)
}
// exp is the expected result for the Row queries that follow.
exp := `{"results":[{"attrs":{},"columns":[1,1300000]}]}` + "\n"
// Verify the data exists on the single node.
if res, err := m0.Query("i", "", `Row(f=1)`); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result: %s", res)
}
// Configure node1
m1 := test.NewCommandNode(false)
m1.Config.Gossip.Port = "0"
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
t.Fatalf("starting second main: %v", err)
}
errc := make(chan error, 1)
go func() {
_, err := m0.API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{})
errc <- err
}()
defer m1.Close()
if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.API.State())
} else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node1 cluster state: %s", m1.API.State())
}
// Verify the data exists on both nodes.
if res, err := m0.Query("i", "", `Row(f=1)`); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result: %s", res)
}
if res, err := m1.Query("i", "", `Row(f=1)`); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result: %s", res)
}
})
t.Run("SkippedShard", func(t *testing.T) {
// Configure node0
m0 := test.MustRunCluster(t, 1)[0]
defer m0.Close()
seed := m0.GossipAddress()
// Create a client for each node.
client0 := m0.Client()
// Create indexes and fields on one node.
if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
} else if err := client0.CreateField(context.Background(), "i", "f"); err != nil {
t.Fatal(err)
}
// Write data on first node. Note that no data is placed on shard 1.
if _, err := m0.Query("i", "", `
Set(1, f=1)
Set(2400000, f=1)
`); err != nil {
t.Fatal(err)
}
// exp is the expected result for the Row queries that follow.
exp := `{"results":[{"attrs":{},"columns":[1,2400000]}]}` + "\n"
// Verify the data exists on the single node.
if res, err := m0.Query("i", "", `Row(f=1)`); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result: %s", res)
}
// Configure node1
m1 := test.NewCommandNode(false)
m1.Config.Gossip.Port = "0"
m1.Config.Gossip.Seeds = []string{seed}
errc := make(chan error, 1)
go func() {
_, err := m0.API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{})
errc <- err
}()
err := m1.Start()
if err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.API.State())
} else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node1 cluster state: %s", m1.API.State())
}
// Verify the data exists on both nodes.
if res, err := m0.Query("i", "", `Row(f=1)`); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result: %s", res)
}
if res, err := m1.Query("i", "", `Row(f=1)`); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result: %s", res)
}
})
}
// Ensure that redundant gossip seeds are used
func TestCluster_GossipMembership(t *testing.T) {
t.Run("Node0Down", func(t *testing.T) {

View file

@ -630,6 +630,53 @@ func TestRemoveNodeAfterItDies(t *testing.T) {
}
}
func TestRemoveConcurrentIndexCreation(t *testing.T) {
cluster := test.MustNewCluster(t, 3)
for _, c := range cluster {
c.Config.Cluster.ReplicaN = 2
}
err := cluster.Start()
if err != nil {
t.Fatalf("starting cluster: %v", err)
}
var wait = true
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond * 1)
}
errc := make(chan error)
go func() {
_, err := cluster[0].API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{})
errc <- err
}()
if _, err := cluster[0].API.RemoveNode(cluster[2].API.Node().ID); err != nil {
t.Fatalf("removing node: %v", err)
}
for i := 0; cluster[0].API.State() != pilosa.ClusterStateNormal; i++ {
time.Sleep(time.Millisecond)
if i > 10 {
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
}
}
hosts := cluster[0].API.Hosts(context.Background())
if len(hosts) != 2 {
t.Fatalf("unexpected hosts: %v", hosts)
}
if err := <-errc; err != nil {
t.Fatalf("error from index creation: %v", err)
}
}
// Ensure program imports timestamps as UTC.
func TestMain_ImportTimestamp(t *testing.T) {
m := test.MustRunCommand()

View file

@ -338,7 +338,7 @@ func (bcast) SendAsync(Message) error {
return nil
}
// SendTo is a test implemenetation of Broadcaster SendTo method.
// SendTo is a test implementation of Broadcaster SendTo method.
func (b bcast) SendTo(to *Node, m Message) error {
switch obj := m.(type) {
case *ResizeInstruction:
@ -349,6 +349,20 @@ func (b bcast) SendTo(to *Node, m Message) error {
case *ResizeInstructionComplete:
coord := b.t.clusterByID(to.ID)
go coord.markResizeInstructionComplete(obj)
case *ClusterStatus:
// Apply the send message to the node.
for _, c := range b.t.Clusters {
if c.Node.ID == to.ID {
c.mergeClusterStatus(obj)
}
}
b.t.mu.RLock()
if obj.State == ClusterStateNormal && b.t.resizing {
close(b.t.resizeDone)
}
b.t.mu.RUnlock()
default:
panic(fmt.Sprintf("message not handled:\n%#v\n", obj))
}
return nil
}