add stateHandler interface to gossip

now it doesn't need to directly reference pilosa.Server
This commit is contained in:
Matt Jaffee 2017-04-17 22:31:32 -05:00 • committed by Travis
parent 2d139d299c
commit 059403e673

View file

@ -13,6 +13,14 @@ import (
"github.com/pilosa/pilosa/internal"
)
// StateHandler specifies two methods which an object must implement to share
// state in the cluster. These are used by the GossipNodeSet to implement the
// LocalState and MergeRemoteState methods of memberlist.Delegate
type StateHandler interface {
LocalState() (proto.Message, error)
HandleRemoteState(proto.Message) error
}
// GossipNodeSet represents a gossip implementation of NodeSet using memberlist
// GossipNodeSet also represents a gossip implementation of pilosa.Broadcaster
// GossipNodeSet also represents an implementation of memberlist.Delegate
@ -22,9 +30,8 @@ type GossipNodeSet struct {
broadcasts *memberlist.TransmitLimitedQueue
server *Server
config *GossipConfig
stateHandler StateHandler
config *GossipConfig
// The writer for any logging.
LogOutput io.Writer
@ -81,7 +88,7 @@ type GossipConfig struct {
}
// NewGossipNodeSet returns a new instance of GossipNodeSet.
func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string, s *Server) *GossipNodeSet {
func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string, sh StateHandler) *GossipNodeSet {
g := &GossipNodeSet{
LogOutput: os.Stderr,
}
@ -98,7 +105,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed
g.config.memberlistConfig.AdvertisePort = gossipPort
g.config.memberlistConfig.Delegate = g
g.server = s
g.stateHandler = sh
return g
}
@ -110,7 +117,7 @@ func (g *GossipNodeSet) SendSync(pb proto.Message) error {
return err
}
mlist := g.server.Cluster.NodeSet.(*GossipNodeSet).memberlist
mlist := g.memberlist
// Direct sends the message directly to every node.
// An error from any node raises an error on the entire operation.
@ -175,7 +182,7 @@ func (g *GossipNodeSet) GetBroadcasts(overhead, limit int) [][]byte {
}
func (g *GossipNodeSet) LocalState(join bool) []byte {
pb, err := g.server.LocalState()
pb, err := g.stateHandler.LocalState()
if err != nil {
g.logger().Printf("error getting local state, err=%s", err)
return []byte{}
@ -197,7 +204,7 @@ func (g *GossipNodeSet) MergeRemoteState(buf []byte, join bool) {
g.logger().Printf("error unmarshalling nodestate data, err=%s", err)
return
}
err := g.server.HandleRemoteState(&pb)
err := g.stateHandler.HandleRemoteState(&pb)
if err != nil {
g.logger().Printf("merge state error: %s", err)
}