From bcb6942c80173bf0ed1a3523451795f21182eb3f Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Tue, 26 Jun 2018 14:36:05 -0500 Subject: [PATCH] continue simplifying memberset and pilosa setup since the gossip MemberSet has access to Server, it wasn't really necessary to pass it a Node object when calling Open on it from Cluster. The end goal is to have it be removed from Cluster entirely, and have it be Opened externally, and this is a step toward that. Exposing Node method on Server doesn't really expose any more than was already there as the same info can be gotten from LocalStatus with a bit of type casting. I figured adding the method was a little cleaner, and we could collapse all the functionality when the dust has settled. The Cluster.open method has been broken into two parts - one of which happens earlier (at NewServer time), and the other will eventually just be "waiting to make sure we've joined the cluster". Right now it's calling Memberset.Open, and then waiting to make sure the cluster has been joined. --- broadcast.go | 4 ++-- cluster.go | 19 +++++++++++++++---- cluster_internal_test.go | 7 ++++--- gossip/gossip.go | 7 ++----- server.go | 35 +++++++++++++++++++++++++++-------- server/server.go | 13 +++++++------ test/pilosa.go | 17 +++++++---------- 7 files changed, 64 insertions(+), 38 deletions(-) diff --git a/broadcast.go b/broadcast.go index 9b894fea2..d102fda22 100644 --- a/broadcast.go +++ b/broadcast.go @@ -27,7 +27,7 @@ import ( type MemberSet interface { // Open starts any network activity implemented by the MemberSet // Node is the local node, used for membership broadcasts. - Open(n *Node) error + Open() error } // StaticMemberSet represents a basic MemberSet for testing. @@ -43,7 +43,7 @@ func NewStaticMemberSet(nodes []*Node) *StaticMemberSet { } // Open implements the MemberSet interface to start network activity, but for a static MemberSet it does nothing. -func (s *StaticMemberSet) Open(n *Node) error { +func (s *StaticMemberSet) Open() error { return nil } diff --git a/cluster.go b/cluster.go index ae34bf03c..75c29ceb8 100644 --- a/cluster.go +++ b/cluster.go @@ -861,7 +861,7 @@ func (h *jmphasher) Hash(key uint64, n int) int { return int(b) } -func (c *Cluster) open() error { +func (c *Cluster) setup() error { // Cluster always comes up in state STARTING until cluster membership is determined. c.state = ClusterStateStarting @@ -876,7 +876,7 @@ func (c *Cluster) open() error { if c.isCoordinator() { err := c.considerTopology() if err != nil { - return fmt.Errorf("considerTopology: %v", err) + return errors.Wrap(err, "considerTopology") } } @@ -885,10 +885,21 @@ func (c *Cluster) open() error { if err != nil { return errors.Wrap(err, "adding local node") } + return nil +} +func (c *Cluster) open() error { + err := c.setup() + if err != nil { + return errors.Wrap(err, "setting up cluster") + } + return c.waitForStarted() +} + +func (c *Cluster) waitForStarted() error { // Open MemberSet communication. - if err := c.MemberSet.Open(c.Node); err != nil { - return fmt.Errorf("opening MemberSet: %v", err) + if err := c.MemberSet.Open(); err != nil { + return errors.Wrap(err, "opening MemberSet") } // If not coordinator then wait for ClusterStatus from coordinator. diff --git a/cluster_internal_test.go b/cluster_internal_test.go index bdd51c047..86e027600 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -25,6 +25,7 @@ import ( "github.com/davecgh/go-spew/spew" "github.com/pilosa/pilosa/internal" + "github.com/pkg/errors" ) // Ensure that fragCombos creates the correct fragment mapping. @@ -591,10 +592,10 @@ func TestCluster_ResizeStates(t *testing.T) { tc.WriteTopology(node.Path, top) // Open TestCluster. - expected := "considerTopology: coordinator node0 is not in topology: [some-other-host]" + expected := "coordinator node0 is not in topology: [some-other-host]" err := tc.Open() - if err == nil || err.Error() != expected { - t.Errorf("did not receive expected error: %s", expected) + if err == nil || errors.Cause(err).Error() != expected { + t.Errorf("did not receive expected error, got: %s", errors.Cause(err).Error()) } // Close TestCluster. diff --git a/gossip/gossip.go b/gossip/gossip.go index 4749c9c58..d67215e85 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -38,7 +38,6 @@ var _ memberlist.Delegate = &GossipMemberSet{} // GossipMemberSet represents a gossip implementation of MemberSet using memberlist. type GossipMemberSet struct { mu sync.RWMutex - node *pilosa.Node memberlist *memberlist.Memberlist handler pilosa.BroadcastHandler @@ -63,7 +62,7 @@ func (g *GossipMemberSet) GetBindAddr() string { } // Open implements the MemberSet interface to start network activity. -func (g *GossipMemberSet) Open(n *pilosa.Node) error { +func (g *GossipMemberSet) Open() error { err := g.gossipEventReceiver.Start(g.pserver) if err != nil { return errors.Wrap(err, "starting event delegate") @@ -72,8 +71,6 @@ func (g *GossipMemberSet) Open(n *pilosa.Node) error { return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()") } - g.node = n - g.mu.Lock() g.memberlist, err = memberlist.Create(g.config.memberlistConfig) g.mu.Unlock() @@ -241,7 +238,7 @@ func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, // NodeMeta implementation of the memberlist.Delegate interface. func (g *GossipMemberSet) NodeMeta(limit int) []byte { - buf, err := proto.Marshal(pilosa.EncodeNode(g.node)) + buf, err := proto.Marshal(pilosa.EncodeNode(g.pserver.Node())) if err != nil { g.Logger.Printf("marshal message error: %s", err) return []byte{} diff --git a/server.go b/server.go index 74bedab5b..27c56be6e 100644 --- a/server.go +++ b/server.go @@ -71,6 +71,7 @@ type Server struct { metricInterval time.Duration diagnosticInterval time.Duration maxWritesPerRequest int + isCoordinator bool primaryTranslateStore TranslateStore @@ -203,6 +204,13 @@ func OptServerClusterDisabled(disabled bool, hosts []string) ServerOption { } } +func OptServerIsCoordinator(is bool) ServerOption { + return func(s *Server) error { + s.isCoordinator = is + return nil + } +} + // NewServer returns a new instance of Server. func NewServer(opts ...ServerOption) (*Server, error) { s := &Server{ @@ -249,6 +257,10 @@ func NewServer(opts ...ServerOption) (*Server, error) { // Get or create NodeID. s.NodeID = s.LoadNodeID() + if s.isCoordinator { + s.Cluster.Coordinator = s.NodeID + } + // Set Cluster Node. node := &Node{ ID: s.NodeID, @@ -271,6 +283,14 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.executor.Cluster = s.Cluster s.executor.TranslateStore = s.translateFile s.executor.MaxWritesPerRequest = s.maxWritesPerRequest + s.Cluster.Broadcaster = s + s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest + s.holder.Broadcaster = s + + err = s.Cluster.setup() + if err != nil { + return nil, errors.Wrap(err, "setting up cluster") + } return s, nil } @@ -290,15 +310,8 @@ func (s *Server) Open() error { return err } - // Cluster settings. - s.Cluster.Broadcaster = s - s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest - - // Initialize Holder. - s.holder.Broadcaster = s - // Open Cluster management. - if err := s.Cluster.open(); err != nil { + if err := s.Cluster.waitForStarted(); err != nil { return fmt.Errorf("opening Cluster: %v", err) } @@ -529,6 +542,12 @@ func (s *Server) SendTo(to *Node, pb proto.Message) error { return s.defaultClient.SendMessage(context.Background(), &to.URI, pb) } +// Node returns the pilosa.Node object. It is used by membership protocols to +// get this node's name(ID), location(URI), and coordinator status. +func (s *Server) Node() *Node { + return s.Cluster.Node +} + // Server implements StatusHandler. // LocalStatus is used to periodically sync information // between nodes. Under normal conditions, nodes should diff --git a/server/server.go b/server/server.go index a6bc1d995..8b958d314 100644 --- a/server/server.go +++ b/server/server.go @@ -247,6 +247,12 @@ func (m *Command) SetupServer() error { primaryTranslateStore = http.NewTranslateStore(m.Config.Translation.PrimaryURL) } + // Set Coordinator. + coordinatorOpt := pilosa.OptServerIsCoordinator(false) + if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 { + coordinatorOpt = pilosa.OptServerIsCoordinator(true) + } + serverOptions := []pilosa.ServerOption{ pilosa.OptServerAntiEntropyInterval(time.Duration(m.Config.AntiEntropy.Interval)), pilosa.OptServerLongQueryTime(time.Duration(m.Config.Cluster.LongQueryTime)), @@ -265,6 +271,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)), pilosa.OptServerPrimaryTranslateStore(primaryTranslateStore), pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts), + coordinatorOpt, } serverOptions = append(serverOptions, m.serverOptions...) @@ -313,12 +320,6 @@ func (m *Command) SetupNetworking() error { } } - // Set Coordinator. - if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 { - m.Server.Cluster.Coordinator = m.Server.NodeID - m.Server.Cluster.Node.IsCoordinator = true - } - gossipMemberSet, err := gossip.NewGossipMemberSet( m.Server.NodeID, m.Server.URI.Host(), diff --git a/test/pilosa.go b/test/pilosa.go index 47b132508..2ae9ae9b0 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -196,13 +196,6 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) ( - SetupNetworking (does the gossip or static stuff) - calls NewTransport - Open server - calls OpenListener */ - - // SetupServer - err = m.SetupServer() - if err != nil { - return seed, err - } - // Open gossip transport to use in SetupServer. transport, err := gossip.NewTransport(host, bindPort, nil) if err != nil { @@ -215,17 +208,21 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) ( } 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 } - m.Server.Cluster.Static = false - go func() { err := m.Handler.Serve() if err != nil {