mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
commit
08590be4c3
14 changed files with 89 additions and 44 deletions
12
broadcast.go
12
broadcast.go
|
|
@ -17,23 +17,27 @@ type NodeSet interface {
|
|||
Open() error
|
||||
}
|
||||
|
||||
// StaticNodeSet represents a basic NodeSet for testing
|
||||
// StaticNodeSet represents a basic NodeSet for testing.
|
||||
type StaticNodeSet struct {
|
||||
nodes []*Node
|
||||
}
|
||||
|
||||
// NewStaticNodeSet creates a statically defined NodeSet.
|
||||
func NewStaticNodeSet() *StaticNodeSet {
|
||||
return &StaticNodeSet{}
|
||||
}
|
||||
|
||||
// Nodes implements the NodeSet interface and returns a list of nodes in the cluster.
|
||||
func (s *StaticNodeSet) Nodes() []*Node {
|
||||
return s.nodes
|
||||
}
|
||||
|
||||
// Open implements the NodeSet interface to start network activity, but for a static NodeSet it does nothing.
|
||||
func (s *StaticNodeSet) Open() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Join sets the NodeSet nodes to the slice of Nodes passed in.
|
||||
func (s *StaticNodeSet) Join(nodes []*Node) error {
|
||||
s.nodes = nodes
|
||||
return nil
|
||||
|
|
@ -49,9 +53,9 @@ func init() {
|
|||
NopBroadcaster = &nopBroadcaster{}
|
||||
}
|
||||
|
||||
// NopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
var NopBroadcaster Broadcaster
|
||||
|
||||
// nopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
type nopBroadcaster struct{}
|
||||
|
||||
// SendSync A no-op implemenetation of Broadcaster SendSync method.
|
||||
|
|
@ -85,8 +89,10 @@ type nopBroadcastReceiver struct{}
|
|||
|
||||
func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil }
|
||||
|
||||
// NopBroadcastReceiver is a no-op implementation of the BroadcastReceiver.
|
||||
var NopBroadcastReceiver = &nopBroadcastReceiver{}
|
||||
|
||||
// Broadcast message types.
|
||||
const (
|
||||
MessageTypeCreateSlice = 1
|
||||
MessageTypeCreateIndex = 2
|
||||
|
|
@ -95,6 +101,7 @@ const (
|
|||
MessageTypeDeleteFrame = 5
|
||||
)
|
||||
|
||||
// MarshalMessage encodes the protobuf message into a byte slice.
|
||||
func MarshalMessage(m proto.Message) ([]byte, error) {
|
||||
var typ uint8
|
||||
switch obj := m.(type) {
|
||||
|
|
@ -118,6 +125,7 @@ func MarshalMessage(m proto.Message) ([]byte, error) {
|
|||
return append([]byte{typ}, buf...), nil
|
||||
}
|
||||
|
||||
// UnmarshalMessage decodes the byte slice into a protobuf message.
|
||||
func UnmarshalMessage(buf []byte) (proto.Message, error) {
|
||||
typ, buf := buf[0], buf[1:]
|
||||
|
||||
|
|
|
|||
16
cache.go
16
cache.go
|
|
@ -53,6 +53,7 @@ func NewLRUCache(maxEntries uint32) *LRUCache {
|
|||
return c
|
||||
}
|
||||
|
||||
// BulkAdd adds a count to the cache unsorted. You should Invalidate after completion.
|
||||
func (c *LRUCache) BulkAdd(id, n uint64) {
|
||||
c.Add(id, n)
|
||||
}
|
||||
|
|
@ -194,6 +195,7 @@ func (c *RankCache) Invalidate() {
|
|||
c.invalidate()
|
||||
}
|
||||
|
||||
// Recalculate rebuilds the cache.
|
||||
func (c *RankCache) Recalculate() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
|
@ -303,19 +305,23 @@ type PairHeap struct {
|
|||
Pairs
|
||||
}
|
||||
|
||||
// Less implemets the Sort interface.
|
||||
// reports whether the element with index i should sort before the element with index j.
|
||||
func (p PairHeap) Less(i, j int) bool { return p.Pairs[i].Count < p.Pairs[j].Count }
|
||||
|
||||
func (h *Pairs) Push(x interface{}) {
|
||||
// Push appends the element onto the Pair slice.
|
||||
func (p *Pairs) Push(x interface{}) {
|
||||
// Push and Pop use pointer receivers because they modify the slice's length,
|
||||
// not just its contents.
|
||||
*h = append(*h, x.(Pair))
|
||||
*p = append(*p, x.(Pair))
|
||||
}
|
||||
|
||||
func (h *Pairs) Pop() interface{} {
|
||||
old := *h
|
||||
// Pop removes the minimum element from the Pair slice.
|
||||
func (p *Pairs) Pop() interface{} {
|
||||
old := *p
|
||||
n := len(old)
|
||||
x := old[n-1]
|
||||
*h = old[0 : n-1]
|
||||
*p = old[0 : n-1]
|
||||
return x
|
||||
}
|
||||
|
||||
|
|
|
|||
29
client.go
29
client.go
|
|
@ -315,6 +315,7 @@ func (c *Client) Import(ctx context.Context, index, frame string, slice uint64,
|
|||
return nil
|
||||
}
|
||||
|
||||
// MarshalImportPayload marshalls the import parameters into a protobuf byte slice.
|
||||
func MarshalImportPayload(index, frame string, slice uint64, bits []Bit) ([]byte, error) {
|
||||
// Separate row and column IDs to reduce allocations.
|
||||
rowIDs := Bits(bits).RowIDs()
|
||||
|
|
@ -982,36 +983,36 @@ func (p Bits) Less(i, j int) bool {
|
|||
}
|
||||
|
||||
// RowIDs returns a slice of all the row IDs.
|
||||
func (a Bits) RowIDs() []uint64 {
|
||||
other := make([]uint64, len(a))
|
||||
for i := range a {
|
||||
other[i] = a[i].RowID
|
||||
func (p Bits) RowIDs() []uint64 {
|
||||
other := make([]uint64, len(p))
|
||||
for i := range p {
|
||||
other[i] = p[i].RowID
|
||||
}
|
||||
return other
|
||||
}
|
||||
|
||||
// ColumnIDs returns a slice of all the column IDs.
|
||||
func (a Bits) ColumnIDs() []uint64 {
|
||||
other := make([]uint64, len(a))
|
||||
for i := range a {
|
||||
other[i] = a[i].ColumnID
|
||||
func (p Bits) ColumnIDs() []uint64 {
|
||||
other := make([]uint64, len(p))
|
||||
for i := range p {
|
||||
other[i] = p[i].ColumnID
|
||||
}
|
||||
return other
|
||||
}
|
||||
|
||||
// Timestamps returns a slice of all the timestamps.
|
||||
func (a Bits) Timestamps() []int64 {
|
||||
other := make([]int64, len(a))
|
||||
for i := range a {
|
||||
other[i] = a[i].Timestamp
|
||||
func (p Bits) Timestamps() []int64 {
|
||||
other := make([]int64, len(p))
|
||||
for i := range p {
|
||||
other[i] = p[i].Timestamp
|
||||
}
|
||||
return other
|
||||
}
|
||||
|
||||
// GroupBySlice returns a map of bits by slice.
|
||||
func (a Bits) GroupBySlice() map[uint64][]Bit {
|
||||
func (p Bits) GroupBySlice() map[uint64][]Bit {
|
||||
m := make(map[uint64][]Bit)
|
||||
for _, bit := range a {
|
||||
for _, bit := range p {
|
||||
slice := bit.ColumnID / SliceWidth
|
||||
m[slice] = append(m[slice], bit)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,8 +13,10 @@ const (
|
|||
|
||||
// DefaultReplicaN is the default number of replicas per partition.
|
||||
DefaultReplicaN = 1
|
||||
)
|
||||
|
||||
// NodeState represents node state returned in /status endpoint for a node in the cluster.
|
||||
// NodeState represents node state returned in /status endpoint for a node in the cluster.
|
||||
const (
|
||||
NodeStateUp = "UP"
|
||||
NodeStateDown = "DOWN"
|
||||
)
|
||||
|
|
@ -125,7 +127,7 @@ func NewCluster() *Cluster {
|
|||
}
|
||||
}
|
||||
|
||||
// NodeSetHosts returns the list of host strings for NodeSet members
|
||||
// NodeSetHosts returns the list of host strings for NodeSet members.
|
||||
func (c *Cluster) NodeSetHosts() []string {
|
||||
if c.NodeSet == nil {
|
||||
return []string{}
|
||||
|
|
@ -152,7 +154,7 @@ func (c *Cluster) NodeStates() map[string]string {
|
|||
return h
|
||||
}
|
||||
|
||||
// State returns the internal ClusterState representation.
|
||||
// Status returns the internal ClusterStatus representation.
|
||||
func (c *Cluster) Status() *internal.ClusterStatus {
|
||||
return &internal.ClusterStatus{
|
||||
Nodes: encodeClusterStatus(c.Nodes),
|
||||
|
|
|
|||
15
config.go
15
config.go
|
|
@ -3,10 +3,16 @@ package pilosa
|
|||
import "time"
|
||||
|
||||
const (
|
||||
// DefaultHost is the default hostname and port to use.
|
||||
DefaultHost = "localhost"
|
||||
DefaultPort = "10101"
|
||||
DefaultClusterType = "static"
|
||||
// DefaultHost is the default hostname to use.
|
||||
DefaultHost = "localhost"
|
||||
|
||||
// DefaultPort is the default port use with the hostname.
|
||||
DefaultPort = "10101"
|
||||
|
||||
// DefaultClusterType sets the node intercommunication method.
|
||||
DefaultClusterType = "static"
|
||||
|
||||
// DefaultInternalPort the port the nodes intercommunicate on.
|
||||
DefaultInternalPort = "14000"
|
||||
)
|
||||
|
||||
|
|
@ -67,6 +73,7 @@ func (d *Duration) UnmarshalText(text []byte) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// MarshalText writes duration value in text format.
|
||||
func (d Duration) MarshalText() (text []byte, err error) {
|
||||
return []byte(d.String()), nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -801,7 +801,7 @@ func (e *Executor) executeSetRowAttrs(ctx context.Context, index string, c *pql.
|
|||
if err != nil {
|
||||
return fmt.Errorf("reading SetRowAttrs() row: %v", err)
|
||||
} else if !ok {
|
||||
return fmt.Errorf("SetRowAttrs() row field '%v' required.", rowLabel)
|
||||
return fmt.Errorf("SetRowAttrs() row field '%v' required", rowLabel)
|
||||
}
|
||||
|
||||
// Copy args and remove reserved fields.
|
||||
|
|
@ -860,7 +860,7 @@ func (e *Executor) executeBulkSetRowAttrs(ctx context.Context, index string, cal
|
|||
if err != nil {
|
||||
return nil, fmt.Errorf("reading SetRowAttrs() row: %v", rowLabel)
|
||||
} else if !ok {
|
||||
return nil, fmt.Errorf("SetRowAttrs row field '%v' required.", rowLabel)
|
||||
return nil, fmt.Errorf("SetRowAttrs row field '%v' required", rowLabel)
|
||||
}
|
||||
|
||||
// Copy args and remove reserved fields.
|
||||
|
|
|
|||
|
|
@ -24,12 +24,13 @@ type GossipNodeSet struct {
|
|||
broadcasts *memberlist.TransmitLimitedQueue
|
||||
|
||||
statusHandler pilosa.StatusHandler
|
||||
config *GossipConfig
|
||||
config *gossipConfig
|
||||
|
||||
// The writer for any logging.
|
||||
LogOutput io.Writer
|
||||
}
|
||||
|
||||
// Nodes implements the NodeSet interface and returns a list of nodes in the cluster.
|
||||
func (g *GossipNodeSet) Nodes() []*pilosa.Node {
|
||||
a := make([]*pilosa.Node, 0, g.memberlist.NumMembers())
|
||||
for _, n := range g.memberlist.Members() {
|
||||
|
|
@ -38,11 +39,13 @@ func (g *GossipNodeSet) Nodes() []*pilosa.Node {
|
|||
return a
|
||||
}
|
||||
|
||||
// Start implements the BroadcastReceiver interface and sets the BroadcastHandler
|
||||
func (g *GossipNodeSet) Start(h pilosa.BroadcastHandler) error {
|
||||
g.handler = h
|
||||
return nil
|
||||
}
|
||||
|
||||
// Open implements the NodeSet interface to start network activity.
|
||||
func (g *GossipNodeSet) Open() error {
|
||||
if g.handler == nil {
|
||||
return fmt.Errorf("opening GossipNodeSet: you must call Start(pilosa.BroadcastHandler) before calling Open()")
|
||||
|
|
@ -75,7 +78,7 @@ func (g *GossipNodeSet) logger() *log.Logger {
|
|||
|
||||
////////////////////////////////////////////////////////////////
|
||||
|
||||
type GossipConfig struct {
|
||||
type gossipConfig struct {
|
||||
gossipSeed string
|
||||
memberlistConfig *memberlist.Config
|
||||
}
|
||||
|
|
@ -87,7 +90,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed
|
|||
}
|
||||
|
||||
//TODO: pull memberlist config from pilosa.cfg file
|
||||
g.config = &GossipConfig{
|
||||
g.config = &gossipConfig{
|
||||
memberlistConfig: memberlist.DefaultLocalConfig(),
|
||||
gossipSeed: gossipSeed,
|
||||
}
|
||||
|
|
@ -103,7 +106,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed
|
|||
return g
|
||||
}
|
||||
|
||||
// SendSync implementation of the Broadcaster interface
|
||||
// SendSync implementation of the Broadcaster interface.
|
||||
func (g *GossipNodeSet) SendSync(pb proto.Message) error {
|
||||
msg, err := pilosa.MarshalMessage(pb)
|
||||
if err != nil {
|
||||
|
|
@ -131,7 +134,7 @@ func (g *GossipNodeSet) SendSync(pb proto.Message) error {
|
|||
return eg.Wait()
|
||||
}
|
||||
|
||||
// SendAsync implementation of the Broadcaster interface
|
||||
// SendAsync implementation of the Broadcaster interface.
|
||||
func (g *GossipNodeSet) SendAsync(pb proto.Message) error {
|
||||
msg, err := pilosa.MarshalMessage(pb)
|
||||
if err != nil {
|
||||
|
|
@ -146,11 +149,13 @@ func (g *GossipNodeSet) SendAsync(pb proto.Message) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// implementation of the memberlist.Delegate interface
|
||||
// NodeMeta implementation of the memberlist.Delegate interface.
|
||||
func (g *GossipNodeSet) NodeMeta(limit int) []byte {
|
||||
return []byte{}
|
||||
}
|
||||
|
||||
// NotifyMsg implementation of the memberlist.Delegate interface
|
||||
// called when a user-data message is received.
|
||||
func (g *GossipNodeSet) NotifyMsg(b []byte) {
|
||||
m, err := pilosa.UnmarshalMessage(b)
|
||||
if err != nil {
|
||||
|
|
@ -163,10 +168,14 @@ func (g *GossipNodeSet) NotifyMsg(b []byte) {
|
|||
}
|
||||
}
|
||||
|
||||
// GetBroadcasts implementation of the memberlist.Delegate interface
|
||||
// called when user data messages can be broadcast.
|
||||
func (g *GossipNodeSet) GetBroadcasts(overhead, limit int) [][]byte {
|
||||
return g.broadcasts.GetBroadcasts(overhead, limit)
|
||||
}
|
||||
|
||||
// LocalState implementation of the memberlist.Delegate interface
|
||||
// sends this Node's state data.
|
||||
func (g *GossipNodeSet) LocalState(join bool) []byte {
|
||||
pb, err := g.statusHandler.LocalStatus()
|
||||
if err != nil {
|
||||
|
|
@ -183,6 +192,8 @@ func (g *GossipNodeSet) LocalState(join bool) []byte {
|
|||
return buf
|
||||
}
|
||||
|
||||
// MergeRemoteState implementation of the memberlist.Delegate interface
|
||||
// receive and process the remote side side's LocalState.
|
||||
func (g *GossipNodeSet) MergeRemoteState(buf []byte, join bool) {
|
||||
// Unmarshal nodestate data.
|
||||
var pb internal.NodeStatus
|
||||
|
|
|
|||
|
|
@ -58,6 +58,7 @@ func NewHandler() *Handler {
|
|||
return handler
|
||||
}
|
||||
|
||||
// NewRouter creates a Gorilla Mux http router.
|
||||
func NewRouter(handler *Handler) *mux.Router {
|
||||
router := mux.NewRouter()
|
||||
router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET")
|
||||
|
|
@ -1315,6 +1316,7 @@ type QueryResponse struct {
|
|||
Err error
|
||||
}
|
||||
|
||||
// MarshalJSON marshals QueryResponse into a JSON-encoded byte slice
|
||||
func (resp *QueryResponse) MarshalJSON() ([]byte, error) {
|
||||
var output struct {
|
||||
Results []interface{} `json:"results,omitempty"`
|
||||
|
|
|
|||
|
|
@ -356,7 +356,7 @@ type HolderSyncer struct {
|
|||
Closing <-chan struct{}
|
||||
}
|
||||
|
||||
// Returns true if the syncer has been marked to close.
|
||||
// IsClosing returns true if the syncer has been marked to close.
|
||||
func (s *HolderSyncer) IsClosing() bool {
|
||||
select {
|
||||
case <-s.Closing:
|
||||
|
|
|
|||
|
|
@ -62,11 +62,11 @@ func (h *HTTPBroadcaster) SendAsync(pb proto.Message) error {
|
|||
|
||||
func (h *HTTPBroadcaster) nodes() ([]*pilosa.Node, error) {
|
||||
if h.server == nil {
|
||||
return nil, errors.New("HTTPBroadcaster has no reference to Server.")
|
||||
return nil, errors.New("HTTPBroadcaster has no reference to Server")
|
||||
}
|
||||
nodeset, ok := h.server.Cluster.NodeSet.(*HTTPNodeSet)
|
||||
if !ok {
|
||||
return nil, errors.New("NodeSet cannot be caste to HTTPNodeSet.")
|
||||
return nil, errors.New("NodeSet cannot be caste to HTTPNodeSet")
|
||||
}
|
||||
return nodeset.Nodes(), nil
|
||||
}
|
||||
|
|
@ -106,12 +106,14 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// HTTPBroadcastReceiver unmarshals incoming messages over HTTP and passes them on to the handler.
|
||||
type HTTPBroadcastReceiver struct {
|
||||
port string
|
||||
handler pilosa.BroadcastHandler
|
||||
logOutput io.Writer
|
||||
}
|
||||
|
||||
// NewHTTPBroadcastReceiver returns a new instance of HTTPBroadcastReceiver.
|
||||
func NewHTTPBroadcastReceiver(port string, logOutput io.Writer) *HTTPBroadcastReceiver {
|
||||
return &HTTPBroadcastReceiver{
|
||||
port: port,
|
||||
|
|
@ -119,6 +121,7 @@ func NewHTTPBroadcastReceiver(port string, logOutput io.Writer) *HTTPBroadcastRe
|
|||
}
|
||||
}
|
||||
|
||||
// Start implements the BroadcastReceiver interface and starts listening for broadcast messages.
|
||||
func (rec *HTTPBroadcastReceiver) Start(b pilosa.BroadcastHandler) error {
|
||||
rec.handler = b
|
||||
go func() {
|
||||
|
|
@ -166,14 +169,17 @@ func NewHTTPNodeSet() *HTTPNodeSet {
|
|||
return &HTTPNodeSet{}
|
||||
}
|
||||
|
||||
// Nodes implements the NodeSet interface and returns a list of nodes in the cluster.
|
||||
func (h *HTTPNodeSet) Nodes() []*pilosa.Node {
|
||||
return h.nodes
|
||||
}
|
||||
|
||||
// Open implements the NodeSet interface to start network activity, but for a HTTPNodeSet it does nothing.
|
||||
func (h *HTTPNodeSet) Open() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Join sets the NodeSet nodes to the slice of Nodes passed in.
|
||||
func (h *HTTPNodeSet) Join(nodes []*pilosa.Node) error {
|
||||
h.nodes = nodes
|
||||
return nil
|
||||
|
|
|
|||
2
index.go
2
index.go
|
|
@ -250,6 +250,7 @@ func (i *Index) MaxSlice() uint64 {
|
|||
return max
|
||||
}
|
||||
|
||||
// SetRemoteMaxSlice sets the remote max slice value received from another node.
|
||||
func (i *Index) SetRemoteMaxSlice(newmax uint64) {
|
||||
i.mu.Lock()
|
||||
defer i.mu.Unlock()
|
||||
|
|
@ -273,6 +274,7 @@ func (i *Index) MaxInverseSlice() uint64 {
|
|||
return max
|
||||
}
|
||||
|
||||
// SetRemoteMaxInverseSlice sets the remote max inverse slice value received from another node.
|
||||
func (i *Index) SetRemoteMaxInverseSlice(v uint64) {
|
||||
i.mu.Lock()
|
||||
defer i.mu.Unlock()
|
||||
|
|
|
|||
|
|
@ -88,7 +88,7 @@ func decodeColumnAttrSet(pb *internal.ColumnAttrSet) *ColumnAttrSet {
|
|||
// TimeFormat is the go-style time format used to parse string dates.
|
||||
const TimeFormat = "2006-01-02T15:04"
|
||||
|
||||
// Restrict name using regex
|
||||
// ValidateName ensures that the name is a valid format.
|
||||
func ValidateName(name string) error {
|
||||
validName := nameRegexp.Match([]byte(name))
|
||||
if validName == false {
|
||||
|
|
|
|||
|
|
@ -281,13 +281,13 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// Server implements StatusHandler.
|
||||
// LocalStatus returns the state of the local node as well as the
|
||||
// holder (indexes/frames) according to the local node.
|
||||
// In a gossip implementation, memberlist.Delegate.LocalState() uses this.
|
||||
// Server implements StatusHandler.
|
||||
func (s *Server) LocalStatus() (proto.Message, error) {
|
||||
if s.Holder == nil {
|
||||
return nil, errors.New("Server.Holder is nil.")
|
||||
return nil, errors.New("Server.Holder is nil")
|
||||
}
|
||||
return &internal.NodeStatus{
|
||||
Host: s.Host,
|
||||
|
|
|
|||
4
stats.go
4
stats.go
|
|
@ -12,7 +12,7 @@ func init() {
|
|||
NopStatsClient = &nopStatsClient{}
|
||||
}
|
||||
|
||||
// Global expvar.
|
||||
// Expvar global expvar map.
|
||||
var Expvar = expvar.NewMap("index")
|
||||
|
||||
// StatsClient represents a client to a stats server.
|
||||
|
|
@ -39,9 +39,9 @@ type StatsClient interface {
|
|||
Timing(name string, value time.Duration)
|
||||
}
|
||||
|
||||
// NopStatsClient represents a client that doesn't do anything.
|
||||
var NopStatsClient StatsClient
|
||||
|
||||
// nopStatsClient represents a client that doesn't do anything.
|
||||
type nopStatsClient struct{}
|
||||
|
||||
func (c *nopStatsClient) Tags() []string { return nil }
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue