diff --git a/etcd/embed.go b/etcd/embed.go index 2cbcce20f..4c306fe5f 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -19,6 +19,7 @@ import ( "context" "fmt" "log" + "net" "path" "strings" "time" @@ -45,6 +46,9 @@ type Options struct { ClusterURL string `toml:"cluster-url"` ClusterName string `toml:"cluster-name"` HeartbeatTTL int64 `toml:"heartbeat-ttl"` + + LPeerSocket []*net.TCPListener + LClientSocket []*net.TCPListener } var ( @@ -124,6 +128,19 @@ func parseOptions(opt Options) *embed.Config { cfg.LPUrls = types.MustNewURLs([]string{opt.LPeerURL}) cfg.APUrls = types.MustNewURLs([]string{opt.APeerURL}) + lps := make([]*net.TCPListener, len(opt.LPeerSocket)) + copy(lps, opt.LPeerSocket) + cfg.LPeerSocket = lps + + lcs := make([]*net.TCPListener, len(opt.LPeerSocket)) + copy(lcs, opt.LClientSocket) + cfg.LClientSocket = lcs + + cfg.Logger = "zap" + cfg.ZapLoggerBuilder = func(*embed.Config) error { + return nil + } + if opt.InitCluster != "" { cfg.InitialCluster = opt.InitCluster cfg.ClusterState = embed.ClusterStateFlagNew diff --git a/go.mod b/go.mod index 5f1edb1b2..766945c7d 100644 --- a/go.mod +++ b/go.mod @@ -2,7 +2,7 @@ module github.com/pilosa/pilosa/v2 replace github.com/hashicorp/memberlist => github.com/pilosa/memberlist v0.1.4-0.20190415211605-f6512523c021 -replace go.etcd.io/etcd => github.com/molecula/etcd v0.0.0-20210108232729-18e95f2f5b93 +replace go.etcd.io/etcd => github.com/molecula/etcd v0.0.0-20210115113447-5d28bda617d2 require ( github.com/CAFxX/gcnotifier v0.0.0-20190112062741-224a280d589d diff --git a/go.sum b/go.sum index f5021bc1b..4fb4fc5e6 100644 --- a/go.sum +++ b/go.sum @@ -247,8 +247,8 @@ github.com/modern-go/reflect2 v1.0.1 h1:9f412s+6RmYXLWZSEzVVgPGK7C2PphHj5RJrvfx9 github.com/modern-go/reflect2 v1.0.1/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0= github.com/molecula/apophenia v0.0.0-20190827192002-68b7a14a478b h1:cZADDaNYM7xn/nklO3g198JerGQjadFuA0ofxBJgK0Y= github.com/molecula/apophenia v0.0.0-20190827192002-68b7a14a478b/go.mod h1:uXd1BiH7xLmgkhVmspdJLENv6uGWrTL/MQX2TN7Yz9s= -github.com/molecula/etcd v0.0.0-20210108232729-18e95f2f5b93 h1:9a+hOGmPrcJEfpK07rzeA0D+F99a+2iha5PfDXGLrbE= -github.com/molecula/etcd v0.0.0-20210108232729-18e95f2f5b93/go.mod h1:yVHk9ub3CSBatqGNg7GRmsnfLWtoW60w4eDYfh7vHDg= +github.com/molecula/etcd v0.0.0-20210115113447-5d28bda617d2 h1:pkzCVLSrFQGVQv3raVGJw6aJCJdIZC/z59tUsSU1Zws= +github.com/molecula/etcd v0.0.0-20210115113447-5d28bda617d2/go.mod h1:1X1h4BZ44WjM0LJof1gKKLap1OA4RsicGCDRtACTkLI= github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223 h1:F9x/1yl3T2AeKLr2AMdilSD8+f9bvMnNN8VS5iDtovc= github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U= github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e h1:fD57ERR4JtEqsWbfPhv4DMiApHyliiK5xCTNVSPiaAs= @@ -367,8 +367,6 @@ github.com/zeebo/blake3 v0.0.4/go.mod h1:YOZo8A49yNqM0X/Y+JmDUZshJWLt1laHsNSn5ny github.com/zeebo/pcg v0.0.0-20181207190024-3cdc6b625a05 h1:4pW5fMvVkrgkMXdvIsVRRTs69DWYA8uNNQsu1stfVKU= github.com/zeebo/pcg v0.0.0-20181207190024-3cdc6b625a05/go.mod h1:Gr+78ptB0MwXxm//LBaEvBiaXY7hXJ6KGe2V32X2F6E= go.etcd.io/bbolt v1.3.2/go.mod h1:IbVyRI1SCnLcuJnV2u8VeU0CEYM7e686BmAb1XKL+uU= -go.etcd.io/bbolt v1.3.3 h1:MUGmc65QhB3pIlaQ5bB4LwqSj6GIonVJXpZiaKNyaKk= -go.etcd.io/bbolt v1.3.3/go.mod h1:IbVyRI1SCnLcuJnV2u8VeU0CEYM7e686BmAb1XKL+uU= go.etcd.io/bbolt v1.3.5 h1:XAzx9gjCb0Rxj7EoqcClPD1d5ZBxZJk0jbuoPHenBt0= go.etcd.io/bbolt v1.3.5/go.mod h1:G5EMThwa9y8QZGBClrRx5EY+Yw9kAhnjy3bSjsnlVTQ= go.opencensus.io v0.21.0/go.mod h1:mSImk1erAIZhrmZN+AvHh14ztQfjbGwt4TtuofqLduU= @@ -453,7 +451,6 @@ golang.org/x/sys v0.0.0-20190502145724-3ef323f4f1fd/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20190507160741-ecd444e8653b/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190606165138-5da285871e9c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190624142023-c5567b49c5d0/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20190826190057-c7b8b68b1456/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191005200804-aed5e4c7ecf9/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191220142924-d4481acd189f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= diff --git a/server/cluster_test.go b/server/cluster_test.go index 99cf6bfbc..75878f493 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -18,6 +18,7 @@ import ( "context" "encoding/json" "fmt" + "net" "net/http" "os" "reflect" @@ -185,8 +186,8 @@ func TestClusterResize_AddNode(t *testing.T) { m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -242,8 +243,8 @@ func TestClusterResize_AddNode(t *testing.T) { m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -298,8 +299,8 @@ func TestClusterResize_AddNode(t *testing.T) { m1 := test.NewCommandNode(t, false) m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -360,8 +361,8 @@ func TestClusterResize_AddNode(t *testing.T) { m1 := test.NewCommandNode(t, false) m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -416,8 +417,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -474,8 +475,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -538,8 +539,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -600,8 +601,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) m1.Config.Gossip.Seeds = []string{seed} - if err := port.GetPorts(func(ports []int) error { - portsCfg := test.GenPortsConfig(test.NewPorts(ports)) + if err := port.GetListeners(func(lsns []*net.TCPListener) error { + portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port m1.Config.DisCo = portsCfg[0].DisCo @@ -661,10 +662,12 @@ func TestCluster_GossipMembership(t *testing.T) { eg.Go(func() error { // Pass invalid seed as first in list m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"} - if err := port.GetPort(func(p int) error { + err := port.GetPort(func(p int) error { m2.Config.Gossip.Port = fmt.Sprintf("%d", p) return m2.Start() - }, 10); err != nil { + }, 10) + + if err != nil { t.Fatalf("starting second main: %v", err) } defer m2.Close() diff --git a/server/server.go b/server/server.go index defe38c25..a25830b8e 100644 --- a/server/server.go +++ b/server/server.go @@ -604,6 +604,10 @@ func (m *Command) Close() error { } } + // prevent the closed sockets from being re-injected into etcd. + m.Config.DisCo.LPeerSocket = nil + m.Config.DisCo.LClientSocket = nil + err := eg.Wait() _ = testhook.Closed(pilosa.NewAuditor(), m, nil) return errors.Wrap(err, "closing everything") diff --git a/test/cluster.go b/test/cluster.go index 41f6d0aad..22714739f 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -19,6 +19,7 @@ import ( "fmt" "io/ioutil" "math" + "net" "path" "strconv" "strings" @@ -269,40 +270,52 @@ 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 - err := port.GetPorts(func(ports []int) error { - portsCfg := GenPortsConfig(NewPorts(ports)) + err := port.GetListeners( - var gossipSeeds []string - for i, cc := range c.Nodes { - i := i - // 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") + func(lsns []*net.TCPListener) (err0 error) { + sliceOfPorts := NewPorts(lsns) + defer func() { + if err0 != nil { + // going to retry. Close the still open listeners + for _, ports := range sliceOfPorts { + _ = ports.Close() + } + } + }() + portsCfg := GenPortsConfig(sliceOfPorts) + + var gossipSeeds []string + for i, cc := range c.Nodes { + i := i + // 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 + + gossipSeeds = append(gossipSeeds, fmt.Sprintf("%s:%s", gossipHost, gossipPort)) } - cc.Config.Gossip.Port = portsCfg[i].Gossip.Port - gossipHost := uri.Host - gossipPort := cc.Config.Gossip.Port + for i, cc := range c.Nodes { + cc := cc + cc.Config.DisCo = portsCfg[i].DisCo + cc.Config.BindGRPC = portsCfg[i].BindGRPC - gossipSeeds = append(gossipSeeds, fmt.Sprintf("%s:%s", gossipHost, gossipPort)) - } + eg.Go(func() error { + fmt.Printf("DISCO CONFIG: %+v\n", cc.Config.DisCo) + cc.Config.Gossip.Seeds = gossipSeeds - for i, cc := range c.Nodes { - cc := cc - cc.Config.DisCo = portsCfg[i].DisCo - cc.Config.BindGRPC = portsCfg[i].BindGRPC + return cc.Start() + }) + } - eg.Go(func() error { - fmt.Printf("DISCO CONFIG: %+v\n", cc.Config.DisCo) - cc.Config.Gossip.Seeds = gossipSeeds + return eg.Wait() + }, 4*len(c.Nodes), 10) - return cc.Start() - }) - } - - return eg.Wait() - }, 4*len(c.Nodes), 10) if err != nil { return err } diff --git a/test/disco.go b/test/disco.go index afcfcf731..b0af11953 100644 --- a/test/disco.go +++ b/test/disco.go @@ -17,18 +17,33 @@ package test import ( "fmt" "io/ioutil" + "net" "strings" "time" "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" ) type Ports struct { - Client, Peer int - Grpc, Gossip int //TODO remove + LsnC *net.TCPListener + PortC int + + LsnP *net.TCPListener + PortP int + + Grpc int + Gossip int //TODO remove +} + +func (ports *Ports) Close() error { + err := ports.LsnC.Close() + err2 := ports.LsnP.Close() + if err != nil { + return err + } + return err2 } //GenPortsConfig creates specific configuration for etcd. @@ -38,9 +53,11 @@ func GenPortsConfig(ports []Ports) []*server.Config { for i := range cfgs { name := fmt.Sprintf("server%d", i) - var lClientURL, lPeerURL string - lClientURL = fmt.Sprintf("http://localhost:%d", ports[i].Client) - lPeerURL = fmt.Sprintf("http://localhost:%d", ports[i].Peer) + lsnC, portC := ports[i].LsnC, ports[i].PortC + lClientURL := fmt.Sprintf("http://localhost:%d", portC) + lsnP, portP := ports[i].LsnP, ports[i].PortP + lPeerURL := fmt.Sprintf("http://localhost:%d", portP) + discoDir := "" if d, err := ioutil.TempDir("/tmp", "disco."); err == nil { discoDir = d @@ -50,22 +67,24 @@ func GenPortsConfig(ports []Ports) []*server.Config { Gossip: gossip.Config{ Port: fmt.Sprint(ports[i].Gossip), }, - BindGRPC: port.ColonZeroString(ports[i].Grpc), + BindGRPC: fmt.Sprintf(":%d", ports[i].Grpc), DisCo: etcd.Options{ - Name: name, - Dir: discoDir, - ClusterName: "bartholemuuuuu", - LClientURL: lClientURL, - AClientURL: lClientURL, - LPeerURL: lPeerURL, - APeerURL: lPeerURL, - HeartbeatTTL: 5 * int64(time.Second), + Name: name, + Dir: discoDir, + ClusterName: "bartholemuuuuu", + LClientURL: lClientURL, + AClientURL: lClientURL, + LPeerURL: lPeerURL, + APeerURL: lPeerURL, + HeartbeatTTL: 5 * int64(time.Second), + LPeerSocket: []*net.TCPListener{lsnP}, + LClientSocket: []*net.TCPListener{lsnC}, }, } clusterURLs[i] = fmt.Sprintf("%s=%s", name, lPeerURL) fmt.Printf("\ndebug test/disco.go: on i=%v, GenPortsConfig Gossip: %v, DisCo.Client: %v, DisCo.Peer: %v, BindGRPC: %v\n", - i, ports[i].Gossip, ports[i].Client, ports[i].Peer, ports[i].Grpc) + i, ports[i].Gossip, portC, portP, ports[i].Grpc) } for i := range cfgs { cfgs[i].DisCo.InitCluster = strings.Join(clusterURLs, ",") @@ -74,15 +93,29 @@ func GenPortsConfig(ports []Ports) []*server.Config { return cfgs } -func NewPorts(ports []int) []Ports { +func NewPorts(lsn []*net.TCPListener) []Ports { var out []Ports - for i := 0; i < len(ports); i = i + 4 { + + n := len(lsn) + ports := make([]int, n) + for i := 0; i < n; i++ { + ports[i] = lsn[i].Addr().(*net.TCPAddr).Port + } + + for i := 0; i < n; i = i + 4 { out = append(out, Ports{ - Client: ports[i], - Peer: ports[i+1], + LsnC: lsn[i], + PortC: ports[i], + LsnP: lsn[i+1], + PortP: ports[i+1], + Grpc: ports[i+2], Gossip: ports[i+3], }) + // make Grpc and Gossip ports available to + // be rebound. + lsn[i+2].Close() + lsn[i+3].Close() } return out diff --git a/test/port/port_mapper.go b/test/port/port_mapper.go index b07151c55..b5c05a6d8 100644 --- a/test/port/port_mapper.go +++ b/test/port/port_mapper.go @@ -27,14 +27,16 @@ func ColonZeroString(port int) string { } func GetPort(wrapper func(int) error, retries int) error { - f := func(ports []int) error { return wrapper(ports[0]) } + f := func(ports []int) error { + return wrapper(ports[0]) + } return GetPorts(f, 1, retries) } func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error { for i := 0; i < retries; i++ { // get all requested ports - listeners := make([]net.Listener, requestedPorts) + listeners := make([]*net.TCPListener, requestedPorts) ports := make([]int, requestedPorts) for i := 0; i < requestedPorts; i++ { l, err := net.Listen("tcp", ":0") @@ -44,7 +46,7 @@ func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error { } ports[i] = l.Addr().(*net.TCPAddr).Port - listeners[i] = l + listeners[i] = l.(*net.TCPListener) } for _, l := range listeners { if err := l.Close(); err != nil { @@ -64,3 +66,32 @@ func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error { return nil } + +func GetListeners(wrapper func([]*net.TCPListener) error, requestedPorts, retries int) error { + for i := 0; i < retries; i++ { + // get all requested ports + listeners := make([]*net.TCPListener, requestedPorts) + ports := make([]int, requestedPorts) + for i := 0; i < requestedPorts; i++ { + l, err := net.Listen("tcp", ":0") + if err != nil { + log.Println("[port_mapper] error getting a free port", err) + return GetListeners(wrapper, requestedPorts, retries-1) + } + + ports[i] = l.Addr().(*net.TCPAddr).Port + listeners[i] = l.(*net.TCPListener) + } + // send to wrapper and check output error + err := wrapper(listeners) + if (err != nil) && (err == syscall.EADDRINUSE || strings.Contains(err.Error(), "address already in use")) { + log.Printf("[port_mapper: %+v] address already in use error calling the wrapper: %v\n", ports, err) + // only retry on address already in use error + continue + } + + return err + } + + return nil +} diff --git a/test/port/port_mapper_test.go b/test/port/port_mapper_test.go deleted file mode 100644 index 29cc8cdd8..000000000 --- a/test/port/port_mapper_test.go +++ /dev/null @@ -1,66 +0,0 @@ -// Copyright 2017 Pilosa Corp. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package port_test - -import ( - "fmt" - "net" - "testing" - - "github.com/pilosa/pilosa/v2/test/port" -) - -func TestPortsAreUnique(t *testing.T) { - t.Skip("do we use this anymore?") - portmap := make(map[int]struct{}) - err := port.GetPorts(func(ports []int) error { - for _, p := range ports { - if _, exists := portmap[p]; exists { - panic(fmt.Sprintf("port %v was already issued!", p)) - } - portmap[p] = struct{}{} - } - - return nil - }, 2000, 3) - if err != nil { - t.Fatal(err) - } -} - -func TestPortsAreUsable(t *testing.T) { - t.Skip("do we use this anymore?") - portmap := make(map[int]struct{}) - err := port.GetPorts(func(ports []int) error { - for _, p := range ports { - if _, exists := portmap[p]; exists { - panic(fmt.Sprintf("port %v was already issued!", p)) - } - - lsn, err := net.Listen("tcp", fmt.Sprintf(":%v", p)) - if err != nil { - panic(err) - } - - portmap[p] = struct{}{} - lsn.Close() - } - - return nil - }, 2000, 3) - if err != nil { - t.Fatal(err) - } -}