disco State

This commit is contained in:
Kuba Podgórski 2021-02-01 21:45:55 +01:00 • committed by Travis
parent 6058fc22e4
commit b611480499
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
3 changed files with 90 additions and 60 deletions

View file

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

View file

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

View file

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