From de7d69234e56a3995bc7ada1876d05909c2f572b Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 18 Jan 2021 10:34:18 -0600 Subject: [PATCH] initial listner --- etcd/embed.go | 17 +++++++++++++++++ server/server.go | 4 ++++ test/disco.go | 27 +++++++++++++++++---------- 3 files changed, 38 insertions(+), 10 deletions(-) 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/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/disco.go b/test/disco.go index afcfcf731..738792262 100644 --- a/test/disco.go +++ b/test/disco.go @@ -17,6 +17,7 @@ package test import ( "fmt" "io/ioutil" + "net" "strings" "time" @@ -39,8 +40,12 @@ func GenPortsConfig(ports []Ports) []*server.Config { 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) + name := fmt.Sprintf("server%d", i) + lsnC, portC := port.MustGetBoundTCPListener() + lClientURL := fmt.Sprintf("http://localhost:%d", portC) + lsnP, portP := port.MustGetBoundTCPListener() + lPeerURL := fmt.Sprintf("http://localhost:%d", portP) + discoDir := "" if d, err := ioutil.TempDir("/tmp", "disco."); err == nil { discoDir = d @@ -52,14 +57,16 @@ func GenPortsConfig(ports []Ports) []*server.Config { }, BindGRPC: port.ColonZeroString(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}, }, }