implement LocalState and internal message reception on Server

greatly simplifies Messenger, to the point of making it basically shell around
MessageBroker. Next stesp are to make the receiving of internal messages more
well-defined and behind an interface, remove/merge Messenger and MessageBroker,
and have the implementations of MessageBroker in separate packages.
This commit is contained in:
Matt Jaffee 2017-04-17 11:06:58 -05:00 committed by Travis
parent 25f3ac92bd
commit 1aca2ce1de
5 changed files with 108 additions and 137 deletions

View file

@ -90,7 +90,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed
type GossipMessageBroker struct {
broadcasts *memberlist.TransmitLimitedQueue
messenger *Messenger
server *Server
// The writer for any logging.
LogOutput io.Writer
@ -103,7 +103,7 @@ func (g *GossipMessageBroker) Send(pb proto.Message, method string) error {
return err
}
mlist := g.messenger.Cluster.NodeSet.(*GossipNodeSet).memberlist
mlist := g.server.Cluster.NodeSet.(*GossipNodeSet).memberlist
// Direct sends the message directly to every node.
// An error from any node raises an error on the entire operation.
@ -136,7 +136,7 @@ func (g *GossipMessageBroker) Send(pb proto.Message, method string) error {
}
func (g *GossipMessageBroker) Receive(pb proto.Message) error {
if err := g.messenger.ReceiveMessage(pb); err != nil {
if err := g.server.ReceiveMessage(pb); err != nil {
return err
}
return nil
@ -164,7 +164,7 @@ func (g *GossipMessageBroker) GetBroadcasts(overhead, limit int) [][]byte {
}
func (g *GossipMessageBroker) LocalState(join bool) []byte {
pb, err := g.messenger.LocalState()
pb, err := g.server.LocalState()
if err != nil {
g.logger().Printf("error getting local state, err=%s", err)
return []byte{}
@ -186,7 +186,7 @@ func (g *GossipMessageBroker) MergeRemoteState(buf []byte, join bool) {
g.logger().Printf("error unmarshalling nodestate data, err=%s", err)
return
}
err := g.messenger.HandleRemoteState(&pb)
err := g.server.HandleRemoteState(&pb)
if err != nil {
g.logger().Printf("merge state error: %s", err)
}
@ -200,15 +200,15 @@ func (g *GossipMessageBroker) logger() *log.Logger {
////////////////////////////////////////////////////////////////
// NewGossipMessageBroker returns a new instance of GossipMessageBroker.
func NewGossipMessageBroker(m *Messenger) *GossipMessageBroker {
func NewGossipMessageBroker(s *Server) *GossipMessageBroker {
g := &GossipMessageBroker{
LogOutput: os.Stderr,
messenger: m,
server: s,
}
g.broadcasts = &memberlist.TransmitLimitedQueue{
NumNodes: func() int {
return g.messenger.Cluster.NodeSet.(*GossipNodeSet).memberlist.NumMembers()
return g.server.Cluster.NodeSet.(*GossipNodeSet).memberlist.NumMembers()
},
RetransmitMult: 3,
}

View file

@ -27,6 +27,7 @@ import (
type Handler struct {
Index *Index
Messenger *Messenger
Server *Server
// Local hostname & cluster configuration.
Host string
@ -213,7 +214,7 @@ func (h *Handler) handlePostMessage(w http.ResponseWriter, r *http.Request) {
return
}
if err := h.Messenger.Broker.Receive(m); err != nil {
if err := h.Server.ReceiveMessage(m); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}

View file

@ -20,15 +20,9 @@ import (
// Messenger represents an internal message handler.
type Messenger struct {
// Broker handles Send/Receive Messages.
// Broker handles Send
Broker MessageBroker
Index *Index
// Local hostname & cluster configuration.
Host string
Cluster *Cluster
// The writer for any logging.
LogOutput io.Writer
}
@ -47,104 +41,12 @@ func (m *Messenger) SendMessage(pb proto.Message, method string) error {
}
return m.Broker.Send(pb, method)
}
func (m *Messenger) ReceiveMessage(pb proto.Message) error {
return m.handleMessage(pb)
}
// handleMessage handles protobuf Messages sent to nodes in the cluster.
func (m *Messenger) handleMessage(pb proto.Message) error {
switch obj := pb.(type) {
case *internal.CreateSliceMessage:
d := m.Index.DB(obj.DB)
if d == nil {
return fmt.Errorf("Local DB not found: %s", obj.DB)
}
d.SetRemoteMaxSlice(obj.Slice)
case *internal.CreateDBMessage:
opt := DBOptions{ColumnLabel: obj.Meta.ColumnLabel}
_, err := m.Index.CreateDB(obj.DB, opt)
if err != nil {
return err
}
case *internal.DeleteDBMessage:
fmt.Println("DELETE:", obj.DB)
if err := m.Index.DeleteDB(obj.DB); err != nil {
return err
}
case *internal.CreateFrameMessage:
db := m.Index.DB(obj.DB)
opt := FrameOptions{RowLabel: obj.Meta.RowLabel}
_, err := db.CreateFrame(obj.Frame, opt)
if err != nil {
return err
}
case *internal.DeleteFrameMessage:
db := m.Index.DB(obj.DB)
if err := db.DeleteFrame(obj.Frame); err != nil {
return err
}
}
return nil
}
// LocalState returns the state of the local node as well as the
// index (dbs/frames) according to the local node.
// In a gossip implementation, memberlist.Delegate.LocalState() uses this.
// It seems odd to have this as part of Messenger, but with the
// exception of Server, it's currenntly the only object with access
// to the necessary information (Host, Index, Cluster).
func (m *Messenger) LocalState() (proto.Message, error) {
if m.Index == nil {
return nil, errors.New("Messenger.Index is nil.")
}
return &internal.NodeState{
Host: m.Host,
State: "OK", // TODO: make this work, pull from m.Cluster.Node
DBs: encodeDBs(m.Index.DBs()),
}, nil
}
// HandleRemoteState receives incoming NodeState from remote nodes.
func (m *Messenger) HandleRemoteState(pb proto.Message) error {
return m.mergeRemoteState(pb.(*internal.NodeState))
}
func (m *Messenger) mergeRemoteState(ns *internal.NodeState) error {
// TODO: update some node state value in the cluster (it should be in cluster.node i guess)
// Create databases that don't exist.
for _, db := range ns.DBs {
opt := DBOptions{
ColumnLabel: db.Meta.ColumnLabel,
TimeQuantum: TimeQuantum(db.Meta.TimeQuantum),
}
d, err := m.Index.CreateDBIfNotExists(db.Name, opt)
if err != nil {
return err
}
// Create frames that don't exist.
for _, f := range db.Frames {
opt := FrameOptions{
RowLabel: f.Meta.RowLabel,
TimeQuantum: TimeQuantum(f.Meta.TimeQuantum),
CacheSize: f.Meta.CacheSize,
}
_, err := d.CreateFrameIfNotExists(f.Name, opt)
if err != nil {
return err
}
}
}
return nil
}
//////////////////////////////////////////////////////////////////
// MessageBroker is an interface for handling incoming/outgoing messages.
type MessageBroker interface {
Send(pb proto.Message, method string) error
Receive(pb proto.Message) error
}
//////////////////////////////////////////////////////////////////
@ -162,21 +64,17 @@ func (c *nopMessageBroker) Send(pb proto.Message, method string) error {
fmt.Println("NOPMessageBroker: Send")
return nil
}
func (c *nopMessageBroker) Receive(pb proto.Message) error {
fmt.Println("NOPMessageBroker: Receive")
return nil
}
//////////////////////////////////////////////////////////////////
// HTTPMessageBroker represents a NodeSet that broadcasts messages over HTTP.
type HTTPMessageBroker struct {
messenger *Messenger
server *Server
}
// NewHTTPMessageBroker returns a new instance of HTTPMessageBroker.
func NewHTTPMessageBroker(m *Messenger) *HTTPMessageBroker {
return &HTTPMessageBroker{messenger: m}
func NewHTTPMessageBroker(s *Server) *HTTPMessageBroker {
return &HTTPMessageBroker{server: s}
}
// Send sends a protobuf message to all nodes simultaneously.
@ -196,7 +94,7 @@ func (h *HTTPMessageBroker) Send(pb proto.Message, method string) error {
var g errgroup.Group
for _, n := range nodes {
// Don't send the message to the local node.
if n.Host == h.messenger.Host {
if n.Host == h.server.Host {
continue
}
node := n
@ -207,19 +105,11 @@ func (h *HTTPMessageBroker) Send(pb proto.Message, method string) error {
return g.Wait()
}
// Receive is called when a node receives a message.
func (h *HTTPMessageBroker) Receive(pb proto.Message) error {
if err := h.messenger.ReceiveMessage(pb); err != nil {
return err
}
return nil
}
func (h *HTTPMessageBroker) nodes() ([]*Node, error) {
if h.messenger == nil {
return nil, errors.New("HTTPMessageBroker has no reference to Messenger.")
if h.server == nil {
return nil, errors.New("HTTPMessageBroker has no reference to Server.")
}
nodeset, ok := h.messenger.Cluster.NodeSet.(*HTTPNodeSet)
nodeset, ok := h.server.Cluster.NodeSet.(*HTTPNodeSet)
if !ok {
return nil, errors.New("NodeSet cannot be caste to HTTPNodeSet.")
}

View file

@ -1,6 +1,7 @@
package pilosa
import (
"errors"
"fmt"
"io"
"io/ioutil"
@ -64,7 +65,7 @@ func NewServer() *Server {
}
s.Handler.Index = s.Index
s.Messenger.Index = s.Index
s.Handler.Server = s // TODO remove
return s
}
@ -111,9 +112,6 @@ func (s *Server) Open() error {
e.Cluster = s.Cluster
// Initialize Messenger.
s.Messenger.Index = s.Index
s.Messenger.Host = s.Host
s.Messenger.Cluster = s.Cluster
s.Messenger.LogOutput = s.LogOutput
// Initialize HTTP handler.
@ -236,6 +234,90 @@ func (s *Server) monitorMaxSlices() {
}
}
// LocalState returns the state of the local node as well as the
// index (dbs/frames) according to the local node.
// In a gossip implementation, memberlist.Delegate.LocalState() uses this.
func (s *Server) LocalState() (proto.Message, error) {
if s.Index == nil {
return nil, errors.New("Messenger.Index is nil.")
}
return &internal.NodeState{
Host: s.Host,
State: "OK", // TODO: make this work, pull from s.Cluster.Node
DBs: encodeDBs(s.Index.DBs()),
}, nil
}
func (s *Server) ReceiveMessage(pb proto.Message) error {
switch obj := pb.(type) {
case *internal.CreateSliceMessage:
d := s.Index.DB(obj.DB)
if d == nil {
return fmt.Errorf("Local DB not found: %s", obj.DB)
}
d.SetRemoteMaxSlice(obj.Slice)
case *internal.CreateDBMessage:
opt := DBOptions{ColumnLabel: obj.Meta.ColumnLabel}
_, err := s.Index.CreateDB(obj.DB, opt)
if err != nil {
return err
}
case *internal.DeleteDBMessage:
fmt.Println("DELETE:", obj.DB)
if err := s.Index.DeleteDB(obj.DB); err != nil {
return err
}
case *internal.CreateFrameMessage:
db := s.Index.DB(obj.DB)
opt := FrameOptions{RowLabel: obj.Meta.RowLabel}
_, err := db.CreateFrame(obj.Frame, opt)
if err != nil {
return err
}
case *internal.DeleteFrameMessage:
db := s.Index.DB(obj.DB)
if err := db.DeleteFrame(obj.Frame); err != nil {
return err
}
}
return nil
}
// HandleRemoteState receives incoming NodeState from remote nodes.
func (s *Server) HandleRemoteState(pb proto.Message) error {
return s.mergeRemoteState(pb.(*internal.NodeState))
}
func (s *Server) mergeRemoteState(ns *internal.NodeState) error {
// TODO: update some node state value in the cluster (it should be in cluster.node i guess)
// Create databases that don't exist.
for _, db := range ns.DBs {
opt := DBOptions{
ColumnLabel: db.Meta.ColumnLabel,
TimeQuantum: TimeQuantum(db.Meta.TimeQuantum),
}
d, err := s.Index.CreateDBIfNotExists(db.Name, opt)
if err != nil {
return err
}
// Create frames that don't exist.
for _, f := range db.Frames {
opt := FrameOptions{
RowLabel: f.Meta.RowLabel,
TimeQuantum: TimeQuantum(f.Meta.TimeQuantum),
CacheSize: f.Meta.CacheSize,
}
_, err := d.CreateFrameIfNotExists(f.Name, opt)
if err != nil {
return err
}
}
}
return nil
}
func checkMaxSlices(hostport string) (map[string]uint64, error) {
// Create HTTP request.
req, err := http.NewRequest("GET", (&url.URL{

View file

@ -94,7 +94,7 @@ func (m *Command) Run(args ...string) (err error) {
if err != nil {
return err
}
m.Server.Messenger = PilosaMessenger(m.Config)
m.Server.Messenger = PilosaMessenger(m.Config, m.Server)
m.Server.Cluster = PilosaCluster(m.Config)
// Associate objects to the MessageBroker based on config.
@ -112,15 +112,13 @@ func (m *Command) Run(args ...string) (err error) {
}
// PilosaMessenger returns a new instance of Messenger based on the config.
func PilosaMessenger(c *pilosa.Config) *pilosa.Messenger {
func PilosaMessenger(c *pilosa.Config, server *pilosa.Server) *pilosa.Messenger {
messenger := pilosa.NewMessenger()
switch c.Cluster.MessengerType {
case "broadcast":
n := pilosa.NewHTTPMessageBroker(messenger)
messenger.Broker = n
messenger.Broker = pilosa.NewHTTPMessageBroker(server)
case "gossip":
n := pilosa.NewGossipMessageBroker(messenger)
messenger.Broker = n
messenger.Broker = pilosa.NewGossipMessageBroker(server)
case "static":
// nop
}