mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-13 08:01:02 +00:00
remove broadcaster methods from gossip- don't use sendAsync anywhere
This commit is contained in:
parent
6ef9b0bb93
commit
56ed9bfbe1
5 changed files with 4 additions and 53 deletions
|
|
@ -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)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
||||
|
|
|
|||
4
view.go
4
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")
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue