From 32d07fc9e1c1fa06f02375e13311b8315c7e5921 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 17 Apr 2017 12:47:26 -0500 Subject: [PATCH] rename MessageBroker -> Broadcaster --- cmd/server.go | 2 +- config.go | 12 ++++++------ db.go | 6 +++--- frame.go | 4 ++-- gossip.go | 38 +++++++++++++++++++------------------- handler.go | 12 ++++++------ handler_test.go | 4 ++-- index.go | 4 ++-- messenger.go | 46 +++++++++++++++++++++++----------------------- server.go | 16 ++++++++-------- server/server.go | 36 ++++++++++++++++++------------------ 11 files changed, 90 insertions(+), 90 deletions(-) diff --git a/cmd/server.go b/cmd/server.go index 27c238eb6..e9d7fcfdf 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -83,7 +83,7 @@ on the configured port.`, flags.DurationVarP((*time.Duration)(&Server.Config.AntiEntropy.Interval), "anti-entropy.interval", "", time.Minute*10, "Interval at which to run anti-entropy routine.") flags.StringVarP(&Server.CPUProfile, "profile.cpu", "", "", "Where to store CPU profile.") flags.DurationVarP(&Server.CPUTime, "profile.cpu-time", "", 30*time.Second, "CPU profile duration.") - flags.StringVarP(&Server.Config.Cluster.MessengerType, "cluster.messenger-type", "", "static", "Type of Messenger to use for inter-host messaging. Choose from [static, broadcast, gossip]") + flags.StringVarP(&Server.Config.Cluster.BroadcasterType, "cluster.broadcaster-type", "", "static", "Type of Broadcaster to use for inter-host messaging. Choose from [static, http, gossip]") flags.StringVarP(&Server.Config.Cluster.Gossip.Seed, "cluster.gossip.seed", "", "", "Host with which to seed the gossip membership.") flags.IntVarP(&Server.Config.Cluster.Gossip.Port, "cluster.gossip.port", "", 0, "Port to which pilosa should bind for gossip.") diff --git a/config.go b/config.go index 150bd1ec9..8b718f4e1 100644 --- a/config.go +++ b/config.go @@ -4,10 +4,10 @@ import "time" const ( // DefaultHost is the default hostname and port to use. - DefaultHost = "localhost" - DefaultPort = "10101" - DefaultMessengerType = "static" - DefaultGossipPort = "14000" + DefaultHost = "localhost" + DefaultPort = "10101" + DefaultBroadcasterType = "static" + DefaultGossipPort = "14000" ) // Config represents the configuration for the command. @@ -17,7 +17,7 @@ type Config struct { Cluster struct { ReplicaN int `toml:"replicas"` - MessengerType string `toml:"messenger-type"` + BroadcasterType string `toml:"broadcaster-type"` Nodes []string `toml:"hosts"` PollingInterval Duration `toml:"polling-interval"` Gossip ConfigGossip `toml:"gossip"` @@ -45,7 +45,7 @@ func NewConfig() *Config { Host: DefaultHost + ":" + DefaultPort, } c.Cluster.ReplicaN = DefaultReplicaN - c.Cluster.MessengerType = DefaultMessengerType + c.Cluster.BroadcasterType = DefaultBroadcasterType c.Cluster.PollingInterval = Duration(DefaultPollingInterval) c.Cluster.Nodes = []string{} c.AntiEntropy.Interval = Duration(DefaultAntiEntropyInterval) diff --git a/db.go b/db.go index 4eb61707e..d2ceaf13a 100644 --- a/db.go +++ b/db.go @@ -43,8 +43,8 @@ type DB struct { // Profile attribute storage and cache profileAttrStore *AttrStore - msgbroker MessageBroker - stats StatsClient + broadcaster Broadcaster + 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.msgbroker = db.msgbroker + f.broadcaster = db.broadcaster return f, nil } diff --git a/frame.go b/frame.go index 1da970efc..a497e6906 100644 --- a/frame.go +++ b/frame.go @@ -38,8 +38,8 @@ type Frame struct { // Bitmap attribute storage and cache bitmapAttrStore *AttrStore - msgbroker MessageBroker - stats StatsClient + broadcaster Broadcaster + stats StatsClient // Frame settings. rowLabel string diff --git a/gossip.go b/gossip.go index 6482a333d..da9c847f5 100644 --- a/gossip.go +++ b/gossip.go @@ -23,7 +23,7 @@ type GossipNodeSet struct { LogOutput io.Writer } -func (g *GossipNodeSet) AttachBroker(mb *GossipMessageBroker) { +func (g *GossipNodeSet) AttachBroadcaster(mb *GossipBroadcaster) { g.config.memberlistConfig.Delegate = mb } @@ -85,9 +85,9 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed //////////////////////////////////////////////////////////////// -// GossipMessageBroker represents a gossip implementation of pilosa.MessageBroker -// GossipMessageBroker also represents an implementation of memberlist.Delegate -type GossipMessageBroker struct { +// GossipBroadcaster represents a gossip implementation of pilosa.Broadcaster +// GossipBroadcaster also represents an implementation of memberlist.Delegate +type GossipBroadcaster struct { broadcasts *memberlist.TransmitLimitedQueue server *Server @@ -96,8 +96,8 @@ type GossipMessageBroker struct { LogOutput io.Writer } -// SendSync implementation of the messenger.MessageBroker interface -func (g *GossipMessageBroker) SendSync(pb proto.Message) error { +// SendSync implementation of the Broadcaster interface +func (g *GossipBroadcaster) SendSync(pb proto.Message) error { msg, err := MarshalMessage(pb) if err != nil { return err @@ -124,8 +124,8 @@ func (g *GossipMessageBroker) SendSync(pb proto.Message) error { return eg.Wait() } -// SendAsync implementation of the messenger.MessageBroker interface -func (g *GossipMessageBroker) SendAsync(pb proto.Message) error { +// SendAsync implementation of the Broadcaster interface +func (g *GossipBroadcaster) SendAsync(pb proto.Message) error { msg, err := MarshalMessage(pb) if err != nil { return err @@ -139,7 +139,7 @@ func (g *GossipMessageBroker) SendAsync(pb proto.Message) error { return nil } -func (g *GossipMessageBroker) Receive(pb proto.Message) error { +func (g *GossipBroadcaster) Receive(pb proto.Message) error { if err := g.server.ReceiveMessage(pb); err != nil { return err } @@ -147,11 +147,11 @@ func (g *GossipMessageBroker) Receive(pb proto.Message) error { } // implementation of the memberlist.Delegate interface -func (g *GossipMessageBroker) NodeMeta(limit int) []byte { +func (g *GossipBroadcaster) NodeMeta(limit int) []byte { return []byte{} } -func (g *GossipMessageBroker) NotifyMsg(b []byte) { +func (g *GossipBroadcaster) NotifyMsg(b []byte) { m, err := UnmarshalMessage(b) if err != nil { g.logger().Printf("unmarshal message error: %s", err) @@ -163,11 +163,11 @@ func (g *GossipMessageBroker) NotifyMsg(b []byte) { } } -func (g *GossipMessageBroker) GetBroadcasts(overhead, limit int) [][]byte { +func (g *GossipBroadcaster) GetBroadcasts(overhead, limit int) [][]byte { return g.broadcasts.GetBroadcasts(overhead, limit) } -func (g *GossipMessageBroker) LocalState(join bool) []byte { +func (g *GossipBroadcaster) LocalState(join bool) []byte { pb, err := g.server.LocalState() if err != nil { g.logger().Printf("error getting local state, err=%s", err) @@ -183,7 +183,7 @@ func (g *GossipMessageBroker) LocalState(join bool) []byte { return buf } -func (g *GossipMessageBroker) MergeRemoteState(buf []byte, join bool) { +func (g *GossipBroadcaster) MergeRemoteState(buf []byte, join bool) { // Unmarshal nodestate data. var pb internal.NodeState if err := proto.Unmarshal(buf, &pb); err != nil { @@ -196,16 +196,16 @@ func (g *GossipMessageBroker) MergeRemoteState(buf []byte, join bool) { } } -// logger returns a logger for the GossipMessageBroker. -func (g *GossipMessageBroker) logger() *log.Logger { +// logger returns a logger for the GossipBroadcaster +func (g *GossipBroadcaster) logger() *log.Logger { return log.New(g.LogOutput, "", log.LstdFlags) } //////////////////////////////////////////////////////////////// -// NewGossipMessageBroker returns a new instance of GossipMessageBroker. -func NewGossipMessageBroker(s *Server) *GossipMessageBroker { - g := &GossipMessageBroker{ +// NewGossipBroadcaster returns a new instance of GossipBroadcaster. +func NewGossipBroadcaster(s *Server) *GossipBroadcaster { + g := &GossipBroadcaster{ LogOutput: os.Stderr, server: s, } diff --git a/handler.go b/handler.go index ed1418f9a..83f969d22 100644 --- a/handler.go +++ b/handler.go @@ -25,9 +25,9 @@ import ( // Handler represents an HTTP handler. type Handler struct { - Index *Index - MsgBroker MessageBroker - Server *Server + Index *Index + Broadcaster Broadcaster + Server *Server // Local hostname & cluster configuration. Host string @@ -370,7 +370,7 @@ func (h *Handler) handlePostDB(w http.ResponseWriter, r *http.Request) { } // Send the delete message to all nodes. - err := h.MsgBroker.SendSync( + err := h.Broadcaster.SendSync( &internal.DeleteDBMessage{ DB: req.DB, }) @@ -514,7 +514,7 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { } // Send the create message to all nodes. - err = h.MsgBroker.SendSync( + err = h.Broadcaster.SendSync( &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.MsgBroker.SendSync( + err := h.Broadcaster.SendSync( &internal.DeleteFrameMessage{ DB: req.DB, Frame: req.Frame, diff --git a/handler_test.go b/handler_test.go index e95e5c0dd..51f48fef7 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.MsgBroker = pilosa.NopMessageBroker + h.Broadcaster = pilosa.NopBroadcaster return h } @@ -828,7 +828,7 @@ func NewServer() *Server { s.Handler.Host = s.Host() // Handler test messages can no-op. - s.Handler.MsgBroker = pilosa.NopMessageBroker + s.Handler.Broadcaster = pilosa.NopBroadcaster // 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 ab659ada1..eb8f4b204 100644 --- a/index.go +++ b/index.go @@ -23,7 +23,7 @@ type Index struct { // Databases by name. dbs map[string]*DB - MsgBroker MessageBroker + Broadcaster Broadcaster // Close management wg sync.WaitGroup closing chan struct{} @@ -246,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.msgbroker = i.MsgBroker + db.broadcaster = i.Broadcaster return db, nil } diff --git a/messenger.go b/messenger.go index 22ba8d61f..0549e35fc 100644 --- a/messenger.go +++ b/messenger.go @@ -15,48 +15,48 @@ import ( "github.com/pilosa/pilosa/internal" ) -// MessageBroker is an interface for handling incoming/outgoing messages. -type MessageBroker interface { +// Broadcaster is an interface for handling incoming/outgoing messages. +type Broadcaster interface { SendSync(pb proto.Message) error SendAsync(pb proto.Message) error } func init() { - NopMessageBroker = &nopMessageBroker{} + NopBroadcaster = &nopBroadcaster{} } -var NopMessageBroker MessageBroker +var NopBroadcaster Broadcaster -// nopMessageBroker represents a MessageBroker that doesn't do anything. -type nopMessageBroker struct{} +// nopBroadcaster represents a Broadcaster that doesn't do anything. +type nopBroadcaster struct{} -// SendSync A no-op implemenetation of MessageBroker SendSync method. -func (c *nopMessageBroker) SendSync(pb proto.Message) error { - fmt.Println("NOPMessageBroker: SendSync") // TODO remove or log properly? +// SendSync A no-op implemenetation of Broadcaster SendSync method. +func (c *nopBroadcaster) SendSync(pb proto.Message) error { + fmt.Println("NOPBroadcaster: SendSync") // TODO remove or log properly? return nil } -// SendAsync A no-op implemenetation of MessageBroker SendAsync method. -func (c *nopMessageBroker) SendAsync(pb proto.Message) error { - fmt.Println("NOPMessageBroker: SendAsync") // TODO remove or log properly? +// SendAsync A no-op implemenetation of Broadcaster SendAsync method. +func (c *nopBroadcaster) SendAsync(pb proto.Message) error { + fmt.Println("NOPBroadcaster: SendAsync") // TODO remove or log properly? return nil } ////////////////////////////////////////////////////////////////// -// HTTPMessageBroker represents a NodeSet that broadcasts messages over HTTP. -type HTTPMessageBroker struct { +// HTTPBroadcaster represents a NodeSet that broadcasts messages over HTTP. +type HTTPBroadcaster struct { server *Server } -// NewHTTPMessageBroker returns a new instance of HTTPMessageBroker. -func NewHTTPMessageBroker(s *Server) *HTTPMessageBroker { - return &HTTPMessageBroker{server: s} +// NewHTTPBroadcaster returns a new instance of HTTPBroadcaster. +func NewHTTPBroadcaster(s *Server) *HTTPBroadcaster { + return &HTTPBroadcaster{server: s} } // SendSync sends a protobuf message to all nodes simultaneously. // It waits for all nodes to respond before the function returns (and returns any errors). -func (h *HTTPMessageBroker) SendSync(pb proto.Message) error { +func (h *HTTPBroadcaster) SendSync(pb proto.Message) error { // Marshal the pb to []byte buf, err := MarshalMessage(pb) if err != nil { @@ -82,15 +82,15 @@ func (h *HTTPMessageBroker) SendSync(pb proto.Message) error { return g.Wait() } -// SendAsync exists to implement the MessageBroker interface, but just calls +// SendAsync exists to implement the Broadcaster interface, but just calls // SendSync. -func (h *HTTPMessageBroker) SendAsync(pb proto.Message) error { +func (h *HTTPBroadcaster) SendAsync(pb proto.Message) error { return h.SendSync(pb) } -func (h *HTTPMessageBroker) nodes() ([]*Node, error) { +func (h *HTTPBroadcaster) nodes() ([]*Node, error) { if h.server == nil { - return nil, errors.New("HTTPMessageBroker has no reference to Server.") + return nil, errors.New("HTTPBroadcaster has no reference to Server.") } nodeset, ok := h.server.Cluster.NodeSet.(*HTTPNodeSet) if !ok { @@ -99,7 +99,7 @@ func (h *HTTPMessageBroker) nodes() ([]*Node, error) { return nodeset.Nodes(), nil } -func (h *HTTPMessageBroker) sendNodeMessage(node *Node, msg []byte) error { +func (h *HTTPBroadcaster) sendNodeMessage(node *Node, msg []byte) error { var client *http.Client client = http.DefaultClient diff --git a/server.go b/server.go index 68f9a6921..2f0ace9ce 100644 --- a/server.go +++ b/server.go @@ -33,9 +33,9 @@ type Server struct { closing chan struct{} // Data storage and HTTP interface. - Index *Index - Handler *Handler - MsgBroker MessageBroker + Index *Index + Handler *Handler + Broadcaster Broadcaster // Cluster configuration. // Host is replaced with actual host after opening if port is ":0". @@ -54,9 +54,9 @@ func NewServer() *Server { s := &Server{ closing: make(chan struct{}), - Index: NewIndex(), - Handler: NewHandler(), - MsgBroker: NopMessageBroker, + Index: NewIndex(), + Handler: NewHandler(), + Broadcaster: NopBroadcaster, AntiEntropyInterval: DefaultAntiEntropyInterval, PollingInterval: DefaultPollingInterval, @@ -112,14 +112,14 @@ func (s *Server) Open() error { e.Cluster = s.Cluster // Initialize HTTP handler. - s.Handler.MsgBroker = s.MsgBroker + s.Handler.Broadcaster = s.Broadcaster s.Handler.Host = s.Host s.Handler.Cluster = s.Cluster s.Handler.Executor = e s.Handler.LogOutput = s.LogOutput // Initialize Index. - s.Index.MsgBroker = s.MsgBroker + s.Index.Broadcaster = s.Broadcaster s.Index.LogOutput = s.LogOutput // Serve HTTP. diff --git a/server/server.go b/server/server.go index 69041ef3a..e880a1362 100644 --- a/server/server.go +++ b/server/server.go @@ -94,11 +94,11 @@ func (m *Command) Run(args ...string) (err error) { if err != nil { return err } - m.Server.MsgBroker = PilosaMessageBroker(m.Config, m.Server) + m.Server.Broadcaster = PilosaBroadcaster(m.Config, m.Server) m.Server.Cluster = PilosaCluster(m.Config) - // Associate objects to the MessageBroker based on config. - AssociateMessageBroker(m.Server, m.Config) + // Associate objects to the Broadcaster based on config. + AssociateBroadcaster(m.Server, m.Config) // Set configuration options. m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) @@ -111,17 +111,17 @@ func (m *Command) Run(args ...string) (err error) { return nil } -// 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": - broker = pilosa.NewHTTPMessageBroker(server) +// PilosaBroadcaster returns a new instance of Broadcaster based on the config. +func PilosaBroadcaster(c *pilosa.Config, server *pilosa.Server) (broadcaster pilosa.Broadcaster) { + switch c.Cluster.BroadcasterType { + case "http": + broadcaster = pilosa.NewHTTPBroadcaster(server) case "gossip": - broker = pilosa.NewGossipMessageBroker(server) + broadcaster = pilosa.NewGossipBroadcaster(server) case "static": - broker = pilosa.NopMessageBroker + broadcaster = pilosa.NopBroadcaster } - return broker + return broadcaster } // PilosaCluster returns a new instance of Cluster based on the config. @@ -134,8 +134,8 @@ func PilosaCluster(c *pilosa.Config) *pilosa.Cluster { } // Setup a Broadcast (over HTTP) or Gossip NodeSet based on config. - switch c.Cluster.MessengerType { - case "broadcast": + switch c.Cluster.BroadcasterType { + case "http": cluster.NodeSet = pilosa.NewHTTPNodeSet() cluster.NodeSet.(*pilosa.HTTPNodeSet).Join(cluster.Nodes) case "gossip": @@ -166,14 +166,14 @@ func PilosaCluster(c *pilosa.Config) *pilosa.Cluster { return cluster } -// AssociateMessageBroker allows an implementation to associate objects to the MessageBroker +// AssociateBroadcaster allows an implementation to associate objects to the Broadcaster // after cluster configuration. -func AssociateMessageBroker(s *pilosa.Server, c *pilosa.Config) { - switch c.Cluster.MessengerType { - case "broadcast": +func AssociateBroadcaster(s *pilosa.Server, c *pilosa.Config) { + switch c.Cluster.BroadcasterType { + case "http": // nop case "gossip": - s.Cluster.NodeSet.(*pilosa.GossipNodeSet).AttachBroker(s.MsgBroker.(*pilosa.GossipMessageBroker)) + s.Cluster.NodeSet.(*pilosa.GossipNodeSet).AttachBroadcaster(s.Broadcaster.(*pilosa.GossipBroadcaster)) case "static": // nop }