mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-11 23:31:03 +00:00
Merge pull request #1421 from jaffee/remove-
remove some unused interfaces and implementations from broadcast.go
This commit is contained in:
commit
f59452a7e7
8 changed files with 85 additions and 165 deletions
2
api.go
2
api.go
|
|
@ -699,7 +699,7 @@ func (api *API) LongQueryTime() time.Duration {
|
|||
if api.Cluster == nil {
|
||||
return 0
|
||||
}
|
||||
return api.Cluster.LongQueryTime
|
||||
return api.Cluster.longQueryTime
|
||||
}
|
||||
|
||||
func (api *API) indexField(indexName string, fieldName string, slice uint64) (*Index, *Field, error) {
|
||||
|
|
|
|||
58
broadcast.go
58
broadcast.go
|
|
@ -23,30 +23,6 @@ import (
|
|||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
// MemberSet represents an interface for Node membership and inter-node communication.
|
||||
type MemberSet interface {
|
||||
// Open starts any network activity implemented by the MemberSet
|
||||
// Node is the local node, used for membership broadcasts.
|
||||
Open() error
|
||||
}
|
||||
|
||||
// StaticMemberSet represents a basic MemberSet for testing.
|
||||
type StaticMemberSet struct {
|
||||
nodes []*Node
|
||||
}
|
||||
|
||||
// NewStaticMemberSet creates a statically defined MemberSet.
|
||||
func NewStaticMemberSet(nodes []*Node) *StaticMemberSet {
|
||||
return &StaticMemberSet{
|
||||
nodes: nodes,
|
||||
}
|
||||
}
|
||||
|
||||
// Open implements the MemberSet interface to start network activity, but for a static MemberSet it does nothing.
|
||||
func (s *StaticMemberSet) Open() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Broadcaster is an interface for broadcasting messages.
|
||||
type Broadcaster interface {
|
||||
SendSync(pb proto.Message) error
|
||||
|
|
@ -56,7 +32,6 @@ type Broadcaster interface {
|
|||
|
||||
func init() {
|
||||
NopBroadcaster = &nopBroadcaster{}
|
||||
NopGossiper = &nopGossiper{}
|
||||
}
|
||||
|
||||
// NopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
|
|
@ -85,39 +60,6 @@ type BroadcastHandler interface {
|
|||
ReceiveMessage(pb proto.Message) error
|
||||
}
|
||||
|
||||
// BroadcastReceiver is the interface for the object which will listen for and
|
||||
// decode broadcast messages before passing them to pilosa to handle. The
|
||||
// implementation of this could be an http server which listens for messages,
|
||||
// gets the protobuf payload, and then passes it to
|
||||
// BroadcastHandler.ReceiveMessage.
|
||||
type BroadcastReceiver interface {
|
||||
// Start starts listening for broadcast messages - it should return
|
||||
// immediately, spawning a goroutine if necessary.
|
||||
Start(BroadcastHandler) error
|
||||
}
|
||||
|
||||
type nopBroadcastReceiver struct{}
|
||||
|
||||
func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil }
|
||||
|
||||
// NopBroadcastReceiver is a no-op implementation of the BroadcastReceiver.
|
||||
var NopBroadcastReceiver = &nopBroadcastReceiver{}
|
||||
|
||||
// Gossiper is an interface for sharing messages via gossip.
|
||||
type Gossiper interface {
|
||||
SendAsync(pb proto.Message) error
|
||||
}
|
||||
|
||||
// NopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
var NopGossiper Gossiper
|
||||
|
||||
type nopGossiper struct{}
|
||||
|
||||
// SendAsync A no-op implementation of Gossiper SendAsync method.
|
||||
func (n *nopGossiper) SendAsync(pb proto.Message) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Broadcast message types.
|
||||
const (
|
||||
messageTypeCreateSlice = iota
|
||||
|
|
|
|||
132
cluster.go
132
cluster.go
|
|
@ -212,7 +212,7 @@ type nodeAction struct {
|
|||
|
||||
// Cluster represents a collection of nodes.
|
||||
type Cluster struct {
|
||||
ID string
|
||||
id string
|
||||
Node *Node
|
||||
Nodes []*Node // TODO phase this out?
|
||||
|
||||
|
|
@ -220,19 +220,16 @@ type Cluster struct {
|
|||
Hasher Hasher
|
||||
|
||||
// The number of partitions in the cluster.
|
||||
PartitionN int
|
||||
partitionN int
|
||||
|
||||
// The number of replicas a partition has.
|
||||
ReplicaN int
|
||||
|
||||
// Threshold for logging long-running queries
|
||||
LongQueryTime time.Duration
|
||||
longQueryTime time.Duration
|
||||
|
||||
// Maximum number of Set() or Clear() commands per request.
|
||||
MaxWritesPerRequest int
|
||||
|
||||
// EventReceiver receives NodeEvents pertaining to node membership.
|
||||
EventReceiver EventReceiver
|
||||
maxWritesPerRequest int
|
||||
|
||||
// Data directory path.
|
||||
Path string
|
||||
|
|
@ -242,8 +239,8 @@ type Cluster struct {
|
|||
Static bool // Static is primarily used for testing in a non-gossip environment.
|
||||
state string
|
||||
Coordinator string
|
||||
Holder *Holder
|
||||
Broadcaster Broadcaster
|
||||
holder *Holder
|
||||
broadcaster Broadcaster
|
||||
|
||||
joiningLeavingNodes chan nodeAction
|
||||
|
||||
|
|
@ -260,7 +257,7 @@ type Cluster struct {
|
|||
wg sync.WaitGroup
|
||||
closing chan struct{}
|
||||
|
||||
Logger Logger
|
||||
logger Logger
|
||||
|
||||
InternalClient InternalClient
|
||||
}
|
||||
|
|
@ -268,10 +265,9 @@ type Cluster struct {
|
|||
// NewCluster returns a new instance of Cluster with defaults.
|
||||
func NewCluster() *Cluster {
|
||||
return &Cluster{
|
||||
Hasher: &jmphasher{},
|
||||
PartitionN: DefaultPartitionN,
|
||||
ReplicaN: 1,
|
||||
EventReceiver: NopEventReceiver,
|
||||
Hasher: &jmphasher{},
|
||||
partitionN: DefaultPartitionN,
|
||||
ReplicaN: 1,
|
||||
|
||||
joiningLeavingNodes: make(chan nodeAction, 10), // buffered channel
|
||||
jobs: make(map[int64]*resizeJob),
|
||||
|
|
@ -280,7 +276,7 @@ func NewCluster() *Cluster {
|
|||
|
||||
InternalClient: NewNopInternalClient(),
|
||||
|
||||
Logger: NopLogger,
|
||||
logger: NopLogger,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -317,7 +313,7 @@ func (c *Cluster) setCoordinator(n *Node) error {
|
|||
_ = c.unprotectedUpdateCoordinator(n)
|
||||
c.mu.Unlock()
|
||||
// Send the update coordinator message to all nodes.
|
||||
err := c.Broadcaster.SendSync(
|
||||
err := c.broadcaster.SendSync(
|
||||
&internal.UpdateCoordinatorMessage{
|
||||
New: EncodeNode(n),
|
||||
})
|
||||
|
|
@ -326,7 +322,7 @@ func (c *Cluster) setCoordinator(n *Node) error {
|
|||
}
|
||||
|
||||
// Broadcast cluster status.
|
||||
return c.Broadcaster.SendSync(c.Status())
|
||||
return c.broadcaster.SendSync(c.Status())
|
||||
}
|
||||
|
||||
// updateCoordinator updates this nodes Coordinator value as well as
|
||||
|
|
@ -358,7 +354,7 @@ func (c *Cluster) unprotectedUpdateCoordinator(n *Node) bool {
|
|||
// addNode adds a node to the Cluster and updates and saves the
|
||||
// new topology.
|
||||
func (c *Cluster) addNode(node *Node) error {
|
||||
c.Logger.Printf("add node %s to cluster on %s", node, c.Node)
|
||||
c.logger.Printf("add node %s to cluster on %s", node, c.Node)
|
||||
|
||||
// If the node being added is the coordinator, set it for this node.
|
||||
if node.IsCoordinator {
|
||||
|
|
@ -409,13 +405,13 @@ func (c *Cluster) nodeIDs() []string {
|
|||
|
||||
func (c *Cluster) setID(id string) {
|
||||
// Don't overwrite ClusterID.
|
||||
if c.ID != "" {
|
||||
if c.id != "" {
|
||||
return
|
||||
}
|
||||
c.ID = id
|
||||
c.id = id
|
||||
|
||||
// Make sure the Topology is updated.
|
||||
c.Topology.ClusterID = c.ID
|
||||
c.Topology.ClusterID = c.id
|
||||
}
|
||||
|
||||
func (c *Cluster) State() string {
|
||||
|
|
@ -436,7 +432,7 @@ func (c *Cluster) setState(state string) {
|
|||
return
|
||||
}
|
||||
|
||||
c.Logger.Printf("change cluster state from %s to %s on %s", c.state, state, c.Node.ID)
|
||||
c.logger.Printf("change cluster state from %s to %s on %s", c.state, state, c.Node.ID)
|
||||
|
||||
var doCleanup bool
|
||||
|
||||
|
|
@ -456,13 +452,13 @@ func (c *Cluster) setState(state string) {
|
|||
if doCleanup {
|
||||
var cleaner HolderCleaner
|
||||
cleaner.Node = c.Node
|
||||
cleaner.Holder = c.Holder
|
||||
cleaner.Holder = c.holder
|
||||
cleaner.Cluster = c
|
||||
cleaner.Closing = c.closing
|
||||
|
||||
// Clean holder.
|
||||
if err := cleaner.CleanHolder(); err != nil {
|
||||
c.Logger.Printf("holder clean error: err=%s", err)
|
||||
c.logger.Printf("holder clean error: err=%s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -478,7 +474,7 @@ func (c *Cluster) setNodeState(state string) error {
|
|||
State: state,
|
||||
}
|
||||
|
||||
c.Logger.Printf("Sending State %s (%s)", state, c.Coordinator)
|
||||
c.logger.Printf("Sending State %s (%s)", state, c.Coordinator)
|
||||
if err := c.sendTo(c.coordinatorNode(), ns); err != nil {
|
||||
return fmt.Errorf("sending node state error: err=%s", err)
|
||||
}
|
||||
|
|
@ -500,7 +496,7 @@ func (c *Cluster) receiveNodeState(nodeID string, state string) error {
|
|||
}
|
||||
|
||||
c.Topology.nodeStates[nodeID] = state
|
||||
c.Logger.Printf("received state %s (%s)", state, nodeID)
|
||||
c.logger.Printf("received state %s (%s)", state, nodeID)
|
||||
|
||||
// Set cluster state to NORMAL.
|
||||
if c.haveTopologyAgreement() && c.allNodesReady() {
|
||||
|
|
@ -513,7 +509,7 @@ func (c *Cluster) receiveNodeState(nodeID string, state string) error {
|
|||
// Status returns the internal ClusterStatus representation.
|
||||
func (c *Cluster) Status() *internal.ClusterStatus {
|
||||
return &internal.ClusterStatus{
|
||||
ClusterID: c.ID,
|
||||
ClusterID: c.id,
|
||||
State: c.state,
|
||||
Nodes: EncodeNodes(c.Nodes),
|
||||
}
|
||||
|
|
@ -715,7 +711,7 @@ func (c *Cluster) fragSources(to *Cluster, idx *Index) (map[string][]*internal.R
|
|||
srcCluster = NewCluster()
|
||||
srcCluster.Nodes = Nodes(c.Nodes).Clone()
|
||||
srcCluster.Hasher = c.Hasher
|
||||
srcCluster.PartitionN = c.PartitionN
|
||||
srcCluster.partitionN = c.partitionN
|
||||
srcCluster.ReplicaN = 1
|
||||
}
|
||||
|
||||
|
|
@ -785,7 +781,7 @@ func (c *Cluster) partition(index string, slice uint64) int {
|
|||
h := fnv.New64a()
|
||||
h.Write([]byte(index))
|
||||
h.Write(buf[:])
|
||||
return int(h.Sum64() % uint64(c.PartitionN))
|
||||
return int(h.Sum64() % uint64(c.partitionN))
|
||||
}
|
||||
|
||||
// sliceNodes returns a list of nodes that own a fragment.
|
||||
|
|
@ -869,7 +865,7 @@ func (c *Cluster) setup() error {
|
|||
return errors.Wrap(err, "loading topology")
|
||||
}
|
||||
|
||||
c.ID = c.Topology.ClusterID
|
||||
c.id = c.Topology.ClusterID
|
||||
|
||||
// Only the coordinator needs to consider the .topology file.
|
||||
if c.isCoordinator() {
|
||||
|
|
@ -910,13 +906,13 @@ func (c *Cluster) waitForStarted() error {
|
|||
Event: uint32(NodeJoin),
|
||||
Node: EncodeNode(c.Node),
|
||||
}
|
||||
if err := c.Broadcaster.SendSync(msg); err != nil {
|
||||
if err := c.broadcaster.SendSync(msg); err != nil {
|
||||
return fmt.Errorf("sending restart NodeJoin: %v", err)
|
||||
}
|
||||
|
||||
c.Logger.Printf("%v wait for joining to complete", c.Node.ID)
|
||||
c.logger.Printf("%v wait for joining to complete", c.Node.ID)
|
||||
<-c.joining
|
||||
c.Logger.Printf("joining has completed")
|
||||
c.logger.Printf("joining has completed")
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
@ -931,7 +927,7 @@ func (c *Cluster) close() error {
|
|||
}
|
||||
|
||||
func (c *Cluster) markAsJoined() {
|
||||
c.Logger.Printf("mark node as joined (received coordinator update)")
|
||||
c.logger.Printf("mark node as joined (received coordinator update)")
|
||||
if !c.joined {
|
||||
c.joined = true
|
||||
close(c.joining)
|
||||
|
|
@ -964,9 +960,9 @@ func (c *Cluster) allNodesReady() bool {
|
|||
func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
|
||||
j, err := c.generateResizeJob(nodeAction)
|
||||
if err != nil {
|
||||
c.Logger.Printf("generateResizeJob error: err=%s", err)
|
||||
c.logger.Printf("generateResizeJob error: err=%s", err)
|
||||
if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil {
|
||||
c.Logger.Printf("setStateAndBroadcast error: err=%s", err)
|
||||
c.logger.Printf("setStateAndBroadcast error: err=%s", err)
|
||||
}
|
||||
return errors.Wrap(err, "setting state")
|
||||
}
|
||||
|
|
@ -980,7 +976,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
|
|||
})
|
||||
|
||||
// Wait for the resizeJob to finish or be aborted.
|
||||
c.Logger.Printf("wait for jobResult")
|
||||
c.logger.Printf("wait for jobResult")
|
||||
jobResult := <-j.result
|
||||
|
||||
// Make sure j.Run() didn't return an error.
|
||||
|
|
@ -988,7 +984,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
|
|||
return errors.Wrap(err, "running job")
|
||||
}
|
||||
|
||||
c.Logger.Printf("received jobResult: %s", jobResult)
|
||||
c.logger.Printf("received jobResult: %s", jobResult)
|
||||
switch jobResult {
|
||||
case resizeJobStateDone:
|
||||
if err := c.completeCurrentJob(resizeJobStateDone); err != nil {
|
||||
|
|
@ -1014,12 +1010,12 @@ func (c *Cluster) setStateAndBroadcast(state string) error {
|
|||
return nil
|
||||
}
|
||||
// Broadcast cluster status changes to the cluster.
|
||||
c.Logger.Printf("broadcasting ClusterStatus: %s", state)
|
||||
return c.Broadcaster.SendSync(c.Status())
|
||||
c.logger.Printf("broadcasting ClusterStatus: %s", state)
|
||||
return c.broadcaster.SendSync(c.Status())
|
||||
}
|
||||
|
||||
func (c *Cluster) sendTo(node *Node, msg proto.Message) error {
|
||||
if err := c.Broadcaster.SendTo(node, msg); err != nil {
|
||||
if err := c.broadcaster.SendTo(node, msg); err != nil {
|
||||
return errors.Wrap(err, "sending")
|
||||
}
|
||||
return nil
|
||||
|
|
@ -1045,7 +1041,7 @@ func (c *Cluster) listenForJoins() {
|
|||
case nodeAction := <-c.joiningLeavingNodes:
|
||||
err := c.handleNodeAction(nodeAction)
|
||||
if err != nil {
|
||||
c.Logger.Printf("handleNodeAction error: err=%s", err)
|
||||
c.logger.Printf("handleNodeAction error: err=%s", err)
|
||||
continue
|
||||
}
|
||||
setNormal = true
|
||||
|
|
@ -1057,7 +1053,7 @@ func (c *Cluster) listenForJoins() {
|
|||
if setNormal {
|
||||
// Put the cluster back to state NORMAL and broadcast.
|
||||
if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil {
|
||||
c.Logger.Printf("setStateAndBroadcast error: err=%s", err)
|
||||
c.logger.Printf("setStateAndBroadcast error: err=%s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1068,7 +1064,7 @@ func (c *Cluster) listenForJoins() {
|
|||
case nodeAction := <-c.joiningLeavingNodes:
|
||||
err := c.handleNodeAction(nodeAction)
|
||||
if err != nil {
|
||||
c.Logger.Printf("handleNodeAction error: err=%s", err)
|
||||
c.logger.Printf("handleNodeAction error: err=%s", err)
|
||||
continue
|
||||
}
|
||||
setNormal = true
|
||||
|
|
@ -1082,7 +1078,7 @@ func (c *Cluster) listenForJoins() {
|
|||
// added/removed. It also saves a reference to the resizeJob in the `jobs` map
|
||||
// for future lookup by JobID.
|
||||
func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) {
|
||||
c.Logger.Printf("generateResizeJob: %v", nodeAction)
|
||||
c.logger.Printf("generateResizeJob: %v", nodeAction)
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
|
|
@ -1090,7 +1086,7 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) {
|
|||
if err != nil {
|
||||
return nil, errors.Wrap(err, "generating job")
|
||||
}
|
||||
c.Logger.Printf("generated resizeJob: %d", j.ID)
|
||||
c.logger.Printf("generated resizeJob: %d", j.ID)
|
||||
|
||||
// Save job in jobs map for future reference.
|
||||
c.jobs[j.ID] = j
|
||||
|
|
@ -1110,13 +1106,13 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) {
|
|||
// the resize instructions to other nodes in the cluster.
|
||||
func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, error) {
|
||||
j := newResizeJob(c.Nodes, nodeAction.node, nodeAction.action)
|
||||
j.Broadcaster = c.Broadcaster
|
||||
j.Broadcaster = c.broadcaster
|
||||
|
||||
// toCluster is a clone of Cluster with the new node added/removed for comparison.
|
||||
toCluster := NewCluster()
|
||||
toCluster.Nodes = Nodes(c.Nodes).Clone()
|
||||
toCluster.Hasher = c.Hasher
|
||||
toCluster.PartitionN = c.PartitionN
|
||||
toCluster.partitionN = c.partitionN
|
||||
toCluster.ReplicaN = c.ReplicaN
|
||||
if nodeAction.action == resizeJobActionRemove {
|
||||
toCluster.removeNodeBasicSorted(nodeAction.node)
|
||||
|
|
@ -1132,7 +1128,7 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob,
|
|||
}
|
||||
|
||||
// Add to multiIndex the instructions for each index.
|
||||
for _, idx := range c.Holder.Indexes() {
|
||||
for _, idx := range c.holder.Indexes() {
|
||||
fragSources, err := c.fragSources(toCluster, idx)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting sources")
|
||||
|
|
@ -1154,7 +1150,7 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob,
|
|||
Node: EncodeNode(toCluster.unprotectedNodeByID(id)),
|
||||
Coordinator: EncodeNode(c.coordinatorNode()),
|
||||
Sources: sources,
|
||||
Schema: c.Holder.EncodeSchema(), // Include the schema to ensure it's in sync on the receiving node.
|
||||
Schema: c.holder.EncodeSchema(), // Include the schema to ensure it's in sync on the receiving node.
|
||||
ClusterStatus: c.Status(),
|
||||
}
|
||||
j.Instructions = append(j.Instructions, instr)
|
||||
|
|
@ -1181,21 +1177,21 @@ func (c *Cluster) completeCurrentJob(state string) error {
|
|||
|
||||
// followResizeInstruction is run by any node that receives a ResizeInstruction.
|
||||
func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) error {
|
||||
c.Logger.Printf("follow resize instruction on %s", c.Node.ID)
|
||||
c.logger.Printf("follow resize instruction on %s", c.Node.ID)
|
||||
// Make sure the cluster status on this node agrees with the Coordinator
|
||||
// before attempting a resize.
|
||||
if err := c.mergeClusterStatus(instr.ClusterStatus); err != nil {
|
||||
return errors.Wrap(err, "merging cluster status")
|
||||
}
|
||||
|
||||
c.Logger.Printf("MergeClusterStatus done, start goroutine")
|
||||
c.logger.Printf("MergeClusterStatus done, start goroutine")
|
||||
|
||||
// The actual resizing runs in a goroutine because we don't want to block
|
||||
// the distribution of other ResizeInstructions to the rest of the cluster.
|
||||
go func() {
|
||||
|
||||
// Make sure the holder has opened.
|
||||
<-c.Holder.opened
|
||||
<-c.holder.opened
|
||||
|
||||
// Prepare the return message.
|
||||
complete := &internal.ResizeInstructionComplete{
|
||||
|
|
@ -1208,19 +1204,19 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err
|
|||
if err := func() error {
|
||||
|
||||
// Sync the schema received in the resize instruction.
|
||||
c.Logger.Printf("Holder ApplySchema")
|
||||
if err := c.Holder.ApplySchema(instr.Schema); err != nil {
|
||||
c.logger.Printf("Holder ApplySchema")
|
||||
if err := c.holder.ApplySchema(instr.Schema); err != nil {
|
||||
return errors.Wrap(err, "applying schema")
|
||||
}
|
||||
|
||||
// Request each source file in ResizeSources.
|
||||
for _, src := range instr.Sources {
|
||||
c.Logger.Printf("get slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
|
||||
c.logger.Printf("get slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
|
||||
|
||||
srcURI := decodeURI(src.Node.URI)
|
||||
|
||||
// Retrieve field.
|
||||
f := c.Holder.Field(src.Index, src.Field)
|
||||
f := c.holder.Field(src.Index, src.Field)
|
||||
if f == nil {
|
||||
return ErrFieldNotFound
|
||||
}
|
||||
|
|
@ -1238,7 +1234,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err
|
|||
}
|
||||
|
||||
// Stream slice from remote node.
|
||||
c.Logger.Printf("retrieve slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
|
||||
c.logger.Printf("retrieve slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
|
||||
rd, err := c.InternalClient.RetrieveSliceFromURI(context.Background(), src.Index, src.Field, src.Slice, srcURI)
|
||||
if err != nil {
|
||||
// For now it is an acceptable error if the fragment is not found
|
||||
|
|
@ -1270,7 +1266,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err
|
|||
}
|
||||
|
||||
if err := c.sendTo(DecodeNode(instr.Coordinator), complete); err != nil {
|
||||
c.Logger.Printf("sending resizeInstructionComplete error: err=%s", err)
|
||||
c.logger.Printf("sending resizeInstructionComplete error: err=%s", err)
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
|
|
@ -1585,10 +1581,10 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) {
|
|||
|
||||
func (c *Cluster) considerTopology() error {
|
||||
// Create ClusterID if one does not already exist.
|
||||
if c.ID == "" {
|
||||
if c.id == "" {
|
||||
u := uuid.NewV4()
|
||||
c.ID = u.String()
|
||||
c.Topology.ClusterID = c.ID
|
||||
c.id = u.String()
|
||||
c.Topology.ClusterID = c.id
|
||||
}
|
||||
|
||||
if c.Static {
|
||||
|
|
@ -1624,7 +1620,7 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error {
|
|||
|
||||
switch e.Event {
|
||||
case NodeJoin:
|
||||
c.Logger.Printf("received NodeJoin event: %v", e)
|
||||
c.logger.Printf("received NodeJoin event: %v", e)
|
||||
// Ignore the event if this is not the coordinator.
|
||||
if !c.isCoordinator() {
|
||||
return nil
|
||||
|
|
@ -1644,7 +1640,7 @@ func (c *Cluster) nodeJoin(node *Node) error {
|
|||
// A host that is not part of the topology can't be added to the STARTING cluster.
|
||||
if !c.Topology.ContainsID(node.ID) {
|
||||
err := fmt.Sprintf("host is not in topology: %s", node.ID)
|
||||
c.Logger.Printf("%v", err)
|
||||
c.logger.Printf("%v", err)
|
||||
return errors.New(err)
|
||||
}
|
||||
|
||||
|
|
@ -1655,7 +1651,7 @@ func (c *Cluster) nodeJoin(node *Node) error {
|
|||
// Only change to normal if there is no existing data. Otherwise,
|
||||
// the coordinator needs to wait to receive READY messages (nodeStates)
|
||||
// from remote nodes before setting the cluster to state NORMAL.
|
||||
if ok, err := c.Holder.HasData(); !ok && err == nil {
|
||||
if ok, err := c.holder.HasData(); !ok && err == nil {
|
||||
// If the result of the previous AddNode completed the joining of nodes
|
||||
// in the topology, then change the state to NORMAL.
|
||||
if c.haveTopologyAgreement() {
|
||||
|
|
@ -1683,7 +1679,7 @@ func (c *Cluster) nodeJoin(node *Node) error {
|
|||
}
|
||||
|
||||
// If the holder does not yet contain data, go ahead and add the node.
|
||||
if ok, err := c.Holder.HasData(); !ok && err == nil {
|
||||
if ok, err := c.holder.HasData(); !ok && err == nil {
|
||||
if err := c.addNode(node); err != nil {
|
||||
return errors.Wrap(err, "adding node")
|
||||
}
|
||||
|
|
@ -1737,7 +1733,7 @@ func (c *Cluster) nodeLeave(node *Node) error {
|
|||
}
|
||||
|
||||
// If the holder does not yet contain data, go ahead and remove the node.
|
||||
if ok, err := c.Holder.HasData(); !ok && err == nil {
|
||||
if ok, err := c.holder.HasData(); !ok && err == nil {
|
||||
if err := c.removeNode(n); err != nil {
|
||||
return errors.Wrap(err, "removing node")
|
||||
}
|
||||
|
|
@ -1759,7 +1755,7 @@ func (c *Cluster) nodeLeave(node *Node) error {
|
|||
func (c *Cluster) mergeClusterStatus(cs *internal.ClusterStatus) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
c.Logger.Printf("merge cluster status: %v", cs)
|
||||
c.logger.Printf("merge cluster status: %v", cs)
|
||||
// Ignore status updates from self (coordinator).
|
||||
if c.unprotectedIsCoordinator() {
|
||||
return nil
|
||||
|
|
|
|||
|
|
@ -341,7 +341,7 @@ func TestCluster_Owners(t *testing.T) {
|
|||
func TestCluster_Partition(t *testing.T) {
|
||||
if err := quick.Check(func(index string, slice uint64, partitionN int) bool {
|
||||
c := NewCluster()
|
||||
c.PartitionN = partitionN
|
||||
c.partitionN = partitionN
|
||||
|
||||
partitionID := c.partition(index, slice)
|
||||
if partitionID < 0 || partitionID >= partitionN {
|
||||
|
|
@ -705,7 +705,7 @@ func TestCluster_ResizeStates(t *testing.T) {
|
|||
|
||||
// Before starting the resize, get the CheckSum to use for
|
||||
// comparison later.
|
||||
node0Field := node0.Holder.Field("i", "f")
|
||||
node0Field := node0.holder.Field("i", "f")
|
||||
node0View := node0Field.View("standard")
|
||||
node0Fragment := node0View.Fragment(1)
|
||||
node0Checksum := node0Fragment.Checksum()
|
||||
|
|
@ -734,7 +734,7 @@ func TestCluster_ResizeStates(t *testing.T) {
|
|||
|
||||
// Bits
|
||||
// Verify that node-1 contains the fragment (i/f/standard/1) transferred from node-0.
|
||||
node1Field := node1.Holder.Field("i", "f")
|
||||
node1Field := node1.holder.Field("i", "f")
|
||||
node1View := node1Field.View("standard")
|
||||
node1Fragment := node1View.Fragment(1)
|
||||
|
||||
|
|
|
|||
18
event.go
18
event.go
|
|
@ -35,21 +35,3 @@ type NodeEvent struct {
|
|||
type EventHandler interface {
|
||||
ReceiveEvent(e *NodeEvent) error
|
||||
}
|
||||
|
||||
// EventReceiver is the interface for the object which will listen for and
|
||||
// decode broadcast messages before passing them to pilosa to handle. The
|
||||
// implementation of this could be an http server which listens for messages,
|
||||
// gets the protobuf payload, and then passes it to
|
||||
// EventHandler.ReceiveMessage.
|
||||
type EventReceiver interface {
|
||||
// Start starts listening for broadcast messages - it should return
|
||||
// immediately, spawning a goroutine if necessary.
|
||||
Start(EventHandler) error
|
||||
}
|
||||
|
||||
type nopEventReceiver struct{}
|
||||
|
||||
func (n *nopEventReceiver) Start(e EventHandler) error { return nil }
|
||||
|
||||
// NopEventReceiver is a no-op implementation of the EventReceiver.
|
||||
var NopEventReceiver = &nopEventReceiver{}
|
||||
|
|
|
|||
|
|
@ -1865,7 +1865,7 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
|
||||
// Generate query with sets & clears, and group the requests to not exceed MaxWritesPerRequest.
|
||||
total := len(set.columnIDs) + len(clear.columnIDs)
|
||||
maxWrites := s.Cluster.MaxWritesPerRequest
|
||||
maxWrites := s.Cluster.maxWritesPerRequest
|
||||
if maxWrites <= 0 {
|
||||
maxWrites = 5000
|
||||
}
|
||||
|
|
|
|||
12
server.go
12
server.go
|
|
@ -123,7 +123,7 @@ func OptServerAntiEntropyInterval(interval time.Duration) ServerOption {
|
|||
|
||||
func OptServerLongQueryTime(dur time.Duration) ServerOption {
|
||||
return func(s *Server) error {
|
||||
s.cluster.LongQueryTime = dur
|
||||
s.cluster.longQueryTime = dur
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
|
@ -246,8 +246,8 @@ func NewServer(opts ...ServerOption) (*Server, error) {
|
|||
s.holder.Stats.SetLogger(s.logger)
|
||||
|
||||
s.cluster.Path = path
|
||||
s.cluster.Logger = s.logger
|
||||
s.cluster.Holder = s.holder
|
||||
s.cluster.logger = s.logger
|
||||
s.cluster.holder = s.holder
|
||||
|
||||
// Initialize translation database.
|
||||
s.translateFile = NewTranslateFile()
|
||||
|
|
@ -282,8 +282,8 @@ func NewServer(opts ...ServerOption) (*Server, error) {
|
|||
s.executor.Cluster = s.cluster
|
||||
s.executor.TranslateStore = s.translateFile
|
||||
s.executor.MaxWritesPerRequest = s.maxWritesPerRequest
|
||||
s.cluster.Broadcaster = s
|
||||
s.cluster.MaxWritesPerRequest = s.maxWritesPerRequest
|
||||
s.cluster.broadcaster = s
|
||||
s.cluster.maxWritesPerRequest = s.maxWritesPerRequest
|
||||
s.holder.Broadcaster = s
|
||||
|
||||
err = s.cluster.setup()
|
||||
|
|
@ -629,7 +629,7 @@ func (s *Server) monitorDiagnostics() {
|
|||
s.diagnostics.Set("NumNodes", len(s.cluster.Nodes))
|
||||
s.diagnostics.Set("NumCPU", runtime.NumCPU())
|
||||
s.diagnostics.Set("NodeID", s.nodeID)
|
||||
s.diagnostics.Set("ClusterID", s.cluster.ID)
|
||||
s.diagnostics.Set("ClusterID", s.cluster.id)
|
||||
s.diagnostics.EnrichWithOSInfo()
|
||||
|
||||
// Flush the diagnostics metrics at startup, then on each tick interval
|
||||
|
|
|
|||
|
|
@ -97,7 +97,7 @@ type commonClusterSettings struct {
|
|||
|
||||
func (t *ClusterCluster) CreateIndex(name string) error {
|
||||
for _, c := range t.Clusters {
|
||||
if _, err := c.Holder.CreateIndexIfNotExists(name, IndexOptions{}); err != nil {
|
||||
if _, err := c.holder.CreateIndexIfNotExists(name, IndexOptions{}); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
|
@ -106,7 +106,7 @@ func (t *ClusterCluster) CreateIndex(name string) error {
|
|||
|
||||
func (t *ClusterCluster) CreateField(index, field string, opt FieldOptions) error {
|
||||
for _, c := range t.Clusters {
|
||||
idx, err := c.Holder.CreateIndexIfNotExists(index, IndexOptions{})
|
||||
idx, err := c.holder.CreateIndexIfNotExists(index, IndexOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -128,7 +128,7 @@ func (t *ClusterCluster) SetBit(index, field string, rowID, colID uint64, x *tim
|
|||
if c == nil {
|
||||
continue
|
||||
}
|
||||
f := c.Holder.Field(index, field)
|
||||
f := c.holder.Field(index, field)
|
||||
if f == nil {
|
||||
return fmt.Errorf("index/field does not exist: %s/%s", index, field)
|
||||
}
|
||||
|
|
@ -227,10 +227,10 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*Cluster, error)
|
|||
c.Hasher = NewTestModHasher()
|
||||
c.Path = path
|
||||
c.Topology = NewTopology()
|
||||
c.Holder = h
|
||||
c.holder = h
|
||||
c.Node = node
|
||||
c.Coordinator = t.common.Nodes[0].ID // the first node is the coordinator
|
||||
c.Broadcaster = t
|
||||
c.broadcaster = t
|
||||
|
||||
// add nodes
|
||||
if saveTopology {
|
||||
|
|
@ -275,7 +275,7 @@ func (t *ClusterCluster) Open() error {
|
|||
if err := c.open(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := c.Holder.Open(); err != nil {
|
||||
if err := c.holder.Open(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := c.setNodeState(NodeStateReady); err != nil {
|
||||
|
|
@ -360,7 +360,7 @@ func (t *ClusterCluster) FollowResizeInstruction(instr *internal.ResizeInstructi
|
|||
destCluster := t.clusterByID(instrNode.ID)
|
||||
|
||||
// Sync the schema received in the resize instruction.
|
||||
if err := destCluster.Holder.ApplySchema(instr.Schema); err != nil {
|
||||
if err := destCluster.holder.ApplySchema(instr.Schema); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
|
@ -368,11 +368,11 @@ func (t *ClusterCluster) FollowResizeInstruction(instr *internal.ResizeInstructi
|
|||
srcNode := DecodeNode(src.Node)
|
||||
srcCluster := t.clusterByID(srcNode.ID)
|
||||
|
||||
srcFragment := srcCluster.Holder.Fragment(src.Index, src.Field, src.View, src.Slice)
|
||||
destFragment := destCluster.Holder.Fragment(src.Index, src.Field, src.View, src.Slice)
|
||||
srcFragment := srcCluster.holder.Fragment(src.Index, src.Field, src.View, src.Slice)
|
||||
destFragment := destCluster.holder.Fragment(src.Index, src.Field, src.View, src.Slice)
|
||||
if destFragment == nil {
|
||||
// Create fragment on destination if it doesn't exist.
|
||||
f := destCluster.Holder.Field(src.Index, src.Field)
|
||||
f := destCluster.holder.Field(src.Index, src.Field)
|
||||
v := f.View(src.View)
|
||||
var err error
|
||||
destFragment, err = v.CreateFragmentIfNotExists(src.Slice)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue