diff --git a/cluster.go b/cluster.go index 9ee35f4fa..204a1016b 100644 --- a/cluster.go +++ b/cluster.go @@ -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 } diff --git a/server.go b/server.go index 5a311918f..3e989e530 100644 --- a/server.go +++ b/server.go @@ -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 { diff --git a/server/cluster_test.go b/server/cluster_test.go index 617bf6cad..a0cd83f1b 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -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) { diff --git a/server/server_test.go b/server/server_test.go index 3b6be7efa..1cc3bd87f 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -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() diff --git a/utils_internal_test.go b/utils_internal_test.go index 016d5276f..bf98f08a8 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -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 }