diff --git a/config.go b/config.go index b74125ce7..150bd1ec9 100644 --- a/config.go +++ b/config.go @@ -1,10 +1,6 @@ package pilosa -import ( - "net" - "strconv" - "time" -) +import "time" const ( // DefaultHost is the default hostname and port to use. @@ -66,79 +62,6 @@ func NewConfigForHosts(hosts []string) *Config { return conf } -// PilosaMessenger returns a new instance of Messenger based on the config. -func (c *Config) PilosaMessenger() *Messenger { - messenger := NewMessenger() - switch c.Cluster.MessengerType { - case "broadcast": - n := NewHTTPMessageBroker() - n.messenger = messenger - messenger.Broker = n - case "gossip": - n := NewGossipMessageBroker() - n.messenger = messenger - messenger.Broker = n - case "static": - // nop - } - return messenger -} - -// PilosaCluster returns a new instance of Cluster based on the config. -func (c *Config) PilosaCluster() *Cluster { - cluster := NewCluster() - cluster.ReplicaN = c.Cluster.ReplicaN - - for _, hostport := range c.Cluster.Nodes { - cluster.Nodes = append(cluster.Nodes, &Node{Host: hostport}) - } - - // Setup a Broadcast (over HTTP) or Gossip NodeSet based on config. - switch c.Cluster.MessengerType { - case "broadcast": - cluster.NodeSet = NewHTTPNodeSet() - cluster.NodeSet.(*HTTPNodeSet).Join(cluster.Nodes) - case "gossip": - gport, err := strconv.Atoi(DefaultGossipPort) - if err != nil { - // what? - } - gossipPort := gport - gossipSeed := DefaultHost - if c.Cluster.Gossip.Port != 0 { - gossipPort = c.Cluster.Gossip.Port - } - if c.Cluster.Gossip.Seed != "" { - gossipSeed = c.Cluster.Gossip.Seed - } - // get the host portion of addr to use for binding - gossipHost, _, err := net.SplitHostPort(c.Host) - if err != nil { - gossipHost = c.Host - } - cluster.NodeSet = NewGossipNodeSet(c.Host, gossipHost, gossipPort, gossipSeed) - case "static": - cluster.NodeSet = NewStaticNodeSet() - default: - cluster.NodeSet = NewStaticNodeSet() - } - - return cluster -} - -// AssociateMessageBroker allows an implementation to associate objects to the MessageBroker -// after cluster configuration. -func (c *Config) AssociateMessageBroker(s *Server) { - switch c.Cluster.MessengerType { - case "broadcast": - // nop - case "gossip": - s.Cluster.NodeSet.(*GossipNodeSet).config.memberlistConfig.Delegate = s.Messenger.Broker.(*GossipMessageBroker) - case "static": - // nop - } -} - // Duration is a TOML wrapper type for time.Duration. type Duration time.Duration diff --git a/gossip.go b/gossip.go index 9b18b45dd..178fa999e 100644 --- a/gossip.go +++ b/gossip.go @@ -23,6 +23,10 @@ type GossipNodeSet struct { LogOutput io.Writer } +func (g *GossipNodeSet) AttachBroker(mb *GossipMessageBroker) { + g.config.memberlistConfig.Delegate = mb +} + func (g *GossipNodeSet) Nodes() []*Node { a := make([]*Node, 0, g.memberlist.NumMembers()) for _, n := range g.memberlist.Members() { @@ -196,9 +200,10 @@ func (g *GossipMessageBroker) logger() *log.Logger { //////////////////////////////////////////////////////////////// // NewGossipMessageBroker returns a new instance of GossipMessageBroker. -func NewGossipMessageBroker() *GossipMessageBroker { +func NewGossipMessageBroker(m *Messenger) *GossipMessageBroker { g := &GossipMessageBroker{ LogOutput: os.Stderr, + messenger: m, } g.broadcasts = &memberlist.TransmitLimitedQueue{ diff --git a/messenger.go b/messenger.go index 03579bee6..1ffbfa51e 100644 --- a/messenger.go +++ b/messenger.go @@ -175,8 +175,8 @@ type HTTPMessageBroker struct { } // NewHTTPMessageBroker returns a new instance of HTTPMessageBroker. -func NewHTTPMessageBroker() *HTTPMessageBroker { - return &HTTPMessageBroker{} +func NewHTTPMessageBroker(m *Messenger) *HTTPMessageBroker { + return &HTTPMessageBroker{messenger: m} } // Send sends a protobuf message to all nodes simultaneously. diff --git a/server/server.go b/server/server.go index ba41f3dd0..11bb055a9 100644 --- a/server/server.go +++ b/server/server.go @@ -9,8 +9,10 @@ import ( "fmt" "io" "math/rand" + "net" "os" "path/filepath" + "strconv" "strings" "time" @@ -92,11 +94,11 @@ func (m *Command) Run(args ...string) (err error) { if err != nil { return err } - m.Server.Messenger = m.Config.PilosaMessenger() - m.Server.Cluster = m.Config.PilosaCluster() + m.Server.Messenger = PilosaMessenger(m.Config) + m.Server.Cluster = PilosaCluster(m.Config) // Associate objects to the MessageBroker based on config. - m.Config.AssociateMessageBroker(m.Server) + AssociateMessageBroker(m.Server, m.Config) // Set configuration options. m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) @@ -109,6 +111,77 @@ func (m *Command) Run(args ...string) (err error) { return nil } +// PilosaMessenger returns a new instance of Messenger based on the config. +func PilosaMessenger(c *pilosa.Config) *pilosa.Messenger { + messenger := pilosa.NewMessenger() + switch c.Cluster.MessengerType { + case "broadcast": + n := pilosa.NewHTTPMessageBroker(messenger) + messenger.Broker = n + case "gossip": + n := pilosa.NewGossipMessageBroker(messenger) + messenger.Broker = n + case "static": + // nop + } + return messenger +} + +// PilosaCluster returns a new instance of Cluster based on the config. +func PilosaCluster(c *pilosa.Config) *pilosa.Cluster { + cluster := pilosa.NewCluster() + cluster.ReplicaN = c.Cluster.ReplicaN + + for _, hostport := range c.Cluster.Nodes { + cluster.Nodes = append(cluster.Nodes, &pilosa.Node{Host: hostport}) + } + + // Setup a Broadcast (over HTTP) or Gossip NodeSet based on config. + switch c.Cluster.MessengerType { + case "broadcast": + cluster.NodeSet = pilosa.NewHTTPNodeSet() + cluster.NodeSet.(*pilosa.HTTPNodeSet).Join(cluster.Nodes) + case "gossip": + gport, err := strconv.Atoi(pilosa.DefaultGossipPort) + if err != nil { + // what? + } + gossipPort := gport + gossipSeed := pilosa.DefaultHost + if c.Cluster.Gossip.Port != 0 { + gossipPort = c.Cluster.Gossip.Port + } + if c.Cluster.Gossip.Seed != "" { + gossipSeed = c.Cluster.Gossip.Seed + } + // get the host portion of addr to use for binding + gossipHost, _, err := net.SplitHostPort(c.Host) + if err != nil { + gossipHost = c.Host + } + cluster.NodeSet = pilosa.NewGossipNodeSet(c.Host, gossipHost, gossipPort, gossipSeed) + case "static": + cluster.NodeSet = pilosa.NewStaticNodeSet() + default: + cluster.NodeSet = pilosa.NewStaticNodeSet() + } + + return cluster +} + +// AssociateMessageBroker allows an implementation to associate objects to the MessageBroker +// after cluster configuration. +func AssociateMessageBroker(s *pilosa.Server, c *pilosa.Config) { + switch c.Cluster.MessengerType { + case "broadcast": + // nop + case "gossip": + s.Cluster.NodeSet.(*pilosa.GossipNodeSet).AttachBroker(s.Messenger.Broker.(*pilosa.GossipMessageBroker)) + case "static": + // nop + } +} + func normalizeHost(host string) (string, error) { if !strings.Contains(host, ":") { host = host + ":"