diff --git a/cluster.go b/cluster.go index a435a20b7..ef67b7a4a 100644 --- a/cluster.go +++ b/cluster.go @@ -263,6 +263,10 @@ func (s *StaticNodeSet) Open() error { return nil } -func (s *StaticNodeSet) Send(pb proto.Message, method string) error { +func (s *StaticNodeSet) SendSync(pb proto.Message) error { + return nil +} + +func (s *StaticNodeSet) SendAsync(pb proto.Message) error { return nil } diff --git a/gossip.go b/gossip.go index 37b8c3b97..6482a333d 100644 --- a/gossip.go +++ b/gossip.go @@ -96,8 +96,8 @@ type GossipMessageBroker struct { LogOutput io.Writer } -// implementation of the messenger.MessageBroker interface -func (g *GossipMessageBroker) Send(pb proto.Message, method string) error { +// SendSync implementation of the messenger.MessageBroker interface +func (g *GossipMessageBroker) SendSync(pb proto.Message) error { msg, err := MarshalMessage(pb) if err != nil { return err @@ -110,28 +110,32 @@ func (g *GossipMessageBroker) Send(pb proto.Message, method string) error { // // Gossip uses the gossip protocol to eventually deliver the message // to every node. - switch method { - case "direct": - var eg errgroup.Group - for _, n := range mlist.Members() { - // Don't send the message to the local node. - if n == mlist.LocalNode() { - continue - } - node := n - eg.Go(func() error { - return mlist.SendToTCP(node, msg) - }) + var eg errgroup.Group + for _, n := range mlist.Members() { + // Don't send the message to the local node. + if n == mlist.LocalNode() { + continue } - return eg.Wait() - case "gossip": - b := &broadcast{ - msg: msg, - notify: nil, - } - g.broadcasts.QueueBroadcast(b) + node := n + eg.Go(func() error { + return mlist.SendToTCP(node, msg) + }) + } + return eg.Wait() +} + +// SendAsync implementation of the messenger.MessageBroker interface +func (g *GossipMessageBroker) SendAsync(pb proto.Message) error { + msg, err := MarshalMessage(pb) + if err != nil { + return err } + b := &broadcast{ + msg: msg, + notify: nil, + } + g.broadcasts.QueueBroadcast(b) return nil } diff --git a/handler.go b/handler.go index 5c015abde..ed1418f9a 100644 --- a/handler.go +++ b/handler.go @@ -370,10 +370,10 @@ func (h *Handler) handlePostDB(w http.ResponseWriter, r *http.Request) { } // Send the delete message to all nodes. - err := h.MsgBroker.Send( + err := h.MsgBroker.SendSync( &internal.DeleteDBMessage{ DB: req.DB, - }, "direct") + }) if err != nil { h.logger().Printf("problem sending DeleteDB message: %s", err) } @@ -514,7 +514,7 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { } // Send the create message to all nodes. - err = h.MsgBroker.Send( + err = h.MsgBroker.SendSync( &internal.CreateFrameMessage{ DB: req.DB, Frame: req.Frame, @@ -522,7 +522,7 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { RowLabel: req.Options.RowLabel, TimeQuantum: string(req.Options.TimeQuantum), }, - }, "direct") + }) if err != nil { h.logger().Printf("problem sending CreateFrame message: %s", err) } @@ -586,11 +586,11 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) { } // Send the delete message to all nodes. - err := h.MsgBroker.Send( + err := h.MsgBroker.SendSync( &internal.DeleteFrameMessage{ DB: req.DB, Frame: req.Frame, - }, "direct") + }) if err != nil { h.logger().Printf("problem sending DeleteFrame message: %s", err) } diff --git a/messenger.go b/messenger.go index fc0c55d77..22ba8d61f 100644 --- a/messenger.go +++ b/messenger.go @@ -17,7 +17,8 @@ import ( // MessageBroker is an interface for handling incoming/outgoing messages. type MessageBroker interface { - Send(pb proto.Message, method string) error + SendSync(pb proto.Message) error + SendAsync(pb proto.Message) error } func init() { @@ -29,8 +30,15 @@ var NopMessageBroker MessageBroker // nopMessageBroker represents a MessageBroker that doesn't do anything. type nopMessageBroker struct{} -func (c *nopMessageBroker) Send(pb proto.Message, method string) error { - fmt.Println("NOPMessageBroker: Send") // TODO remove or log properly? +// 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? + 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? return nil } @@ -46,9 +54,9 @@ func NewHTTPMessageBroker(s *Server) *HTTPMessageBroker { return &HTTPMessageBroker{server: s} } -// Send sends a protobuf message to all nodes simultaneously. +// 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) Send(pb proto.Message, method string) error { +func (h *HTTPMessageBroker) SendSync(pb proto.Message) error { // Marshal the pb to []byte buf, err := MarshalMessage(pb) if err != nil { @@ -74,6 +82,12 @@ func (h *HTTPMessageBroker) Send(pb proto.Message, method string) error { return g.Wait() } +// SendAsync exists to implement the MessageBroker interface, but just calls +// SendSync. +func (h *HTTPMessageBroker) SendAsync(pb proto.Message) error { + return h.SendSync(pb) +} + func (h *HTTPMessageBroker) nodes() ([]*Node, error) { if h.server == nil { return nil, errors.New("HTTPMessageBroker has no reference to Server.")