mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
various cluster test fixups/cleanups
Some cluster tests failed sporadically. In order to fix them, I introduced some debugging-related functionality, which revealed several new bugs that were actually existing bugs we just happened not to hit in testing. This combines various fixes. We start with "make the nodes used in testing have distinct names based on the test case name", which lets us discover that we are leaking clusters, which continue to sit around talking with each other. That in turn causes significantly higher load on access to ephemeral ports, which causes sporadic failures when we shut a node down and try to restart it, but something else has gotten assigned its ephemeral port number since then. Part of the fix is to try to rebind on port 0 if an attempt to bind to a specified port over 32k fails. This is a guess; the actual ephemeral port range could be 16k+, 32k+, or 48k+, or just about anything else really, but it seems reasonable in practice. There were bugs in the oft-repeated loops to await the cluster achieving a given state, and it could hang forever if it didn't, so we add a timeout and a standard function on the test.Cluster type to handle that. Note that the timeout seems irrelevant; in every case I've tried, a timeout of 0 is fine because the node start doesn't complete until the cluster state has changed. Add a method to test.Command to run a query, expecting a specific result. Also clean up some of the formatting and generation of queries, and allow parameterized (badly) queries. This lets us fix a subtle bug, which is that test cases were depending on assumptions about shardwidths. Also improve the diagnostic output from some of these functions so test failures are more comprehensible. But actually that dependency on shardwidths was ALSO revealing a genuine underlying bug, which is that a node resize did not correctly propagate the schema to a new node if there was no data present on shards that node would own. We now also have a test case that hits that (or would, if we hadn't fixed it). Add comments explaining the server options parameters for MustNewCluster and MustRunCluster. Also, we implement the ReadFrom and WriteTo behaviors for InMemTranslateStore, without which some of the cluster resize tests fail. Props to the comment for specifically stating that they wouldn't work if that happened, which probably saved me several hours of debugging. The implementations may not be robust, but InMemTranslateStore is intended to be used only in lightweight and transient testing.
This commit is contained in:
parent
ef8b054367
commit
364b533ead
8 changed files with 278 additions and 209 deletions
12
cluster.go
12
cluster.go
|
|
@ -1431,11 +1431,11 @@ func (c *cluster) unprotectedGenerateResizeJobByAction(nodeAction nodeAction) (*
|
|||
}
|
||||
|
||||
for _, node := range toCluster.nodes {
|
||||
// If a host doesn't need to request data, mark it as complete.
|
||||
if len(fragmentSourcesByNode[node.ID]) == 0 && len(translationSourcesByNode[node.ID]) == 0 {
|
||||
j.IDs[node.ID] = true
|
||||
continue
|
||||
}
|
||||
// We may send a resize instruction that has no sources that
|
||||
// the node needs to read from -- for instance, if there's no
|
||||
// data in any fragments it would process. But it still needs
|
||||
// to get the NodeStatus to pick up the schema so it knows
|
||||
// about existing indexes.
|
||||
instr := &ResizeInstruction{
|
||||
JobID: j.ID,
|
||||
Node: toCluster.unprotectedNodeByID(node.ID),
|
||||
|
|
@ -2410,7 +2410,7 @@ func (c *cluster) translateIndexIDs(ctx context.Context, indexName string, ids [
|
|||
}
|
||||
|
||||
func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idSet map[uint64]struct{}) (map[uint64]string, error) {
|
||||
idMap := make(map[uint64]string)
|
||||
idMap := make(map[uint64]string, len(idSet))
|
||||
|
||||
index := c.holder.Index(indexName)
|
||||
if index == nil {
|
||||
|
|
|
|||
|
|
@ -199,22 +199,20 @@ func TestClusterResize_AddNode(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
col := pilosa.ShardWidth + 20
|
||||
|
||||
// Write data on first node.
|
||||
if _, err := m0.Query(t, "i", "", `
|
||||
if _, err := m0.Queryf(t, "i", "", `
|
||||
Set(1, f=1)
|
||||
Set(1300000, f=1)
|
||||
`); err != nil {
|
||||
Set(%d, f=1)
|
||||
`, col); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// exp is the expected result for the Row queries that follow.
|
||||
exp := `{"results":[{"attrs":{},"columns":[1,1300000]}]}` + "\n"
|
||||
exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col)
|
||||
|
||||
// Verify the data exists on the single node.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
|
||||
// Configure node1
|
||||
m1 := test.NewCommandNode(false)
|
||||
|
|
@ -233,16 +231,57 @@ func TestClusterResize_AddNode(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify the data exists on both nodes.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
m1.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
})
|
||||
t.Run("OneShard", 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 res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
} else if err := client0.CreateField(context.Background(), "i", "f"); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
|
||||
// Write data on first node.
|
||||
if _, err := m0.Query(t, "i", "", `
|
||||
Set(1, f=1)
|
||||
`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// exp is the expected result for the Row queries that follow.
|
||||
exp := `{"results":[{"attrs":{},"columns":[1]}]}`
|
||||
|
||||
// Verify the data exists on the single node.
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
|
||||
// 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())
|
||||
}
|
||||
|
||||
// Verify the data exists on both nodes.
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
m1.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
})
|
||||
t.Run("SkippedShard", func(t *testing.T) {
|
||||
// Configure node0
|
||||
|
|
@ -261,23 +300,21 @@ func TestClusterResize_AddNode(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
col := pilosa.ShardWidth*2 + 20
|
||||
|
||||
// Write data on first node. Note that no data is placed on shard 1.
|
||||
if _, err := m0.Query(t, "i", "", `
|
||||
if _, err := m0.Queryf(t, "i", "", `
|
||||
Set(1, f=1)
|
||||
Set(2400000, f=1)
|
||||
`); err != nil {
|
||||
Set(%d, f=1)
|
||||
`, col); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// exp is the expected result for the Row queries that follow.
|
||||
exp := `{"results":[{"attrs":{},"columns":[1,2400000]}]}` + "\n"
|
||||
exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col)
|
||||
|
||||
// Verify the data exists on the single node.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
|
||||
// Configure node1
|
||||
m1 := test.NewCommandNode(false)
|
||||
|
|
@ -296,16 +333,8 @@ func TestClusterResize_AddNode(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify the data exists on both nodes.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
m1.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -371,23 +400,21 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
col := pilosa.ShardWidth + 20
|
||||
|
||||
// Write data on first node.
|
||||
if _, err := m0.Query(t, "i", "", `
|
||||
if _, err := m0.Queryf(t, "i", "", `
|
||||
Set(1, f=1)
|
||||
Set(1300000, f=1)
|
||||
`); err != nil {
|
||||
Set(%d, f=1)
|
||||
`, col); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// exp is the expected result for the Row queries that follow.
|
||||
exp := `{"results":[{"attrs":{},"columns":[1,1300000]}]}` + "\n"
|
||||
exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col)
|
||||
|
||||
// Verify the data exists on the single node.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
|
||||
// Configure node1
|
||||
m1 := test.NewCommandNode(false)
|
||||
|
|
@ -411,16 +438,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify the data exists on both nodes.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
m1.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
})
|
||||
t.Run("SkippedShard", func(t *testing.T) {
|
||||
// Configure node0
|
||||
|
|
@ -439,23 +458,21 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
col := pilosa.ShardWidth*2 + 20
|
||||
|
||||
// Write data on first node. Note that no data is placed on shard 1.
|
||||
if _, err := m0.Query(t, "i", "", `
|
||||
if _, err := m0.Queryf(t, "i", "", `
|
||||
Set(1, f=1)
|
||||
Set(2400000, f=1)
|
||||
`); err != nil {
|
||||
Set(%d, f=1)
|
||||
`, col); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// exp is the expected result for the Row queries that follow.
|
||||
exp := `{"results":[{"attrs":{},"columns":[1,2400000]}]}` + "\n"
|
||||
exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col)
|
||||
|
||||
// Verify the data exists on the single node.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
|
||||
// Configure node1
|
||||
m1 := test.NewCommandNode(false)
|
||||
|
|
@ -479,16 +496,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify the data exists on both nodes.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
m1.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
})
|
||||
t.Run("WithIndexKeys", func(t *testing.T) {
|
||||
// Configure node0
|
||||
|
|
@ -516,14 +525,10 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
}
|
||||
|
||||
// exp is the expected result for the Row queries that follow.
|
||||
exp := `{"results":[{"attrs":{},"columns":[],"keys":["col2","col1"]}]}` + "\n"
|
||||
exp := `{"results":[{"attrs":{},"columns":[],"keys":["col2","col1"]}]}`
|
||||
|
||||
// Verify the data exists on the single node.
|
||||
if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
|
||||
// Configure node1
|
||||
m1 := test.NewCommandNode(false)
|
||||
|
|
@ -545,15 +550,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
} else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) {
|
||||
t.Fatalf("unexpected node1 cluster state: %s", m1.API.State())
|
||||
}
|
||||
|
||||
// Verify the data exists on both nodes.
|
||||
for i, node := range []*test.Command{m0, m1} {
|
||||
if res, err := node.Query(t, "i", "", `Row(f=1)`); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("node%d expected: %s, but got: %s", i, exp, res)
|
||||
}
|
||||
}
|
||||
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
m1.QueryExpect(t, "i", "", `Row(f=1)`, exp)
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1290,6 +1290,7 @@ func TestCluster_TranslateStore(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("starting cluster 0: %v", err)
|
||||
}
|
||||
defer cluster[0].Close()
|
||||
|
||||
test.Do(t, "POST", cluster[0].URL()+"/index/i0", "{\"options\": {\"keys\": true}}")
|
||||
}
|
||||
|
|
@ -1306,6 +1307,7 @@ func TestClusterTranslator(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("starting cluster 0: %v", err)
|
||||
}
|
||||
defer cluster[0].Close()
|
||||
cluster[1] = test.NewCommandNode(false,
|
||||
server.OptCommandServerOptions(
|
||||
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
|
||||
|
|
@ -1318,6 +1320,7 @@ func TestClusterTranslator(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("starting cluster 1: %v", err)
|
||||
}
|
||||
defer cluster[1].Close()
|
||||
|
||||
test.Do(t, "POST", cluster[0].URL()+"/index/i0", "{\"options\": {\"keys\": true}}")
|
||||
test.Do(t, "POST", cluster[0].URL()+"/index/i0/field/f0", "{\"options\": {\"keys\": true}}")
|
||||
|
|
|
|||
|
|
@ -406,6 +406,17 @@ func (m *Command) setupNetworking() error {
|
|||
// get the host portion of addr to use for binding
|
||||
gossipHost := m.listenURI.Host
|
||||
m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger())
|
||||
if err != nil && gossipPort >= 32768 {
|
||||
// In testing, we sometimes try to reuse an ephemeral port.
|
||||
// Which probably works. If it doesn't, this test will take
|
||||
// about a minute longer because we'll come back in from a
|
||||
// new port. See also the gossip config in gossip/gossip.go.
|
||||
// TODO: Maybe make that more configurable here.
|
||||
m.logger.Printf("ephemeral port %d already occupied, switching to :0 (%v)", gossipPort, err)
|
||||
m.Config.Gossip.Port = "0"
|
||||
gossipPort = 0
|
||||
m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger())
|
||||
}
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "getting transport")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -630,15 +630,9 @@ func TestClusteringNodesReplica1(t *testing.T) {
|
|||
cluster := test.MustRunCluster(t, 3)
|
||||
defer cluster.Close()
|
||||
|
||||
var wait = true
|
||||
for wait {
|
||||
wait = false
|
||||
for _, node := range cluster {
|
||||
if node.API.State() != pilosa.ClusterStateNormal {
|
||||
wait = true
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond * 1)
|
||||
err := cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
|
||||
if err := cluster[2].Command.Close(); err != nil {
|
||||
|
|
@ -665,14 +659,9 @@ func TestClusteringNodesReplica1(t *testing.T) {
|
|||
t.Fatalf("restarting node 2: %v", err)
|
||||
}
|
||||
|
||||
for wait {
|
||||
wait = false
|
||||
for _, node := range cluster {
|
||||
if node.API.State() != pilosa.ClusterStateNormal {
|
||||
wait = true
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
err = cluster.AwaitState(pilosa.ClusterStateNormal, 200*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("resuming normal operations: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -685,24 +674,20 @@ func TestClusteringNodesReplica2(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
defer cluster.Close()
|
||||
|
||||
var wait = true
|
||||
for wait {
|
||||
wait = false
|
||||
for _, node := range cluster {
|
||||
if node.API.State() != pilosa.ClusterStateNormal {
|
||||
wait = true
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond * 1)
|
||||
err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
|
||||
if err := cluster[2].Command.Close(); err != nil {
|
||||
t.Fatalf("closing third node: %v", err)
|
||||
}
|
||||
|
||||
if cluster[0].API.State() != pilosa.ClusterStateDegraded {
|
||||
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
|
||||
err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("after closing first server: %v", err)
|
||||
}
|
||||
|
||||
// confirm that cluster keeps accepting queries if replication > 1
|
||||
|
|
@ -715,8 +700,9 @@ func TestClusteringNodesReplica2(t *testing.T) {
|
|||
t.Fatalf("closing 2nd node: %v", err)
|
||||
}
|
||||
|
||||
if cluster[0].API.State() != pilosa.ClusterStateStarting {
|
||||
t.Fatalf("expected state to be Starting, but got %s", cluster[0].API.State())
|
||||
err = cluster.AwaitCoordinatorState(pilosa.ClusterStateStarting, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("after closing second server: %v", err)
|
||||
}
|
||||
|
||||
if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") {
|
||||
|
|
@ -739,8 +725,9 @@ func TestClusteringNodesReplica2(t *testing.T) {
|
|||
t.Fatalf("restarting node 2: %v", err)
|
||||
}
|
||||
|
||||
if cluster[0].API.State() != pilosa.ClusterStateDegraded {
|
||||
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
|
||||
err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("after restarting first server: %v", err)
|
||||
}
|
||||
|
||||
// Create new main with the same config.
|
||||
|
|
@ -756,19 +743,12 @@ func TestClusteringNodesReplica2(t *testing.T) {
|
|||
|
||||
// Run new program.
|
||||
if err := cluster[1].Start(); err != nil {
|
||||
t.Fatalf("restarting node 2: %v", err)
|
||||
t.Fatalf("restarting node 1: %v", err)
|
||||
}
|
||||
|
||||
defer cluster.Close()
|
||||
|
||||
for wait {
|
||||
wait = false
|
||||
for _, node := range cluster {
|
||||
if node.API.State() != pilosa.ClusterStateNormal {
|
||||
wait = true
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
err = cluster.AwaitState(pilosa.ClusterStateNormal, 200*time.Microsecond)
|
||||
if err != nil {
|
||||
t.Fatalf("resuming normal operations: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -781,32 +761,38 @@ func TestRemoveNodeAfterItDies(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
// The anonymous function is necessary so that the slice
|
||||
// passed to Close() as a receiver is the modified value
|
||||
// of cluster, because we're removing the last entry from it
|
||||
// below.
|
||||
defer func() {
|
||||
cluster.Close()
|
||||
}()
|
||||
|
||||
var wait = true
|
||||
for wait {
|
||||
wait = false
|
||||
for _, node := range cluster {
|
||||
if node.API.State() != pilosa.ClusterStateNormal {
|
||||
wait = true
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond * 1)
|
||||
err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
|
||||
if err := cluster[2].Command.Close(); err != nil {
|
||||
// prevent double-closing cluster[2] from the deferred Close above
|
||||
disabled, cluster := cluster[2], cluster[:2]
|
||||
|
||||
if err := disabled.Command.Close(); err != nil {
|
||||
t.Fatalf("closing third node: %v", err)
|
||||
}
|
||||
|
||||
if cluster[0].API.State() != pilosa.ClusterStateDegraded {
|
||||
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
|
||||
err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
|
||||
if _, err := cluster[0].API.RemoveNode(cluster[2].API.Node().ID); err != nil {
|
||||
if _, err := cluster[0].API.RemoveNode(disabled.API.Node().ID); err != nil {
|
||||
t.Fatalf("removing failed node: %v", err)
|
||||
}
|
||||
|
||||
if cluster[0].API.State() != pilosa.ClusterStateNormal {
|
||||
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
|
||||
err = cluster.AwaitCoordinatorState(pilosa.ClusterStateNormal, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("removing disabled node: %v", err)
|
||||
}
|
||||
|
||||
hosts := cluster[0].API.Hosts(context.Background())
|
||||
|
|
@ -824,16 +810,10 @@ func TestRemoveConcurrentIndexCreation(t *testing.T) {
|
|||
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)
|
||||
defer cluster.Close()
|
||||
err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
|
||||
errc := make(chan error)
|
||||
|
|
@ -846,11 +826,9 @@ func TestRemoveConcurrentIndexCreation(t *testing.T) {
|
|||
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())
|
||||
}
|
||||
err = cluster.AwaitCoordinatorState(pilosa.ClusterStateNormal, 100*time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
|
||||
hosts := cluster[0].API.Hosts(context.Background())
|
||||
|
|
|
|||
|
|
@ -16,9 +16,9 @@ package test
|
|||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"path"
|
||||
"runtime"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
|
|
@ -140,10 +140,53 @@ func (c Cluster) Close() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// MustNewCluster creates a new cluster
|
||||
// AwaitState waits for the cluster coordinator (assumed to be the first
|
||||
// node) to reach a specified state.
|
||||
func (c Cluster) AwaitCoordinatorState(expectedState string, timeout time.Duration) error {
|
||||
if len(c) < 1 {
|
||||
return errors.New("can't await coordinator state on an empty cluster")
|
||||
}
|
||||
return c[:1].AwaitState(expectedState, timeout)
|
||||
}
|
||||
|
||||
// ExceptionalState returns an error if any node in the cluster is not
|
||||
// in the expected state.
|
||||
func (c Cluster) ExceptionalState(expectedState string) error {
|
||||
for _, node := range c {
|
||||
state := node.API.State()
|
||||
if state != expectedState {
|
||||
return fmt.Errorf("node %q: state %s", node.ID(), state)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AwaitState waits for the whole cluster to reach a specified state.
|
||||
func (c Cluster) AwaitState(expectedState string, timeout time.Duration) (err error) {
|
||||
if len(c) < 1 {
|
||||
return errors.New("can't await state of an empty cluster")
|
||||
}
|
||||
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 = c.ExceptionalState(expectedState); err == nil {
|
||||
return err
|
||||
}
|
||||
time.Sleep(1 * time.Millisecond)
|
||||
}
|
||||
return fmt.Errorf("waited %v for cluster to reach state %q: %v",
|
||||
elapsed, expectedState, err)
|
||||
}
|
||||
|
||||
// MustNewCluster creates a new cluster. If opts contains only one
|
||||
// slice of command options, those options are used with every node.
|
||||
// If it is empty, default options are used. Otherwise, it must contain size
|
||||
// slices of command options, which are used with corresponding nodes.
|
||||
func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster {
|
||||
tb.Helper()
|
||||
c, err := newCluster(size, opts...)
|
||||
c, err := newCluster(tb, size, opts...)
|
||||
if err != nil {
|
||||
tb.Fatalf("new cluster: %v", err)
|
||||
}
|
||||
|
|
@ -151,7 +194,7 @@ func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu
|
|||
}
|
||||
|
||||
// newCluster creates a new cluster
|
||||
func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
if size == 0 {
|
||||
return nil, errors.New("cluster must contain at least one node")
|
||||
}
|
||||
|
|
@ -160,26 +203,7 @@ func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
|||
}
|
||||
|
||||
cluster := make(Cluster, size)
|
||||
// try to find a Test function to use as the "name" for our node.
|
||||
name := "node"
|
||||
callers := make([]uintptr, 10)
|
||||
n := runtime.Callers(2, callers)
|
||||
callers = callers[:n]
|
||||
for _, pc := range callers {
|
||||
fn := runtime.FuncForPC(pc)
|
||||
if fn != nil {
|
||||
fnName := fn.Name()
|
||||
sections := strings.Split(fnName, ".")
|
||||
if len(sections) > 1 {
|
||||
fnName = sections[2]
|
||||
}
|
||||
if strings.HasPrefix(fnName, "Test") {
|
||||
name = "test" + fnName[4:]
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
_ = name
|
||||
name := tb.Name()
|
||||
for i := 0; i < size; i++ {
|
||||
var commandOpts []server.CommandOption
|
||||
if len(opts) > 0 {
|
||||
|
|
@ -197,8 +221,8 @@ func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
|||
}
|
||||
|
||||
// runCluster creates and starts a new cluster
|
||||
func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
cluster, err := newCluster(size, opts...)
|
||||
func runCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
cluster, err := newCluster(tb, size, opts...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "new cluster")
|
||||
}
|
||||
|
|
@ -209,7 +233,8 @@ func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
|||
return cluster, nil
|
||||
}
|
||||
|
||||
// MustRunCluster creates and starts a new cluster
|
||||
// MustRunCluster creates and starts a new cluster. The opts parameter
|
||||
// is slightly magical; see MustNewCluster.
|
||||
func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster {
|
||||
// We want tests to default to using the in-memory translate store, so we
|
||||
// prepend opts with that functional option. If a different translate store
|
||||
|
|
@ -217,7 +242,7 @@ func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu
|
|||
opts = prependOpts(opts)
|
||||
|
||||
tb.Helper()
|
||||
c, err := runCluster(size, opts...)
|
||||
c, err := runCluster(tb, size, opts...)
|
||||
if err != nil {
|
||||
tb.Fatalf("run cluster: %v", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -195,6 +195,9 @@ func (m *Command) MustRecalculateCaches(tb testing.TB) {
|
|||
// URL returns the base URL string for accessing the running program.
|
||||
func (m *Command) URL() string { return m.API.Node().URI.String() }
|
||||
|
||||
// ID returns the node ID used by the running program.
|
||||
func (m *Command) ID() string { return m.API.Node().ID }
|
||||
|
||||
// Client returns a client to connect to the program.
|
||||
func (m *Command) Client() *http.InternalClient {
|
||||
return m.Server.InternalClient().(*http.InternalClient)
|
||||
|
|
@ -202,13 +205,41 @@ func (m *Command) Client() *http.InternalClient {
|
|||
|
||||
// Query executes a query against the program through the HTTP API.
|
||||
func (m *Command) Query(t *testing.T, index, rawQuery, query string) (string, error) {
|
||||
resp := Do(t, "POST", m.URL()+fmt.Sprintf("/index/%s/query?", index)+rawQuery, query)
|
||||
resp := Do(t, "POST", fmt.Sprintf("%s/index/%s/query?%s", m.URL(), index, rawQuery), query)
|
||||
if resp.StatusCode != gohttp.StatusOK {
|
||||
return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body)
|
||||
}
|
||||
return resp.Body, nil
|
||||
}
|
||||
|
||||
// Queryf is like Query, but with a format string.
|
||||
func (m *Command) Queryf(t *testing.T, index, rawQuery, query string, params ...interface{}) (string, error) {
|
||||
query = fmt.Sprintf(query, params...)
|
||||
resp := Do(t, "POST", fmt.Sprintf("%s/index/%s/query?%s", m.URL(), index, rawQuery), query)
|
||||
if resp.StatusCode != gohttp.StatusOK {
|
||||
return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body)
|
||||
}
|
||||
return resp.Body, nil
|
||||
}
|
||||
|
||||
// QueryExpect executes a query against the program through the HTTP API, and
|
||||
// confirms that it got an expected response.
|
||||
func (m *Command) QueryExpect(t *testing.T, index, rawQuery, query string, expected string) {
|
||||
resp := Do(t, "POST", fmt.Sprintf("%s/index/%s/query?%s", m.URL(), index, rawQuery), query)
|
||||
if resp.StatusCode != gohttp.StatusOK {
|
||||
t.Fatalf("invalid status from %s: %d, body=%q", m.ID(), resp.StatusCode, resp.Body)
|
||||
}
|
||||
last := len(resp.Body) - 1
|
||||
// Trim trailing newline so we don't need it to be present in the expected data.
|
||||
if last >= 0 && resp.Body[last] == '\n' {
|
||||
resp.Body = resp.Body[:last]
|
||||
}
|
||||
|
||||
if resp.Body != expected {
|
||||
t.Fatalf("node %s, query %q: expected response %s, got %s", m.ID(), query, expected, resp.Body)
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Command) QueryProtobuf(indexName string, query string) (*pilosa.QueryResponse, error) {
|
||||
var ser proto.Serializer
|
||||
queryReq := &pilosa.QueryRequest{
|
||||
|
|
@ -261,8 +292,6 @@ func (m *Command) RecalculateCaches(t *testing.T) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
////////////////////////////////////////////////////////////////////////////////////
|
||||
|
||||
// Do executes http.Do() with an http.NewRequest().
|
||||
func Do(t *testing.T, method, urlStr string, body string) *httpResponse {
|
||||
t.Helper()
|
||||
|
|
|
|||
47
translate.go
47
translate.go
|
|
@ -16,7 +16,9 @@ package pilosa
|
|||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"sync"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
|
|
@ -414,20 +416,43 @@ func (s *InMemTranslateStore) EntryReader(ctx context.Context, offset uint64) (T
|
|||
return newInMemTranslateEntryReader(ctx, s, offset), nil
|
||||
}
|
||||
|
||||
// WriteTo ensures that the TranslateStore implements io.WriterTo.
|
||||
// It's not important that this be implemented. It would really
|
||||
// only be necessary if we wanted to test cluster resizing while using
|
||||
// an in-memory translate store.
|
||||
// WriteTo implements io.WriterTo. It's not efficient or careful, but we
|
||||
// don't expect to use InMemTranslateStore much, it's mostly there to
|
||||
// avoid disk load during testing.
|
||||
func (s *InMemTranslateStore) WriteTo(w io.Writer) (int64, error) {
|
||||
return 0, nil // TODO: try to use ErrNotImplemented
|
||||
bytes, err := json.Marshal(s.keysByID)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
n, err := w.Write(bytes)
|
||||
return int64(n), err
|
||||
}
|
||||
|
||||
// ReadFrom ensures that the TranslateStore implements io.ReaderFrom.
|
||||
// It's not important that this be implemented. It would really
|
||||
// only be necessary if we wanted to test cluster resizing while using
|
||||
// an in-memory translate store.
|
||||
func (s *InMemTranslateStore) ReadFrom(r io.Reader) (int64, error) {
|
||||
return 0, nil // TODO: try to use ErrNotImplemented
|
||||
// ReadFrom implements io.ReaderFrom. It's not efficient or careful, but we
|
||||
// don't expect to use InMemTranslateStore much, it's mostly there to
|
||||
// avoid disk load during testing.
|
||||
func (s *InMemTranslateStore) ReadFrom(r io.Reader) (count int64, err error) {
|
||||
var bytes []byte
|
||||
bytes, err = ioutil.ReadAll(r)
|
||||
count = int64(len(bytes))
|
||||
if err != nil {
|
||||
return count, err
|
||||
}
|
||||
var keysByID map[uint64]string
|
||||
err = json.Unmarshal(bytes, &keysByID)
|
||||
if err != nil {
|
||||
return count, err
|
||||
}
|
||||
s.maxID = 0
|
||||
s.keysByID = keysByID
|
||||
s.idsByKey = make(map[string]uint64, len(s.keysByID))
|
||||
for k, v := range s.keysByID {
|
||||
s.idsByKey[v] = k
|
||||
if k > s.maxID {
|
||||
s.maxID = k
|
||||
}
|
||||
}
|
||||
return count, nil
|
||||
}
|
||||
|
||||
// MaxID returns the highest identifier in the store.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue