add more AwaitState calls in the tests

This commit is contained in:
Travis 2021-02-01 21:21:06 -06:00
parent 4811958de4
commit 629bfa3ac8
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
3 changed files with 59 additions and 26 deletions

View file

@ -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

View file

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

View file

@ -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
}