From 7c2f6a9bae2f1de825a9ecbcc9d448d2e9488b48 Mon Sep 17 00:00:00 2001 From: Travis Date: Wed, 7 Dec 2016 14:09:08 -0600 Subject: [PATCH] adds a basic gossip implementation using hashicorp/memberlist Changes Gossiper to NodeSet. Adds StaticNodeSet (for testing) and GossipNodeSet (for memberlist) implementations. Changes NodeSet interface to return a pilosa-specific generic instead of `NodeSet interface`. Removes NumMembers from interface (which is specific to memberlist). Implements a Messenger interface with which to send inter-node messages via NodeSet. The Pilosa implementation occurs in the GossipNodeSet. Implements the Messenger as an object on Server, Handler, and Index. Uses HealthStatus constants. Removes commented-out code. for gossip, make sure to bind to both host and port, and advertise those as well adjust messenger to work with the db schema logic add dependencies: hashicorp/memberlist, golang.org/x/sync add dependency: golang.org/x/net Uses `errgroup` to handle errors from broadcast messages. Marshals message one time instead of once for every node. Adds error handling for some errors that were being swallowed. --- cluster.go | 184 ++++++++++++++++++++++++++++++++++++++++++++++- cluster_test.go | 40 +++++++++++ config.go | 50 +++++++++++-- db.go | 5 +- executor_test.go | 4 +- frame.go | 21 +++++- glide.lock | 26 ++++++- glide.yaml | 5 ++ gossip.go | 165 ++++++++++++++++++++++++++++++++++++++++++ handler.go | 54 +++++++++++++- index.go | 28 +++++++- index_test.go | 27 +++++++ messenger.go | 73 +++++++++++++++++++ server.go | 27 +++---- 14 files changed, 682 insertions(+), 27 deletions(-) create mode 100644 gossip.go create mode 100644 messenger.go diff --git a/cluster.go b/cluster.go index beb33496b..8268fe508 100644 --- a/cluster.go +++ b/cluster.go @@ -1,8 +1,17 @@ package pilosa import ( + "bytes" "encoding/binary" + "fmt" "hash/fnv" + "io/ioutil" + "net/http" + "net/url" + + "golang.org/x/sync/errgroup" + + "github.com/gogo/protobuf/proto" ) const ( @@ -11,6 +20,10 @@ const ( // DefaultReplicaN is the default number of replicas per partition. DefaultReplicaN = 1 + + // HealthStatus is the return value of the /health endpoint for a node in the cluster. + HealthStatusUp = "UP" + HealthStatusDown = "DOWN" ) // Node represents a node in the cluster. @@ -81,7 +94,8 @@ func (a Nodes) Clone() []*Node { // Cluster represents a collection of nodes. type Cluster struct { - Nodes []*Node + Nodes []*Node + NodeSet NodeSet // Hashing algorithm used to assign partitions to nodes. Hasher Hasher @@ -102,6 +116,33 @@ func NewCluster() *Cluster { } } +// NodeSetHosts returns the list of host strings for NodeSet members +func (c *Cluster) NodeSetHosts() []string { + if c.NodeSet == nil { + return []string{} + } + a := make([]string, 0, len(c.NodeSet.Nodes())) + for _, m := range c.NodeSet.Nodes() { + a = append(a, m.Host) + } + return a +} + +// Health returns a list of nodes in the cluster along with each node's state (UP/DOWN). +func (c *Cluster) Health() map[string]string { + h := make(map[string]string) + for _, n := range c.Nodes { + h[n.Host] = HealthStatusDown + } + // we are assuming that NodeSetHosts is a subset of c.Nodes + for _, m := range c.NodeSetHosts() { + if _, ok := h[m]; ok { + h[m] = HealthStatusUp + } + } + return h +} + // NodeByHost returns a node reference by host. func (c *Cluster) NodeByHost(host string) *Node { for _, n := range c.Nodes { @@ -157,6 +198,21 @@ func (c *Cluster) PartitionNodes(partitionID int) []*Node { return nodes } +// NodeSet represents an interface to maintaining Node state. +type NodeSet interface { + // Returns a list of all Nodes in the cluster + Nodes() []*Node + + // Attempts to join a cluster having `nodes` as its existing members + Join(nodes []*Node) (int, error) + + // Open starts any network activity implemented by the NodeSet + Open() error + + // SetMessageHandler provides the NodeSet with a function to call on ReceiveMessage + SetMessageHandler(f func(proto.Message) error) +} + // Hasher represents an interface to hash integers into buckets. type Hasher interface { // Hashes the key into a number between [0,N). @@ -179,3 +235,129 @@ func (h *jmphasher) Hash(key uint64, n int) int { } return int(b) } + +// HTTPNodeSet represents a NodeSet that broadcasts messages over HTTP. +type HTTPNodeSet struct { + nodes []*Node + messageHandler func(m proto.Message) error +} + +// NewHTTPNodeSet returns a new instance of HTTPNodeSet. +func NewHTTPNodeSet() *HTTPNodeSet { + return &HTTPNodeSet{} +} + +func (h *HTTPNodeSet) Nodes() []*Node { + return h.nodes +} + +func (h *HTTPNodeSet) Join(nodes []*Node) (int, error) { + h.nodes = nodes + return 0, nil +} + +func (h *HTTPNodeSet) Open() error { + return nil +} + +// SendMessage asyncronously broadcasts a protobuf message to all nodes. +func (h *HTTPNodeSet) SendMessage(pb proto.Message) error { + + // Marshal the pb to []byte + buf, err := MarshalMessage(pb) + if err != nil { + return err + } + + var g errgroup.Group + for _, n := range h.nodes { + node := n + g.Go(func() error { + return h.sendNodeMessage(node, buf) + }) + } + return g.Wait() +} + +// ReceiveMessage is called when a node recieves a message. +func (h *HTTPNodeSet) ReceiveMessage(pb proto.Message) error { + return h.messageHandler(pb) +} + +func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { + var client *http.Client + client = http.DefaultClient + + // Create HTTP request. + req, err := http.NewRequest("POST", (&url.URL{ + Scheme: "http", + Host: node.Host, + Path: "/message", + }).String(), bytes.NewReader(msg)) + if err != nil { + return err + } + + // Require protobuf encoding. + req.Header.Set("Content-Type", "application/x-protobuf") + + // Send request to remote node. + resp, err := client.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + + // Read response into buffer. + body, err := ioutil.ReadAll(resp.Body) + + if err != nil { + return err + } + + // Check status code. + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("invalid status: code=%d, err=%s", resp.StatusCode, body) + } + + return nil +} + +// SetMessageHandler provides the Messenger with a function to handle incoming messages. +func (h *HTTPNodeSet) SetMessageHandler(f func(proto.Message) error) { + h.messageHandler = f +} + +// StaticNodeSet represents a basic NodeSet for testing +type StaticNodeSet struct { + Messenger + nodes []*Node +} + +func NewStaticNodeSet() *StaticNodeSet { + return &StaticNodeSet{} +} + +func (s *StaticNodeSet) Nodes() []*Node { + return s.nodes +} + +func (s *StaticNodeSet) Join(nodes []*Node) (int, error) { + s.nodes = nodes + return 0, nil +} + +func (s *StaticNodeSet) Open() error { + return nil +} + +func (s *StaticNodeSet) SetMessageHandler(f func(proto.Message) error) { + return +} + +func (s *StaticNodeSet) SendMessage(pb proto.Message) error { + return nil +} +func (s *StaticNodeSet) ReceiveMessage(pb proto.Message) error { + return nil +} diff --git a/cluster_test.go b/cluster_test.go index f5ff29edf..2e9b01d89 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -77,6 +77,46 @@ func TestHasher(t *testing.T) { } } +// Ensure that an empty cluster returns a valid (empty) NodeSet +func TestCluster_NodeSetHosts(t *testing.T) { + + c := pilosa.Cluster{} + + if h := c.NodeSetHosts(); !reflect.DeepEqual(h, []string{}) { + t.Fatalf("unexpected slice of hosts: %s", h) + } +} + +// Ensure cluster can compare its Nodes and Members +func TestCluster_Health(t *testing.T) { + c := pilosa.Cluster{ + Nodes: []*pilosa.Node{ + {Host: "serverA:1000"}, + {Host: "serverB:1000"}, + {Host: "serverC:1000"}, + }, + NodeSet: &pilosa.StaticNodeSet{}, + } + + j, err := c.NodeSet.Join([]*pilosa.Node{ + &pilosa.Node{Host: "serverA:1000"}, + &pilosa.Node{Host: "serverC:1000"}, + &pilosa.Node{Host: "serverD:1000"}, + }) + if err != nil { + t.Fatalf("unexpected gossiper nodes: %s", j) + } + + // Verify a DOWN node is reported, and extraneous nodes are ignored + if a := c.Health(); !reflect.DeepEqual(a, map[string]string{ + "serverA:1000": "UP", + "serverB:1000": "DOWN", + "serverC:1000": "UP", + }) { + t.Fatalf("unexpected health: %s", spew.Sdump(a)) + } +} + // NewCluster returns a cluster with n nodes and uses a mod-based hasher. func NewCluster(n int) *pilosa.Cluster { c := pilosa.NewCluster() diff --git a/config.go b/config.go index 7d2efc5fe..4c0554ad4 100644 --- a/config.go +++ b/config.go @@ -1,11 +1,15 @@ package pilosa -import "time" +import ( + "net" + "time" +) const ( // DefaultHost is the default hostname and port to use. - DefaultHost = "localhost" - DefaultPort = "10101" + DefaultHost = "localhost" + DefaultPort = "10101" + DefaultGossipPort = 14000 ) // Config represents the configuration for the command. @@ -14,9 +18,11 @@ type Config struct { Host string `toml:"host"` Cluster struct { - ReplicaN int `toml:"replicas"` - Nodes []string `toml:"hosts"` - PollingInterval Duration `toml:"polling-interval"` + ReplicaN int `toml:"replicas"` + MessengerType string `toml:"messenger-type"` + Nodes []string `toml:"hosts"` + PollingInterval Duration `toml:"polling-interval"` + Gossip *ConfigGossip `toml:"gossip"` } `toml:"cluster"` Plugins struct { @@ -30,6 +36,15 @@ type Config struct { LogPath string `toml:"log-path"` } +type ConfigNode struct { + Host string `toml:"host"` +} + +type ConfigGossip struct { + Port int `toml:"port"` + Seed string `toml:"seed"` +} + // NewConfig returns an instance of Config with default options. func NewConfig() *Config { c := &Config{ @@ -61,6 +76,29 @@ func (c *Config) PilosaCluster() *Cluster { cluster.Nodes = append(cluster.Nodes, &Node{Host: hostport}) } + // Setup a Broadcast (over HTTP) or Gossip NodeSet based on config. + if c.Cluster.MessengerType == "broadcast" { + cluster.NodeSet = NewHTTPNodeSet() + cluster.NodeSet.Join(cluster.Nodes) + } else if (c.Cluster.MessengerType == "gossip") && (c.Cluster.Gossip != nil) { + gossipPort := DefaultGossipPort + gossipSeed := DefaultHost + if c.Cluster.Gossip.Port != 0 { + gossipPort = c.Cluster.Gossip.Port + } + if c.Cluster.Gossip.Seed != "" { + gossipSeed = c.Cluster.Gossip.Seed + } + // get the host portion of addr to use for binding + gossipHost, _, err := net.SplitHostPort(c.Host) + if err != nil { + gossipHost = c.Host + } + cluster.NodeSet = NewGossipNodeSet(c.Host, gossipHost, gossipPort, gossipSeed) + } else { + cluster.NodeSet = NewStaticNodeSet() + } + return cluster } diff --git a/db.go b/db.go index 3abff2a9c..ca03188a3 100644 --- a/db.go +++ b/db.go @@ -43,7 +43,8 @@ type DB struct { // Profile attribute storage and cache profileAttrStore *AttrStore - stats StatsClient + messenger Messenger + stats StatsClient LogOutput io.Writer } @@ -67,6 +68,7 @@ func NewDB(path, name string) (*DB, error) { columnLabel: DefaultColumnLabel, + messenger: NopMessenger, stats: NopStatsClient, LogOutput: ioutil.Discard, }, nil @@ -416,6 +418,7 @@ func (db *DB) newFrame(path, name string) (*Frame, error) { } f.LogOutput = db.LogOutput f.stats = db.stats.WithTags(fmt.Sprintf("frame:%s", name)) + f.messenger = db.messenger return f, nil } diff --git a/executor_test.go b/executor_test.go index be005e732..7a0c95622 100644 --- a/executor_test.go +++ b/executor_test.go @@ -470,7 +470,9 @@ func TestExecutor_Execute_Remote_Bitmap(t *testing.T) { t.Fatalf("unexpected db: %s", db) } else if query.String() != `Bitmap(frame="f", id=10)` { t.Fatalf("unexpected query: %s", query.String()) - } else if !reflect.DeepEqual(slices, []uint64{0}) { //TODO: this is incorrect because the calling node doesn't know about slice 2 + // NOTE: while the following is technically incorrect (it should be {0, 2}) because the calling node doesn't know about slice 2 yet, + // we are ok with this and assuming that the calling node will become aware of slice 2 via inter-node messaging + } else if !reflect.DeepEqual(slices, []uint64{0}) { t.Fatalf("unexpected slices: %+v", slices) } diff --git a/frame.go b/frame.go index be949b19d..b299d2ac2 100644 --- a/frame.go +++ b/frame.go @@ -38,7 +38,8 @@ type Frame struct { // Bitmap attribute storage and cache bitmapAttrStore *AttrStore - stats StatsClient + messenger Messenger + stats StatsClient // Frame settings. rowLabel string @@ -66,7 +67,8 @@ func NewFrame(path, db, name string) (*Frame, error) { views: make(map[string]*View), bitmapAttrStore: NewAttrStore(filepath.Join(path, ".data")), - stats: NopStatsClient, + messenger: NopMessenger, + stats: NopStatsClient, rowLabel: DefaultRowLabel, cacheType: DefaultCacheType, @@ -409,6 +411,21 @@ func (f *Frame) CreateViewIfNotExists(name string) (*View, error) { view.BitmapAttrStore = f.bitmapAttrStore f.views[view.Name()] = view + // 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 + */ + return view, nil } diff --git a/glide.lock b/glide.lock index 7fffeca5e..e4b6ce617 100644 --- a/glide.lock +++ b/glide.lock @@ -1,6 +1,10 @@ hash: 743e8f978eb4ad8f80a2ab71b05caebbf50b6769b71aa457bc4f144fef8c6595 updated: 2017-04-18T15:33:39.035615802-05:00 imports: +- name: github.com/armon/go-metrics + version: 97c69685293dce4c0a2d0b19535179bbc976e4d2 +- name: github.com/aws/aws-sdk-go + version: 819b71cf8430e434c1eee7e7e8b0f2b8870be899 - name: github.com/boltdb/bolt version: 4b1ebc1869ad66568b313d0dc410e2be72670dda - name: github.com/BurntSushi/toml @@ -24,13 +28,21 @@ imports: subpackages: - lru - name: github.com/golang/protobuf - version: 888eb0692c857ec880338addf316bd662d5e630e + version: 8ee79997227bf9b34611aee7946ae64735e6fd93 subpackages: - proto - name: github.com/gorilla/context version: 08b5f424b9271eedf6f9f0ce86cb9396ed337a42 - name: github.com/gorilla/mux version: 392c28fe23e1c45ddba891b0320b3b5df220beea +- name: github.com/hashicorp/errwrap + version: 7554cd9344cec97297fa6649b055a8c98c2a1e55 +- name: github.com/hashicorp/go-msgpack + version: fa3f63826f7c23912c15263591e65d54d080b458 + subpackages: + - codec +- name: github.com/hashicorp/go-multierror + version: ed905158d87462226a13fe39ddf685ea65f1c11f - name: github.com/hashicorp/hcl version: 630949a3c5fa3c613328e1b8256052cbc2327c9b subpackages: @@ -42,10 +54,14 @@ imports: - json/parser - json/scanner - json/token +- name: github.com/hashicorp/memberlist + version: 9800c50ab79c002353852a9b1095e9591b161513 - name: github.com/inconshreveable/mousetrap version: 76626ae9c91c4f2a10f34cad8ce83ea42c93bb75 - name: github.com/magiconair/properties version: b3b15ef068fd0b17ddf408a23669f20811d194d2 +- name: github.com/miekg/dns + version: ca336a1f95a6b89be9c250df26c7a41742eb4a6f - name: github.com/mitchellh/mapstructure version: db1efb556f84b25a0a13a04aad883943538ad2e0 - name: github.com/pelletier/go-buffruneio @@ -68,6 +84,14 @@ imports: version: 9ff6c6923cfffbcd502984b8e0c80539a94968b7 - name: github.com/spf13/viper version: 7538d73b4eb9511d85a9f1dfef202eeb8ac260f4 +- name: golang.org/x/net + version: 60c41d1de8da134c05b7b40154a9a82bf5b7edb9 + subpackages: + - context +- name: golang.org/x/sync + version: 450f422ab23cf9881c94e2db30cac0eb1b7cf80c + subpackages: + - errgroup - name: golang.org/x/sys version: c200b10b5d5e122be351b67af224adc6128af5bf subpackages: diff --git a/glide.yaml b/glide.yaml index d4615b0ba..f317a0a1b 100644 --- a/glide.yaml +++ b/glide.yaml @@ -31,3 +31,8 @@ import: - package: github.com/spf13/viper - package: github.com/gorilla/mux version: ^1.3.0 +- package: github.com/aws/aws-sdk-go + version: ^1.6.10 +- package: github.com/hashicorp/memberlist +- package: golang.org/x/sync +- package: golang.org/x/net diff --git a/gossip.go b/gossip.go new file mode 100644 index 000000000..79d301a97 --- /dev/null +++ b/gossip.go @@ -0,0 +1,165 @@ +package pilosa + +import ( + "io" + "log" + "os" + + "github.com/gogo/protobuf/proto" + "github.com/hashicorp/memberlist" +) + +// 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 + + config *GossipConfig + + messageHandler func(m 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 = append(a, &Node{Host: n.Name}) + } + return a +} + +func (g *GossipNodeSet) Join(nodes []*Node) (int, error) { + return g.Memberlist.Join(Nodes(nodes).Hosts()) +} + +func (g *GossipNodeSet) Open() error { + ml, err := memberlist.Create(g.config.memberlistConfig) + if err != nil { + return err + } + 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{ + NumNodes: func() int { + return g.Memberlist.NumMembers() + }, + RetransmitMult: 3, + } + return nil +} + +func (g *GossipNodeSet) SetMessageHandler(f func(proto.Message) error) { + g.messageHandler = f +} + +// implementation of the messenger.Messenger interface +func (g *GossipNodeSet) SendMessage(pb proto.Message) error { + msg, err := MarshalMessage(pb) + if err != nil { + return err + } + + b := &broadcast{ + msg: msg, + notify: nil, + } + g.Broadcasts.QueueBroadcast(b) + return nil +} + +func (g *GossipNodeSet) ReceiveMessage(pb proto.Message) error { + err := g.messageHandler(pb) + if err != nil { + return err + } + return nil +} + +// implementation of the memberlist.Delegate interface +func (g *GossipNodeSet) NodeMeta(limit int) []byte { + return []byte{} +} + +func (g *GossipNodeSet) NotifyMsg(b []byte) { + m, err := UnmarshalMessage(b) + if err != nil { + g.logger().Printf("unmarshal message error: %s", err) + return + } + if err := g.ReceiveMessage(m); err != nil { + g.logger().Printf("receive message error: %s", err) + return + } +} + +func (g *GossipNodeSet) GetBroadcasts(overhead, limit int) [][]byte { + return g.Broadcasts.GetBroadcasts(overhead, limit) +} + +func (g *GossipNodeSet) LocalState(join bool) []byte { + return []byte{} +} + +func (g *GossipNodeSet) MergeRemoteState(buf []byte, join bool) { + return +} + +// logger returns a logger for the GossipNodeSet. +func (g *GossipNodeSet) logger() *log.Logger { + return log.New(g.LogOutput, "", log.LstdFlags) +} + +// broadcast represents an implementation of memberlist.Broadcast +type broadcast struct { + msg []byte + notify chan<- struct{} +} + +func (b *broadcast) Invalidates(other memberlist.Broadcast) bool { + return false +} + +func (b *broadcast) Message() []byte { + return b.msg +} + +func (b *broadcast) Finished() { + if b.notify != nil { + close(b.notify) + } +} + +//////////////////////////////////////////////////////////////// + +type GossipConfig struct { + gossipSeed string + memberlistConfig *memberlist.Config +} + +// NewGossipNodeSet returns a new instance of GossipNodeSet. +func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string) *GossipNodeSet { + g := &GossipNodeSet{ + LogOutput: os.Stderr, + } + + //TODO: pull memberlist config from pilosa.cfg file + g.config = &GossipConfig{ + memberlistConfig: memberlist.DefaultLocalConfig(), + gossipSeed: gossipSeed, + } + g.config.memberlistConfig.Name = name + g.config.memberlistConfig.BindAddr = gossipHost + 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 0e25c3f2f..d3c5a28cc 100644 --- a/handler.go +++ b/handler.go @@ -25,7 +25,8 @@ import ( // Handler represents an HTTP handler. type Handler struct { - Index *Index + Index *Index + Messenger Messenger // Local hostname & cluster configuration. Host string @@ -49,6 +50,7 @@ type Handler struct { func NewHandler() *Handler { handler := &Handler{ LogOutput: os.Stderr, + Messenger: NopMessenger, } handler.Router = NewRouter(handler) return handler @@ -111,10 +113,23 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { } } +// handleGetStatus handles GET /status requests. +func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) { + if err := json.NewEncoder(w).Encode(getStatusResponse{ + Health: h.Cluster.Health(), + }); err != nil { + h.logger().Printf("write status response error: %s", err) + } +} + type getSchemaResponse struct { DBs []*DBInfo `json:"dbs"` } +type getStatusResponse struct { + Health map[string]string `json:"health"` +} + // handlePostQuery handles /query requests. func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { dbName := mux.Vars(r)["db"] @@ -177,6 +192,36 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } } +// handlePostMessage handles /message requests. +func (h *Handler) handlePostMessage(w http.ResponseWriter, r *http.Request) { + // Verify that request is only communicating over protobufs. + if r.Header.Get("Content-Type") != "application/x-protobuf" { + http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType) + return + } + + // Read entire body. + body, err := ioutil.ReadAll(r.Body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + // Unmarshal message to specific proto type. + m, err := UnmarshalMessage(body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + if err := h.Messenger.ReceiveMessage(m); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + return +} + func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) { var ms map[string]uint64 if inverse, _ := strconv.ParseBool(r.URL.Query().Get("inverse")); inverse { @@ -324,6 +369,13 @@ func (h *Handler) handlePostDB(w http.ResponseWriter, r *http.Request) { return } + // Send the delete message to all nodes. + // NOTE: this calls a second DeleteDB on the local node + h.Messenger.SendMessage( + &internal.DeleteDBMessage{ + DB: req.DB, + }) + // Encode response. if err := json.NewEncoder(w).Encode(postDBResponse{}); err != nil { h.logger().Printf("response encoding error: %s", err) diff --git a/index.go b/index.go index 7597984ba..01409a277 100644 --- a/index.go +++ b/index.go @@ -11,6 +11,9 @@ import ( "sort" "sync" "time" + + "github.com/gogo/protobuf/proto" + "github.com/pilosa/pilosa/internal" ) // DefaultCacheFlushInterval is the default value for Fragment.CacheFlushInterval. @@ -23,6 +26,8 @@ type Index struct { // Databases by name. dbs map[string]*DB + Messenger Messenger + // Close management wg sync.WaitGroup closing chan struct{} @@ -45,7 +50,8 @@ func NewIndex() *Index { dbs: make(map[string]*DB), closing: make(chan struct{}, 0), - Stats: NopStatsClient, + Messenger: NopMessenger, + Stats: NopStatsClient, CacheFlushInterval: DefaultCacheFlushInterval, @@ -243,6 +249,7 @@ func (i *Index) newDB(path, name string) (*DB, error) { } db.LogOutput = i.LogOutput db.stats = i.Stats.WithTags(fmt.Sprintf("db:%s", db.Name())) + db.messenger = i.Messenger return db, nil } @@ -338,6 +345,25 @@ func (i *Index) flushCaches() { } } +// HandleMessage handles protobuf Messages broadcasted to nodes in the +// cluster from the Cluster's NodeSet. +func (i *Index) HandleMessage(pb proto.Message) error { + switch obj := pb.(type) { + case *internal.CreateSliceMessage: + d := i.DB(obj.DB) + if d == nil { + return fmt.Errorf("Local DB not found: %s", obj.DB) + } + d.SetRemoteMaxSlice(obj.Slice) + case *internal.DeleteDBMessage: + err := i.DeleteDB(obj.DB) + if err != nil { + return err + } + } + return nil +} + func (i *Index) logger() *log.Logger { return log.New(i.LogOutput, "", log.LstdFlags) } // IndexSyncer is an active anti-entropy tool that compares the local index diff --git a/index_test.go b/index_test.go index 4a1be4cf2..50ce1979e 100644 --- a/index_test.go +++ b/index_test.go @@ -9,6 +9,7 @@ import ( "testing" "github.com/pilosa/pilosa" + "github.com/pilosa/pilosa/internal" "github.com/pilosa/pilosa/pql" ) @@ -162,6 +163,32 @@ func TestIndexSyncer_SyncIndex(t *testing.T) { } } +// Ensure index can handle Messenger messages. +func TestIndex_HandleMessage(t *testing.T) { + // Create a local index. + idx0 := MustOpenIndex() + defer idx0.Close() + + idx0.MustCreateDBIfNotExists("d", pilosa.DBOptions{}) + + msg0 := &internal.CreateSliceMessage{ + DB: "d", + Slice: 8, + } + idx0.HandleMessage(msg0) + if ms := idx0.MaxSlices(); !reflect.DeepEqual(ms, map[string]uint64{"d": 8}) { + t.Fatalf("unexpected max slice: %s", ms) + } + + msg1 := &internal.DeleteDBMessage{ + DB: "d", + } + idx0.HandleMessage(msg1) + if ms := idx0.MaxSlices(); !reflect.DeepEqual(ms, map[string]uint64{}) { + t.Fatalf("unexpected delete db: %s", ms) + } +} + // Index is a test wrapper for pilosa.Index. type Index struct { *pilosa.Index diff --git a/messenger.go b/messenger.go new file mode 100644 index 000000000..115627ae4 --- /dev/null +++ b/messenger.go @@ -0,0 +1,73 @@ +package pilosa + +import ( + "fmt" + "reflect" + + "github.com/gogo/protobuf/proto" + "github.com/pilosa/pilosa/internal" +) + +func init() { + NopMessenger = &nopMessenger{} +} + +var NopMessenger Messenger + +// nopMessenger represents a Messenger that doesn't do anything. +type nopMessenger struct{} + +func (c *nopMessenger) SendMessage(pb proto.Message) error { + fmt.Println("NOPMessenger: Send") + return nil +} +func (c *nopMessenger) ReceiveMessage(pb proto.Message) error { + fmt.Println("NOPMessenger: Receive") + return nil +} + +type Messenger interface { + SendMessage(pb proto.Message) error + ReceiveMessage(pb proto.Message) error +} + +const ( + MessageTypeCreateSlice = 1 + MessageTypeDeleteDB = 2 +) + +func MarshalMessage(m proto.Message) ([]byte, error) { + var typ uint8 + switch obj := m.(type) { + case *internal.CreateSliceMessage: + typ = MessageTypeCreateSlice + case *internal.DeleteDBMessage: + typ = MessageTypeDeleteDB + default: + return nil, fmt.Errorf("message type not implemented for marshalling: %s", reflect.TypeOf(obj)) + } + buf, err := proto.Marshal(m) + if err != nil { + return nil, err + } + return append([]byte{typ}, buf...), nil +} + +func UnmarshalMessage(buf []byte) (proto.Message, error) { + typ, buf := buf[0], buf[1:] + + var m proto.Message + switch typ { + case MessageTypeCreateSlice: + m = &internal.CreateSliceMessage{} + case MessageTypeDeleteDB: + m = &internal.DeleteDBMessage{} + default: + return nil, fmt.Errorf("invalid message type: %d", typ) + } + + if err := proto.Unmarshal(buf, m); err != nil { + return nil, err + } + return m, nil +} diff --git a/server.go b/server.go index 55702acf9..f4cef42f6 100644 --- a/server.go +++ b/server.go @@ -32,8 +32,9 @@ type Server struct { closing chan struct{} // Data storage and HTTP interface. - Index *Index - Handler *Handler + Index *Index + Handler *Handler + Messenger Messenger // Cluster configuration. // Host is replaced with actual host after opening if port is ":0". @@ -52,8 +53,9 @@ func NewServer() *Server { s := &Server{ closing: make(chan struct{}), - Index: NewIndex(), - Handler: NewHandler(), + Index: NewIndex(), + Handler: NewHandler(), + Messenger: NopMessenger, AntiEntropyInterval: DefaultAntiEntropyInterval, PollingInterval: DefaultPollingInterval, @@ -96,6 +98,11 @@ func (s *Server) Open() error { return err } + // Open NodeSet communication + if err := s.Cluster.NodeSet.Open(); err != nil { + return err + } + // Create executor for executing queries. e := NewExecutor() e.Index = s.Index @@ -202,21 +209,15 @@ func (s *Server) monitorMaxSlices() { if s.Host != node.Host { maxSlices, _ := checkMaxSlices(node.Host) for db, newmax := range maxSlices { - // if we don't know about a db locally, create it - // so that the /schema endpoint can report it + // if we don't know about a db locally, log an error because + // db's should be created and synced prior to slice creation if localdb := s.Index.DB(db); localdb != nil { if newmax > oldmaxslices[db] { oldmaxslices[db] = newmax localdb.SetRemoteMaxSlice(newmax) } } else { - d := s.Index.DB(db) - if d == nil { - s.logger().Printf("Local DB not found: %s", db) - return - } - oldmaxslices[db] = newmax - d.SetRemoteMaxSlice(newmax) + s.logger().Printf("Local DB not found: %s", db) } } }