mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 20:07:51 +00:00
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.
This commit is contained in:
parent
a5f9236f3a
commit
bcb6942c80
7 changed files with 64 additions and 38 deletions
|
|
@ -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
|
||||
}
|
||||
|
||||
|
|
|
|||
19
cluster.go
19
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.
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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{}
|
||||
|
|
|
|||
35
server.go
35
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
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue