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.
This commit is contained in:
Travis 2017-03-17 17:19:42 -05:00
parent 2b35df879c
commit 8fa45649b6
9 changed files with 291 additions and 70 deletions

View file

@ -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 {

36
db.go
View file

@ -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.

View file

@ -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.

View file

@ -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

View file

@ -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)

View file

@ -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
}

View file

@ -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)
}

View file

@ -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.

View file

@ -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)