From 1aca2ce1de6c10045301e41102110a0aa3efcdb7 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 17 Apr 2017 11:06:58 -0500 Subject: [PATCH] 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. --- gossip.go | 16 +++--- handler.go | 3 +- messenger.go | 126 +++-------------------------------------------- server.go | 90 +++++++++++++++++++++++++++++++-- server/server.go | 10 ++-- 5 files changed, 108 insertions(+), 137 deletions(-) diff --git a/gossip.go b/gossip.go index 178fa999e..59abef70d 100644 --- a/gossip.go +++ b/gossip.go @@ -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, } diff --git a/handler.go b/handler.go index 0a7cb278d..7616fd2d8 100644 --- a/handler.go +++ b/handler.go @@ -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 } diff --git a/messenger.go b/messenger.go index 1ffbfa51e..3672a407a 100644 --- a/messenger.go +++ b/messenger.go @@ -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.") } diff --git a/server.go b/server.go index a93d29f57..7dd35fa8f 100644 --- a/server.go +++ b/server.go @@ -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{ diff --git a/server/server.go b/server/server.go index 11bb055a9..0b61bf405 100644 --- a/server/server.go +++ b/server/server.go @@ -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 }