mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
simplify gossipEventReceiver - no longer needs separate Start method
also remove some dead code
This commit is contained in:
parent
965bd08225
commit
75c1440eb5
1 changed files with 11 additions and 42 deletions
|
|
@ -55,12 +55,7 @@ type GossipMemberSet struct {
|
|||
}
|
||||
|
||||
// Open implements the MemberSet interface to start network activity.
|
||||
func (g *GossipMemberSet) Open() error {
|
||||
err := g.gossipEventReceiver.Start(g.pserver)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "starting event delegate")
|
||||
}
|
||||
|
||||
func (g *GossipMemberSet) Open() (err error) {
|
||||
g.mu.Lock()
|
||||
g.memberlist, err = memberlist.Create(g.config.memberlistConfig)
|
||||
g.mu.Unlock()
|
||||
|
|
@ -163,11 +158,9 @@ func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSet
|
|||
return nil, errors.Wrap(err, "executing option")
|
||||
}
|
||||
}
|
||||
ger := newGossipEventReceiver(g.logger)
|
||||
ger := newGossipEventReceiver(g.logger, s)
|
||||
g.gossipEventReceiver = ger
|
||||
|
||||
g.handler = s
|
||||
|
||||
if g.transport == nil {
|
||||
port, err := strconv.Atoi(cfg.Port)
|
||||
if err != nil {
|
||||
|
|
@ -299,15 +292,18 @@ type gossipEventReceiver struct {
|
|||
ch chan memberlist.NodeEvent
|
||||
eventHandler pilosa.EventHandler
|
||||
|
||||
Logger *log.Logger
|
||||
logger *log.Logger
|
||||
}
|
||||
|
||||
// newGossipEventReceiver returns a new instance of GossipEventReceiver.
|
||||
func newGossipEventReceiver(logger *log.Logger) *gossipEventReceiver {
|
||||
return &gossipEventReceiver{
|
||||
ch: make(chan memberlist.NodeEvent, 1),
|
||||
Logger: logger,
|
||||
func newGossipEventReceiver(logger *log.Logger, pserver pilosa.EventHandler) *gossipEventReceiver {
|
||||
ger := &gossipEventReceiver{
|
||||
ch: make(chan memberlist.NodeEvent, 1),
|
||||
logger: logger,
|
||||
eventHandler: pserver,
|
||||
}
|
||||
go ger.listen()
|
||||
return ger
|
||||
}
|
||||
|
||||
func (g *gossipEventReceiver) NotifyJoin(n *memberlist.Node) {
|
||||
|
|
@ -322,13 +318,6 @@ func (g *gossipEventReceiver) NotifyUpdate(n *memberlist.Node) {
|
|||
g.ch <- memberlist.NodeEvent{memberlist.NodeUpdate, n}
|
||||
}
|
||||
|
||||
// Start implements the pilosa.EventReceiver interface and sets the EventHandler.
|
||||
func (g *gossipEventReceiver) Start(h pilosa.EventHandler) error {
|
||||
g.eventHandler = h
|
||||
go g.listen()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (g *gossipEventReceiver) listen() {
|
||||
var nodeEventType pilosa.NodeEventType
|
||||
for {
|
||||
|
|
@ -356,31 +345,11 @@ func (g *gossipEventReceiver) listen() {
|
|||
Node: node,
|
||||
}
|
||||
if err := g.eventHandler.ReceiveEvent(ne); err != nil {
|
||||
g.Logger.Printf("receive event error: %s", err)
|
||||
g.logger.Printf("receive event error: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// broadcast represents an implementation of memberlist.Broadcast
|
||||
type broadcast struct {
|
||||
msg []byte
|
||||
notify chan<- struct{}
|
||||
}
|
||||
|
||||
func (b *broadcast) Invalidates(other memberlist.Broadcast) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (b *broadcast) Message() []byte {
|
||||
return b.msg
|
||||
}
|
||||
|
||||
func (b *broadcast) Finished() {
|
||||
if b.notify != nil {
|
||||
close(b.notify)
|
||||
}
|
||||
}
|
||||
|
||||
// Transport is a gossip transport for binding to a port.
|
||||
type Transport struct {
|
||||
//memberlist.Transport
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue