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)