mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 09:05:55 +00:00
Merge pull request #458 from seebs/eaddrinuse
Fix very-sporadic EADDRINUSE failures in CI testing (and related cluster test issues)
This commit is contained in:
commit
ab6a3aff10
8 changed files with 478 additions and 377 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())
|
||||
|
|
|
|||
253
test/cluster.go
253
test/cluster.go
|
|
@ -14,7 +14,260 @@
|
|||
|
||||
package test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"path"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
"github.com/pilosa/pilosa/v2/server"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
// modHasher represents a simple, mod-based hashing.
|
||||
type ModHasher struct{}
|
||||
|
||||
func (*ModHasher) Hash(key uint64, n int) int { return int(key) % n }
|
||||
|
||||
// Cluster represents a Pilosa cluster (multiple Command instances)
|
||||
type Cluster []*Command
|
||||
|
||||
// Query executes an API.Query through one of the cluster's node's API. It fails
|
||||
// the test if there is an error.
|
||||
func (c Cluster) Query(t testing.TB, index, query string) pilosa.QueryResponse {
|
||||
t.Helper()
|
||||
if len(c) == 0 {
|
||||
t.Fatal("must have at least one node in cluster to query")
|
||||
}
|
||||
|
||||
return c[0].QueryAPI(t, &pilosa.QueryRequest{Index: index, Query: query})
|
||||
}
|
||||
|
||||
func (c Cluster) ImportBits(t testing.TB, index, field string, rowcols [][2]uint64) {
|
||||
t.Helper()
|
||||
byShard := make(map[uint64][][2]uint64)
|
||||
for _, rowcol := range rowcols {
|
||||
shard := rowcol[1] / pilosa.ShardWidth
|
||||
byShard[shard] = append(byShard[shard], rowcol)
|
||||
}
|
||||
|
||||
for shard, bits := range byShard {
|
||||
rowIDs := make([]uint64, len(bits))
|
||||
colIDs := make([]uint64, len(bits))
|
||||
for i, bit := range bits {
|
||||
rowIDs[i] = bit[0]
|
||||
colIDs[i] = bit[1]
|
||||
}
|
||||
nodes, err := c[0].API.ShardNodes(context.Background(), index, shard)
|
||||
if err != nil {
|
||||
t.Fatalf("getting shard nodes: %v", err)
|
||||
}
|
||||
// TODO won't be necessary to do all nodes once that works hits
|
||||
// (travis) this TODO is not clear to me, but I think it's
|
||||
// suggesting that elsewhere we would support importing to a
|
||||
// single node, regardless of where the data ends up.
|
||||
for _, node := range nodes {
|
||||
for _, com := range c {
|
||||
if com.API.Node().ID != node.ID {
|
||||
continue
|
||||
}
|
||||
err := com.API.Import(context.Background(), &pilosa.ImportRequest{
|
||||
Index: index,
|
||||
Field: field,
|
||||
Shard: shard,
|
||||
RowIDs: rowIDs,
|
||||
ColumnIDs: colIDs,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("importing data: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// CreateField creates the index (if necessary) and field specified.
|
||||
func (c Cluster) CreateField(t testing.TB, index string, iopts pilosa.IndexOptions, field string, fopts ...pilosa.FieldOption) *pilosa.Field {
|
||||
t.Helper()
|
||||
idx, err := c[0].API.CreateIndex(context.Background(), index, iopts)
|
||||
if err != nil && !strings.Contains(err.Error(), "index already exists") {
|
||||
t.Fatalf("creating index: %v", err)
|
||||
} else if err != nil { // index exists
|
||||
idx, err = c[0].API.Index(context.Background(), index)
|
||||
if err != nil {
|
||||
t.Fatalf("getting index: %v", err)
|
||||
}
|
||||
}
|
||||
if idx.Options() != iopts {
|
||||
t.Logf("existing index options:\n%v\ndon't match given opts:\n%v\n in pilosa/test.Cluster.CreateField", idx.Options(), iopts)
|
||||
}
|
||||
|
||||
f, err := c[0].API.CreateField(context.Background(), index, field, fopts...)
|
||||
// we'll assume the field doesn't exist because checking if the options
|
||||
// match seems painful.
|
||||
if err != nil {
|
||||
t.Fatalf("creating field: %v", err)
|
||||
}
|
||||
return f
|
||||
}
|
||||
|
||||
// Start runs a Cluster
|
||||
func (c Cluster) Start() error {
|
||||
var gossipSeeds = make([]string, len(c))
|
||||
for i, cc := range c {
|
||||
cc.Config.Gossip.Port = "0"
|
||||
cc.Config.Gossip.Seeds = gossipSeeds[:i]
|
||||
if err := cc.Start(); err != nil {
|
||||
return errors.Wrapf(err, "starting server %d", i)
|
||||
}
|
||||
gossipSeeds[i] = cc.GossipAddress()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Stop stops a Cluster
|
||||
func (c Cluster) Close() error {
|
||||
for i, cc := range c {
|
||||
if err := cc.Close(); err != nil {
|
||||
return errors.Wrapf(err, "stopping server %d", i)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// 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(tb, size, opts...)
|
||||
if err != nil {
|
||||
tb.Fatalf("new cluster: %v", err)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// newCluster creates a new cluster
|
||||
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")
|
||||
}
|
||||
if len(opts) != size && len(opts) != 0 && len(opts) != 1 {
|
||||
return nil, errors.New("Slice of CommandOptions must be of length 0, 1, or equal to the number of cluster nodes")
|
||||
}
|
||||
|
||||
cluster := make(Cluster, size)
|
||||
name := tb.Name()
|
||||
for i := 0; i < size; i++ {
|
||||
var commandOpts []server.CommandOption
|
||||
if len(opts) > 0 {
|
||||
commandOpts = opts[i%len(opts)]
|
||||
}
|
||||
m := NewCommandNode(i == 0, commandOpts...)
|
||||
err := ioutil.WriteFile(path.Join(m.Config.DataDir, ".id"), []byte(name+"_"+strconv.Itoa(i)), 0600)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "writing node id")
|
||||
}
|
||||
cluster[i] = m
|
||||
}
|
||||
|
||||
return cluster, nil
|
||||
}
|
||||
|
||||
// runCluster creates and starts a new cluster
|
||||
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")
|
||||
}
|
||||
|
||||
if err = cluster.Start(); err != nil {
|
||||
return nil, errors.Wrap(err, "starting cluster")
|
||||
}
|
||||
return cluster, nil
|
||||
}
|
||||
|
||||
// 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
|
||||
// has been specified, it will override this one.
|
||||
opts = prependOpts(opts)
|
||||
|
||||
tb.Helper()
|
||||
c, err := runCluster(tb, size, opts...)
|
||||
if err != nil {
|
||||
tb.Fatalf("run cluster: %v", err)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// prependOpts applies prependTestServerOpts to each of the ops (one per
|
||||
// node, or one for the entire cluser).
|
||||
func prependOpts(opts [][]server.CommandOption) [][]server.CommandOption {
|
||||
if len(opts) == 0 {
|
||||
opts = [][]server.CommandOption{
|
||||
prependTestServerOpts([]server.CommandOption{}),
|
||||
}
|
||||
} else {
|
||||
for i := range opts {
|
||||
opts[i] = prependTestServerOpts(opts[i])
|
||||
}
|
||||
}
|
||||
return opts
|
||||
}
|
||||
|
||||
// prependTestServerOpts prepends opts with the OpenInMemTranslateStore.
|
||||
func prependTestServerOpts(opts []server.CommandOption) []server.CommandOption {
|
||||
defaultOpts := []server.CommandOption{
|
||||
server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore), pilosa.OptServerNodeDownRetries(5, 100*time.Millisecond)),
|
||||
}
|
||||
return append(defaultOpts, opts...)
|
||||
}
|
||||
|
|
|
|||
231
test/pilosa.go
231
test/pilosa.go
|
|
@ -21,9 +21,7 @@ import (
|
|||
"io/ioutil"
|
||||
gohttp "net/http"
|
||||
"os"
|
||||
"path"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
|
@ -32,7 +30,6 @@ import (
|
|||
"github.com/pilosa/pilosa/v2/encoding/proto"
|
||||
"github.com/pilosa/pilosa/v2/http"
|
||||
"github.com/pilosa/pilosa/v2/server"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
////////////////////////////////////////////////////////////////////////////////////
|
||||
|
|
@ -198,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)
|
||||
|
|
@ -205,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{
|
||||
|
|
@ -264,201 +292,6 @@ func (m *Command) RecalculateCaches(t *testing.T) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// Cluster represents a Pilosa cluster (multiple Command instances)
|
||||
type Cluster []*Command
|
||||
|
||||
// Query executes an API.Query through one of the cluster's node's API. It fails
|
||||
// the test if there is an error.
|
||||
func (c Cluster) Query(t testing.TB, index, query string) pilosa.QueryResponse {
|
||||
t.Helper()
|
||||
if len(c) == 0 {
|
||||
t.Fatal("must have at least one node in cluster to query")
|
||||
}
|
||||
|
||||
return c[0].QueryAPI(t, &pilosa.QueryRequest{Index: index, Query: query})
|
||||
}
|
||||
|
||||
func (c Cluster) ImportBits(t testing.TB, index, field string, rowcols [][2]uint64) {
|
||||
t.Helper()
|
||||
byShard := make(map[uint64][][2]uint64)
|
||||
for _, rowcol := range rowcols {
|
||||
shard := rowcol[1] / pilosa.ShardWidth
|
||||
byShard[shard] = append(byShard[shard], rowcol)
|
||||
}
|
||||
|
||||
for shard, bits := range byShard {
|
||||
rowIDs := make([]uint64, len(bits))
|
||||
colIDs := make([]uint64, len(bits))
|
||||
for i, bit := range bits {
|
||||
rowIDs[i] = bit[0]
|
||||
colIDs[i] = bit[1]
|
||||
}
|
||||
nodes, err := c[0].API.ShardNodes(context.Background(), index, shard)
|
||||
if err != nil {
|
||||
t.Fatalf("getting shard nodes: %v", err)
|
||||
}
|
||||
// TODO won't be necessary to do all nodes once that works hits
|
||||
// (travis) this TODO is not clear to me, but I think it's
|
||||
// suggesting that elsewhere we would support importing to a
|
||||
// single node, regardless of where the data ends up.
|
||||
for _, node := range nodes {
|
||||
for _, com := range c {
|
||||
if com.API.Node().ID != node.ID {
|
||||
continue
|
||||
}
|
||||
err := com.API.Import(context.Background(), &pilosa.ImportRequest{
|
||||
Index: index,
|
||||
Field: field,
|
||||
Shard: shard,
|
||||
RowIDs: rowIDs,
|
||||
ColumnIDs: colIDs,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("importing data: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// CreateField creates the index (if necessary) and field specified.
|
||||
func (c Cluster) CreateField(t testing.TB, index string, iopts pilosa.IndexOptions, field string, fopts ...pilosa.FieldOption) *pilosa.Field {
|
||||
t.Helper()
|
||||
idx, err := c[0].API.CreateIndex(context.Background(), index, iopts)
|
||||
if err != nil && !strings.Contains(err.Error(), "index already exists") {
|
||||
t.Fatalf("creating index: %v", err)
|
||||
} else if err != nil { // index exists
|
||||
idx, err = c[0].API.Index(context.Background(), index)
|
||||
if err != nil {
|
||||
t.Fatalf("getting index: %v", err)
|
||||
}
|
||||
}
|
||||
if idx.Options() != iopts {
|
||||
t.Logf("existing index options:\n%v\ndon't match given opts:\n%v\n in pilosa/test.Cluster.CreateField", idx.Options(), iopts)
|
||||
}
|
||||
|
||||
f, err := c[0].API.CreateField(context.Background(), index, field, fopts...)
|
||||
// we'll assume the field doesn't exist because checking if the options
|
||||
// match seems painful.
|
||||
if err != nil {
|
||||
t.Fatalf("creating field: %v", err)
|
||||
}
|
||||
return f
|
||||
}
|
||||
|
||||
// Start runs a Cluster
|
||||
func (c Cluster) Start() error {
|
||||
var gossipSeeds = make([]string, len(c))
|
||||
for i, cc := range c {
|
||||
cc.Config.Gossip.Port = "0"
|
||||
cc.Config.Gossip.Seeds = gossipSeeds[:i]
|
||||
if err := cc.Start(); err != nil {
|
||||
return errors.Wrapf(err, "starting server %d", i)
|
||||
}
|
||||
gossipSeeds[i] = cc.GossipAddress()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Stop stops a Cluster
|
||||
func (c Cluster) Close() error {
|
||||
for i, cc := range c {
|
||||
if err := cc.Close(); err != nil {
|
||||
return errors.Wrapf(err, "stopping server %d", i)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// MustNewCluster creates a new cluster
|
||||
func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster {
|
||||
tb.Helper()
|
||||
c, err := newCluster(size, opts...)
|
||||
if err != nil {
|
||||
tb.Fatalf("new cluster: %v", err)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// newCluster creates a new cluster
|
||||
func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
if size == 0 {
|
||||
return nil, errors.New("cluster must contain at least one node")
|
||||
}
|
||||
if len(opts) != size && len(opts) != 0 && len(opts) != 1 {
|
||||
return nil, errors.New("Slice of CommandOptions must be of length 0, 1, or equal to the number of cluster nodes")
|
||||
}
|
||||
|
||||
cluster := make(Cluster, size)
|
||||
for i := 0; i < size; i++ {
|
||||
var commandOpts []server.CommandOption
|
||||
if len(opts) > 0 {
|
||||
commandOpts = opts[i%len(opts)]
|
||||
}
|
||||
m := NewCommandNode(i == 0, commandOpts...)
|
||||
err := ioutil.WriteFile(path.Join(m.Config.DataDir, ".id"), []byte("node"+strconv.Itoa(i)), 0600)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "writing node id")
|
||||
}
|
||||
cluster[i] = m
|
||||
}
|
||||
|
||||
return cluster, nil
|
||||
}
|
||||
|
||||
// runCluster creates and starts a new cluster
|
||||
func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
cluster, err := newCluster(size, opts...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "new cluster")
|
||||
}
|
||||
|
||||
if err = cluster.Start(); err != nil {
|
||||
return nil, errors.Wrap(err, "starting cluster")
|
||||
}
|
||||
return cluster, nil
|
||||
}
|
||||
|
||||
// MustRunCluster creates and starts a new cluster
|
||||
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
|
||||
// has been specified, it will override this one.
|
||||
opts = prependOpts(opts)
|
||||
|
||||
tb.Helper()
|
||||
c, err := runCluster(size, opts...)
|
||||
if err != nil {
|
||||
tb.Fatalf("run cluster: %v", err)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// prependOpts applies prependTestServerOpts to each of the ops (one per
|
||||
// node, or one for the entire cluser).
|
||||
func prependOpts(opts [][]server.CommandOption) [][]server.CommandOption {
|
||||
if len(opts) == 0 {
|
||||
opts = [][]server.CommandOption{
|
||||
prependTestServerOpts([]server.CommandOption{}),
|
||||
}
|
||||
} else {
|
||||
for i := range opts {
|
||||
opts[i] = prependTestServerOpts(opts[i])
|
||||
}
|
||||
}
|
||||
return opts
|
||||
}
|
||||
|
||||
// prependTestServerOpts prepends opts with the OpenInMemTranslateStore.
|
||||
func prependTestServerOpts(opts []server.CommandOption) []server.CommandOption {
|
||||
defaultOpts := []server.CommandOption{
|
||||
server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore), pilosa.OptServerNodeDownRetries(5, 100*time.Millisecond)),
|
||||
}
|
||||
return append(defaultOpts, opts...)
|
||||
}
|
||||
|
||||
////////////////////////////////////////////////////////////////////////////////////
|
||||
|
||||
// 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