remove messenger "wrapper" around messagebroker

all the functionality aside from the broker has been removed
This commit is contained in:
Matt Jaffee 2017-04-17 12:14:11 -05:00 committed by Travis
parent 1aca2ce1de
commit 5079ee7a22
10 changed files with 27 additions and 68 deletions

View file

@ -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
}

4
db.go
View file

@ -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
}

View file

@ -38,7 +38,7 @@ type Frame struct {
// Bitmap attribute storage and cache
bitmapAttrStore *AttrStore
messenger *Messenger
msgbroker MessageBroker
stats StatsClient
// Frame settings.

View file

@ -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 {

View file

@ -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,

View file

@ -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()

View file

@ -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
}

View file

@ -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
}

View file

@ -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,

View file

@ -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
}