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.
This commit is contained in:
Matt Jaffee 2018-06-28 17:16:08 -05:00
parent 12a49c3e14
commit 91454a5cd0
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
8 changed files with 18 additions and 99 deletions

2
Gopkg.lock generated
View file

@ -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

View file

@ -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

View file

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

View file

@ -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

View file

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

View file

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

View file

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

View file

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