diff --git a/test/cluster.go b/test/cluster.go index 88ab423a3..5c6a1ffc8 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -245,21 +245,24 @@ func (c *Cluster) CreateField(t testing.TB, index string, iopts pilosa.IndexOpti // Start runs a Cluster func (c *Cluster) Start() error { var eg errgroup.Group - // seedCh is a channel of host:port values to use - // as gossip seeds during startup. - seedCh := make(chan string, len(c.Nodes)) - for i, cc := range c.Nodes { - i := i - cc := cc - eg.Go(func() error { - // get the bind uri to use as the host portion of the gossip seed. - uri, err := pilosa.AddressWithDefaults(cc.Config.Bind) - if err != nil { - return errors.Wrap(err, "processing bind address") - } + err := port.GetPorts(func(ports []int) error { + portsCfg := GenPortsConfig(NewPorts(ports)) - if err := port.GetPort(func(p int) error { - cc.Config.Gossip.Port = fmt.Sprint(p) + // seedCh is a channel of host:port values to use + // as gossip seeds during startup. + seedCh := make(chan string, len(c.Nodes)) + for i, cc := range c.Nodes { + i := i + cc.Config.DisCo = portsCfg[i].DisCo + cc.Config.BindGRPC = portsCfg[i].BindGRPC + eg.Go(func() error { + // get the bind uri to use as the host portion of the gossip seed. + uri, err := pilosa.AddressWithDefaults(cc.Config.Bind) + if err != nil { + return errors.Wrap(err, "processing bind address") + } + + cc.Config.Gossip.Port = portsCfg[i].Gossip.Port gossipHost := uri.Host gossipPort := cc.Config.Gossip.Port @@ -280,18 +283,16 @@ func (c *Cluster) Start() error { } return nil - }, 10); err != nil { - return errors.Wrap(err, "getting gossip port") - } + }) + // fixes race on gossip: time.Sleep(time.Second) + } - return nil - }) - // fixes race on gossip: time.Sleep(time.Second) - } - err := eg.Wait() + return eg.Wait() + }, 4*len(c.Nodes), 10) if err != nil { return err } + return c.AwaitState(pilosa.ClusterStateNormal, 10*time.Second) } @@ -393,7 +394,7 @@ func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (*Clust return nil, errors.New("cluster must contain at least one node") } - opts = appendOpts(opts, GenDisCoConfig(size)) + //opts = appendOpts(opts, GenDisCoConfig(size)) 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") diff --git a/test/disco.go b/test/disco.go index 98650d902..c316a1f40 100644 --- a/test/disco.go +++ b/test/disco.go @@ -19,40 +19,41 @@ import ( "strings" "github.com/pilosa/pilosa/v2/etcd" + "github.com/pilosa/pilosa/v2/gossip" "github.com/pilosa/pilosa/v2/server" "github.com/pilosa/pilosa/v2/test/port" ) -//GenDisCoConfig creates specific configuration for etcd. -func GenDisCoConfig(clusterSize int) []*server.Config { - cfgs := make([]*server.Config, clusterSize) +type Ports struct { + Client, Peer int + Grpc, Gossip int //TODO remove +} - clusterURLs := make([]string, clusterSize) +//GenPortsConfig creates specific configuration for etcd. +func GenPortsConfig(ports []Ports) []*server.Config { + cfgs := make([]*server.Config, len(ports)) + clusterURLs := make([]string, len(ports)) for i := range cfgs { name := fmt.Sprintf("server%d", i) var lClientURL, lPeerURL string - err := port.GetPorts(func(ports []int) error { - lClientURL = fmt.Sprintf("http://localhost:%d", ports[0]) - lPeerURL = fmt.Sprintf("http://localhost:%d", ports[1]) + lClientURL = fmt.Sprintf("http://localhost:%d", ports[i].Client) + lPeerURL = fmt.Sprintf("http://localhost:%d", ports[i].Peer) - cfgs[i] = &server.Config{ - BindGRPC: port.ColonZeroString(ports[2]), - DisCo: etcd.Options{ - Name: name, - Dir: "", - ClusterName: "bartholemuuuuu", - LClientURL: lClientURL, - AClientURL: lClientURL, - LPeerURL: lPeerURL, - APeerURL: lPeerURL, - }, - } - - return nil - }, 3, 10) - if err != nil { - panic(err) + cfgs[i] = &server.Config{ + Gossip: gossip.Config{ + Port: fmt.Sprint(ports[i].Gossip), + }, + BindGRPC: port.ColonZeroString(ports[i].Grpc), + DisCo: etcd.Options{ + Name: name, + Dir: "", + ClusterName: "bartholemuuuuu", + LClientURL: lClientURL, + AClientURL: lClientURL, + LPeerURL: lPeerURL, + APeerURL: lPeerURL, + }, } clusterURLs[i] = fmt.Sprintf("%s=%s", name, lPeerURL) @@ -64,3 +65,17 @@ func GenDisCoConfig(clusterSize int) []*server.Config { return cfgs } + +func NewPorts(ports []int) []Ports { + var out []Ports + for i := 0; i < len(ports); i = i + 4 { + out = append(out, Ports{ + Client: ports[i], + Peer: ports[i+1], + Grpc: ports[i+2], + Gossip: ports[i+3], + }) + } + + return out +}