From c312bc13165f608c1ed50f3a78cce1e54665258f Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 27 Jun 2018 07:18:16 -0500 Subject: [PATCH 1/3] simplify arguments to NewGossipMemberSet --- gossip/gossip.go | 9 +++++---- server/server.go | 2 -- 2 files changed, 5 insertions(+), 6 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index d67215e85..0355260d6 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -161,7 +161,8 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption { } // NewGossipMemberSet returns a new instance of GossipMemberSet based on options. -func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) { +func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) { + host := s.Node().URI.Host() g := &GossipMemberSet{ Logger: pilosa.NopLogger, } @@ -206,11 +207,11 @@ func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, // memberlist config conf := memberlist.DefaultWANConfig() conf.Transport = g.transport.Net - conf.Name = name - conf.BindAddr = host + conf.Name = s.Node().ID + conf.BindAddr = s.Node().URI.Host() conf.BindPort = port conf.AdvertisePort = port - conf.AdvertiseAddr = hostToIP(host) + conf.AdvertiseAddr = hostToIP(s.Node().URI.Host()) // conf.TCPTimeout = time.Duration(cfg.StreamTimeout) conf.SuspicionMult = cfg.SuspicionMult diff --git a/server/server.go b/server/server.go index 8b958d314..614a4efcf 100644 --- a/server/server.go +++ b/server/server.go @@ -321,8 +321,6 @@ func (m *Command) SetupNetworking() error { } gossipMemberSet, err := gossip.NewGossipMemberSet( - m.Server.NodeID, - m.Server.URI.Host(), m.Config.Gossip, m.Server, gossip.WithLogger(m.logger.Logger()), From 41bdfceb577fdfdcf9ce89f2afcfa9397b50d39a Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 27 Jun 2018 07:22:32 -0500 Subject: [PATCH 2/3] remove MemberSet from cluster, Open in server package --- cluster.go | 13 +++---------- server/server.go | 3 +-- utils_internal_test.go | 1 - 3 files changed, 4 insertions(+), 13 deletions(-) diff --git a/cluster.go b/cluster.go index 75c29ceb8..e10aae84e 100644 --- a/cluster.go +++ b/cluster.go @@ -212,10 +212,9 @@ type nodeAction struct { // Cluster represents a collection of nodes. type Cluster struct { - ID string - Node *Node - Nodes []*Node // TODO phase this out? - MemberSet MemberSet + ID string + Node *Node + Nodes []*Node // TODO phase this out? // Hashing algorithm used to assign partitions to nodes. Hasher Hasher @@ -897,11 +896,6 @@ func (c *Cluster) open() error { } func (c *Cluster) waitForStarted() error { - // Open MemberSet communication. - if err := c.MemberSet.Open(); err != nil { - return errors.Wrap(err, "opening MemberSet") - } - // If not coordinator then wait for ClusterStatus from coordinator. if !c.isCoordinator() { // In the case where a node has been restarted and memberlist has @@ -1821,6 +1815,5 @@ func (c *Cluster) setStatic(hosts []string) error { } c.Nodes = append(c.Nodes, &Node{URI: *uri}) } - c.MemberSet = NewStaticMemberSet(c.Nodes) return nil } diff --git a/server/server.go b/server/server.go index 614a4efcf..cfb06fc3a 100644 --- a/server/server.go +++ b/server/server.go @@ -329,8 +329,7 @@ func (m *Command) SetupNetworking() error { if err != nil { return errors.Wrap(err, "getting memberset") } - m.Server.Cluster.MemberSet = gossipMemberSet - return nil + return errors.Wrap(gossipMemberSet.Open(), "opening gossip memberset") } // Close shuts down the server. diff --git a/utils_internal_test.go b/utils_internal_test.go index ea98875c4..20ff2b70b 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -228,7 +228,6 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*Cluster, error) c.Path = path c.Topology = NewTopology() c.Holder = h - c.MemberSet = NewStaticMemberSet(c.Nodes) c.Node = node c.Coordinator = t.common.Nodes[0].ID // the first node is the coordinator c.Broadcaster = t From fbe035ef25a6ae36648013657935d172e23a7739 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 27 Jun 2018 10:46:37 -0500 Subject: [PATCH 3/3] simplify test cluster setup by exposing gossip transport on server.Command --- server/cluster_test.go | 173 ++++++++++++----------------------------- server/server.go | 22 +++--- test/pilosa.go | 84 +++----------------- 3 files changed, 75 insertions(+), 204 deletions(-) diff --git a/server/cluster_test.go b/server/cluster_test.go index 73de82279..db21ce784 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -24,10 +24,9 @@ import ( "testing" "time" - "golang.org/x/sync/errgroup" - "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/test" + "golang.org/x/sync/errgroup" ) // Ensure program can send/receive broadcast messages. @@ -132,76 +131,34 @@ func TestClusterResize_EmptyNode(t *testing.T) { // Ensure that a cluster of empty nodes comes up in a NORMAL state. func TestClusterResize_EmptyNodes(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) - defer m0.Close() + clus := test.MustRunMainWithCluster(t, 2) + defer clus[0].Close() + defer clus[1].Close() - gossipHost := "localhost" - gossipPort := 0 - seed, err := m0.RunWithTransport(gossipHost, gossipPort, []string{}) - if err != nil { - t.Fatal(err) - } - - // Configure node1 - m1 := test.NewMainWithCluster(false) - defer m1.Close() - - seed, err = m1.RunWithTransport(gossipHost, gossipPort, []string{seed}) - if err != nil { - t.Fatal(err) - } - - if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if clus[0].Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", clus[0].Server.Cluster.State()) + } else if clus[1].Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", clus[1].Server.Cluster.State()) } } // Ensure that adding a node correctly resizes the cluster. func TestClusterResize_AddNode(t *testing.T) { t.Run("NoData", func(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) - defer m0.Close() + clus := test.MustRunMainWithCluster(t, 2) - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } - - // Configure node1 - m1 := test.NewMainWithCluster(false) - defer m1.Close() - - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) - } - - if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if !checkClusterState(clus[0].Server.Cluster, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", clus[0].Server.Cluster.State()) + } else if !checkClusterState(clus[1].Server.Cluster, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", clus[1].Server.Cluster.State()) } }) t.Run("WithIndex", func(t *testing.T) { // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() // Create a client for each node. client0 := m0.Client() @@ -215,19 +172,13 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewMainWithCluster(false) - defer m1.Close() - - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) } + defer m1.Close() if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) @@ -236,19 +187,14 @@ func TestClusterResize_AddNode(t *testing.T) { } }) t.Run("ContinuousSlices", func(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() // Create a client for each node. client0 := m0.Client() - //client1 := m1.Client() // Create indexes and fields on one node. if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { @@ -267,19 +213,13 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewMainWithCluster(false) - defer m1.Close() - - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) } + defer m1.Close() if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) @@ -288,19 +228,14 @@ func TestClusterResize_AddNode(t *testing.T) { } }) t.Run("SkippedSlice", func(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() // Create a client for each node. client0 := m0.Client() - //client1 := m1.Client() // Create indexes and fields on one node. if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { @@ -319,19 +254,13 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewMainWithCluster(false) - defer m1.Close() - - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) } + defer m1.Close() if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) @@ -345,37 +274,37 @@ func TestClusterResize_AddNode(t *testing.T) { func TestCluster_GossipMembership(t *testing.T) { t.Run("Node0Down", func(t *testing.T) { // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() + + var eg errgroup.Group // Configure node1 m1 := test.NewMainWithCluster(false) defer m1.Close() - - var eg errgroup.Group eg.Go(func() error { + m1.Config.Gossip.Port = "0" // Pass invalid seed as first in list - _, err := m1.RunWithTransport("localhost", 0, []string{"http://localhost:8765", seed}) + m1.Config.Gossip.Seeds = []string{"http://localhost:8765", seed} + err := m1.Start() if err != nil { - return err + t.Fatalf("starting second main: %v", err) } return nil }) - // Configure node2 + // Configure node1 m2 := test.NewMainWithCluster(false) defer m2.Close() - eg.Go(func() error { - // Pass invalid seed as last in list - _, err := m2.RunWithTransport("localhost", 0, []string{seed, "http://localhost:8765"}) + m2.Config.Gossip.Port = "0" + // Pass invalid seed as first in list + m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"} + err := m2.Start() if err != nil { - return err + t.Fatalf("starting second main: %v", err) } return nil }) diff --git a/server/server.go b/server/server.go index cfb06fc3a..164dc94e3 100644 --- a/server/server.go +++ b/server/server.go @@ -60,7 +60,7 @@ type Command struct { Config *Config // Gossip transport - GossipTransport *gossip.Transport + gossipTransport *gossip.Transport // Standard input/output *pilosa.CmdIO @@ -310,21 +310,16 @@ func (m *Command) SetupNetworking() error { // get the host portion of addr to use for binding gossipHost := m.Server.URI.Host() - var transport *gossip.Transport - if m.GossipTransport != nil { - transport = m.GossipTransport - } else { - transport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) - if err != nil { - return errors.Wrap(err, "getting transport") - } + m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) + if err != nil { + return errors.Wrap(err, "getting transport") } gossipMemberSet, err := gossip.NewGossipMemberSet( m.Config.Gossip, m.Server, gossip.WithLogger(m.logger.Logger()), - gossip.WithTransport(transport), + gossip.WithTransport(m.gossipTransport), ) if err != nil { return errors.Wrap(err, "getting memberset") @@ -332,6 +327,13 @@ func (m *Command) SetupNetworking() error { return errors.Wrap(gossipMemberSet.Open(), "opening gossip memberset") } +// GossipTransport allows a caller to return the gossip transport created when +// setting up the GossipMemberSet. This is useful if one needs to determine the +// allocated ephemeral port programmatically. (usually used in tests) +func (m *Command) GossipTransport() *gossip.Transport { + return m.gossipTransport +} + // Close shuts down the server. func (m *Command) Close() error { var logErr error diff --git a/test/pilosa.go b/test/pilosa.go index 2ae9ae9b0..ece7c27d6 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -19,14 +19,12 @@ import ( "fmt" "io" "io/ioutil" - "log" gohttp "net/http" "os" "strings" "testing" "time" - "github.com/pilosa/pilosa/gossip" "github.com/pilosa/pilosa/http" "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/toml" @@ -59,6 +57,13 @@ func OptAllowedOrigins(origins []string) server.CommandOption { } } +// GossipAddress returns the address on which gossip is listening after a Main +// has been setup. Useful to pass as a seed to other nodes when creating and +// testing clusters. +func (m *Main) GossipAddress() string { + return m.GossipTransport().URI.String() +} + // NewMain returns a new instance of Main with a temporary data directory and random port. func NewMain(opts ...server.CommandOption) *Main { path, err := ioutil.TempDir("", "pilosa-") @@ -116,25 +121,20 @@ func runMainWithCluster(size int, opts ...[]server.CommandOption) ([]*Main, erro } mains := make([]*Main, size) - - gossipHost := "localhost" - gossipPort := 0 - var err error var gossipSeeds = make([]string, size) - for i := 0; i < size; i++ { var commandOpts []server.CommandOption if len(opts) > 0 { commandOpts = opts[i%len(opts)] } m := NewMainWithCluster(i == 0, commandOpts...) - m.Config.Cluster.Disabled = false + m.Config.Gossip.Port = "0" + m.Config.Gossip.Seeds = gossipSeeds[:i] - gossipSeeds[i], err = m.RunWithTransport(gossipHost, gossipPort, gossipSeeds[:i]) - if err != nil { - return nil, errors.Wrap(err, "RunWithTransport") + if err := m.Start(); err != nil { + return nil, errors.Wrapf(err, "Starting server %d", i) } - + gossipSeeds[i] = m.GossipTransport().URI.String() mains[i] = m } @@ -179,66 +179,6 @@ func (m *Main) Reopen() error { return nil } -// RunWithTransport runs Main and returns the dynamically allocated gossip port. -func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (seed string, err error) { - defer close(m.Started) - - /* - TEST: - - SetupServer (just static settings from config) - - OpenListener (sets Server.Name to use in gossip) - - NewTransport (gossip) - - SetupNetworking (does the gossip or static stuff) - uses Server.Name - - Open server - - PRODUCTION: - - SetupServer (just static settings from config) - - SetupNetworking (does the gossip or static stuff) - calls NewTransport - - Open server - calls OpenListener - */ - // Open gossip transport to use in SetupServer. - transport, err := gossip.NewTransport(host, bindPort, nil) - if err != nil { - return seed, err - } - m.GossipTransport = transport - - if len(joinSeeds) != 0 { - m.Config.Gossip.Seeds = joinSeeds - } else { - m.Config.Gossip.Seeds = []string{transport.URI.String()} - } - seed = transport.URI.String() - - // SetupServer - m.Config.Cluster.Disabled = false - err = m.SetupServer() - if err != nil { - return seed, err - } - - // SetupNetworking - err = m.SetupNetworking() - if err != nil { - return seed, err - } - - go func() { - err := m.Handler.Serve() - if err != nil { - log.Printf("Handler serve error: %v", err) - } - }() - - // Initialize server. - err = m.Server.Open() - if err != nil { - return seed, err - } - - return seed, nil -} - // URL returns the base URL string for accessing the running program. func (m *Main) URL() string { return "http://" + m.Server.Addr().String() }