diff --git a/broadcast.go b/broadcast.go index 12f2b0bfe..57f0bb496 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 ( messageTypeCreateShard = iota diff --git a/broadcast_test.go b/broadcast_test.go index 5c97df0d0..415228718 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 21c91e528..ce0214537 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 0355260d6..4f562af5e 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 @@ -51,26 +50,11 @@ type GossipMemberSet struct { logger *log.Logger transport *Transport - 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) + gossipEventReceiver *gossipEventReceiver } // 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") - } - if g.handler == nil { - return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()") - } - +func (g *GossipMemberSet) Open() (err error) { g.mu.Lock() g.memberlist, err = memberlist.Create(g.config.memberlistConfig) g.mu.Unlock() @@ -173,11 +157,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 { @@ -255,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 } @@ -300,46 +282,42 @@ 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 + eventHandler *pilosa.Server - 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, +// newGossipEventReceiver returns a new instance of GossipEventReceiver. +func newGossipEventReceiver(logger *log.Logger, pserver *pilosa.Server) *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) { +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 { - g.eventHandler = h - go g.listen() - return nil -} - -func (g *GossipEventReceiver) listen() { +func (g *gossipEventReceiver) listen() { var nodeEventType pilosa.NodeEventType for { e := <-g.ch @@ -359,38 +337,17 @@ 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 := &pilosa.NodeEvent{ - Event: nodeEventType, - Node: node, + ne := &internal.NodeEventMessage{ + Event: uint32(nodeEventType), + Node: &n, } - if err := g.eventHandler.ReceiveEvent(ne); err != nil { - g.Logger.Printf("receive event error: %s", err) + if err := g.eventHandler.ReceiveMessage(ne); err != nil { + 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 diff --git a/server.go b/server.go index b04661f53..457a9123a 100644 --- a/server.go +++ b/server.go @@ -41,8 +41,7 @@ const ( // Ensure Server implements interfaces. var _ Broadcaster = &Server{} -var _ BroadcastHandler = &Server{} -var _ StatusHandler = &Server{} +var _ MemberServer = &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. @@ -707,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 { @@ -733,15 +722,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) { @@ -753,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 7843ed443..6a84d168d 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, }