From 8db7a78a937b66ce4d48f20cd1e6cfe578a390bb Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Tue, 23 Oct 2018 11:10:24 -0500 Subject: [PATCH] unpushed WIP on cluster-tests --- gossip/gossip.go | 28 ++++++++++++++++++++---- server.go | 16 ++++++++++---- server/config.go | 6 +++++ server/server.go | 6 +++++ server/server_test.go | 3 ++- test/pilosa.go | 51 ++++++++++++++++++++++++++++++++++++++++++- uri.go | 17 +++++++++++---- 7 files changed, 113 insertions(+), 14 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index 789ed50b0..d03016792 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -205,8 +205,18 @@ func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*mem conf.Name = api.Node().ID conf.BindAddr = api.Node().URI.Host conf.BindPort = port + if cfg.AdvertisePort != "" { + port, err = strconv.Atoi(cfg.Port) + if err != nil { + return nil, fmt.Errorf("convert advertise port: %s", err) + } + } conf.AdvertisePort = port - conf.AdvertiseAddr = hostToIP(api.Node().URI.Host) + if cfg.AdvertiseHost != "" { + conf.AdvertiseAddr = cfg.AdvertiseHost + } else { + conf.AdvertiseAddr = hostToIP(api.Node().URI.Host) + } // conf.TCPTimeout = time.Duration(cfg.StreamTimeout) conf.SuspicionMult = cfg.SuspicionMult @@ -447,10 +457,20 @@ func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) { // Config holds toml-friendly memberlist configuration. type Config struct { + // Host is the host gossip will bind to. If left blank it will be set to the + // host from Pilosa. + Host string `toml:"host"` // Port indicates the port to which pilosa should bind for internal state sharing. - Port string `toml:"port"` - Seeds []string `toml:"seeds"` - Key string `toml:"key"` + Port string `toml:"port"` + // AdvertiseHost is the hostname or IP other nodes should use to connect to + // this host. If left blank, the value for Host will be used. This is useful + // in some proxy and NAT scenarios. + AdvertiseHost string `toml:"advertise-host` + // AdvertisePort is the port other nodes will use to connect to this one. + // Behaves like AdvertiseHost. + AdvertisePort string `toml:"advertise-port"` + Seeds []string `toml:"seeds"` + Key string `toml:"key"` // StreamTimeout is the timeout for establishing a stream connection with // a remote node for a full state sync, and for stream read and write // operations. Maps to memberlist TCPTimeout. diff --git a/server.go b/server.go index a4cb8fc8b..36190e315 100644 --- a/server.go +++ b/server.go @@ -62,6 +62,7 @@ type Server struct { // nolint: maligned nodeID string uri URI + advertiseURI URI antiEntropyInterval time.Duration metricInterval time.Duration diagnosticInterval time.Duration @@ -73,7 +74,7 @@ type Server struct { // nolint: maligned dataDir string } -// TODO: have this return an interface for Holder instead of concrete object? +// TODO (2.0): have this return an interface for Holder instead of concrete object? func (s *Server) Holder() *Holder { return s.holder } @@ -160,6 +161,13 @@ func OptServerInternalClient(c InternalClient) ServerOption { } } +func OptServerAdvertiseURI(u *URI) ServerOption { + return func(s *Server) error { + s.advertiseURI = *u + return nil + } +} + // DEPRECATED func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption { return func(s *Server) error { @@ -296,7 +304,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { // Set Cluster Node. node := &Node{ ID: s.nodeID, - URI: s.uri, + URI: s.advertiseURI, IsCoordinator: s.cluster.Coordinator == s.nodeID, } s.cluster.Node = node @@ -581,7 +589,7 @@ func (s *Server) SendSync(m Message) error { node := node s.logger.Printf("SendSync to: %s", node.URI) // Don't forward the message to ourselves. - if s.uri == node.URI { + if s.advertiseURI == node.URI { continue } @@ -676,7 +684,7 @@ func (s *Server) monitorDiagnostics() { s.diagnostics.Logger = s.logger s.diagnostics.SetVersion(Version) - s.diagnostics.Set("Host", s.uri.Host) + s.diagnostics.Set("Host", s.advertiseURI.Host) s.diagnostics.Set("Cluster", strings.Join(s.cluster.nodeIDs(), ",")) s.diagnostics.Set("NumNodes", len(s.cluster.nodes)) s.diagnostics.Set("NumCPU", runtime.NumCPU()) diff --git a/server/config.go b/server/config.go index d255e4bb6..210a2e298 100644 --- a/server/config.go +++ b/server/config.go @@ -36,9 +36,15 @@ type Config struct { // DataDir is the directory where Pilosa stores both indexed data and // running state such as cluster topology information. DataDir string `toml:"data-dir"` + // Bind is the host:port on which Pilosa will listen. Bind string `toml:"bind"` + // Advertise is the host:port that this node will report as its address to + // others. If left blank (the default), this will be set to the bind address + // once it is listening. + Advertise string `toml:"advertise"` + // MaxWritesPerRequest limits the number of mutating commands that can be in // a single request to the server. This includes Set, Clear, // SetRowAttrs & SetColumnAttrs. diff --git a/server/server.go b/server/server.go index 010eb2f8a..8521abdec 100644 --- a/server/server.go +++ b/server/server.go @@ -248,6 +248,11 @@ func (m *Command) SetupServer() error { uri.SetPort(uint16(m.ln.Addr().(*net.TCPAddr).Port)) } + advertURI, err := pilosa.NewURIFromAddressWithDefault(m.Config.Advertise, uri) + if err != nil { + return errors.Wrapf(err, "processing avertise address '%s'", m.Config.Advertise) + } + c := http.GetHTTPClient(TLSConfig) // Primary store configuration is handled automatically now. @@ -276,6 +281,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerGCNotifier(gcnotify.NewActiveGCNotifier()), pilosa.OptServerStatsClient(statsClient), pilosa.OptServerURI(uri), + pilosa.OptServerAdvertiseURI(advertURI), pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)), pilosa.OptServerPrimaryTranslateStoreFunc(http.NewTranslateStore), pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts), diff --git a/server/server_test.go b/server/server_test.go index aa202a54c..7a5e85572 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -405,7 +405,6 @@ func TestClusteringNodesReplica1(t *testing.T) { // Create new main with the same config. config := cluster[2].Command.Config config.Translation.MapSize = 100000 - // config.Bind = cluster[2].API.Node().URI.HostPort() // this isn't necessary, but makes the test run way faster config.Gossip.Port = strconv.Itoa(int(cluster[2].Command.GossipTransport().URI.Port)) @@ -413,6 +412,8 @@ func TestClusteringNodesReplica1(t *testing.T) { cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr) cluster[2].Command.Config = config + time.Sleep(time.Second * 40) + // Run new program. if err := cluster[2].Start(); err != nil { t.Fatalf("restarting node 2: %v", err) diff --git a/test/pilosa.go b/test/pilosa.go index 27331d53f..46424e09d 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -19,11 +19,13 @@ import ( "context" "fmt" "io/ioutil" + "math/rand" gohttp "net/http" "os" "path" "strconv" "strings" + "sync/atomic" "testing" "time" @@ -31,8 +33,37 @@ import ( "github.com/pilosa/pilosa/http" "github.com/pilosa/pilosa/server" "github.com/pkg/errors" + + "github.com/Shopify/toxiproxy" + tox "github.com/Shopify/toxiproxy/client" ) +type portAllocator struct { + port *uint32 +} + +var ports *portAllocator + +// Next generates a new port and returns it +func (n *portAllocator) Next() (nextPort uint32) { + return atomic.AddUint32(n.port, 1) +} + +var proxy string + +func init() { + r := rand.New(rand.NewSource(int64(time.Now().Nanosecond()))) + p := uint32(r.Uint32()%40000 + 20000) + ports = &portAllocator{ + port: &p, + } + + tport := strconv.Itoa(int(ports.Next())) + proxy = "localhost:" + tport + go toxiproxy.NewServer().Listen("localhost", tport) + +} + //////////////////////////////////////////////////////////////////////////////////// // Command represents a test wrapper for server.Command. type Command struct { @@ -232,6 +263,8 @@ func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu // newCluster creates a new cluster func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { + tclient := tox.NewClient(proxy) + if size == 0 { return nil, errors.New("cluster must contain at least one node") } @@ -245,8 +278,24 @@ func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { if len(opts) > 0 { commandOpts = opts[i%len(opts)] } + aport := strconv.Itoa(int(ports.Next())) + name := "node" + strconv.Itoa(i) m := NewCommandNode(i == 0, commandOpts...) - err := ioutil.WriteFile(path.Join(m.Config.DataDir, ".id"), []byte("node"+strconv.Itoa(i)), 0600) + m.Config.Bind = "localhost:" + strconv.Itoa(int(ports.Next())) + m.Config.Advertise = "localhost:" + aport + p, err := tclient.CreateProxy(name+aport, m.Config.Advertise, m.Config.Bind) + if err != nil { + return nil, errors.Wrap(err, "setting up toxiproxy") + } + m.Config.Gossip.Port = strconv.Itoa(int(ports.Next())) + aport = strconv.Itoa(int(ports.Next())) + m.Config.Gossip.AdvertisePort = aport + p, err = tclient.CreateProxy(name+"-gossip"+aport, "localhost:"+m.Config.Gossip.AdvertisePort, "localhost:"+m.Config.Gossip.Port) + if err != nil { + return nil, errors.Wrap(err, "setting up toxiproxy for gossip") + } + fmt.Println(p) + err = ioutil.WriteFile(path.Join(m.Config.DataDir, ".id"), []byte(name), 0600) if err != nil { return nil, errors.Wrap(err, "writing node id") } diff --git a/uri.go b/uri.go index 231457691..7cf985f0b 100644 --- a/uri.go +++ b/uri.go @@ -82,6 +82,10 @@ func NewURIFromAddress(address string) (*URI, error) { return parseAddress(address) } +func NewURIFromAddressWithDefault(address string, base *URI) (*URI, error) { + return parseAddressWithDefault(address, base) +} + // setScheme sets the scheme of this URI. func (u *URI) setScheme(scheme string) error { m := schemeRegexp.FindStringSubmatch(scheme) @@ -154,20 +158,20 @@ func (u URI) Type() string { return "URI" } -func parseAddress(address string) (uri *URI, err error) { +func parseAddressWithDefault(address string, def *URI) (uri *URI, err error) { m := addressRegexp.FindStringSubmatch(address) if m == nil { return nil, errors.New("invalid address") } - scheme := "http" + scheme := def.Scheme if m[2] != "" { scheme = m[2] } - host := "localhost" + host := def.Host if m[3] != "" { host = m[3] } - var port = 10101 + var port = int(def.Port) if m[5] != "" { port, err = strconv.Atoi(m[5]) if err != nil { @@ -185,6 +189,11 @@ func parseAddress(address string) (uri *URI, err error) { return uri, nil } +func parseAddress(address string) (uri *URI, err error) { + u, err := parseAddressWithDefault(address, defaultURI()) + return u, err +} + // MarshalJSON marshals URI into a JSON-encoded byte slice. func (u *URI) MarshalJSON() ([]byte, error) { var output struct {