split MessageBroker.Send into SendSync and SendAsync

This commit is contained in:
Matt Jaffee 2017-04-17 12:30:02 -05:00 • committed by Travis
parent 5079ee7a22
commit 4b8c2630e3
4 changed files with 55 additions and 33 deletions

View file

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

View file

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

View file

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

View file

@ -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.")