mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 20:37:52 +00:00
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
This commit is contained in:
parent
8e2e9254f7
commit
b0f1fc7523
21 changed files with 544 additions and 498 deletions
10
cache.go
10
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 {
|
||||
|
|
|
|||
124
cluster.go
124
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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.")
|
||||
|
||||
|
|
|
|||
57
config.go
57
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
|
||||
|
||||
|
|
|
|||
8
db.go
8
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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
32
frame.go
32
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"`
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
2
glide.lock
generated
2
glide.lock
generated
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -35,4 +35,3 @@ import:
|
|||
version: ^1.6.10
|
||||
- package: github.com/hashicorp/memberlist
|
||||
- package: golang.org/x/sync
|
||||
- package: golang.org/x/net
|
||||
|
|
|
|||
305
gossip.go
305
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)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
11
handler.go
11
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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
*/
|
||||
|
|
|
|||
44
index.go
44
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
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
262
messenger.go
262
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
|
||||
|
|
|
|||
67
server.go
67
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.
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
4
view.go
4
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,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue