simplify test cluster setup by exposing gossip transport on server.Command

This commit is contained in:
Matt Jaffee 2018-06-27 10:46:37 -05:00
parent 41bdfceb57
commit fbe035ef25
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
3 changed files with 75 additions and 204 deletions

View file

@ -24,10 +24,9 @@ import (
"testing"
"time"
"golang.org/x/sync/errgroup"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/test"
"golang.org/x/sync/errgroup"
)
// Ensure program can send/receive broadcast messages.
@ -132,76 +131,34 @@ func TestClusterResize_EmptyNode(t *testing.T) {
// Ensure that a cluster of empty nodes comes up in a NORMAL state.
func TestClusterResize_EmptyNodes(t *testing.T) {
// Configure node0
m0 := test.NewMainWithCluster(true)
defer m0.Close()
clus := test.MustRunMainWithCluster(t, 2)
defer clus[0].Close()
defer clus[1].Close()
gossipHost := "localhost"
gossipPort := 0
seed, err := m0.RunWithTransport(gossipHost, gossipPort, []string{})
if err != nil {
t.Fatal(err)
}
// Configure node1
m1 := test.NewMainWithCluster(false)
defer m1.Close()
seed, err = m1.RunWithTransport(gossipHost, gossipPort, []string{seed})
if err != nil {
t.Fatal(err)
}
if m0.Server.Cluster.State() != pilosa.ClusterStateNormal {
t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State())
} else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal {
t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State())
if clus[0].Server.Cluster.State() != pilosa.ClusterStateNormal {
t.Fatalf("unexpected node0 cluster state: %s", clus[0].Server.Cluster.State())
} else if clus[1].Server.Cluster.State() != pilosa.ClusterStateNormal {
t.Fatalf("unexpected node1 cluster state: %s", clus[1].Server.Cluster.State())
}
}
// Ensure that adding a node correctly resizes the cluster.
func TestClusterResize_AddNode(t *testing.T) {
t.Run("NoData", func(t *testing.T) {
// Configure node0
m0 := test.NewMainWithCluster(true)
defer m0.Close()
clus := test.MustRunMainWithCluster(t, 2)
seed, err := m0.RunWithTransport("localhost", 0, []string{})
if err != nil {
t.Fatal(err)
}
// Configure node1
m1 := test.NewMainWithCluster(false)
defer m1.Close()
var eg errgroup.Group
eg.Go(func() error {
_, err = m1.RunWithTransport("localhost", 0, []string{seed})
if err != nil {
return err
}
return nil
})
if err := eg.Wait(); err != nil {
t.Fatal(err)
}
if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State())
} else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State())
if !checkClusterState(clus[0].Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", clus[0].Server.Cluster.State())
} else if !checkClusterState(clus[1].Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node1 cluster state: %s", clus[1].Server.Cluster.State())
}
})
t.Run("WithIndex", func(t *testing.T) {
// Configure node0
m0 := test.NewMainWithCluster(true)
m0 := test.MustRunMainWithCluster(t, 1)[0]
defer m0.Close()
seed, err := m0.RunWithTransport("localhost", 0, []string{})
if err != nil {
t.Fatal(err)
}
seed := m0.GossipAddress()
// Create a client for each node.
client0 := m0.Client()
@ -215,19 +172,13 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewMainWithCluster(false)
defer m1.Close()
var eg errgroup.Group
eg.Go(func() error {
_, err = m1.RunWithTransport("localhost", 0, []string{seed})
if err != nil {
return err
}
return nil
})
if err := eg.Wait(); err != nil {
t.Fatal(err)
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.Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State())
@ -236,19 +187,14 @@ func TestClusterResize_AddNode(t *testing.T) {
}
})
t.Run("ContinuousSlices", func(t *testing.T) {
// Configure node0
m0 := test.NewMainWithCluster(true)
m0 := test.MustRunMainWithCluster(t, 1)[0]
defer m0.Close()
seed, err := m0.RunWithTransport("localhost", 0, []string{})
if err != nil {
t.Fatal(err)
}
seed := m0.GossipAddress()
// Create a client for each node.
client0 := m0.Client()
//client1 := m1.Client()
// Create indexes and fields on one node.
if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
@ -267,19 +213,13 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewMainWithCluster(false)
defer m1.Close()
var eg errgroup.Group
eg.Go(func() error {
_, err = m1.RunWithTransport("localhost", 0, []string{seed})
if err != nil {
return err
}
return nil
})
if err := eg.Wait(); err != nil {
t.Fatal(err)
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.Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State())
@ -288,19 +228,14 @@ func TestClusterResize_AddNode(t *testing.T) {
}
})
t.Run("SkippedSlice", func(t *testing.T) {
// Configure node0
m0 := test.NewMainWithCluster(true)
m0 := test.MustRunMainWithCluster(t, 1)[0]
defer m0.Close()
seed, err := m0.RunWithTransport("localhost", 0, []string{})
if err != nil {
t.Fatal(err)
}
seed := m0.GossipAddress()
// Create a client for each node.
client0 := m0.Client()
//client1 := m1.Client()
// Create indexes and fields on one node.
if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
@ -319,19 +254,13 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewMainWithCluster(false)
defer m1.Close()
var eg errgroup.Group
eg.Go(func() error {
_, err = m1.RunWithTransport("localhost", 0, []string{seed})
if err != nil {
return err
}
return nil
})
if err := eg.Wait(); err != nil {
t.Fatal(err)
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.Server.Cluster, pilosa.ClusterStateNormal, 1000) {
t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State())
@ -345,37 +274,37 @@ func TestClusterResize_AddNode(t *testing.T) {
func TestCluster_GossipMembership(t *testing.T) {
t.Run("Node0Down", func(t *testing.T) {
// Configure node0
m0 := test.NewMainWithCluster(true)
m0 := test.MustRunMainWithCluster(t, 1)[0]
defer m0.Close()
seed, err := m0.RunWithTransport("localhost", 0, []string{})
if err != nil {
t.Fatal(err)
}
seed := m0.GossipAddress()
var eg errgroup.Group
// Configure node1
m1 := test.NewMainWithCluster(false)
defer m1.Close()
var eg errgroup.Group
eg.Go(func() error {
m1.Config.Gossip.Port = "0"
// Pass invalid seed as first in list
_, err := m1.RunWithTransport("localhost", 0, []string{"http://localhost:8765", seed})
m1.Config.Gossip.Seeds = []string{"http://localhost:8765", seed}
err := m1.Start()
if err != nil {
return err
t.Fatalf("starting second main: %v", err)
}
return nil
})
// Configure node2
// Configure node1
m2 := test.NewMainWithCluster(false)
defer m2.Close()
eg.Go(func() error {
// Pass invalid seed as last in list
_, err := m2.RunWithTransport("localhost", 0, []string{seed, "http://localhost:8765"})
m2.Config.Gossip.Port = "0"
// Pass invalid seed as first in list
m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"}
err := m2.Start()
if err != nil {
return err
t.Fatalf("starting second main: %v", err)
}
return nil
})

View file

@ -60,7 +60,7 @@ type Command struct {
Config *Config
// Gossip transport
GossipTransport *gossip.Transport
gossipTransport *gossip.Transport
// Standard input/output
*pilosa.CmdIO
@ -310,21 +310,16 @@ func (m *Command) SetupNetworking() error {
// get the host portion of addr to use for binding
gossipHost := m.Server.URI.Host()
var transport *gossip.Transport
if m.GossipTransport != nil {
transport = m.GossipTransport
} else {
transport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger())
if err != nil {
return errors.Wrap(err, "getting transport")
}
m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger())
if err != nil {
return errors.Wrap(err, "getting transport")
}
gossipMemberSet, err := gossip.NewGossipMemberSet(
m.Config.Gossip,
m.Server,
gossip.WithLogger(m.logger.Logger()),
gossip.WithTransport(transport),
gossip.WithTransport(m.gossipTransport),
)
if err != nil {
return errors.Wrap(err, "getting memberset")
@ -332,6 +327,13 @@ func (m *Command) SetupNetworking() error {
return errors.Wrap(gossipMemberSet.Open(), "opening gossip memberset")
}
// GossipTransport allows a caller to return the gossip transport created when
// setting up the GossipMemberSet. This is useful if one needs to determine the
// allocated ephemeral port programmatically. (usually used in tests)
func (m *Command) GossipTransport() *gossip.Transport {
return m.gossipTransport
}
// Close shuts down the server.
func (m *Command) Close() error {
var logErr error

View file

@ -19,14 +19,12 @@ import (
"fmt"
"io"
"io/ioutil"
"log"
gohttp "net/http"
"os"
"strings"
"testing"
"time"
"github.com/pilosa/pilosa/gossip"
"github.com/pilosa/pilosa/http"
"github.com/pilosa/pilosa/server"
"github.com/pilosa/pilosa/toml"
@ -59,6 +57,13 @@ func OptAllowedOrigins(origins []string) server.CommandOption {
}
}
// GossipAddress returns the address on which gossip is listening after a Main
// has been setup. Useful to pass as a seed to other nodes when creating and
// testing clusters.
func (m *Main) GossipAddress() string {
return m.GossipTransport().URI.String()
}
// NewMain returns a new instance of Main with a temporary data directory and random port.
func NewMain(opts ...server.CommandOption) *Main {
path, err := ioutil.TempDir("", "pilosa-")
@ -116,25 +121,20 @@ func runMainWithCluster(size int, opts ...[]server.CommandOption) ([]*Main, erro
}
mains := make([]*Main, size)
gossipHost := "localhost"
gossipPort := 0
var err error
var gossipSeeds = make([]string, size)
for i := 0; i < size; i++ {
var commandOpts []server.CommandOption
if len(opts) > 0 {
commandOpts = opts[i%len(opts)]
}
m := NewMainWithCluster(i == 0, commandOpts...)
m.Config.Cluster.Disabled = false
m.Config.Gossip.Port = "0"
m.Config.Gossip.Seeds = gossipSeeds[:i]
gossipSeeds[i], err = m.RunWithTransport(gossipHost, gossipPort, gossipSeeds[:i])
if err != nil {
return nil, errors.Wrap(err, "RunWithTransport")
if err := m.Start(); err != nil {
return nil, errors.Wrapf(err, "Starting server %d", i)
}
gossipSeeds[i] = m.GossipTransport().URI.String()
mains[i] = m
}
@ -179,66 +179,6 @@ func (m *Main) Reopen() error {
return nil
}
// RunWithTransport runs Main and returns the dynamically allocated gossip port.
func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (seed string, err error) {
defer close(m.Started)
/*
TEST:
- SetupServer (just static settings from config)
- OpenListener (sets Server.Name to use in gossip)
- NewTransport (gossip)
- SetupNetworking (does the gossip or static stuff) - uses Server.Name
- Open server
PRODUCTION:
- SetupServer (just static settings from config)
- SetupNetworking (does the gossip or static stuff) - calls NewTransport
- Open server - calls OpenListener
*/
// Open gossip transport to use in SetupServer.
transport, err := gossip.NewTransport(host, bindPort, nil)
if err != nil {
return seed, err
}
m.GossipTransport = transport
if len(joinSeeds) != 0 {
m.Config.Gossip.Seeds = joinSeeds
} else {
m.Config.Gossip.Seeds = []string{transport.URI.String()}
}
seed = transport.URI.String()
// SetupServer
m.Config.Cluster.Disabled = false
err = m.SetupServer()
if err != nil {
return seed, err
}
// SetupNetworking
err = m.SetupNetworking()
if err != nil {
return seed, err
}
go func() {
err := m.Handler.Serve()
if err != nil {
log.Printf("Handler serve error: %v", err)
}
}()
// Initialize server.
err = m.Server.Open()
if err != nil {
return seed, err
}
return seed, nil
}
// URL returns the base URL string for accessing the running program.
func (m *Main) URL() string { return "http://" + m.Server.Addr().String() }