diff --git a/cluster.go b/cluster.go index 13f620b88..a435a20b7 100644 --- a/cluster.go +++ b/cluster.go @@ -248,7 +248,6 @@ func (h *HTTPNodeSet) Join(nodes []*Node) error { // StaticNodeSet represents a basic NodeSet for testing type StaticNodeSet struct { - Messenger nodes []*Node } @@ -264,9 +263,6 @@ func (s *StaticNodeSet) Open() error { return nil } -func (s *StaticNodeSet) SendMessage(pb proto.Message, method string) error { - return nil -} -func (s *StaticNodeSet) ReceiveMessage(pb proto.Message) error { +func (s *StaticNodeSet) Send(pb proto.Message, method string) error { return nil } diff --git a/db.go b/db.go index cc27e18c9..4eb61707e 100644 --- a/db.go +++ b/db.go @@ -43,7 +43,7 @@ type DB struct { // Profile attribute storage and cache profileAttrStore *AttrStore - messenger *Messenger + msgbroker MessageBroker stats StatsClient LogOutput io.Writer @@ -417,7 +417,7 @@ func (db *DB) newFrame(path, name string) (*Frame, error) { } f.LogOutput = db.LogOutput f.stats = db.stats.WithTags(fmt.Sprintf("frame:%s", name)) - f.messenger = db.messenger + f.msgbroker = db.msgbroker return f, nil } diff --git a/frame.go b/frame.go index c5a40b1d1..1da970efc 100644 --- a/frame.go +++ b/frame.go @@ -38,7 +38,7 @@ type Frame struct { // Bitmap attribute storage and cache bitmapAttrStore *AttrStore - messenger *Messenger + msgbroker MessageBroker stats StatsClient // Frame settings. diff --git a/gossip.go b/gossip.go index 59abef70d..37b8c3b97 100644 --- a/gossip.go +++ b/gossip.go @@ -96,7 +96,7 @@ type GossipMessageBroker struct { LogOutput io.Writer } -// implementation of the messenger.Messenger interface +// implementation of the messenger.MessageBroker interface func (g *GossipMessageBroker) Send(pb proto.Message, method string) error { msg, err := MarshalMessage(pb) if err != nil { diff --git a/handler.go b/handler.go index 7616fd2d8..5c015abde 100644 --- a/handler.go +++ b/handler.go @@ -26,7 +26,7 @@ import ( // Handler represents an HTTP handler. type Handler struct { Index *Index - Messenger *Messenger + MsgBroker MessageBroker Server *Server // Local hostname & cluster configuration. @@ -370,7 +370,7 @@ func (h *Handler) handlePostDB(w http.ResponseWriter, r *http.Request) { } // Send the delete message to all nodes. - err := h.Messenger.SendMessage( + err := h.MsgBroker.Send( &internal.DeleteDBMessage{ DB: req.DB, }, "direct") @@ -514,7 +514,7 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { } // Send the create message to all nodes. - err = h.Messenger.SendMessage( + err = h.MsgBroker.Send( &internal.CreateFrameMessage{ DB: req.DB, Frame: req.Frame, @@ -586,7 +586,7 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) { } // Send the delete message to all nodes. - err := h.Messenger.SendMessage( + err := h.MsgBroker.Send( &internal.DeleteFrameMessage{ DB: req.DB, Frame: req.Frame, diff --git a/handler_test.go b/handler_test.go index 92bcc7601..e95e5c0dd 100644 --- a/handler_test.go +++ b/handler_test.go @@ -794,7 +794,7 @@ func NewHandler() *Handler { h.Handler.LogOutput = ioutil.Discard // Handler test messages can no-op. - h.Messenger = pilosa.NewMessenger() + h.MsgBroker = pilosa.NopMessageBroker return h } @@ -828,8 +828,7 @@ func NewServer() *Server { s.Handler.Host = s.Host() // Handler test messages can no-op. - s.Handler.Messenger = pilosa.NewMessenger() - + s.Handler.MsgBroker = pilosa.NopMessageBroker // Create a default cluster on the handler s.Handler.Cluster = NewCluster(1) s.Handler.Cluster.Nodes[0].Host = s.Host() diff --git a/index.go b/index.go index 5e9f5187e..ab659ada1 100644 --- a/index.go +++ b/index.go @@ -23,8 +23,7 @@ type Index struct { // Databases by name. dbs map[string]*DB - Messenger *Messenger - + MsgBroker MessageBroker // Close management wg sync.WaitGroup closing chan struct{} @@ -247,7 +246,7 @@ func (i *Index) newDB(path, name string) (*DB, error) { } db.LogOutput = i.LogOutput db.stats = i.Stats.WithTags(fmt.Sprintf("db:%s", db.Name())) - db.messenger = i.Messenger + db.msgbroker = i.MsgBroker return db, nil } diff --git a/messenger.go b/messenger.go index 3672a407a..fc0c55d77 100644 --- a/messenger.go +++ b/messenger.go @@ -4,11 +4,9 @@ import ( "bytes" "errors" "fmt" - "io" "io/ioutil" "net/http" "net/url" - "os" "reflect" "golang.org/x/sync/errgroup" @@ -17,40 +15,11 @@ import ( "github.com/pilosa/pilosa/internal" ) -// Messenger represents an internal message handler. -type Messenger struct { - - // Broker handles Send - Broker MessageBroker - - // The writer for any logging. - LogOutput io.Writer -} - -// NewMessenger returns a new instance of Messenger with a default logger. -func NewMessenger() *Messenger { - return &Messenger{ - Broker: NopMessageBroker, - LogOutput: os.Stderr, - } -} - -func (m *Messenger) SendMessage(pb proto.Message, method string) error { - if m.Broker == nil { - return errors.New("Messenger.Broker is not defined.") - } - return m.Broker.Send(pb, method) -} - -////////////////////////////////////////////////////////////////// - // MessageBroker is an interface for handling incoming/outgoing messages. type MessageBroker interface { Send(pb proto.Message, method string) error } -////////////////////////////////////////////////////////////////// - func init() { NopMessageBroker = &nopMessageBroker{} } @@ -61,7 +30,7 @@ var NopMessageBroker MessageBroker type nopMessageBroker struct{} func (c *nopMessageBroker) Send(pb proto.Message, method string) error { - fmt.Println("NOPMessageBroker: Send") + fmt.Println("NOPMessageBroker: Send") // TODO remove or log properly? return nil } diff --git a/server.go b/server.go index 7dd35fa8f..68f9a6921 100644 --- a/server.go +++ b/server.go @@ -35,7 +35,7 @@ type Server struct { // Data storage and HTTP interface. Index *Index Handler *Handler - Messenger *Messenger + MsgBroker MessageBroker // Cluster configuration. // Host is replaced with actual host after opening if port is ":0". @@ -56,7 +56,7 @@ func NewServer() *Server { Index: NewIndex(), Handler: NewHandler(), - Messenger: NewMessenger(), + MsgBroker: NopMessageBroker, AntiEntropyInterval: DefaultAntiEntropyInterval, PollingInterval: DefaultPollingInterval, @@ -111,18 +111,15 @@ func (s *Server) Open() error { e.Host = s.Host e.Cluster = s.Cluster - // Initialize Messenger. - s.Messenger.LogOutput = s.LogOutput - // Initialize HTTP handler. - s.Handler.Messenger = s.Messenger + s.Handler.MsgBroker = s.MsgBroker s.Handler.Host = s.Host s.Handler.Cluster = s.Cluster s.Handler.Executor = e s.Handler.LogOutput = s.LogOutput // Initialize Index. - s.Index.Messenger = s.Messenger + s.Index.MsgBroker = s.MsgBroker s.Index.LogOutput = s.LogOutput // Serve HTTP. @@ -239,7 +236,7 @@ func (s *Server) monitorMaxSlices() { // In a gossip implementation, memberlist.Delegate.LocalState() uses this. func (s *Server) LocalState() (proto.Message, error) { if s.Index == nil { - return nil, errors.New("Messenger.Index is nil.") + return nil, errors.New("Server.Index is nil.") } return &internal.NodeState{ Host: s.Host, diff --git a/server/server.go b/server/server.go index 0b61bf405..69041ef3a 100644 --- a/server/server.go +++ b/server/server.go @@ -94,7 +94,7 @@ func (m *Command) Run(args ...string) (err error) { if err != nil { return err } - m.Server.Messenger = PilosaMessenger(m.Config, m.Server) + m.Server.MsgBroker = PilosaMessageBroker(m.Config, m.Server) m.Server.Cluster = PilosaCluster(m.Config) // Associate objects to the MessageBroker based on config. @@ -111,18 +111,17 @@ 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, server *pilosa.Server) *pilosa.Messenger { - messenger := pilosa.NewMessenger() +// PilosaMessageBroker returns a new instance of MessageBroker based on the config. +func PilosaMessageBroker(c *pilosa.Config, server *pilosa.Server) (broker pilosa.MessageBroker) { switch c.Cluster.MessengerType { case "broadcast": - messenger.Broker = pilosa.NewHTTPMessageBroker(server) + broker = pilosa.NewHTTPMessageBroker(server) case "gossip": - messenger.Broker = pilosa.NewGossipMessageBroker(server) + broker = pilosa.NewGossipMessageBroker(server) case "static": - // nop + broker = pilosa.NopMessageBroker } - return messenger + return broker } // PilosaCluster returns a new instance of Cluster based on the config. @@ -174,7 +173,7 @@ func AssociateMessageBroker(s *pilosa.Server, c *pilosa.Config) { case "broadcast": // nop case "gossip": - s.Cluster.NodeSet.(*pilosa.GossipNodeSet).AttachBroker(s.Messenger.Broker.(*pilosa.GossipMessageBroker)) + s.Cluster.NodeSet.(*pilosa.GossipNodeSet).AttachBroker(s.MsgBroker.(*pilosa.GossipMessageBroker)) case "static": // nop }