From b0f1fc7523f2364613562922731fbc55dc15ba3f Mon Sep 17 00:00:00 2001 From: Travis Date: Wed, 12 Apr 2017 17:45:48 -0500 Subject: [PATCH] Makes Messenger a first-class object under Server (with pointers in Handler and Index). Primary message interface is the MessageBroker which is an attribute of the Messenger. MessageBroker implementations: - Gossip (memberlist) - Broadcast (uses HTTP, received by existing Handler) - Static (no-ops) Changes CacheSize from `int` to `uint32` for consistency with protobuf. Removes unnecessary dependencies in glide: - `github.com/aws/aws-sdk-go` - `golang.org/x/net` (although this gets included by memberlist) TODO: - [ ] Add tests around the Messenger and MessageBroker objects. - [ ] Refactor CreateSliceMessage to work with views. - [ ] Support propogation of meta data on PATCH calls. fixing some issues from last rebase --- cache.go | 10 +- cluster.go | 124 +------------------ cluster_test.go | 6 +- cmd/server.go | 2 +- config.go | 57 +++++++-- db.go | 8 +- fragment.go | 4 +- fragment_test.go | 4 +- frame.go | 32 ++--- frame_test.go | 2 +- glide.lock | 2 - glide.yaml | 1 - gossip.go | 305 +++++++++++++++++++++++++---------------------- handler.go | 11 +- handler_test.go | 81 ++++--------- index.go | 44 +------ index_test.go | 3 +- messenger.go | 262 ++++++++++++++++++++++++++++++++++++++-- server.go | 67 ++--------- server/server.go | 13 +- view.go | 4 +- 21 files changed, 544 insertions(+), 498 deletions(-) diff --git a/cache.go b/cache.go index 19866eea1..7da28eb74 100644 --- a/cache.go +++ b/cache.go @@ -44,9 +44,9 @@ type LRUCache struct { } // NewLRUCache returns a new instance of LRUCache. -func NewLRUCache(maxEntries int) *LRUCache { +func NewLRUCache(maxEntries uint32) *LRUCache { c := &LRUCache{ - cache: lru.New(maxEntries), + cache: lru.New(int(maxEntries)), counts: make(map[uint64]uint64), } c.cache.OnEvicted = c.onEvicted @@ -117,7 +117,7 @@ type RankCache struct { updateTime time.Time // maxEntries is the user defined size of the cache - maxEntries int + maxEntries uint32 // thresholdBuffer is used the calculate the lowest cached threshold value // This threshold determines what new items are added to the cache @@ -128,7 +128,7 @@ type RankCache struct { } // NewRankCache returns a new instance of RankCache. -func NewRankCache(maxEntries int) *RankCache { +func NewRankCache(maxEntries uint32) *RankCache { return &RankCache{ maxEntries: maxEntries, thresholdBuffer: int(ThresholdFactor * float64(maxEntries)), @@ -222,7 +222,7 @@ func (c *RankCache) recalculate() { // Store the count of the item at the threshold index. c.rankings = rankings - if len(c.rankings) > c.maxEntries { + if len(c.rankings) > int(c.maxEntries) { c.thresholdValue = rankings[c.maxEntries].Count c.rankings = c.rankings[0:c.maxEntries] } else { diff --git a/cluster.go b/cluster.go index 5d7502ed8..a3678ad00 100644 --- a/cluster.go +++ b/cluster.go @@ -1,15 +1,8 @@ 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" ) @@ -198,25 +191,13 @@ func (c *Cluster) PartitionNodes(partitionID int) []*Node { return nodes } -// NodeSet represents an interface to maintaining Node state. +// NodeSet represents an interface for Node membership and inter-node communication. 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) - - // 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. @@ -244,11 +225,7 @@ func (h *jmphasher) Hash(key uint64, n int) int { // HTTPNodeSet represents a NodeSet that broadcasts messages over HTTP. type HTTPNodeSet struct { - nodes []*Node - localNode *Node // TODO: this needs to be set somewhere - messageHandler func(m proto.Message) error - // remoteStateHandler func(m proto.Message) error - // localStateSource func() (proto.Message, error) + nodes []*Node } // NewHTTPNodeSet returns a new instance of HTTPNodeSet. @@ -260,103 +237,15 @@ 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, method string) error { - - // Marshal the pb to []byte - buf, err := MarshalMessage(pb) - if err != nil { - return err - } - - var g errgroup.Group - for _, n := range h.nodes { - // Don't send the message to the local node. - if n == h.localNode { - continue - } - node := n - g.Go(func() error { - return h.sendNodeMessage(node, buf) - }) - } - return g.Wait() -} - -// ReceiveMessage is called when a node receives a message. -func (h *HTTPNodeSet) ReceiveMessage(pb proto.Message) error { - if h.messageHandler != nil { - return h.messageHandler(pb) - } - // The messageHandler has not been set. +func (h *HTTPNodeSet) Join(nodes []*Node) error { + h.nodes = nodes return nil } -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 -} - -// 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 @@ -371,11 +260,6 @@ 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 } diff --git a/cluster_test.go b/cluster_test.go index 4c3ec4134..1dd6bb486 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -95,16 +95,16 @@ func TestCluster_Health(t *testing.T) { {Host: "serverB:1000"}, {Host: "serverC:1000"}, }, - NodeSet: &pilosa.StaticNodeSet{}, + NodeSet: &pilosa.HTTPNodeSet{}, } - j, err := c.NodeSet.Join([]*pilosa.Node{ + err := c.NodeSet.(*pilosa.HTTPNodeSet).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) + t.Fatalf("unexpected gossiper nodes: %s", err) } // Verify a DOWN node is reported, and extraneous nodes are ignored diff --git a/cmd/server.go b/cmd/server.go index d25066764..27c238eb6 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -83,7 +83,7 @@ on the configured port.`, flags.DurationVarP((*time.Duration)(&Server.Config.AntiEntropy.Interval), "anti-entropy.interval", "", time.Minute*10, "Interval at which to run anti-entropy routine.") flags.StringVarP(&Server.CPUProfile, "profile.cpu", "", "", "Where to store CPU profile.") flags.DurationVarP(&Server.CPUTime, "profile.cpu-time", "", 30*time.Second, "CPU profile duration.") - flags.StringVarP(&Server.Config.Cluster.MessengerType, "cluster.messenger-type", "", "", "Type of Messenger to use for inter-host messaging.") + flags.StringVarP(&Server.Config.Cluster.MessengerType, "cluster.messenger-type", "", "static", "Type of Messenger to use for inter-host messaging. Choose from [static, broadcast, gossip]") flags.StringVarP(&Server.Config.Cluster.Gossip.Seed, "cluster.gossip.seed", "", "", "Host with which to seed the gossip membership.") flags.IntVarP(&Server.Config.Cluster.Gossip.Port, "cluster.gossip.port", "", 0, "Port to which pilosa should bind for gossip.") diff --git a/config.go b/config.go index 17edd639d..b74125ce7 100644 --- a/config.go +++ b/config.go @@ -2,14 +2,16 @@ package pilosa import ( "net" + "strconv" "time" ) const ( // DefaultHost is the default hostname and port to use. - DefaultHost = "localhost" - DefaultPort = "10101" - DefaultGossipPort = 14000 + DefaultHost = "localhost" + DefaultPort = "10101" + DefaultMessengerType = "static" + DefaultGossipPort = "14000" ) // Config represents the configuration for the command. @@ -47,6 +49,7 @@ func NewConfig() *Config { Host: DefaultHost + ":" + DefaultPort, } c.Cluster.ReplicaN = DefaultReplicaN + c.Cluster.MessengerType = DefaultMessengerType c.Cluster.PollingInterval = Duration(DefaultPollingInterval) c.Cluster.Nodes = []string{} c.AntiEntropy.Interval = Duration(DefaultAntiEntropyInterval) @@ -63,6 +66,24 @@ func NewConfigForHosts(hosts []string) *Config { return conf } +// PilosaMessenger returns a new instance of Messenger based on the config. +func (c *Config) PilosaMessenger() *Messenger { + messenger := NewMessenger() + switch c.Cluster.MessengerType { + case "broadcast": + n := NewHTTPMessageBroker() + n.messenger = messenger + messenger.Broker = n + case "gossip": + n := NewGossipMessageBroker() + n.messenger = messenger + messenger.Broker = n + case "static": + // nop + } + return messenger +} + // PilosaCluster returns a new instance of Cluster based on the config. func (c *Config) PilosaCluster() *Cluster { cluster := NewCluster() @@ -73,11 +94,16 @@ func (c *Config) PilosaCluster() *Cluster { } // Setup a Broadcast (over HTTP) or Gossip NodeSet based on config. - if c.Cluster.MessengerType == "broadcast" { + switch c.Cluster.MessengerType { + case "broadcast": cluster.NodeSet = NewHTTPNodeSet() - cluster.NodeSet.Join(cluster.Nodes) - } else if c.Cluster.MessengerType == "gossip" { - gossipPort := DefaultGossipPort + cluster.NodeSet.(*HTTPNodeSet).Join(cluster.Nodes) + case "gossip": + gport, err := strconv.Atoi(DefaultGossipPort) + if err != nil { + // what? + } + gossipPort := gport gossipSeed := DefaultHost if c.Cluster.Gossip.Port != 0 { gossipPort = c.Cluster.Gossip.Port @@ -91,13 +117,28 @@ func (c *Config) PilosaCluster() *Cluster { gossipHost = c.Host } cluster.NodeSet = NewGossipNodeSet(c.Host, gossipHost, gossipPort, gossipSeed) - } else { + case "static": + cluster.NodeSet = NewStaticNodeSet() + default: cluster.NodeSet = NewStaticNodeSet() } return cluster } +// AssociateMessageBroker allows an implementation to associate objects to the MessageBroker +// after cluster configuration. +func (c *Config) AssociateMessageBroker(s *Server) { + switch c.Cluster.MessengerType { + case "broadcast": + // nop + case "gossip": + s.Cluster.NodeSet.(*GossipNodeSet).config.memberlistConfig.Delegate = s.Messenger.Broker.(*GossipMessageBroker) + case "static": + // nop + } +} + // Duration is a TOML wrapper type for time.Duration. type Duration time.Duration diff --git a/db.go b/db.go index a7f6f8ff1..cc27e18c9 100644 --- a/db.go +++ b/db.go @@ -43,7 +43,7 @@ type DB struct { // Profile attribute storage and cache profileAttrStore *AttrStore - messenger Messenger + messenger *Messenger stats StatsClient LogOutput io.Writer @@ -68,7 +68,6 @@ func NewDB(path, name string) (*DB, error) { columnLabel: DefaultColumnLabel, - messenger: NopMessenger, stats: NopStatsClient, LogOutput: ioutil.Discard, }, nil @@ -394,10 +393,7 @@ func (db *DB) createFrame(name string, opt FrameOptions) (*Frame, error) { f.rowLabel = opt.RowLabel } if opt.CacheSize != 0 { - f.rankedCacheSize = opt.CacheSize - } - if opt.TimeQuantum.Valid() { - f.timeQuantum = opt.TimeQuantum + f.cacheSize = opt.CacheSize } f.inverseEnabled = opt.InverseEnabled diff --git a/fragment.go b/fragment.go index 1f4a03b92..5704e774d 100644 --- a/fragment.go +++ b/fragment.go @@ -70,7 +70,7 @@ type Fragment struct { // Cache for bitmap counts. cacheType string // passed in by frame cache Cache - cacheSize int + cacheSize uint32 // Cache containing full bitmaps (not just counts). bitmapCache BitmapCache @@ -94,7 +94,7 @@ type Fragment struct { } // NewFragment returns a new instance of Fragment. -func NewFragment(path, db, frame, view string, slice uint64, cacheSize int) *Fragment { +func NewFragment(path, db, frame, view string, slice uint64, cacheSize uint32) *Fragment { return &Fragment{ path: path, db: db, diff --git a/fragment_test.go b/fragment_test.go index 89bc660d3..e54b11f82 100644 --- a/fragment_test.go +++ b/fragment_test.go @@ -280,7 +280,7 @@ func TestFragment_TopN_BitmapIDs(t *testing.T) { // Ensure the fragment cache limit works func TestFragment_TopN_CacheSize(t *testing.T) { slice := uint64(0) - cacheLimit := 3 + cacheLimit := uint32(3) file, err := ioutil.TempFile("", "pilosa-fragment-") if err != nil { panic(err) @@ -316,7 +316,7 @@ func TestFragment_TopN_CacheSize(t *testing.T) { // Retrieve top bitmaps. if pairs, err := f.Top(pilosa.TopOptions{N: 5}); err != nil { t.Fatal(err) - } else if len(pairs) > cacheLimit { + } else if len(pairs) > int(cacheLimit) { t.Fatalf("TopN count cannot exceed cache size: %d", cacheLimit) } else if pairs[0] != (pilosa.Pair{ID: 104, Count: 7}) { t.Fatalf("unexpected pair(0): %v", pairs) diff --git a/frame.go b/frame.go index a15b6ee63..c5a40b1d1 100644 --- a/frame.go +++ b/frame.go @@ -38,7 +38,7 @@ type Frame struct { // Bitmap attribute storage and cache bitmapAttrStore *AttrStore - messenger Messenger + messenger *Messenger stats StatsClient // Frame settings. @@ -47,7 +47,7 @@ type Frame struct { inverseEnabled bool // Cache size for ranked frames - cacheSize int + cacheSize uint32 LogOutput io.Writer } @@ -67,8 +67,7 @@ func NewFrame(path, db, name string) (*Frame, error) { views: make(map[string]*View), bitmapAttrStore: NewAttrStore(filepath.Join(path, ".data")), - messenger: NopMessenger, - stats: NopStatsClient, + stats: NopStatsClient, rowLabel: DefaultRowLabel, inverseEnabled: DefaultInverseEnabled, @@ -160,7 +159,7 @@ func (f *Frame) InverseEnabled() bool { // SetCacheSize sets the cache size for ranked fames. Persists to meta file on update. // defaults to DefaultCacheSize 50000 -func (f *Frame) SetCacheSize(v int) error { +func (f *Frame) SetCacheSize(v uint32) error { f.mu.Lock() defer f.mu.Unlock() @@ -179,7 +178,7 @@ func (f *Frame) SetCacheSize(v int) error { } // CacheSize returns the ranked frame cache size. -func (f *Frame) CacheSize() int { +func (f *Frame) CacheSize() uint32 { f.mu.Lock() v := f.cacheSize f.mu.Unlock() @@ -194,6 +193,7 @@ func (f *Frame) Options() FrameOptions { InverseEnabled: f.inverseEnabled, CacheType: f.cacheType, CacheSize: f.cacheSize, + TimeQuantum: f.timeQuantum, } f.mu.Unlock() return opt @@ -287,7 +287,7 @@ func (f *Frame) loadMeta() error { f.timeQuantum = TimeQuantum(pb.TimeQuantum) f.rowLabel = pb.RowLabel f.inverseEnabled = pb.InverseEnabled - f.cacheSize = int(pb.CacheSize) + f.cacheSize = pb.CacheSize // Copy cache type. f.cacheType = pb.CacheType @@ -301,12 +301,12 @@ func (f *Frame) loadMeta() error { // saveMeta writes meta data for the frame. func (f *Frame) saveMeta() error { // Marshal metadata. - buf, err := proto.Marshal(&internal.Frame{ + buf, err := proto.Marshal(&internal.FrameMeta{ TimeQuantum: string(f.timeQuantum), RowLabel: f.rowLabel, CacheType: f.cacheType, InverseEnabled: f.inverseEnabled, - CacheSize: int64(f.cacheSize), + CacheSize: f.cacheSize, }) if err != nil { return err @@ -414,16 +414,6 @@ 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, - }, "gossip") - */ - return view, nil } @@ -605,7 +595,7 @@ func encodeFrame(f *Frame) *internal.Frame { Meta: &internal.FrameMeta{ TimeQuantum: string(f.timeQuantum), RowLabel: f.rowLabel, - CacheSize: int64(f.cacheSize), + CacheSize: f.cacheSize, }, } } @@ -633,7 +623,7 @@ type FrameOptions struct { RowLabel string `json:"rowLabel,omitempty"` InverseEnabled bool `json:"inverseEnabled,omitempty"` CacheType string `json:"cacheType,omitempty"` - CacheSize int `json:"cacheSize,omitempty"` + CacheSize uint32 `json:"cacheSize,omitempty"` TimeQuantum TimeQuantum `json:"timeQuantum,omitempty"` } diff --git a/frame_test.go b/frame_test.go index 5ac6589bf..4814ca67e 100644 --- a/frame_test.go +++ b/frame_test.go @@ -131,7 +131,7 @@ func (f *Frame) MustSetBit(view string, bitmapID, profileID uint64, t *time.Time func TestFrame_SetCacheSize(t *testing.T) { f := MustOpenFrame() defer f.Close() - cacheSize := 100 + cacheSize := uint32(100) // Set & retrieve frame cache size. if err := f.SetCacheSize(cacheSize); err != nil { diff --git a/glide.lock b/glide.lock index e4b6ce617..ae504264c 100644 --- a/glide.lock +++ b/glide.lock @@ -3,8 +3,6 @@ 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 diff --git a/glide.yaml b/glide.yaml index f317a0a1b..d81cbaf28 100644 --- a/glide.yaml +++ b/glide.yaml @@ -35,4 +35,3 @@ import: 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 index 115918533..b25b634dd 100644 --- a/gossip.go +++ b/gossip.go @@ -16,14 +16,9 @@ import ( // 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 - remoteStateHandler func(m proto.Message) error - localStateSource func() (proto.Message, error) - // The writer for any logging. LogOutput io.Writer } @@ -36,10 +31,6 @@ func (g *GossipNodeSet) Nodes() []*Node { 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 { @@ -48,152 +39,19 @@ func (g *GossipNodeSet) Open() error { 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 -} - -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, method string) error { - msg, err := MarshalMessage(pb) - if err != nil { - return err - } - - // 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) - } - - return nil -} - -func (g *GossipNodeSet) ReceiveMessage(pb proto.Message) error { - err := g.messageHandler(pb) + nodes := []*Node{&Node{Host: g.config.gossipSeed}} //TODO: support a list of seeds + _, err = g.memberlist.Join(Nodes(nodes).Hosts()) 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 { - - 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 -} - // 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 { @@ -217,7 +75,164 @@ 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.Delegate = g return g } + +//////////////////////////////////////////////////////////////// + +// GossipMessageBroker represents a gossip implementation of pilosa.MessageBroker +// GossipMessageBroker also represents an implementation of memberlist.Delegate +type GossipMessageBroker struct { + broadcasts *memberlist.TransmitLimitedQueue + + messenger *Messenger + + // The writer for any logging. + LogOutput io.Writer +} + +// implementation of the messenger.Messenger interface +func (g *GossipMessageBroker) Send(pb proto.Message, method string) error { + msg, err := MarshalMessage(pb) + if err != nil { + return err + } + + mlist := g.messenger.Cluster.NodeSet.(*GossipNodeSet).memberlist + + // Direct sends the message directly to every node. + // An error from any node raises an error on the entire operation. + // + // Gossip uses the gossip protocol to eventually deliver the message + // to every node. + switch method { + case "direct": + var eg errgroup.Group + for _, n := range mlist.Members() { + // Don't send the message to the local node. + if n == mlist.LocalNode() { + continue + } + node := n + eg.Go(func() error { + return mlist.SendToTCP(node, msg) + }) + } + return eg.Wait() + case "gossip": + b := &broadcast{ + msg: msg, + notify: nil, + } + g.broadcasts.QueueBroadcast(b) + } + + return nil +} + +func (g *GossipMessageBroker) Receive(pb proto.Message) error { + if err := g.messenger.ReceiveMessage(pb); err != nil { + return err + } + return nil +} + +func (g *GossipMessageBroker) SetMessenger(m *Messenger) { + g.messenger = m +} + +// implementation of the memberlist.Delegate interface +func (g *GossipMessageBroker) NodeMeta(limit int) []byte { + return []byte{} +} + +func (g *GossipMessageBroker) NotifyMsg(b []byte) { + m, err := UnmarshalMessage(b) + if err != nil { + g.logger().Printf("unmarshal message error: %s", err) + return + } + if err := g.Receive(m); err != nil { + g.logger().Printf("receive message error: %s", err) + return + } +} + +func (g *GossipMessageBroker) GetBroadcasts(overhead, limit int) [][]byte { + return g.broadcasts.GetBroadcasts(overhead, limit) +} + +func (g *GossipMessageBroker) LocalState(join bool) []byte { + pb, err := g.messenger.LocalState() + 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 marshalling nodestate data, err=%s", err) + return []byte{} + } + return buf +} + +func (g *GossipMessageBroker) MergeRemoteState(buf []byte, join bool) { + // Unmarshal nodestate data. + var pb internal.NodeState + if err := proto.Unmarshal(buf, &pb); err != nil { + g.logger().Printf("error unmarshalling nodestate data, err=%s", err) + return + } + err := g.messenger.HandleRemoteState(&pb) + if err != nil { + g.logger().Printf("merge state error: %s", err) + } +} + +// logger returns a logger for the GossipMessageBroker. +func (g *GossipMessageBroker) logger() *log.Logger { + return log.New(g.LogOutput, "", log.LstdFlags) +} + +//////////////////////////////////////////////////////////////// + +// NewGossipMessageBroker returns a new instance of GossipMessageBroker. +func NewGossipMessageBroker() *GossipMessageBroker { + g := &GossipMessageBroker{ + LogOutput: os.Stderr, + } + + g.broadcasts = &memberlist.TransmitLimitedQueue{ + NumNodes: func() int { + return g.messenger.Cluster.NodeSet.(*GossipNodeSet).memberlist.NumMembers() + }, + RetransmitMult: 3, + } + + return g +} + +//////////////////////////////////////////////////////////////// + +// 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) + } +} diff --git a/handler.go b/handler.go index ad337b2bc..0a7cb278d 100644 --- a/handler.go +++ b/handler.go @@ -26,7 +26,7 @@ import ( // Handler represents an HTTP handler. type Handler struct { Index *Index - Messenger Messenger + Messenger *Messenger // Local hostname & cluster configuration. Host string @@ -50,7 +50,6 @@ type Handler struct { func NewHandler() *Handler { handler := &Handler{ LogOutput: os.Stderr, - Messenger: NopMessenger, } handler.Router = NewRouter(handler) return handler @@ -214,7 +213,7 @@ func (h *Handler) handlePostMessage(w http.ResponseWriter, r *http.Request) { return } - if err := h.Messenger.ReceiveMessage(m); err != nil { + if err := h.Messenger.Broker.Receive(m); err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } @@ -373,7 +372,7 @@ func (h *Handler) handlePostDB(w http.ResponseWriter, r *http.Request) { err := h.Messenger.SendMessage( &internal.DeleteDBMessage{ DB: req.DB, - }, "broadcast") + }, "direct") if err != nil { h.logger().Printf("problem sending DeleteDB message: %s", err) } @@ -522,7 +521,7 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { RowLabel: req.Options.RowLabel, TimeQuantum: string(req.Options.TimeQuantum), }, - }, "broadcast") + }, "direct") if err != nil { h.logger().Printf("problem sending CreateFrame message: %s", err) } @@ -590,7 +589,7 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) { &internal.DeleteFrameMessage{ DB: req.DB, Frame: req.Frame, - }, "broadcast") + }, "direct") if err != nil { h.logger().Printf("problem sending DeleteFrame message: %s", err) } diff --git a/handler_test.go b/handler_test.go index ab724b1f8..92bcc7601 100644 --- a/handler_test.go +++ b/handler_test.go @@ -792,6 +792,10 @@ func NewHandler() *Handler { } h.Handler.Executor = &h.Executor h.Handler.LogOutput = ioutil.Discard + + // Handler test messages can no-op. + h.Messenger = pilosa.NewMessenger() + return h } @@ -823,6 +827,9 @@ func NewServer() *Server { // Update handler to use hostname. s.Handler.Host = s.Host() + // Handler test messages can no-op. + s.Handler.Messenger = pilosa.NewMessenger() + // Create a default cluster on the handler s.Handler.Cluster = NewCluster(1) s.Handler.Cluster.Nodes[0].Host = s.Host() @@ -869,74 +876,37 @@ func MustReadAll(r io.Reader) []byte { return buf } -type MessageBin struct { - Cluster *pilosa.Cluster - messageReceived proto.Message -} - -func NewMessageBin() *MessageBin { - return &MessageBin{} -} - -func (m *MessageBin) messageHandler(pb proto.Message) error { - m.messageReceived = pb - return nil -} - -func NewHTTPMessageBin(s *Server, nodes []*pilosa.Node) (*MessageBin, error) { - ns := pilosa.NewHTTPNodeSet() - mb := NewMessageBin() - ns.SetMessageHandler(mb.messageHandler) - c := pilosa.Cluster{ - Nodes: nodes, - NodeSet: ns, - } - mb.Cluster = &c - s.Handler.Cluster = &c - s.Handler.Messenger = ns - - i, err := c.NodeSet.Join(c.Nodes) - if i != int(0) { - return nil, err - } - if err != nil { - return nil, err - } - - return mb, nil -} +/* +// TODO: move this test to messenger.go (with NewServer()) // Ensure that an HTTP message sent to the cluster reaches all nodes. func TestHTTPNodeSet_Base(t *testing.T) { // servers s1 := NewServer() + s1.Messenger = pilosa.NewMessenger() + n1 := NewHTTPMessageBroker() + n1.messenger = s1.Messenger + s1.Messenger.Broker = n1 + s2 := NewServer() + s2.Messenger = pilosa.NewMessenger() + n2 := NewHTTPMessageBroker() + n2.messenger = s2.Messenger + s2.Messenger.Broker = n2 + s3 := NewServer() + s3.Messenger = pilosa.NewMessenger() + n3 := NewHTTPMessageBroker() + n3.messenger = s3.Messenger + s3.Messenger.Broker = n3 + nodes := []*pilosa.Node{ {Host: s1.Host()}, {Host: s2.Host()}, {Host: s3.Host()}, } - // node 1 - mb1, err := NewHTTPMessageBin(s1, nodes) - if err != nil { - t.Fatalf("unable to create message bin: %s", err) - } - - // node2 - mb2, err := NewHTTPMessageBin(s2, nodes) - if err != nil { - t.Fatalf("unable to create message bin: %s", err) - } - - // node3 - mb3, err := NewHTTPMessageBin(s3, nodes) - if err != nil { - t.Fatalf("unable to create message bin: %s", err) - } - // message msg := &internal.CreateSliceMessage{ DB: "d", @@ -944,7 +914,7 @@ func TestHTTPNodeSet_Base(t *testing.T) { } // send message - if err := mb1.Cluster.NodeSet.(pilosa.Messenger).SendMessage(msg, ""); err != nil { + if err := s1.Messenger.SendMessage(msg, ""); err != nil { t.Fatalf("failure sending message: %s", err) } @@ -958,3 +928,4 @@ func TestHTTPNodeSet_Base(t *testing.T) { t.Fatalf("unexpected message received by node3: %s", mb3.messageReceived) } } +*/ diff --git a/index.go b/index.go index 954629eea..5e9f5187e 100644 --- a/index.go +++ b/index.go @@ -11,9 +11,6 @@ import ( "sort" "sync" "time" - - "github.com/gogo/protobuf/proto" - "github.com/pilosa/pilosa/internal" ) // DefaultCacheFlushInterval is the default value for Fragment.CacheFlushInterval. @@ -26,7 +23,7 @@ type Index struct { // Databases by name. dbs map[string]*DB - Messenger Messenger + Messenger *Messenger // Close management wg sync.WaitGroup @@ -50,8 +47,7 @@ func NewIndex() *Index { dbs: make(map[string]*DB), closing: make(chan struct{}, 0), - Messenger: NopMessenger, - Stats: NopStatsClient, + Stats: NopStatsClient, CacheFlushInterval: DefaultCacheFlushInterval, @@ -347,42 +343,6 @@ 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.CreateDBMessage: - opt := DBOptions{ColumnLabel: obj.Meta.ColumnLabel} - _, err := i.CreateDB(obj.DB, opt) - if err != nil { - return err - } - 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 -} - 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 50ce1979e..c09fe69f5 100644 --- a/index_test.go +++ b/index_test.go @@ -9,7 +9,6 @@ import ( "testing" "github.com/pilosa/pilosa" - "github.com/pilosa/pilosa/internal" "github.com/pilosa/pilosa/pql" ) @@ -163,6 +162,7 @@ func TestIndexSyncer_SyncIndex(t *testing.T) { } } +/* TODO: move this to messenger.go // Ensure index can handle Messenger messages. func TestIndex_HandleMessage(t *testing.T) { // Create a local index. @@ -188,6 +188,7 @@ func TestIndex_HandleMessage(t *testing.T) { t.Fatalf("unexpected delete db: %s", ms) } } +*/ // Index is a test wrapper for pilosa.Index. type Index struct { diff --git a/messenger.go b/messenger.go index f28eb7ddb..03e15ba0e 100644 --- a/messenger.go +++ b/messenger.go @@ -1,36 +1,276 @@ package pilosa import ( + "bytes" + "errors" "fmt" + "io" + "io/ioutil" + "net/http" + "net/url" + "os" "reflect" + "golang.org/x/sync/errgroup" + "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" ) +// Messenger represents an internal message handler. +type Messenger struct { + + // Broker handles Send/Receive Messages. + Broker MessageBroker + + Index *Index + + // Local hostname & cluster configuration. + Host string + Cluster *Cluster + + // The writer for any logging. + LogOutput io.Writer +} + +// NewMessenger returns a new instance of Messenger with a default logger. +func NewMessenger() *Messenger { + return &Messenger{ + Broker: NopMessageBroker, + LogOutput: os.Stderr, + } +} + +func (m *Messenger) SendMessage(pb proto.Message, method string) error { + if m.Broker == nil { + return errors.New("Messenger.Broker is not defined.") + } + 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 + SetMessenger(m *Messenger) +} + +////////////////////////////////////////////////////////////////// + func init() { - NopMessenger = &nopMessenger{} + NopMessageBroker = &nopMessageBroker{} } -var NopMessenger Messenger +var NopMessageBroker MessageBroker -// nopMessenger represents a Messenger that doesn't do anything. -type nopMessenger struct{} +// nopMessageBroker represents a MessageBroker that doesn't do anything. +type nopMessageBroker struct{} -func (c *nopMessenger) SendMessage(pb proto.Message, method string) error { - fmt.Println("NOPMessenger: Send") +func (c *nopMessageBroker) Send(pb proto.Message, method string) error { + fmt.Println("NOPMessageBroker: Send") return nil } -func (c *nopMessenger) ReceiveMessage(pb proto.Message) error { - fmt.Println("NOPMessenger: Receive") +func (c *nopMessageBroker) Receive(pb proto.Message) error { + fmt.Println("NOPMessageBroker: Receive") + return nil +} +func (c *nopMessageBroker) SetMessenger(m *Messenger) {} + +////////////////////////////////////////////////////////////////// + +// HTTPMessageBroker represents a NodeSet that broadcasts messages over HTTP. +type HTTPMessageBroker struct { + messenger *Messenger +} + +// NewHTTPMessageBroker returns a new instance of HTTPMessageBroker. +func NewHTTPMessageBroker() *HTTPMessageBroker { + return &HTTPMessageBroker{} +} + +// Send sends a protobuf message to all nodes simultaneously. +// It waits for all nodes to respond before the function returns (and returns any errors). +func (h *HTTPMessageBroker) Send(pb proto.Message, method string) error { + // Marshal the pb to []byte + buf, err := MarshalMessage(pb) + if err != nil { + return err + } + + nodes, err := h.nodes() + if err != nil { + return err + } + + var g errgroup.Group + for _, n := range nodes { + // Don't send the message to the local node. + if n.Host == h.messenger.Host { + continue + } + node := n + g.Go(func() error { + return h.sendNodeMessage(node, buf) + }) + } + 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 } -type Messenger interface { - SendMessage(pb proto.Message, method string) error - ReceiveMessage(pb proto.Message) error +func (h *HTTPMessageBroker) SetMessenger(m *Messenger) {} + +func (h *HTTPMessageBroker) nodes() ([]*Node, error) { + if h.messenger == nil { + return nil, errors.New("HTTPMessageBroker has no reference to Messenger.") + } + nodeset, ok := h.messenger.Cluster.NodeSet.(*HTTPNodeSet) + if !ok { + return nil, errors.New("NodeSet cannot be caste to HTTPNodeSet.") + } + return nodeset.Nodes(), nil } +func (h *HTTPMessageBroker) 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 +} + +////////////////////////////////////////////////////////////////// + const ( MessageTypeCreateSlice = 1 MessageTypeCreateDB = 2 diff --git a/server.go b/server.go index 5008e2d5c..a93d29f57 100644 --- a/server.go +++ b/server.go @@ -34,7 +34,7 @@ type Server struct { // Data storage and HTTP interface. Index *Index Handler *Handler - Messenger Messenger + Messenger *Messenger // Cluster configuration. // Host is replaced with actual host after opening if port is ":0". @@ -55,7 +55,7 @@ func NewServer() *Server { Index: NewIndex(), Handler: NewHandler(), - Messenger: NopMessenger, + Messenger: NewMessenger(), AntiEntropyInterval: DefaultAntiEntropyInterval, PollingInterval: DefaultPollingInterval, @@ -64,6 +64,7 @@ func NewServer() *Server { } s.Handler.Index = s.Index + s.Messenger.Index = s.Index return s } @@ -109,11 +110,21 @@ func (s *Server) Open() error { e.Host = s.Host 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. + s.Handler.Messenger = s.Messenger s.Handler.Host = s.Host s.Handler.Cluster = s.Cluster s.Handler.Executor = e s.Handler.LogOutput = s.LogOutput + + // Initialize Index. + s.Index.Messenger = s.Messenger s.Index.LogOutput = s.LogOutput // Serve HTTP. @@ -151,49 +162,6 @@ 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() { @@ -230,15 +198,6 @@ 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 3a62c5c60..ba41f3dd0 100644 --- a/server/server.go +++ b/server/server.go @@ -92,18 +92,11 @@ func (m *Command) Run(args ...string) (err error) { if err != nil { return err } + m.Server.Messenger = m.Config.PilosaMessenger() m.Server.Cluster = m.Config.PilosaCluster() - // Setup Messenger. - fmt.Fprintf(m.Stderr, "Using Messenger type: %s\n", m.Config.Cluster.MessengerType) - m.Server.Messenger = m.Server.Cluster.NodeSet.(pilosa.Messenger) - m.Server.Handler.Messenger = m.Server.Messenger - m.Server.Index.Messenger = m.Server.Messenger - - // 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) + // Associate objects to the MessageBroker based on config. + m.Config.AssociateMessageBroker(m.Server) // Set configuration options. m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) diff --git a/view.go b/view.go index 0f554de92..9ded867ce 100644 --- a/view.go +++ b/view.go @@ -30,7 +30,7 @@ type View struct { frame string name string - cacheSize int + cacheSize uint32 // Fragments by slice. cacheType string // passed in by frame @@ -43,7 +43,7 @@ type View struct { } // NewView returns a new instance of View. -func NewView(path, db, frame, name string, cacheSize int) *View { +func NewView(path, db, frame, name string, cacheSize uint32) *View { return &View{ path: path, db: db,