diff --git a/gossip/gossip.go b/gossip/gossip.go index 58c344b39..3afdabe31 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -26,6 +26,7 @@ import ( "github.com/hashicorp/memberlist" "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" + "github.com/pkg/errors" ) // Ensure GossipNodeSet implements interfaces. @@ -76,7 +77,7 @@ func (g *GossipNodeSet) Open() error { } ml, err := memberlist.Create(g.config.memberlistConfig) if err != nil { - return err + return errors.Wrap(err, "creating memberlist") } g.memberlist = ml g.broadcasts = &memberlist.TransmitLimitedQueue{ @@ -90,7 +91,7 @@ func (g *GossipNodeSet) Open() error { nodes := []*pilosa.Node{&pilosa.Node{Scheme: "gossip", Host: g.config.gossipSeed}} //TODO: support a list of seeds err = g.joinWithRetry(pilosa.Nodes(nodes).Hosts()) if err != nil { - return err + return errors.Wrap(err, "joinWithRetry") } return nil } diff --git a/test/pilosa.go b/test/pilosa.go index cb3f7a07a..9c541a6b8 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -8,41 +8,113 @@ import ( "testing" "github.com/pilosa/pilosa/server" + "github.com/pkg/errors" ) func MustNewRunningServer(t *testing.T) *server.Command { - s := server.NewCommand(&bytes.Buffer{}, ioutil.Discard, ioutil.Discard) - s.Config.Bind = ":0" - port := strconv.Itoa(MustOpenPort(t)) - s.Config.GossipPort = port - s.Config.GossipSeed = "localhost:" + port - td, err := ioutil.TempDir("", "") + s, err := newServer() if err != nil { - t.Fatalf("error creating temp data directory: %v", err) + t.Fatalf("getting new server: %v", err) } - s.Config.DataDir = td + err = s.Run() if err != nil { - t.Fatalf("error running new pilosa server: %v", err) + t.Fatalf("running new pilosa server: %v", err) } return s } -func MustOpenPort(t *testing.T) int { +func newServer() (*server.Command, error) { + s := server.NewCommand(&bytes.Buffer{}, ioutil.Discard, ioutil.Discard) + + port, err := openPort() + if err != nil { + return nil, errors.Wrap(err, "getting port") + } + s.Config.Bind = "localhost:" + strconv.Itoa(port) + + gport, err := openPort() + if err != nil { + return nil, errors.Wrap(err, "getting gossip port") + } + s.Config.GossipPort = strconv.Itoa(gport) + + s.Config.GossipSeed = "localhost:" + s.Config.GossipPort + s.Config.Cluster.Type = "gossip" + td, err := ioutil.TempDir("", "") + if err != nil { + return nil, errors.Wrap(err, "temp dir") + } + s.Config.DataDir = td + return s, nil +} + +func openPort() (int, error) { addr, err := net.ResolveTCPAddr("tcp", ":0") if err != nil { - t.Fatalf("resolving new port addr: %v", err) + return 0, errors.Wrap(err, "resolving new port addr") } - l, err := net.ListenTCP("tcp", addr) if err != nil { - t.Fatalf("listening to get new port: %v", err) + return 0, errors.Wrap(err, "listening to get new port") } - defer func() { - err := l.Close() - if err != nil { - t.Logf("error closing listener in MustOpenPort: %v", err) - } - }() - return l.Addr().(*net.TCPAddr).Port + port := l.Addr().(*net.TCPAddr).Port + err = l.Close() + if err != nil { + return port, errors.Wrap(err, "closing listener") + } + return port, nil + +} + +func MustOpenPort(t *testing.T) int { + port, err := openPort() + if err != nil { + t.Fatalf("allocating new port: %v", err) + } + return port +} + +type Cluster struct { + Servers []*server.Command +} + +func MustNewServerCluster(t *testing.T, size int) *Cluster { + cluster, err := NewServerCluster(size) + if err != nil { + t.Fatalf("new cluster: %v", err) + } + return cluster +} + +func NewServerCluster(size int) (cluster *Cluster, err error) { + cluster = &Cluster{ + Servers: make([]*server.Command, size), + } + hosts := make([]string, size) + for i := 0; i < size; i++ { + s, err := newServer() + if err != nil { + return nil, errors.Wrap(err, "new server") + } + cluster.Servers[i] = s + hosts[i] = s.Config.Bind + s.Config.GossipSeed = cluster.Servers[0].Config.GossipSeed + + } + + for _, s := range cluster.Servers { + s.Config.Cluster.Hosts = hosts + } + for i, s := range cluster.Servers { + err := s.Run() + if err != nil { + for j := 0; j <= i; j++ { + cluster.Servers[j].Close() + } + return nil, errors.Wrapf(err, "starting server %d of %d. Config: %#v", i+1, size, s.Config) + } + } + + return cluster, nil } diff --git a/test/pilosa_test.go b/test/pilosa_test.go new file mode 100644 index 000000000..78090a896 --- /dev/null +++ b/test/pilosa_test.go @@ -0,0 +1,49 @@ +package test_test + +import ( + "net/http" + "testing" + + "encoding/json" + + "github.com/pilosa/pilosa/test" +) + +func TestNewCluster(t *testing.T) { + cluster := test.MustNewServerCluster(t, 3) + response, err := http.Get("http://" + cluster.Servers[0].Server.Addr().String() + "/status") + if err != nil { + t.Fatalf("getting schema: %v", err) + } + dec := json.NewDecoder(response.Body) + a := StatusResp{} + err = dec.Decode(&a) + if err != nil { + t.Fatalf("decoding status response: %v", err) + } + + bytes, err := json.MarshalIndent(a, "", " ") + if err != nil { + t.Fatalf("encoding: %v", err) + } + + if len(a.Status.Nodes) != 3 { + t.Fatalf("wrong number of nodes in status: %s", bytes) + } + + for i, node := range a.Status.Nodes { + if node.State != "UP" { + t.Fatalf("node %d should be up but is %s", i, node.State) + } + } +} + +type StatusResp struct { + Status struct { + Nodes []struct { + Host string + Schema string + State string + } + } `json:"status"` +}