From b6114804993aaa8388763b668160e0b78139c673 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Mon, 1 Feb 2021 21:45:55 +0100 Subject: [PATCH] disco State --- cluster.go | 9 ++-- server/cluster_test.go | 120 ++++++++++++++++++++++++----------------- server/server_test.go | 21 +++++--- 3 files changed, 90 insertions(+), 60 deletions(-) diff --git a/cluster.go b/cluster.go index a3c9b5d6a..d02a70417 100644 --- a/cluster.go +++ b/cluster.go @@ -44,10 +44,11 @@ import ( const ( // ClusterState represents the state returned in the /status endpoint. - ClusterStateStarting = "STARTING" - ClusterStateDegraded = "DEGRADED" // cluster is running but we've lost some # of hosts >0 but < replicaN - ClusterStateNormal = "NORMAL" - ClusterStateResizing = "RESIZING" + ClusterStateStarting = disco.ClusterStateStarting + ClusterStateDegraded = disco.ClusterStateDegraded // cluster is running but we've lost some # of hosts >0 but < replicaN + ClusterStateNormal = disco.ClusterStateNormal + ClusterStateResizing = disco.ClusterStateResizing + ClusterStateDown = disco.ClusterStateDown // NodeState represents the state of a node during startup. nodeStateReady = "READY" diff --git a/server/cluster_test.go b/server/cluster_test.go index da6cf611f..f4d22b08d 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -121,8 +121,9 @@ func TestClusterResize_EmptyNode(t *testing.T) { m0 := test.RunCommand(t) defer m0.Close() - if m0.API.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected cluster state: %s", m0.API.State()) + state0, err := m0.API.State() + if err != nil || state0 != pilosa.ClusterStateNormal { + t.Fatalf("unexpected cluster state: %s, error: %v", state0, err) } } @@ -131,10 +132,12 @@ func TestClusterResize_EmptyNodes(t *testing.T) { clus := test.MustRunCluster(t, 2) defer clus.Close() - if clus.GetNode(0).API.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", clus.GetNode(0).API.State()) - } else if clus.GetNode(1).API.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", clus.GetNode(1).API.State()) + state0, err0 := clus.GetNode(0).API.State() + state1, err1 := clus.GetNode(1).API.State() + if err0 != nil || state0 != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || state1 != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } } @@ -157,10 +160,12 @@ func TestClusterResize_AddNode(t *testing.T) { clus := test.MustRunCluster(t, 2) defer clus.Close() - if !test.CheckClusterState(clus.GetNode(0), pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", clus.GetNode(0).API.State()) - } else if !test.CheckClusterState(clus.GetNode(1), pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", clus.GetNode(1).API.State()) + state0, err0 := clus.GetNode(0).API.State() + state1, err1 := clus.GetNode(1).API.State() + if err0 != nil || !test.CheckClusterState(clus.GetNode(0), pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(clus.GetNode(1), pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } }) t.Run("WithIndex", func(t *testing.T) { @@ -200,10 +205,12 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error; %v", state1, err1) } }) t.Run("ContinuousShards", func(t *testing.T) { @@ -259,10 +266,12 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } // Verify the data exists on both nodes. @@ -317,10 +326,12 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } // Verify the data exists on both nodes. @@ -382,10 +393,12 @@ func TestClusterResize_AddNode(t *testing.T) { defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } // Verify the data exists on both nodes. @@ -438,10 +451,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } if err := <-errc; err != nil { @@ -503,10 +518,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { }() defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } // Verify the data exists on both nodes. @@ -570,10 +587,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } defer m1.Close() - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } // Verify the data exists on both nodes. @@ -633,10 +652,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { t.Fatalf("starting second main: %v", err) } - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) } m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) @@ -693,12 +714,15 @@ func TestCluster_GossipMembership(t *testing.T) { t.Fatal(err) } - if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) - } else if !test.CheckClusterState(m2, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node2 cluster state: %s", m2.API.State()) + state0, err0 := m0.API.State() + state1, err1 := m1.API.State() + state2, err2 := m2.API.State() + if err0 != nil || !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0) + } else if err1 != nil || !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1) + } else if err2 != nil || !test.CheckClusterState(m2, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node2 cluster state: %s, error: %v", state2, err2) } numNodes := len(m0.API.Hosts(context.Background())) diff --git a/server/server_test.go b/server/server_test.go index c80044457..4952f479e 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -32,6 +32,7 @@ import ( "time" "github.com/pilosa/pilosa/v2" + "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/http" "github.com/pilosa/pilosa/v2/pql" "github.com/pilosa/pilosa/v2/roaring" @@ -630,7 +631,7 @@ func TestClusteringNodesReplica1(t *testing.T) { cluster := test.MustRunCluster(t, 3) defer cluster.Close() - if err := cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond); err != nil { + if err := cluster.AwaitState(string(disco.ClusterStateNormal), 100*time.Millisecond); err != nil { t.Fatalf("starting cluster: %v", err) } @@ -638,12 +639,12 @@ func TestClusteringNodesReplica1(t *testing.T) { t.Fatalf("closing third node: %v", err) } - if err := cluster.AwaitCoordinatorState(pilosa.ClusterStateStarting, 30*time.Second); err != nil { + if err := cluster.AwaitCoordinatorState(string(disco.ClusterStateDown), 60*time.Second); err != nil { t.Fatalf("starting cluster: %v", err) } // confirm that cluster stops accepting queries after one node closes - if _, err := cluster.GetCoordinator().API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { + if _, err := cluster.GetCoordinator().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) } } @@ -659,7 +660,7 @@ func TestClusteringNodesReplica2(t *testing.T) { } defer cluster.Close() - err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond) + err = cluster.AwaitState(string(disco.ClusterStateDown), 100*time.Millisecond) if err != nil { t.Fatalf("starting cluster: %v", err) } @@ -670,7 +671,7 @@ func TestClusteringNodesReplica2(t *testing.T) { t.Fatalf("closing third node: %v", err) } - err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 30*time.Second) + err = cluster.AwaitCoordinatorState(string(disco.ClusterStateDegraded), 30*time.Second) if err != nil { t.Fatalf("after closing first server: %v", err) } @@ -946,7 +947,7 @@ func TestClusterQueriesAfterRestart(t *testing.T) { err = cmd1.Command.Close() if err != nil { - t.Fatalf("closing node0: %v", err) + t.Fatalf("closing node1: %v", err) } // confirm that cluster stops accepting queries after one node closes @@ -964,10 +965,14 @@ func TestClusterQueriesAfterRestart(t *testing.T) { cmd1.Command.Config = config err = cmd1.Start() if err != nil { - t.Fatalf("reopening node 0: %v", err) + t.Fatalf("reopening node 1: %v", err) } - for cmd1.API.State() != pilosa.ClusterStateNormal { + state1, err1 := cmd1.API.State() + if err1 != nil { + t.Fatalf("getting state foor node 1: %v", err) + } + for state1 != pilosa.ClusterStateNormal { time.Sleep(time.Millisecond) }