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 }