From 5e48023565f2856485273dc91ce1a779c4d58a47 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 27 Jun 2018 14:48:07 -0500 Subject: [PATCH 1/3] remove some unused interfaces and implementations from broadcast.go --- broadcast.go | 58 ---------------------------------------------------- 1 file changed, 58 deletions(-) diff --git a/broadcast.go b/broadcast.go index d102fda22..b7b13fe04 100644 --- a/broadcast.go +++ b/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 From d6029e64fc46a5de6a2492aa82dc6f46aa01f3f5 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 27 Jun 2018 15:33:53 -0500 Subject: [PATCH 2/3] remove EventReceiver - not used anymore --- cluster.go | 10 +++------- event.go | 18 ------------------ 2 files changed, 3 insertions(+), 25 deletions(-) diff --git a/cluster.go b/cluster.go index e10aae84e..e1ea2ada7 100644 --- a/cluster.go +++ b/cluster.go @@ -231,9 +231,6 @@ type Cluster struct { // Maximum number of Set() or Clear() commands per request. MaxWritesPerRequest int - // EventReceiver receives NodeEvents pertaining to node membership. - EventReceiver EventReceiver - // Data directory path. Path string Topology *Topology @@ -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), diff --git a/event.go b/event.go index 5df69361b..69229b1e4 100644 --- a/event.go +++ b/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{} From 648cfe7ad0be05b91871d60fbadf98f7bea47b03 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 27 Jun 2018 15:39:12 -0500 Subject: [PATCH 3/3] unexport Cluster fields which could be automatically unexported --- api.go | 2 +- cluster.go | 124 +++++++++++++++++++-------------------- cluster_internal_test.go | 6 +- fragment.go | 2 +- server.go | 12 ++-- utils_internal_test.go | 20 +++---- 6 files changed, 83 insertions(+), 83 deletions(-) diff --git a/api.go b/api.go index b455f4301..5a576878c 100644 --- a/api.go +++ b/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) { diff --git a/cluster.go b/cluster.go index e1ea2ada7..de728472a 100644 --- a/cluster.go +++ b/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,16 +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 + maxWritesPerRequest int // Data directory path. Path string @@ -239,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 @@ -257,7 +257,7 @@ type Cluster struct { wg sync.WaitGroup closing chan struct{} - Logger Logger + logger Logger InternalClient InternalClient } @@ -266,7 +266,7 @@ type Cluster struct { func NewCluster() *Cluster { return &Cluster{ Hasher: &jmphasher{}, - PartitionN: DefaultPartitionN, + partitionN: DefaultPartitionN, ReplicaN: 1, joiningLeavingNodes: make(chan nodeAction, 10), // buffered channel @@ -276,7 +276,7 @@ func NewCluster() *Cluster { InternalClient: NewNopInternalClient(), - Logger: NopLogger, + logger: NopLogger, } } @@ -313,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), }) @@ -322,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 @@ -354,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 { @@ -405,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 { @@ -432,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 @@ -452,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) } } } @@ -474,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) } @@ -496,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() { @@ -509,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), } @@ -711,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 } @@ -781,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. @@ -865,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() { @@ -906,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 @@ -927,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) @@ -960,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") } @@ -976,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. @@ -984,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 { @@ -1010,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 @@ -1041,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 @@ -1053,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) } } @@ -1064,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 @@ -1078,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() @@ -1086,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 @@ -1106,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) @@ -1128,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") @@ -1150,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) @@ -1177,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{ @@ -1204,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 } @@ -1234,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 @@ -1266,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 @@ -1581,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 { @@ -1620,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 @@ -1640,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) } @@ -1651,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() { @@ -1679,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") } @@ -1733,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") } @@ -1755,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 diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 86e027600..b53180b39 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -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) diff --git a/fragment.go b/fragment.go index 28978d337..3853dd7cd 100644 --- a/fragment.go +++ b/fragment.go @@ -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 } diff --git a/server.go b/server.go index 403e19a0a..c07720a54 100644 --- a/server.go +++ b/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 diff --git a/utils_internal_test.go b/utils_internal_test.go index 20ff2b70b..5e19a83b9 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -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)