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, }