From c4d3bde50198802c3df44369c165eacd62b0b14a Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Tue, 18 Apr 2017 08:41:23 -0500 Subject: [PATCH] move gossip code to separate package --- gossip.go => gossip/gossip.go | 23 ++++++++++++----------- server/server.go | 3 ++- 2 files changed, 14 insertions(+), 12 deletions(-) rename gossip.go => gossip/gossip.go (90%) diff --git a/gossip.go b/gossip/gossip.go similarity index 90% rename from gossip.go rename to gossip/gossip.go index 664b91aa7..050435441 100644 --- a/gossip.go +++ b/gossip/gossip.go @@ -1,4 +1,4 @@ -package pilosa +package gossip import ( "fmt" @@ -10,6 +10,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/hashicorp/memberlist" + "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" ) @@ -26,7 +27,7 @@ type StateHandler interface { // GossipNodeSet also represents an implementation of memberlist.Delegate type GossipNodeSet struct { memberlist *memberlist.Memberlist - handler BroadcastHandler + handler pilosa.BroadcastHandler broadcasts *memberlist.TransmitLimitedQueue @@ -37,15 +38,15 @@ type GossipNodeSet struct { LogOutput io.Writer } -func (g *GossipNodeSet) Nodes() []*Node { - a := make([]*Node, 0, g.memberlist.NumMembers()) +func (g *GossipNodeSet) Nodes() []*pilosa.Node { + a := make([]*pilosa.Node, 0, g.memberlist.NumMembers()) for _, n := range g.memberlist.Members() { - a = append(a, &Node{Host: n.Name}) + a = append(a, &pilosa.Node{Host: n.Name}) } return a } -func (g *GossipNodeSet) Start(h BroadcastHandler) error { +func (g *GossipNodeSet) Start(h pilosa.BroadcastHandler) error { g.handler = h return nil } @@ -61,8 +62,8 @@ func (g *GossipNodeSet) Open() error { g.memberlist = ml // attach to gossip seed node - nodes := []*Node{&Node{Host: g.config.gossipSeed}} //TODO: support a list of seeds - _, err = g.memberlist.Join(Nodes(nodes).Hosts()) + nodes := []*pilosa.Node{&pilosa.Node{Host: g.config.gossipSeed}} //TODO: support a list of seeds + _, err = g.memberlist.Join(pilosa.Nodes(nodes).Hosts()) if err != nil { return err } @@ -112,7 +113,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed // SendSync implementation of the Broadcaster interface func (g *GossipNodeSet) SendSync(pb proto.Message) error { - msg, err := MarshalMessage(pb) + msg, err := pilosa.MarshalMessage(pb) if err != nil { return err } @@ -140,7 +141,7 @@ func (g *GossipNodeSet) SendSync(pb proto.Message) error { // SendAsync implementation of the Broadcaster interface func (g *GossipNodeSet) SendAsync(pb proto.Message) error { - msg, err := MarshalMessage(pb) + msg, err := pilosa.MarshalMessage(pb) if err != nil { return err } @@ -166,7 +167,7 @@ func (g *GossipNodeSet) NodeMeta(limit int) []byte { } func (g *GossipNodeSet) NotifyMsg(b []byte) { - m, err := UnmarshalMessage(b) + m, err := pilosa.UnmarshalMessage(b) if err != nil { g.logger().Printf("unmarshal message error: %s", err) return diff --git a/server/server.go b/server/server.go index 9d00eabfb..80c0457f8 100644 --- a/server/server.go +++ b/server/server.go @@ -17,6 +17,7 @@ import ( "time" "github.com/pilosa/pilosa" + "github.com/pilosa/pilosa/gossip" ) func init() { @@ -143,7 +144,7 @@ func (m *Command) SetupServer() error { if err != nil { gossipHost = m.Config.Host } - gossipNodeSet := pilosa.NewGossipNodeSet(m.Config.Host, gossipHost, gossipPort, gossipSeed, m.Server) + gossipNodeSet := gossip.NewGossipNodeSet(m.Config.Host, gossipHost, gossipPort, gossipSeed, m.Server) m.Server.Cluster.NodeSet = gossipNodeSet m.Server.Broadcaster = gossipNodeSet m.Server.BroadcastReceiver = gossipNodeSet