From e459c9a77b292d7079b042cd0b9e86d9987928e4 Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 29 Jan 2021 14:11:51 -0600 Subject: [PATCH] change Config.DisCo to Config.Etcd --- ctl/server.go | 20 ++++++++++---------- server/cluster_test.go | 16 ++++++++-------- server/config.go | 35 +++++++++++++++++------------------ server/server.go | 16 +++++++--------- test/cluster.go | 3 +-- test/disco.go | 6 +++--- 6 files changed, 46 insertions(+), 50 deletions(-) diff --git a/ctl/server.go b/ctl/server.go index 0d814cb86..e279230d3 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -73,16 +73,16 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.DurationVarP((*time.Duration)(&srv.Config.Gossip.Interval), "gossip.interval", "", (time.Duration)(srv.Config.Gossip.Interval), "Interval between sending messages that need to be gossiped that haven't piggybacked on probing messages.") flags.DurationVarP((*time.Duration)(&srv.Config.Gossip.ToTheDeadTime), "gossip.to-the-dead-time", "", (time.Duration)(srv.Config.Gossip.ToTheDeadTime), "Interval after which a node has died that we will still try to gossip to it.") - // DisCo - flags.StringVarP(&srv.Config.DisCo.Name, "disco.name", "", srv.Config.DisCo.Name, "Name of node in DisCo.") - flags.StringVarP(&srv.Config.DisCo.Dir, "disco.dir", "", srv.Config.DisCo.Dir, "Directory to use for DisCo.") - flags.StringVarP(&srv.Config.DisCo.LClientURL, "disco.listen-client-addr", "", srv.Config.DisCo.LClientURL, "Listen client address.") - flags.StringVarP(&srv.Config.DisCo.AClientURL, "disco.advertise-client-addr", "", srv.Config.DisCo.AClientURL, "Advertise client address.") - flags.StringVarP(&srv.Config.DisCo.LPeerURL, "disco.listen-peer-addr", "", srv.Config.DisCo.LPeerURL, "Listen peer address.") - flags.StringVarP(&srv.Config.DisCo.APeerURL, "disco.advertise-peer-addr", "", srv.Config.DisCo.APeerURL, "Advertise peer address.") - flags.StringVarP(&srv.Config.DisCo.ClusterURL, "disco.cluster-url", "", srv.Config.DisCo.ClusterURL, "Cluster URL to join.") - flags.StringVarP(&srv.Config.DisCo.ClusterName, "disco.cluster-name", "", srv.Config.DisCo.ClusterName, "Cluster name.") - flags.StringVarP(&srv.Config.DisCo.InitCluster, "disco.initial-cluster", "", srv.Config.DisCo.InitCluster, "Initial cluster name1=apurl1,name2=apurl2") + // Etcd + flags.StringVarP(&srv.Config.Etcd.Name, "etcd.name", "", srv.Config.Etcd.Name, "Name of node in Etcd.") + flags.StringVarP(&srv.Config.Etcd.Dir, "etcd.dir", "", srv.Config.Etcd.Dir, "Directory to use for Etcd.") + flags.StringVarP(&srv.Config.Etcd.LClientURL, "etcd.listen-client-addr", "", srv.Config.Etcd.LClientURL, "Listen client address.") + flags.StringVarP(&srv.Config.Etcd.AClientURL, "etcd.advertise-client-addr", "", srv.Config.Etcd.AClientURL, "Advertise client address.") + flags.StringVarP(&srv.Config.Etcd.LPeerURL, "etcd.listen-peer-addr", "", srv.Config.Etcd.LPeerURL, "Listen peer address.") + flags.StringVarP(&srv.Config.Etcd.APeerURL, "etcd.advertise-peer-addr", "", srv.Config.Etcd.APeerURL, "Advertise peer address.") + flags.StringVarP(&srv.Config.Etcd.ClusterURL, "etcd.cluster-url", "", srv.Config.Etcd.ClusterURL, "Cluster URL to join.") + flags.StringVarP(&srv.Config.Etcd.ClusterName, "etcd.cluster-name", "", srv.Config.Etcd.ClusterName, "Cluster name.") + flags.StringVarP(&srv.Config.Etcd.InitCluster, "etcd.initial-cluster", "", srv.Config.Etcd.InitCluster, "Initial cluster name1=apurl1,name2=apurl2") // AntiEntropy flags.DurationVarP((*time.Duration)(&srv.Config.AntiEntropy.Interval), "anti-entropy.interval", "", (time.Duration)(srv.Config.AntiEntropy.Interval), "Interval at which to run anti-entropy routine.") diff --git a/server/cluster_test.go b/server/cluster_test.go index 973c59960..a590da062 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -189,7 +189,7 @@ func TestClusterResize_AddNode(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC return m1.Start() @@ -246,7 +246,7 @@ func TestClusterResize_AddNode(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC return m1.Start() @@ -302,7 +302,7 @@ func TestClusterResize_AddNode(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC return m1.Start() @@ -364,7 +364,7 @@ func TestClusterResize_AddNode(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC return m1.Start() @@ -420,7 +420,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC return m1.Start() }, 4, 10); err != nil { @@ -478,7 +478,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC return m1.Start() }, 4, 10); err != nil { @@ -542,7 +542,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC errc := make(chan error, 1) @@ -604,7 +604,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) m1.Config.Gossip.Port = portsCfg[0].Gossip.Port - m1.Config.DisCo = portsCfg[0].DisCo + m1.Config.Etcd = portsCfg[0].Etcd m1.Config.BindGRPC = portsCfg[0].BindGRPC errc := make(chan error, 1) diff --git a/server/config.go b/server/config.go index f374e5db2..105f45c5c 100644 --- a/server/config.go +++ b/server/config.go @@ -130,8 +130,8 @@ type Config struct { LongQueryTime toml.Duration `toml:"long-query-time"` } `toml:"cluster"` - // DisCo config is based on embedded etcd. - DisCo petcd.Options `toml:"disco"` + // Etcd config is based on embedded etcd. + Etcd petcd.Options `toml:"etcd"` LongQueryTime toml.Duration `toml:"long-query-time"` // Gossip config is based around memberlist.Config. @@ -225,24 +225,24 @@ type Config struct { // We disallow zero because the tests need to be using from the pre-allocated // block of ports maintained by the pilosa/test/port port-mapper. func (c *Config) MustValidate() { - err := c.Validate() + err := c.validate() if err != nil { panic(err) } } -func (c *Config) Validate() error { - fmt.Printf("Validate() called on Config = '%#v'\n", c) +// validate ... +func (c *Config) validate() error { hostPort := []string{ "Bind", c.Bind, // :10101 "BindGRPC", c.BindGRPC, // :20101 "Advertise", c.Advertise, // on hp = 'http://localhost:63002' "AdvertiseGRPC", c.AdvertiseGRPC, // on hp = 'http://localhost:63003' - "DisCo.LClientURL", c.DisCo.LClientURL, // on hp = ':14000' - //c.DisCo.AClientURL, // hardcoded to same as LClientURL - "DisCo.LPeerURL", c.DisCo.LPeerURL, // ":" - //c.DisCo.APeerURL, // hardcoded to same as LPeerURL - "DisCo.ClusterURL", c.DisCo.ClusterURL, + "Etcd.LClientURL", c.Etcd.LClientURL, // on hp = ':14000' + //c.Etcd.AClientURL, // hardcoded to same as LClientURL + "Etcd.LPeerURL", c.Etcd.LPeerURL, // ":" + //c.Etcd.APeerURL, // hardcoded to same as LPeerURL + "Etcd.ClusterURL", c.Etcd.ClusterURL, "Gossip.Port", fmt.Sprintf(":%v", c.Gossip.Port), "Gossip.AdvertisePort", fmt.Sprintf(":%v", c.Gossip.AdvertisePort), "Postgres.Bind", c.Postgres.Bind, @@ -265,7 +265,6 @@ func (c *Config) Validate() error { continue } - fmt.Printf(" on name = '%v', hp = '%v'\n", name, hp) hp = strings.TrimPrefix(hp, "http://") hp = strings.TrimPrefix(hp, "https://") splt := strings.Split(hp, ":") @@ -356,13 +355,13 @@ func NewConfig() *Config { c.Postgres.WriteTimeout = toml.Duration(10 * time.Second) // we don't really need a connection limit - c.DisCo.AClientURL = "http://localhost:10301" - c.DisCo.LClientURL = "http://localhost:10301" - c.DisCo.APeerURL = "http://localhost:10401" - c.DisCo.LPeerURL = "http://localhost:10401" - c.DisCo.Dir = "" - c.DisCo.Name = "nodeName" - c.DisCo.ClusterName = "clusterName" + c.Etcd.AClientURL = "http://localhost:10301" + c.Etcd.LClientURL = "http://localhost:10301" + c.Etcd.APeerURL = "http://localhost:10401" + c.Etcd.LPeerURL = "http://localhost:10401" + c.Etcd.Dir = "" + c.Etcd.Name = "nodeName" + c.Etcd.ClusterName = "clusterName" return c } diff --git a/server/server.go b/server/server.go index 6495e9cb8..527404007 100644 --- a/server/server.go +++ b/server/server.go @@ -23,7 +23,6 @@ import ( "bytes" "context" "crypto/tls" - "fmt" "io" "io/ioutil" "log" @@ -122,8 +121,7 @@ func OptCommandConfig(config *Config) CommandOption { return func(c *Command) error { defer c.Config.MustValidate() if c.Config != nil { - c.Config.DisCo = config.DisCo - fmt.Printf("setting c.ConfigDisCo to '%#v'", config.DisCo) + c.Config.Etcd = config.Etcd return nil } c.Config = config @@ -395,16 +393,16 @@ func (m *Command) SetupServer() error { coordinatorOpt = pilosa.OptServerIsCoordinator(true) } - // If a DisCo.Dir is not provided, nest a default under the pilosa data dir. - if m.Config.DisCo.Dir == "" { + // If an Etcd.Dir is not provided, nest a default under the pilosa data dir. + if m.Config.Etcd.Dir == "" { path, err := expandDirName(m.Config.DataDir) if err != nil { return errors.Wrapf(err, "expanding directory name: %s", m.Config.DataDir) } - m.Config.DisCo.Dir = filepath.Join(path, pilosa.DefaultDiscoDir) + m.Config.Etcd.Dir = filepath.Join(path, pilosa.DefaultDiscoDir) } - e := petcd.NewEtcd(m.Config.DisCo, m.Config.Cluster.ReplicaN) + e := petcd.NewEtcd(m.Config.Etcd, m.Config.Cluster.ReplicaN) discoOpt := pilosa.OptServerDisCo(e, e, e, e, e, e, e) serverOptions := []pilosa.ServerOption{ @@ -583,8 +581,8 @@ 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 + m.Config.Etcd.LPeerSocket = nil + m.Config.Etcd.LClientSocket = nil err := eg.Wait() _ = testhook.Closed(pilosa.NewAuditor(), m, nil) diff --git a/test/cluster.go b/test/cluster.go index 9031e766f..445565ed2 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -422,11 +422,10 @@ func (c *Cluster) Start() error { for i, cc := range c.Nodes { cc := cc - cc.Config.DisCo = portsCfg[i].DisCo + cc.Config.Etcd = portsCfg[i].Etcd cc.Config.BindGRPC = portsCfg[i].BindGRPC eg.Go(func() error { - fmt.Printf("DISCO CONFIG: %+v\n", cc.Config.DisCo) cc.Config.Gossip.Seeds = gossipSeeds return cc.Start() diff --git a/test/disco.go b/test/disco.go index b0af11953..936026995 100644 --- a/test/disco.go +++ b/test/disco.go @@ -68,7 +68,7 @@ func GenPortsConfig(ports []Ports) []*server.Config { Port: fmt.Sprint(ports[i].Gossip), }, BindGRPC: fmt.Sprintf(":%d", ports[i].Grpc), - DisCo: etcd.Options{ + Etcd: etcd.Options{ Name: name, Dir: discoDir, ClusterName: "bartholemuuuuu", @@ -83,11 +83,11 @@ func GenPortsConfig(ports []Ports) []*server.Config { } 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", + fmt.Printf("\ndebug test/disco.go: on i=%v, GenPortsConfig Gossip: %v, Etcd.Client: %v, Etcd.Peer: %v, BindGRPC: %v\n", i, ports[i].Gossip, portC, portP, ports[i].Grpc) } for i := range cfgs { - cfgs[i].DisCo.InitCluster = strings.Join(clusterURLs, ",") + cfgs[i].Etcd.InitCluster = strings.Join(clusterURLs, ",") } return cfgs