mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 17:15:56 +00:00
Merge pull request #1428 from jaffee/more-gossip-stuff
More gossip stuff
This commit is contained in:
commit
0fae552577
7 changed files with 41 additions and 179 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
13
event.go
13
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
29
server.go
29
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue