diff --git a/cluster.go b/cluster.go index b0bc00548..d6982915d 100644 --- a/cluster.go +++ b/cluster.go @@ -570,13 +570,13 @@ func (c *cluster) Nodes() []*topology.Node { // Set node states and IsPrimary. for _, node := range nodes { node.IsCoordinator = node.ID == primaryNode.ID - // s, err := c.stator.NodeState(context.Background(), node.ID) - // if err != nil { - // node.State = nodeStateDown - // continue - // } - // node.State = string(s) + s, err := c.stator.NodeState(context.Background(), node.ID) + if err != nil { + node.State = nodeStateDown + continue + } + node.State = string(s) } return nodes diff --git a/server/server_test.go b/server/server_test.go index 298be37eb..c1bf488f2 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -108,6 +108,10 @@ func TestMain_Set_Quick(t *testing.T) { t.Fatal(err) } + if err := m.AwaitState(string(pilosa.ClusterStateNormal), 10*time.Second); err != nil { + t.Fatalf("restarting cluster: %v", err) + } + // Validate data after reopening. for field, fieldSet := range SetCommands(cmds).Fields() { for id, columnIDs := range fieldSet { @@ -187,6 +191,10 @@ func TestMain_SetRowAttrs(t *testing.T) { t.Fatal(err) } + if err := m.AwaitState(string(pilosa.ClusterStateNormal), 10*time.Second); err != nil { + t.Fatalf("restarting cluster: %v", err) + } + // Query rows after reopening. if res, err := m.Query(t, "i", "columnAttrs=true", `Row(x=1)`); err != nil { t.Fatal(err) @@ -243,6 +251,10 @@ func TestMain_SetColumnAttrs(t *testing.T) { t.Fatal(err) } + if err := m.AwaitState(string(pilosa.ClusterStateNormal), 10*time.Second); err != nil { + t.Fatalf("restarting cluster: %v", err) + } + // Query row after reopening. if res, err := m.Query(t, "i", "columnAttrs=true", `Row(x=1)`); err != nil { t.Fatal(err) @@ -650,7 +662,10 @@ func TestClusteringNodesReplica1(t *testing.T) { } func TestClusteringNodesReplica2(t *testing.T) { - cluster := test.MustNewCluster(t, 3) + // Because this test shuts down 2 nodes, it needs to start as a 5-node + // cluster in order to retain enough available nodes for raft leader + // election. + cluster := test.MustNewCluster(t, 5) for _, c := range cluster.Nodes { c.Config.Cluster.ReplicaN = 2 } @@ -660,11 +675,6 @@ func TestClusteringNodesReplica2(t *testing.T) { } defer cluster.Close() - err = cluster.AwaitState(string(disco.ClusterStateDown), 100*time.Millisecond) - if err != nil { - t.Fatalf("starting cluster: %v", err) - } - coord, others := cluster.GetCoordinator(), cluster.GetNonCoordinators() if err := others[0].Close(); err != nil { @@ -676,22 +686,25 @@ func TestClusteringNodesReplica2(t *testing.T) { t.Fatalf("after closing first server: %v", err) } - // confirm that cluster keeps accepting queries if replication > 1 - if _, err := coord.API.CreateIndex(context.Background(), "anewindex", pilosa.IndexOptions{}); err != nil { - t.Fatalf("got unexpected error creating index: %v", err) - } + // We no longer support mutations or schema changes when the cluster is in + // state DEGRADED, so this test doesn't apply anymore. + // + // // confirm that cluster keeps accepting queries if replication > 1 + // if _, err := coord.API.CreateIndex(context.Background(), "anewindex", pilosa.IndexOptions{}); err != nil { + // t.Fatalf("got unexpected error creating index: %v", err) + // } // confirm that cluster stops accepting queries if 2 nodes fail and replication == 2 if err := others[1].Close(); err != nil { t.Fatalf("closing 2nd node: %v", err) } - err = cluster.AwaitCoordinatorState(string(pilosa.ClusterStateStarting), 30*time.Second) + err = cluster.AwaitCoordinatorState(string(pilosa.ClusterStateDown), 30*time.Second) if err != nil { t.Fatalf("after closing second server: %v", err) } - if _, err := coord.API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { + if _, err := coord.API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state DOWN") { t.Fatalf("got unexpected error querying an incomplete cluster: %v", err) } } @@ -1201,20 +1214,15 @@ func TestClusterCreatedAtRace(t *testing.T) { cluster := test.MustRunCluster(t, 4) defer cluster.Close() - err := cluster.AwaitState(string(pilosa.ClusterStateNormal), 100*time.Millisecond) - if err != nil { - t.Fatalf("starting cluster: %v", err) - } - for _, com := range cluster.Nodes { nodes := com.API.Hosts(context.Background()) for _, n := range nodes { - if n.State != "READY" { - t.Fatalf("unexpected node state after upping cluster: %v", nodes) // server_test.go:1245: unexpected node state after upping cluster: [Node:http://localhost:43075:READY:TestClusterCreatedAtRace/run-0__0 Node:http://localhost:42301:READY:TestClusterCreatedAtRace/run-0__1 Node:http://localhost:42031:DOWN:TestClusterCreatedAtRace/run-0__2 Node:http://localhost:43671:READY:TestClusterCreatedAtRace/run-0__3] + if n.State != string(disco.NodeStateStarted) { + t.Fatalf("unexpected node state (%s) after upping cluster: %v", n.State, nodes) } } } - _, err = cluster.Nodes[0].API.CreateIndex(context.Background(), "anindex", pilosa.IndexOptions{}) + _, err := cluster.Nodes[0].API.CreateIndex(context.Background(), "anindex", pilosa.IndexOptions{}) if err != nil && errors.Cause(err).Error() != pilosa.ErrIndexExists.Error() { t.Fatal(err) } diff --git a/test/pilosa.go b/test/pilosa.go index 1a6635151..55fc13bb1 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -394,3 +394,28 @@ func RetryUntil(timeout time.Duration, fn func() error) (err error) { } } } + +// AwaitState waits for the whole cluster to reach a specified state. +func (m *Command) AwaitState(expectedState string, timeout time.Duration) (err error) { + startTime := time.Now() + var elapsed time.Duration + for elapsed = 0; elapsed <= timeout; elapsed = time.Since(startTime) { + // Counterintuitive: We're returning if the err *is* nil, + // meaning we've reached the expected state. + if err = m.exceptionalState(expectedState); err == nil { + return err + } + time.Sleep(1 * time.Millisecond) + } + return fmt.Errorf("waited %v for command to reach state %q: %v", + elapsed, expectedState, err) +} + +// exceptionalState returns an error if the node is not in the expected state. +func (m *Command) exceptionalState(expectedState string) error { + state, err := m.API.State() + if err != nil || state != expectedState { + return fmt.Errorf("node %q: state %s: err %v", m.ID(), state, err) + } + return nil +}