change Config.DisCo to Config.Etcd

This commit is contained in:
Travis 2021-01-29 14:11:51 -06:00
parent 9a70e15cb3
commit e459c9a77b
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
6 changed files with 46 additions and 50 deletions

View file

@ -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.")

View file

@ -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)

View file

@ -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
}

View file

@ -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)

View file

@ -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()

View file

@ -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