From 13e08ac3a546edce3da167bd394cfc246a712624 Mon Sep 17 00:00:00 2001 From: Travis Date: Tue, 21 Mar 2017 16:06:22 -0500 Subject: [PATCH] add tests for `HTTPNodeSet` --- cluster.go | 6 +++- cluster_test.go | 6 ++-- gossip.go | 3 -- handler_test.go | 90 +++++++++++++++++++++++++++++++++++++++++++++++ messenger_test.go | 37 +++++++++++++++++++ 5 files changed, 135 insertions(+), 7 deletions(-) create mode 100644 messenger_test.go diff --git a/cluster.go b/cluster.go index 44cbe90cd..5d7502ed8 100644 --- a/cluster.go +++ b/cluster.go @@ -294,7 +294,11 @@ func (h *HTTPNodeSet) SendMessage(pb proto.Message, method string) error { // ReceiveMessage is called when a node receives a message. func (h *HTTPNodeSet) ReceiveMessage(pb proto.Message) error { - return h.messageHandler(pb) + if h.messageHandler != nil { + return h.messageHandler(pb) + } + // The messageHandler has not been set. + return nil } func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { diff --git a/cluster_test.go b/cluster_test.go index 2e9b01d89..4c3ec4134 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -109,9 +109,9 @@ func TestCluster_Health(t *testing.T) { // Verify a DOWN node is reported, and extraneous nodes are ignored if a := c.Health(); !reflect.DeepEqual(a, map[string]string{ - "serverA:1000": "UP", - "serverB:1000": "DOWN", - "serverC:1000": "UP", + "serverA:1000": pilosa.HealthStatusUp, + "serverB:1000": pilosa.HealthStatusDown, + "serverC:1000": pilosa.HealthStatusUp, }) { t.Fatalf("unexpected health: %s", spew.Sdump(a)) } diff --git a/gossip.go b/gossip.go index e2b3ddcde..115918533 100644 --- a/gossip.go +++ b/gossip.go @@ -1,7 +1,6 @@ package pilosa import ( - "fmt" "io" "log" "os" @@ -124,8 +123,6 @@ func (g *GossipNodeSet) NodeMeta(limit int) []byte { } func (g *GossipNodeSet) NotifyMsg(b []byte) { - loc := g.memberlist.LocalNode() - fmt.Println("Received Msg:", loc) m, err := UnmarshalMessage(b) if err != nil { g.logger().Printf("unmarshal message error: %s", err) diff --git a/handler_test.go b/handler_test.go index c19a83a93..ab724b1f8 100644 --- a/handler_test.go +++ b/handler_test.go @@ -868,3 +868,93 @@ func MustReadAll(r io.Reader) []byte { } return buf } + +type MessageBin struct { + Cluster *pilosa.Cluster + messageReceived proto.Message +} + +func NewMessageBin() *MessageBin { + return &MessageBin{} +} + +func (m *MessageBin) messageHandler(pb proto.Message) error { + m.messageReceived = pb + return nil +} + +func NewHTTPMessageBin(s *Server, nodes []*pilosa.Node) (*MessageBin, error) { + ns := pilosa.NewHTTPNodeSet() + mb := NewMessageBin() + ns.SetMessageHandler(mb.messageHandler) + c := pilosa.Cluster{ + Nodes: nodes, + NodeSet: ns, + } + mb.Cluster = &c + s.Handler.Cluster = &c + s.Handler.Messenger = ns + + i, err := c.NodeSet.Join(c.Nodes) + if i != int(0) { + return nil, err + } + if err != nil { + return nil, err + } + + return mb, nil +} + +// Ensure that an HTTP message sent to the cluster reaches all nodes. +func TestHTTPNodeSet_Base(t *testing.T) { + + // servers + s1 := NewServer() + s2 := NewServer() + s3 := NewServer() + nodes := []*pilosa.Node{ + {Host: s1.Host()}, + {Host: s2.Host()}, + {Host: s3.Host()}, + } + + // node 1 + mb1, err := NewHTTPMessageBin(s1, nodes) + if err != nil { + t.Fatalf("unable to create message bin: %s", err) + } + + // node2 + mb2, err := NewHTTPMessageBin(s2, nodes) + if err != nil { + t.Fatalf("unable to create message bin: %s", err) + } + + // node3 + mb3, err := NewHTTPMessageBin(s3, nodes) + if err != nil { + t.Fatalf("unable to create message bin: %s", err) + } + + // message + msg := &internal.CreateSliceMessage{ + DB: "d", + Slice: 8, + } + + // send message + if err := mb1.Cluster.NodeSet.(pilosa.Messenger).SendMessage(msg, ""); err != nil { + t.Fatalf("failure sending message: %s", err) + } + + if !reflect.DeepEqual(mb1.messageReceived, msg) { + t.Fatalf("unexpected message received by node1: %s", mb1.messageReceived) + } + if !reflect.DeepEqual(mb2.messageReceived, msg) { + t.Fatalf("unexpected message received by node2: %s", mb2.messageReceived) + } + if !reflect.DeepEqual(mb3.messageReceived, msg) { + t.Fatalf("unexpected message received by node3: %s", mb3.messageReceived) + } +} diff --git a/messenger_test.go b/messenger_test.go new file mode 100644 index 000000000..ee2d720a9 --- /dev/null +++ b/messenger_test.go @@ -0,0 +1,37 @@ +package pilosa_test + +import ( + "reflect" + "testing" + + "github.com/gogo/protobuf/proto" + "github.com/pilosa/pilosa" + "github.com/pilosa/pilosa/internal" +) + +// Ensure a message can be marshaled and unmarshaled. +func TestMessage_Marshal(t *testing.T) { + + testMessageMarshal(t, &internal.CreateSliceMessage{ + DB: "d", + Slice: 8, + }) + + testMessageMarshal(t, &internal.DeleteDBMessage{ + DB: "d", + }) +} + +func testMessageMarshal(t *testing.T, m proto.Message) { + marshalled, err := pilosa.MarshalMessage(m) + if err != nil { + t.Fatal(err) + } + unmarshalled, err := pilosa.UnmarshalMessage(marshalled) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(unmarshalled, m) { + t.Fatalf("unexpected message marshalling: %s", unmarshalled) + } +}