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 {