From 8fa45649b678745e6eb5e8d9262415426c3500d1 Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 17 Mar 2017 17:19:42 -0500 Subject: [PATCH] WIP: Added support for NodeState. Modified the gossip implementation to optionally send messages directly (instead of over gossip). Added support for local and remote NodeState, as well as function to merge remote into local state. Refactored Protobuf models to include `DBMeta` and `FrameMeta` as well as support for NodeState. --- cluster.go | 38 +++++++++++++++---- db.go | 36 +++++++++++++++--- frame.go | 48 ++++++++++++++++-------- gossip.go | 95 +++++++++++++++++++++++++++++++++++++++--------- handler.go | 32 ++++++++++++++-- index.go | 30 ++++++++------- messenger.go | 26 +++++++++---- server.go | 52 ++++++++++++++++++++++++++ server/server.go | 4 +- 9 files changed, 291 insertions(+), 70 deletions(-) diff --git a/cluster.go b/cluster.go index 061f76c62..bbdc8f51f 100644 --- a/cluster.go +++ b/cluster.go @@ -211,6 +211,12 @@ type NodeSet interface { // SetMessageHandler provides the NodeSet with a function to call on ReceiveMessage SetMessageHandler(f func(proto.Message) error) + + // SetRemoteStateHandler provides the function to call on MergeRemoteState + SetRemoteStateHandler(f func(proto.Message) error) + + // SetLocalStateSource provides the function to get the current node's local state. + SetLocalStateSource(f func() (proto.Message, error)) } // Hasher represents an interface to hash integers into buckets. @@ -240,6 +246,8 @@ func (h *jmphasher) Hash(key uint64, n int) int { type HTTPNodeSet struct { nodes []*Node messageHandler func(m proto.Message) error + // remoteStateHandler func(m proto.Message) error + // localStateSource func() (proto.Message, error) } // NewHTTPNodeSet returns a new instance of HTTPNodeSet. @@ -261,7 +269,7 @@ func (h *HTTPNodeSet) Open() error { } // SendMessage asyncronously broadcasts a protobuf message to all nodes. -func (h *HTTPNodeSet) SendMessage(pb proto.Message) error { +func (h *HTTPNodeSet) SendMessage(pb proto.Message, method string) error { // Marshal the pb to []byte buf, err := MarshalMessage(pb) @@ -279,13 +287,12 @@ func (h *HTTPNodeSet) SendMessage(pb proto.Message) error { return g.Wait() } -// ReceiveMessage is called when a node recieves a message. +// ReceiveMessage is called when a node receives a message. func (h *HTTPNodeSet) ReceiveMessage(pb proto.Message) error { return h.messageHandler(pb) } func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { - fmt.Println("sendNodeMessage:", node.Host) var client *http.Client client = http.DefaultClient @@ -303,9 +310,7 @@ func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { req.Header.Set("Content-Type", "application/x-protobuf") // Send request to remote node. - fmt.Println("Send") resp, err := client.Do(req) - fmt.Println("Got back") if err != nil { return err } @@ -317,7 +322,6 @@ func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { if err != nil { return err } - fmt.Println("code:", resp.StatusCode) // Check status code. if resp.StatusCode != http.StatusOK { @@ -332,6 +336,18 @@ func (h *HTTPNodeSet) SetMessageHandler(f func(proto.Message) error) { h.messageHandler = f } +// SetRemoteStateHandler provides the Messenger with a function to merge remote state. +func (h *HTTPNodeSet) SetRemoteStateHandler(f func(proto.Message) error) { + // not implemented + // h.remoteStateHandler = f +} + +// SetLocalStateSource currently no-ops. +func (h *HTTPNodeSet) SetLocalStateSource(f func() (proto.Message, error)) { + // not implemented + // h.localStateSource = f +} + // StaticNodeSet represents a basic NodeSet for testing type StaticNodeSet struct { Messenger @@ -359,7 +375,15 @@ func (s *StaticNodeSet) SetMessageHandler(f func(proto.Message) error) { return } -func (s *StaticNodeSet) SendMessage(pb proto.Message) error { +func (s *StaticNodeSet) SetRemoteStateHandler(f func(proto.Message) error) { + return +} + +func (s *StaticNodeSet) SetLocalStateSource(f func() (proto.Message, error)) { + return +} + +func (s *StaticNodeSet) SendMessage(pb proto.Message, method string) error { return nil } func (s *StaticNodeSet) ReceiveMessage(pb proto.Message) error { diff --git a/db.go b/db.go index ca03188a3..a7f6f8ff1 100644 --- a/db.go +++ b/db.go @@ -173,7 +173,7 @@ func (db *DB) openFrames() error { // loadMeta reads meta data for the database, if any. func (db *DB) loadMeta() error { - var pb internal.DB + var pb internal.DBMeta // Read data from meta file. buf, err := ioutil.ReadFile(filepath.Join(db.path, ".meta")) @@ -199,7 +199,7 @@ func (db *DB) loadMeta() error { // saveMeta writes meta data for the database. func (db *DB) saveMeta() error { // Marshal metadata. - buf, err := proto.Marshal(&internal.DB{ + buf, err := proto.Marshal(&internal.DBMeta{ TimeQuantum: string(db.timeQuantum), ColumnLabel: db.columnLabel, }) @@ -251,10 +251,10 @@ func (db *DB) MaxSlice() uint64 { return max } -func (db *DB) SetRemoteMaxSlice(v uint64) { +func (db *DB) SetRemoteMaxSlice(newmax uint64) { db.mu.Lock() defer db.mu.Unlock() - db.remoteMaxSlice = v + db.remoteMaxSlice = newmax } // MaxInverseSlice returns the max inverse slice in the database according to this node. @@ -396,6 +396,9 @@ func (db *DB) createFrame(name string, opt FrameOptions) (*Frame, error) { if opt.CacheSize != 0 { f.rankedCacheSize = opt.CacheSize } + if opt.TimeQuantum.Valid() { + f.timeQuantum = opt.TimeQuantum + } f.inverseEnabled = opt.InverseEnabled if err := f.saveMeta(); err != nil { @@ -509,9 +512,32 @@ func MergeSchemas(a, b []*DBInfo) []*DBInfo { return dbs } +// encodeDBs converts a into its internal representation. +func encodeDBs(a []*DB) []*internal.DB { + other := make([]*internal.DB, len(a)) + for i := range a { + other[i] = encodeDB(a[i]) + } + return other +} + +// encodeDB converts d into its internal representation. +func encodeDB(d *DB) *internal.DB { + return &internal.DB{ + Name: d.name, + Meta: &internal.DBMeta{ + ColumnLabel: d.columnLabel, + TimeQuantum: string(d.timeQuantum), + }, + MaxSlice: d.remoteMaxSlice, + Frames: encodeFrames(d.Frames()), + } +} + // DBOptions represents options to set when initializing a db. type DBOptions struct { - ColumnLabel string `json:"columnLabel,omitempty"` + ColumnLabel string `json:"columnLabel,omitempty"` + TimeQuantum TimeQuantum `json:"timeQuantum,omitempty"` } // hasTime returns true if a contains a non-nil time. diff --git a/frame.go b/frame.go index b299d2ac2..51cb2c41b 100644 --- a/frame.go +++ b/frame.go @@ -264,7 +264,7 @@ func (f *Frame) openViews() error { // loadMeta reads meta data for the frame, if any. func (f *Frame) loadMeta() error { - var pb internal.Frame + var pb internal.FrameMeta // Read data from meta file. buf, err := ioutil.ReadFile(filepath.Join(f.path, ".meta")) @@ -413,17 +413,12 @@ func (f *Frame) CreateViewIfNotExists(name string) (*View, error) { // TODO: this needs to be refactored for views /* - // Send a MaxSlice message - f.messenger.SendMessage( - &internal.CreateSliceMessage{ - DB: f.db, - Slice: slice, - }) - - frag.BitmapAttrStore = f.bitmapAttrStore - - // Save to lookup. - f.fragments[slice] = frag + // Send a MaxSlice message + f.messenger.SendMessage( + &internal.CreateSliceMessage{ + DB: f.db, + Slice: slice, + }, "gossip") */ return view, nil @@ -591,6 +586,26 @@ func (f *Frame) Import(bitmapIDs, profileIDs []uint64, timestamps []*time.Time) return nil } +// encodeFrames converts a into its internal representation. +func encodeFrames(a []*Frame) []*internal.Frame { + other := make([]*internal.Frame, len(a)) + for i := range a { + other[i] = encodeFrame(a[i]) + } + return other +} + +// encodeFrame converts f into its internal representation. +func encodeFrame(f *Frame) *internal.Frame { + return &internal.Frame{ + Name: f.name, + Meta: &internal.FrameMeta{ + TimeQuantum: string(f.timeQuantum), + RowLabel: f.rowLabel, + }, + } +} + type frameSlice []*Frame func (p frameSlice) Swap(i, j int) { p[i], p[j] = p[j], p[i] } @@ -611,10 +626,11 @@ func (p frameInfoSlice) Less(i, j int) bool { return p[i].Name < p[j].Name } // FrameOptions represents options to set when initializing a frame. type FrameOptions struct { - RowLabel string `json:"rowLabel,omitempty"` - InverseEnabled bool `json:"inverseEnabled,omitempty"` - CacheType string `json:"cacheType,omitempty"` - CacheSize int `json:"cacheSize,omitempty"` + RowLabel string `json:"rowLabel,omitempty"` + InverseEnabled bool `json:"inverseEnabled,omitempty"` + CacheType string `json:"cacheType,omitempty"` + CacheSize int `json:"cacheSize,omitempty"` + TimeQuantum TimeQuantum `json:"timeQuantum,omitempty"` } // importBitSet represents slices of row and column ids. diff --git a/gossip.go b/gossip.go index 79d301a97..e2b3ddcde 100644 --- a/gossip.go +++ b/gossip.go @@ -1,38 +1,44 @@ package pilosa import ( + "fmt" "io" "log" "os" + "golang.org/x/sync/errgroup" + "github.com/gogo/protobuf/proto" "github.com/hashicorp/memberlist" + "github.com/pilosa/pilosa/internal" ) // GossipNodeSet represents a gossip implementation of NodeSet using memberlist // GossipNodeSet also represents an implementation of memberlist.Delegate type GossipNodeSet struct { - Memberlist *memberlist.Memberlist - Broadcasts *memberlist.TransmitLimitedQueue + memberlist *memberlist.Memberlist + broadcasts *memberlist.TransmitLimitedQueue config *GossipConfig - messageHandler func(m proto.Message) error + messageHandler func(m proto.Message) error + remoteStateHandler func(m proto.Message) error + localStateSource func() (proto.Message, error) // The writer for any logging. LogOutput io.Writer } func (g *GossipNodeSet) Nodes() []*Node { - a := make([]*Node, 0, g.Memberlist.NumMembers()) - for _, n := range g.Memberlist.Members() { + a := make([]*Node, 0, g.memberlist.NumMembers()) + for _, n := range g.memberlist.Members() { a = append(a, &Node{Host: n.Name}) } return a } func (g *GossipNodeSet) Join(nodes []*Node) (int, error) { - return g.Memberlist.Join(Nodes(nodes).Hosts()) + return g.memberlist.Join(Nodes(nodes).Hosts()) } func (g *GossipNodeSet) Open() error { @@ -40,14 +46,14 @@ func (g *GossipNodeSet) Open() error { if err != nil { return err } - g.Memberlist = ml + g.memberlist = ml // attach to gossip seed node g.Join([]*Node{&Node{Host: g.config.gossipSeed}}) //TODO: support a list of seeds - g.Broadcasts = &memberlist.TransmitLimitedQueue{ + g.broadcasts = &memberlist.TransmitLimitedQueue{ NumNodes: func() int { - return g.Memberlist.NumMembers() + return g.memberlist.NumMembers() }, RetransmitMult: 3, } @@ -58,18 +64,49 @@ func (g *GossipNodeSet) SetMessageHandler(f func(proto.Message) error) { g.messageHandler = f } +func (g *GossipNodeSet) SetRemoteStateHandler(f func(proto.Message) error) { + g.remoteStateHandler = f +} + +func (g *GossipNodeSet) SetLocalStateSource(f func() (proto.Message, error)) { + g.localStateSource = f +} + // implementation of the messenger.Messenger interface -func (g *GossipNodeSet) SendMessage(pb proto.Message) error { +func (g *GossipNodeSet) SendMessage(pb proto.Message, method string) error { msg, err := MarshalMessage(pb) if err != nil { return err } - b := &broadcast{ - msg: msg, - notify: nil, + // Broadcast asyncronously sends the message directly to each node. + // An error from any node raises an error on the entire operation. + // This is a blocking operation. + // + // Gossip uses the gossip protocol to eventually deliver the message + // to every node. + switch method { + case "broadcast": + var eg errgroup.Group + for _, n := range g.memberlist.Members() { + // Don't send the message to the local node. + if n == g.memberlist.LocalNode() { + continue + } + node := n + eg.Go(func() error { + return g.memberlist.SendToTCP(node, msg) + }) + } + return eg.Wait() + case "gossip": + b := &broadcast{ + msg: msg, + notify: nil, + } + g.broadcasts.QueueBroadcast(b) } - g.Broadcasts.QueueBroadcast(b) + return nil } @@ -87,6 +124,8 @@ 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) @@ -99,14 +138,37 @@ func (g *GossipNodeSet) NotifyMsg(b []byte) { } func (g *GossipNodeSet) GetBroadcasts(overhead, limit int) [][]byte { - return g.Broadcasts.GetBroadcasts(overhead, limit) + return g.broadcasts.GetBroadcasts(overhead, limit) } func (g *GossipNodeSet) LocalState(join bool) []byte { - return []byte{} + + pb, err := g.localStateSource() + if err != nil { + g.logger().Printf("error getting local state, err=%s", err) + return []byte{} + } + + // Marshal nodestate data to bytes. + buf, err := proto.Marshal(pb) + if err != nil { + g.logger().Printf("error marshaling nodestate data, err=%s", err) + return []byte{} + } + return buf } func (g *GossipNodeSet) MergeRemoteState(buf []byte, join bool) { + // Unmarshal nodestate data. + var pb internal.NodeState + if err := proto.Unmarshal(buf, &pb); err != nil { + g.logger().Printf("error unmarshaling nodestate data, err=%s", err) + return + } + err := g.remoteStateHandler(&pb) + if err != nil { + g.logger().Printf("merge state error: %s", err) + } return } @@ -158,7 +220,6 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed g.config.memberlistConfig.BindPort = gossipPort g.config.memberlistConfig.AdvertiseAddr = gossipHost g.config.memberlistConfig.AdvertisePort = gossipPort - g.config.memberlistConfig.GossipNodes = 1 g.config.memberlistConfig.Delegate = g return g diff --git a/handler.go b/handler.go index d3c5a28cc..ad337b2bc 100644 --- a/handler.go +++ b/handler.go @@ -370,11 +370,13 @@ func (h *Handler) handlePostDB(w http.ResponseWriter, r *http.Request) { } // Send the delete message to all nodes. - // NOTE: this calls a second DeleteDB on the local node - h.Messenger.SendMessage( + err := h.Messenger.SendMessage( &internal.DeleteDBMessage{ DB: req.DB, - }) + }, "broadcast") + if err != nil { + h.logger().Printf("problem sending DeleteDB message: %s", err) + } // Encode response. if err := json.NewEncoder(w).Encode(postDBResponse{}); err != nil { @@ -511,6 +513,20 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { return } + // Send the create message to all nodes. + err = h.Messenger.SendMessage( + &internal.CreateFrameMessage{ + DB: req.DB, + Frame: req.Frame, + Meta: &internal.FrameMeta{ + RowLabel: req.Options.RowLabel, + TimeQuantum: string(req.Options.TimeQuantum), + }, + }, "broadcast") + if err != nil { + h.logger().Printf("problem sending CreateFrame message: %s", err) + } + // Encode response. if err := json.NewEncoder(w).Encode(postFrameResponse{}); err != nil { h.logger().Printf("response encoding error: %s", err) @@ -569,6 +585,16 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) { return } + // Send the delete message to all nodes. + err := h.Messenger.SendMessage( + &internal.DeleteFrameMessage{ + DB: req.DB, + Frame: req.Frame, + }, "broadcast") + if err != nil { + h.logger().Printf("problem sending DeleteFrame message: %s", err) + } + // Encode response. if err := json.NewEncoder(w).Encode(deleteFrameResponse{}); err != nil { h.logger().Printf("response encoding error: %s", err) diff --git a/index.go b/index.go index f7269d670..954629eea 100644 --- a/index.go +++ b/index.go @@ -194,7 +194,6 @@ func (i *Index) CreateDB(name string, opt DBOptions) (*DB, error) { // Ensure db doesn't already exist. if i.dbs[name] != nil { - fmt.Println("ErrDatabaseExists: 2") return nil, ErrDatabaseExists } return i.createDB(name, opt) @@ -236,18 +235,12 @@ func (i *Index) createDB(name string, opt DBOptions) (*DB, error) { // Update options. db.SetColumnLabel(opt.ColumnLabel) + db.SetTimeQuantum(opt.TimeQuantum) i.dbs[db.Name()] = db i.Stats.Count("dbN", 1) - // Send a CreateDB message - i.Messenger.SendMessage( - &internal.CreateDBMessage{ - DB: name, - ColumnLabel: opt.ColumnLabel, - }) - return db, nil } @@ -364,17 +357,28 @@ func (i *Index) HandleMessage(pb proto.Message) error { return fmt.Errorf("Local DB not found: %s", obj.DB) } d.SetRemoteMaxSlice(obj.Slice) - case *internal.DeleteDBMessage: - err := i.DeleteDB(obj.DB) + case *internal.CreateDBMessage: + opt := DBOptions{ColumnLabel: obj.Meta.ColumnLabel} + _, err := i.CreateDB(obj.DB, opt) if err != nil { return err } - case *internal.CreateDBMessage: - opt := DBOptions{ColumnLabel: obj.ColumnLabel} - _, err := i.CreateDB(obj.DB, opt) + case *internal.DeleteDBMessage: + if err := i.DeleteDB(obj.DB); err != nil { + return err + } + case *internal.CreateFrameMessage: + db := i.DB(obj.DB) + opt := FrameOptions{RowLabel: obj.Meta.RowLabel} + _, err := db.CreateFrame(obj.Frame, opt) if err != nil { return err } + case *internal.DeleteFrameMessage: + db := i.DB(obj.DB) + if err := db.DeleteFrame(obj.Frame); err != nil { + return err + } } return nil } diff --git a/messenger.go b/messenger.go index 5e2d128e2..f28eb7ddb 100644 --- a/messenger.go +++ b/messenger.go @@ -17,7 +17,7 @@ var NopMessenger Messenger // nopMessenger represents a Messenger that doesn't do anything. type nopMessenger struct{} -func (c *nopMessenger) SendMessage(pb proto.Message) error { +func (c *nopMessenger) SendMessage(pb proto.Message, method string) error { fmt.Println("NOPMessenger: Send") return nil } @@ -27,14 +27,16 @@ func (c *nopMessenger) ReceiveMessage(pb proto.Message) error { } type Messenger interface { - SendMessage(pb proto.Message) error + SendMessage(pb proto.Message, method string) error ReceiveMessage(pb proto.Message) error } const ( MessageTypeCreateSlice = 1 - MessageTypeDeleteDB = 2 - MessageTypeCreateDB = 3 + MessageTypeCreateDB = 2 + MessageTypeDeleteDB = 3 + MessageTypeCreateFrame = 4 + MessageTypeDeleteFrame = 5 ) func MarshalMessage(m proto.Message) ([]byte, error) { @@ -42,10 +44,14 @@ func MarshalMessage(m proto.Message) ([]byte, error) { switch obj := m.(type) { case *internal.CreateSliceMessage: typ = MessageTypeCreateSlice - case *internal.DeleteDBMessage: - typ = MessageTypeDeleteDB case *internal.CreateDBMessage: typ = MessageTypeCreateDB + case *internal.DeleteDBMessage: + typ = MessageTypeDeleteDB + case *internal.CreateFrameMessage: + typ = MessageTypeCreateFrame + case *internal.DeleteFrameMessage: + typ = MessageTypeDeleteFrame default: return nil, fmt.Errorf("message type not implemented for marshalling: %s", reflect.TypeOf(obj)) } @@ -63,10 +69,14 @@ func UnmarshalMessage(buf []byte) (proto.Message, error) { switch typ { case MessageTypeCreateSlice: m = &internal.CreateSliceMessage{} - case MessageTypeDeleteDB: - m = &internal.DeleteDBMessage{} case MessageTypeCreateDB: m = &internal.CreateDBMessage{} + case MessageTypeDeleteDB: + m = &internal.DeleteDBMessage{} + case MessageTypeCreateFrame: + m = &internal.CreateFrameMessage{} + case MessageTypeDeleteFrame: + m = &internal.DeleteFrameMessage{} default: return nil, fmt.Errorf("invalid message type: %d", typ) } diff --git a/server.go b/server.go index cde64c3ca..0df1182ba 100644 --- a/server.go +++ b/server.go @@ -153,6 +153,49 @@ func (s *Server) Addr() net.Addr { return s.ln.Addr() } +// LocalState returns the state of the local node as well as the +// index (dbs/frames) according to the local node. +func (s *Server) LocalState() (proto.Message, error) { + // TODO: are there errors to handle? + pb := encodeLocalState(s) + return pb, nil +} + +// HandleRemoteState provides the current, local state. +// In a gossip implementation, memberlist.Delegate.LocalState() uses this. +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), + } + _, err := d.CreateFrameIfNotExists(f.Name, opt) + if err != nil { + return err + } + } + } + + return nil +} + func (s *Server) logger() *log.Logger { return log.New(s.LogOutput, "", log.LstdFlags) } func (s *Server) monitorAntiEntropy() { @@ -189,6 +232,15 @@ func (s *Server) monitorAntiEntropy() { } } +// encodeLocalState converts s into its internal representation. +func encodeLocalState(s *Server) *internal.NodeState { + return &internal.NodeState{ + Host: s.Host, + State: "OK", // TODO: make this work, pull from cluster.Node + DBs: encodeDBs(s.Index.DBs()), + } +} + // monitorMaxSlices periodically pulls the highest slice from each node in the cluster. func (s *Server) monitorMaxSlices() { // Ignore if only one node in the cluster. diff --git a/server/server.go b/server/server.go index da83da52c..3a62c5c60 100644 --- a/server/server.go +++ b/server/server.go @@ -100,8 +100,10 @@ func (m *Command) Run(args ...string) (err error) { m.Server.Handler.Messenger = m.Server.Messenger m.Server.Index.Messenger = m.Server.Messenger - // Set message handler. + // Set message and state handlers. m.Server.Cluster.NodeSet.SetMessageHandler(m.Server.Index.HandleMessage) + m.Server.Cluster.NodeSet.SetRemoteStateHandler(m.Server.HandleRemoteState) + m.Server.Cluster.NodeSet.SetLocalStateSource(m.Server.LocalState) // Set configuration options. m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval)