mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
add tests for HTTPNodeSet
This commit is contained in:
parent
3abfea140b
commit
13e08ac3a5
5 changed files with 135 additions and 7 deletions
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
37
messenger_test.go
Normal file
37
messenger_test.go
Normal file
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue