From 09adc3b46182d986b6b50a9b39d79587a0540e24 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Thu, 28 Jun 2018 07:54:56 -0500 Subject: [PATCH 1/6] unexport gossipEventRecevier in gossip.go --- gossip/gossip.go | 23 ++++++++++------------- 1 file changed, 10 insertions(+), 13 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index 0355260d6..f5297031c 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -51,7 +51,7 @@ type GossipMemberSet struct { logger *log.Logger transport *Transport - gossipEventReceiver *GossipEventReceiver + gossipEventReceiver *gossipEventReceiver } // GetBindAddr returns the gossip bind address based on config and auto bind port. @@ -67,9 +67,6 @@ func (g *GossipMemberSet) Open() error { if err != nil { return errors.Wrap(err, "starting event delegate") } - if g.handler == nil { - return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()") - } g.mu.Lock() g.memberlist, err = memberlist.Create(g.config.memberlistConfig) @@ -300,12 +297,12 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) { } } -// GossipEventReceiver is used to enable an application to receive +// gossipEventReceiver is used to enable an application to receive // events about joins and leaves over a channel. // // Care must be taken that events are processed in a timely manner from // the channel, since this delegate will block until an event can be sent. -type GossipEventReceiver struct { +type gossipEventReceiver struct { ch chan memberlist.NodeEvent eventHandler pilosa.EventHandler @@ -313,33 +310,33 @@ type GossipEventReceiver struct { } // NewGossipEventReceiver returns a new instance of GossipEventReceiver. -func NewGossipEventReceiver(logger *log.Logger) *GossipEventReceiver { - return &GossipEventReceiver{ +func NewGossipEventReceiver(logger *log.Logger) *gossipEventReceiver { + return &gossipEventReceiver{ ch: make(chan memberlist.NodeEvent, 1), Logger: logger, } } -func (g *GossipEventReceiver) NotifyJoin(n *memberlist.Node) { +func (g *gossipEventReceiver) NotifyJoin(n *memberlist.Node) { g.ch <- memberlist.NodeEvent{memberlist.NodeJoin, n} } -func (g *GossipEventReceiver) NotifyLeave(n *memberlist.Node) { +func (g *gossipEventReceiver) NotifyLeave(n *memberlist.Node) { g.ch <- memberlist.NodeEvent{memberlist.NodeLeave, n} } -func (g *GossipEventReceiver) NotifyUpdate(n *memberlist.Node) { +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 { +func (g *gossipEventReceiver) Start(h pilosa.EventHandler) error { g.eventHandler = h go g.listen() return nil } -func (g *GossipEventReceiver) listen() { +func (g *gossipEventReceiver) listen() { var nodeEventType pilosa.NodeEventType for { e := <-g.ch From 965bd0822504bf448a78b29b51c77b8bf799340e Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Thu, 28 Jun 2018 11:00:45 -0500 Subject: [PATCH 2/6] unexport newGossipEventReceiver and remove some dead code --- gossip/gossip.go | 13 +++---------- 1 file changed, 3 insertions(+), 10 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index f5297031c..40af378c3 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -54,13 +54,6 @@ type GossipMemberSet struct { gossipEventReceiver *gossipEventReceiver } -// GetBindAddr returns the gossip bind address based on config and auto bind port. -// This method is currently only used in a test scenario where a second node needs -// the auto-bind address of the first node to use as its gossip seed. -func (g *GossipMemberSet) GetBindAddr() string { - return fmt.Sprintf("%s:%d", g.config.memberlistConfig.BindAddr, g.config.memberlistConfig.BindPort) -} - // Open implements the MemberSet interface to start network activity. func (g *GossipMemberSet) Open() error { err := g.gossipEventReceiver.Start(g.pserver) @@ -170,7 +163,7 @@ func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSet return nil, errors.Wrap(err, "executing option") } } - ger := NewGossipEventReceiver(g.logger) + ger := newGossipEventReceiver(g.logger) g.gossipEventReceiver = ger g.handler = s @@ -309,8 +302,8 @@ type gossipEventReceiver struct { Logger *log.Logger } -// NewGossipEventReceiver returns a new instance of GossipEventReceiver. -func NewGossipEventReceiver(logger *log.Logger) *gossipEventReceiver { +// newGossipEventReceiver returns a new instance of GossipEventReceiver. +func newGossipEventReceiver(logger *log.Logger) *gossipEventReceiver { return &gossipEventReceiver{ ch: make(chan memberlist.NodeEvent, 1), Logger: logger, From 75c1440eb5fa513275610c6cc7ff1dda6fa9ae30 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Thu, 28 Jun 2018 11:03:56 -0500 Subject: [PATCH 3/6] simplify gossipEventReceiver - no longer needs separate Start method also remove some dead code --- gossip/gossip.go | 53 ++++++++++-------------------------------------- 1 file changed, 11 insertions(+), 42 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index 40af378c3..3bfccaaac 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -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 From 12a49c3e143f80cfb17abc200c4795756e14e380 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Thu, 28 Jun 2018 14:45:20 -0500 Subject: [PATCH 4/6] remove ClusterStatus method and StatusHandler interface gossipEventReceiver uses ReceiveMessage instead of ReceiveEvent --- gossip/gossip.go | 14 +++++++------- server.go | 15 --------------- 2 files changed, 7 insertions(+), 22 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index 3bfccaaac..477fdb5ae 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -290,13 +290,13 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) { // the channel, since this delegate will block until an event can be sent. type gossipEventReceiver struct { ch chan memberlist.NodeEvent - eventHandler pilosa.EventHandler + eventHandler *pilosa.Server logger *log.Logger } // newGossipEventReceiver returns a new instance of GossipEventReceiver. -func newGossipEventReceiver(logger *log.Logger, pserver pilosa.EventHandler) *gossipEventReceiver { +func newGossipEventReceiver(logger *log.Logger, pserver *pilosa.Server) *gossipEventReceiver { ger := &gossipEventReceiver{ ch: make(chan memberlist.NodeEvent, 1), logger: logger, @@ -338,13 +338,13 @@ func (g *gossipEventReceiver) listen() { if err := proto.Unmarshal(e.Node.Meta, &n); err != nil { panic("failed to unmarshal event node meta data") } - node := pilosa.DecodeNode(&n) + // node := pilosa.DecodeNode(&n) - ne := &pilosa.NodeEvent{ - Event: nodeEventType, - Node: node, + ne := &internal.NodeEventMessage{ + Event: uint32(nodeEventType), + Node: &n, } - if err := g.eventHandler.ReceiveEvent(ne); err != nil { + if err := g.eventHandler.ReceiveMessage(ne); err != nil { g.logger.Printf("receive event error: %s", err) } } diff --git a/server.go b/server.go index c07720a54..083ad4e2a 100644 --- a/server.go +++ b/server.go @@ -42,7 +42,6 @@ const ( // Ensure Server implements interfaces. var _ Broadcaster = &Server{} var _ BroadcastHandler = &Server{} -var _ StatusHandler = &Server{} // Server represents a holder wrapped by a running HTTP server. type Server struct { @@ -557,11 +556,6 @@ func (s *Server) LocalStatus() (proto.Message, error) { return &ns, nil } -// ClusterStatus returns the ClusterState and NodeSet for the cluster. -func (s *Server) ClusterStatus() (proto.Message, error) { - return s.cluster.Status(), nil -} - // HandleRemoteStatus receives incoming NodeStatus from remote nodes. func (s *Server) HandleRemoteStatus(pb proto.Message) error { // Ignore NodeStatus messages until the cluster is in a Normal state. @@ -733,15 +727,6 @@ func countOpenFiles() (int, error) { } } -// StatusHandler specifies the methods which an object must implement to share -// state in the cluster. These are used by the GossipMemberSet to implement the -// LocalState and MergeRemoteState methods of memberlist.Delegate -type StatusHandler interface { - LocalStatus() (proto.Message, error) - ClusterStatus() (proto.Message, error) - HandleRemoteStatus(proto.Message) error -} - func expandDirName(path string) (string, error) { prefix := "~" + string(filepath.Separator) if strings.HasPrefix(path, prefix) { From 91454a5cd0a8756e945ed9001cb8341c58fea416 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Thu, 28 Jun 2018 17:16:08 -0500 Subject: [PATCH 5/6] collapse *handler interfaces into MemberServer gossip now takes a single "MemberServer" which is implemented by server. Several interfaces have been removed. MemberServer contains ReceiveMessage which is a superset of the functionality of ReceiveEvent, LocalStatus and HandleRemoteStatus are all that's left of StatusHandler - ClusterStatus was not used and is gone. The Node() method is actually a subset of LocalStatus() functionality. Maybe we should break up localstatus or remove Node... not sure. Remove BroadcastReceiver test which was a bit silly. NodeEvent can now be unexported, and is. --- Gopkg.lock | 2 +- broadcast.go | 6 ---- broadcast_test.go | 69 ------------------------------------------ cluster.go | 6 ++-- event.go | 13 ++------ gossip/gossip.go | 5 ++- server.go | 14 +++++---- utils_internal_test.go | 2 +- 8 files changed, 18 insertions(+), 99 deletions(-) diff --git a/Gopkg.lock b/Gopkg.lock index d3a12ef8a..33187bfe2 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -310,6 +310,6 @@ [solve-meta] analyzer-name = "dep" analyzer-version = 1 - inputs-digest = "40bd9c0a1a403580ad77f9ae84e81a97da1d1622b3f620bd000271c52b50b8b5" + inputs-digest = "da6d02118ca77527c4ff00e9522880032fc052fb39bc8efe6c76602857c8c84e" solver-name = "gps-cdcl" solver-version = 1 diff --git a/broadcast.go b/broadcast.go index b7b13fe04..f97438479 100644 --- a/broadcast.go +++ b/broadcast.go @@ -54,12 +54,6 @@ func (c *nopBroadcaster) SendTo(to *Node, pb proto.Message) error { return nil } -// BroadcastHandler is the interface for the pilosa object which knows how to -// handle broadcast messages. (Hint: this is implemented by pilosa.Server) -type BroadcastHandler interface { - ReceiveMessage(pb proto.Message) error -} - // Broadcast message types. const ( messageTypeCreateSlice = iota diff --git a/broadcast_test.go b/broadcast_test.go index 23bd962a7..ab4bcaec5 100644 --- a/broadcast_test.go +++ b/broadcast_test.go @@ -15,16 +15,12 @@ package pilosa_test import ( - "bytes" "reflect" "testing" - "io/ioutil" - "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" - "github.com/pilosa/pilosa/server" ) // Ensure a message can be marshaled and unmarshaled. @@ -53,68 +49,3 @@ func testMessageMarshal(t *testing.T, m proto.Message) { t.Fatalf("unexpected message marshalling: %s", unmarshalled) } } - -// Ensure that BroadcastReceiver can register a BroadcastHandler. -func TestBroadcast_BroadcastReceiver(t *testing.T) { - t.Skip("broadcast receiver") - path, err := ioutil.TempDir("", "pilosa-") - if err != nil { - panic(err) - } - com := server.NewCommand(bytes.NewBuffer([]byte{}), ioutil.Discard, ioutil.Discard) - com.Config.Bind = "localhost:0" - com.Config.DataDir = path - err = com.SetupServer() // this test shouldn't need to import pilosa/server just to set up the Server, but it really shouldn't need to setup the Server at all. The Server should not be the implementation of Broadcast* TODO - if err != nil { - t.Fatalf("setting up server: %v", err) - } - // s := com.Server - - // sbr := NewSimpleBroadcastReceiver() - // sbh := NewSimpleBroadcastHandler() - - // s.BroadcastReceiver = sbr - // s.BroadcastReceiver.Start(sbh) - - // msg := &internal.DeleteIndexMessage{ - // Index: "i", - // } - - // s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg) - - // // Make sure the message received is what was sentd - // if !reflect.DeepEqual(sbh.receivedMessage, msg) { - // t.Fatalf("unexpected message: %s", sbh.receivedMessage) - // } -} - -type SimpleBroadcastReceiver struct { - broadcastHandler pilosa.BroadcastHandler -} - -func NewSimpleBroadcastReceiver() *SimpleBroadcastReceiver { - return &SimpleBroadcastReceiver{} -} - -func (r *SimpleBroadcastReceiver) Start(h pilosa.BroadcastHandler) error { - r.broadcastHandler = h - return nil -} - -func (r *SimpleBroadcastReceiver) Receive(pb proto.Message) error { - r.broadcastHandler.ReceiveMessage(pb) - return nil -} - -type SimpleBroadcastHandler struct { - receivedMessage proto.Message -} - -func NewSimpleBroadcastHandler() *SimpleBroadcastHandler { - return &SimpleBroadcastHandler{} -} - -func (h *SimpleBroadcastHandler) ReceiveMessage(pb proto.Message) error { - h.receivedMessage = pb.(proto.Message) - return nil -} diff --git a/cluster.go b/cluster.go index de728472a..48428a5d5 100644 --- a/cluster.go +++ b/cluster.go @@ -108,8 +108,8 @@ func DecodeNode(node *internal.Node) *Node { } } -func DecodeNodeEvent(ne *internal.NodeEventMessage) *NodeEvent { - return &NodeEvent{ +func DecodeNodeEvent(ne *internal.NodeEventMessage) *nodeEvent { + return &nodeEvent{ Event: NodeEventType(ne.Event), Node: DecodeNode(ne.Node), } @@ -1612,7 +1612,7 @@ func (c *Cluster) considerTopology() error { } // ReceiveEvent represents an implementation of EventHandler. -func (c *Cluster) ReceiveEvent(e *NodeEvent) error { +func (c *Cluster) ReceiveEvent(e *nodeEvent) error { // Ignore events sent from this node. if e.Node.ID == c.Node.ID { return nil diff --git a/event.go b/event.go index 69229b1e4..aa1e0e890 100644 --- a/event.go +++ b/event.go @@ -14,8 +14,7 @@ package pilosa -// NodeEventType are the types of events that can be sent from the -// ChannelEventDelegate. +// NodeEventType are the types of node events. type NodeEventType int const ( @@ -24,14 +23,8 @@ const ( NodeUpdate ) -// NodeEvent is a single event related to node activity in the cluster. -type NodeEvent struct { +// nodeEvent is a single event related to node activity in the cluster. +type nodeEvent struct { Event NodeEventType Node *Node } - -// EventHandler is the interface for the pilosa object which knows how to -// handle broadcast messages. (Hint: this is implemented by pilosa.Server) -type EventHandler interface { - ReceiveEvent(e *NodeEvent) error -} diff --git a/gossip/gossip.go b/gossip/gossip.go index 477fdb5ae..2200efad7 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -39,11 +39,10 @@ var _ memberlist.Delegate = &GossipMemberSet{} type GossipMemberSet struct { mu sync.RWMutex memberlist *memberlist.Memberlist - handler pilosa.BroadcastHandler broadcasts *memberlist.TransmitLimitedQueue - pserver *pilosa.Server + pserver pilosa.MemberServer config *gossipConfig Logger pilosa.Logger @@ -238,7 +237,7 @@ func (g *GossipMemberSet) NotifyMsg(b []byte) { g.Logger.Printf("unmarshal message error: %s", err) return } - if err := g.handler.ReceiveMessage(m); err != nil { + if err := g.pserver.ReceiveMessage(m); err != nil { g.Logger.Printf("receive message error: %s", err) return } diff --git a/server.go b/server.go index 083ad4e2a..c1b787e04 100644 --- a/server.go +++ b/server.go @@ -41,7 +41,7 @@ const ( // Ensure Server implements interfaces. var _ Broadcaster = &Server{} -var _ BroadcastHandler = &Server{} +var _ MemberServer = &Server{} // Server represents a holder wrapped by a running HTTP server. type Server struct { @@ -701,11 +701,6 @@ func (s *Server) monitorRuntime() { } } -// ReceiveEvent implements the EventHandler interface. -func (s *Server) ReceiveEvent(e *NodeEvent) error { - return s.cluster.ReceiveEvent(e) -} - // countOpenFiles on operating systems that support lsof. func countOpenFiles() (int, error) { switch runtime.GOOS { @@ -738,3 +733,10 @@ func expandDirName(path string) (string, error) { } return path, nil } + +type MemberServer interface { + ReceiveMessage(proto.Message) error + LocalStatus() (proto.Message, error) + HandleRemoteStatus(proto.Message) error + Node() *Node +} diff --git a/utils_internal_test.go b/utils_internal_test.go index 5e19a83b9..8dbce71bc 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -162,7 +162,7 @@ func (t *ClusterCluster) addNode() error { // Send NodeJoin event to coordinator. if id > 0 { coord := t.Clusters[0] - ev := &NodeEvent{ + ev := &nodeEvent{ Event: NodeJoin, Node: c.Node, } From d2d84d2649b5dc92c0f22b6cd9bfffacfe5d8876 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Fri, 29 Jun 2018 06:52:29 -0500 Subject: [PATCH 6/6] remove commented code --- gossip/gossip.go | 1 - 1 file changed, 1 deletion(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index 2200efad7..4f562af5e 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -337,7 +337,6 @@ func (g *gossipEventReceiver) listen() { if err := proto.Unmarshal(e.Node.Meta, &n); err != nil { panic("failed to unmarshal event node meta data") } - // node := pilosa.DecodeNode(&n) ne := &internal.NodeEventMessage{ Event: uint32(nodeEventType),