diff --git a/cluster.go b/cluster.go index 4ddaf3a87..8dc3a60d1 100644 --- a/cluster.go +++ b/cluster.go @@ -910,7 +910,7 @@ func (c *Cluster) open() error { Event: uint32(NodeJoin), Node: EncodeNode(c.Node), } - if err := c.Broadcaster.SendAsync(msg); err != nil { + if err := c.Broadcaster.SendSync(msg); err != nil { return fmt.Errorf("sending restart NodeJoin: %v", err) } diff --git a/gossip/gossip.go b/gossip/gossip.go index dfcb2f759..9154e9fc7 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -24,8 +24,6 @@ import ( "sync" "time" - "golang.org/x/sync/errgroup" - "github.com/gogo/protobuf/proto" "github.com/hashicorp/memberlist" "github.com/pilosa/pilosa" @@ -36,7 +34,6 @@ import ( // Ensure GossipMemberSet implements interfaces. var _ pilosa.BroadcastReceiver = &GossipMemberSet{} -var _ pilosa.Gossiper = &GossipMemberSet{} var _ memberlist.Delegate = &GossipMemberSet{} // GossipMemberSet represents a gossip implementation of MemberSet using memberlist. @@ -240,49 +237,6 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe return g, nil } -// SendSync implementation of the Broadcaster interface. -func (g *GossipMemberSet) SendSync(pb proto.Message) error { - msg, err := pilosa.MarshalMessage(pb) - if err != nil { - return fmt.Errorf("marshal message: %s", err) - } - - mlist := g.memberlist - - // Direct sends the message directly to every node. - // An error from any node raises an error on the entire operation. - // - // Gossip uses the gossip protocol to eventually deliver the message - // to every node. - 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) - }) - } - return eg.Wait() -} - -// SendAsync implementation of the Gossiper interface. -func (g *GossipMemberSet) SendAsync(pb proto.Message) error { - msg, err := pilosa.MarshalMessage(pb) - if err != nil { - return fmt.Errorf("marshal message: %s", err) - } - - b := &broadcast{ - msg: msg, - notify: nil, - } - g.broadcasts.QueueBroadcast(b) - return nil -} - // NodeMeta implementation of the memberlist.Delegate interface. func (g *GossipMemberSet) NodeMeta(limit int) []byte { buf, err := proto.Marshal(pilosa.EncodeNode(g.node)) diff --git a/server.go b/server.go index 4a453cef5..dd26a2c0b 100644 --- a/server.go +++ b/server.go @@ -62,7 +62,6 @@ type Server struct { handler Handler Broadcaster Broadcaster BroadcastReceiver BroadcastReceiver - Gossiper Gossiper systemInfo SystemInfo gcNotifier GCNotifier NewAttrStore func(string) AttrStore @@ -547,7 +546,7 @@ func (s *Server) SendSync(pb proto.Message) error { // SendAsync represents an implementation of Broadcaster. func (s *Server) SendAsync(pb proto.Message) error { - return s.Gossiper.SendAsync(pb) + return ErrNotImplemented } // SendTo represents an implementation of Broadcaster. diff --git a/server/server.go b/server/server.go index 0c4184d60..b6bf6e864 100644 --- a/server/server.go +++ b/server/server.go @@ -268,7 +268,6 @@ func (m *Command) SetupNetworking() error { m.Server.Broadcaster = pilosa.NopBroadcaster m.Server.Cluster.MemberSet = pilosa.NewStaticMemberSet(m.Server.Cluster.Nodes) m.Server.BroadcastReceiver = pilosa.NopBroadcastReceiver - m.Server.Gossiper = pilosa.NopGossiper return nil } @@ -313,7 +312,6 @@ func (m *Command) SetupNetworking() error { m.Server.Cluster.MemberSet = gossipMemberSet m.Server.Broadcaster = m.Server m.Server.BroadcastReceiver = gossipMemberSet - m.Server.Gossiper = gossipMemberSet return nil } diff --git a/view.go b/view.go index b21c3f15a..428fc6f54 100644 --- a/view.go +++ b/view.go @@ -237,13 +237,13 @@ func (v *View) createFragmentIfNotExists(slice uint64) (*Fragment, error) { v.maxSlice = slice // Send the create slice message to all nodes. - err := v.broadcaster.SendAsync( + err := v.broadcaster.SendSync( &internal.CreateSliceMessage{ Index: v.index, Slice: slice, }) if err != nil { - return nil, errors.Wrap(err, "sending message") + return nil, errors.Wrap(err, "sending createslice message") } }