mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 09:05:55 +00:00
Merge pull request #1417 from jaffee/extracting-memberset
continue simplifying memberset and pilosa setup
This commit is contained in:
commit
a0f21064ed
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