From 7801b81b10c2b3ef0fbc2c308c1f7fa9ab1e758e Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 2 Jul 2018 07:56:51 -0500 Subject: [PATCH] unexport cluster (gorename) --- api.go | 2 +- cluster.go | 120 +++++++++++++++++++-------------------- cluster_internal_test.go | 6 +- executor.go | 2 +- fragment.go | 2 +- holder.go | 4 +- server.go | 2 +- utils_internal_test.go | 8 +-- 8 files changed, 73 insertions(+), 73 deletions(-) diff --git a/api.go b/api.go index 0f19ef19d..bc3840ee5 100644 --- a/api.go +++ b/api.go @@ -36,7 +36,7 @@ import ( // wrapped by a handler which provides an external interface (e.g. HTTP). type API struct { Holder *Holder - Cluster *Cluster + Cluster *cluster server *Server } diff --git a/cluster.go b/cluster.go index ce0214537..920c06d8d 100644 --- a/cluster.go +++ b/cluster.go @@ -210,8 +210,8 @@ type nodeAction struct { action string } -// Cluster represents a collection of nodes. -type Cluster struct { +// cluster represents a collection of nodes. +type cluster struct { id string Node *Node Nodes []*Node // TODO phase this out? @@ -263,8 +263,8 @@ type Cluster struct { } // NewCluster returns a new instance of Cluster with defaults. -func NewCluster() *Cluster { - return &Cluster{ +func NewCluster() *cluster { + return &cluster{ Hasher: &jmphasher{}, partitionN: DefaultPartitionN, ReplicaN: 1, @@ -281,18 +281,18 @@ func NewCluster() *Cluster { } // coordinatorNode returns the coordinator node. -func (c *Cluster) coordinatorNode() *Node { +func (c *cluster) coordinatorNode() *Node { return c.unprotectedNodeByID(c.Coordinator) } // isCoordinator is true if this node is the coordinator. -func (c *Cluster) isCoordinator() bool { +func (c *cluster) isCoordinator() bool { c.mu.RLock() defer c.mu.RUnlock() return c.unprotectedIsCoordinator() } -func (c *Cluster) unprotectedIsCoordinator() bool { +func (c *cluster) unprotectedIsCoordinator() bool { return c.Coordinator == c.Node.ID } @@ -300,7 +300,7 @@ func (c *Cluster) unprotectedIsCoordinator() bool { // Coordinator. In response to this, the current node // will consider itself coordinator and update the other // nodes with its version of Cluster.Status. -func (c *Cluster) setCoordinator(n *Node) error { +func (c *cluster) setCoordinator(n *Node) error { c.mu.Lock() // Verify that the new Coordinator value matches // this node. @@ -329,13 +329,13 @@ func (c *Cluster) setCoordinator(n *Node) error { // changing the corresponding node's IsCoordinator value // to true, and sets all other nodes to false. Returns true if the value // changed. -func (c *Cluster) updateCoordinator(n *Node) bool { +func (c *cluster) updateCoordinator(n *Node) bool { c.mu.Lock() defer c.mu.Unlock() return c.unprotectedUpdateCoordinator(n) } -func (c *Cluster) unprotectedUpdateCoordinator(n *Node) bool { +func (c *cluster) unprotectedUpdateCoordinator(n *Node) bool { var changed bool if c.Coordinator != n.ID { c.Coordinator = n.ID @@ -353,7 +353,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 { +func (c *cluster) addNode(node *Node) error { 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. @@ -380,7 +380,7 @@ func (c *Cluster) addNode(node *Node) error { // removeNode removes a node from the Cluster and updates and saves the // new topology. -func (c *Cluster) removeNode(node *Node) error { +func (c *cluster) removeNode(node *Node) error { // remove from cluster if !c.removeNodeBasicSorted(node) { return nil @@ -399,11 +399,11 @@ func (c *Cluster) removeNode(node *Node) error { } // nodeIDs returns the list of IDs in the cluster. -func (c *Cluster) nodeIDs() []string { +func (c *cluster) nodeIDs() []string { return Nodes(c.Nodes).IDs() } -func (c *Cluster) setID(id string) { +func (c *cluster) setID(id string) { // Don't overwrite ClusterID. if c.id != "" { return @@ -414,19 +414,19 @@ func (c *Cluster) setID(id string) { c.Topology.ClusterID = c.id } -func (c *Cluster) State() string { +func (c *cluster) State() string { c.mu.RLock() defer c.mu.RUnlock() return c.state } -func (c *Cluster) SetState(state string) { +func (c *cluster) SetState(state string) { c.mu.Lock() c.setState(state) c.mu.Unlock() } -func (c *Cluster) setState(state string) { +func (c *cluster) setState(state string) { // Ignore cases where the state hasn't changed. if state == c.state { return @@ -463,7 +463,7 @@ func (c *Cluster) setState(state string) { } } -func (c *Cluster) setNodeState(state string) error { +func (c *cluster) setNodeState(state string) error { if c.isCoordinator() { return c.receiveNodeState(c.Node.ID, state) } @@ -485,7 +485,7 @@ func (c *Cluster) setNodeState(state string) error { // receiveNodeState sets node state in Topology in order for the // Coordinator to keep track of, during startup, which nodes have // finished opening their Holder. -func (c *Cluster) receiveNodeState(nodeID string, state string) error { +func (c *cluster) receiveNodeState(nodeID string, state string) error { if !c.isCoordinator() { return nil } @@ -507,7 +507,7 @@ func (c *Cluster) receiveNodeState(nodeID string, state string) error { } // Status returns the internal ClusterStatus representation. -func (c *Cluster) Status() *internal.ClusterStatus { +func (c *cluster) Status() *internal.ClusterStatus { return &internal.ClusterStatus{ ClusterID: c.id, State: c.state, @@ -515,14 +515,14 @@ func (c *Cluster) Status() *internal.ClusterStatus { } } -func (c *Cluster) nodeByID(id string) *Node { +func (c *cluster) nodeByID(id string) *Node { c.mu.RLock() defer c.mu.RUnlock() return c.unprotectedNodeByID(id) } // unprotectedNodeByID returns a node reference by ID. -func (c *Cluster) unprotectedNodeByID(id string) *Node { +func (c *cluster) unprotectedNodeByID(id string) *Node { for _, n := range c.Nodes { if n.ID == id { return n @@ -532,7 +532,7 @@ func (c *Cluster) unprotectedNodeByID(id string) *Node { } // nodePositionByID returns the position of the node in slice c.Nodes. -func (c *Cluster) nodePositionByID(nodeID string) int { +func (c *cluster) nodePositionByID(nodeID string) int { for i, n := range c.Nodes { if n.ID == nodeID { return i @@ -543,7 +543,7 @@ func (c *Cluster) nodePositionByID(nodeID string) int { // addNodeBasicSorted adds a node to the cluster, sorted by id. // Returns a pointer to the node and true if the node was added. -func (c *Cluster) addNodeBasicSorted(node *Node) bool { +func (c *cluster) addNodeBasicSorted(node *Node) bool { n := c.unprotectedNodeByID(node.ID) if n != nil { return false @@ -559,7 +559,7 @@ func (c *Cluster) addNodeBasicSorted(node *Node) bool { // removeNodeBasicSorted removes a node from the cluster, maintaining // the sort order. Returns true if the node was removed. -func (c *Cluster) removeNodeBasicSorted(node *Node) bool { +func (c *cluster) removeNodeBasicSorted(node *Node) bool { i := c.nodePositionByID(node.ID) if i < 0 { return false @@ -613,7 +613,7 @@ func (a viewsByField) addView(field, view string) { a[field] = append(a[field], view) } -func (c *Cluster) fragsByHost(idx *Index) fragsByHost { +func (c *cluster) fragsByHost(idx *Index) fragsByHost { // fieldViews is a map of field to slice of views. fieldViews := make(viewsByField) @@ -628,7 +628,7 @@ func (c *Cluster) fragsByHost(idx *Index) fragsByHost { // fragCombos returns a map (by uri) of lists of fragments for a given index // by creating every combination of field/view specified in `fieldViews` up to maxShard. -func (c *Cluster) fragCombos(idx string, maxShard uint64, fieldViews viewsByField) fragsByHost { +func (c *cluster) fragCombos(idx string, maxShard uint64, fieldViews viewsByField) fragsByHost { t := make(fragsByHost) for i := uint64(0); i <= maxShard; i++ { nodes := c.shardNodes(idx, i) @@ -647,7 +647,7 @@ func (c *Cluster) fragCombos(idx string, maxShard uint64, fieldViews viewsByFiel // diff compares c with another cluster and determines if a node is being // added or removed. An error is returned for any case other than where // exactly one node is added or removed. -func (c *Cluster) diff(other *Cluster) (action string, nodeID string, err error) { +func (c *cluster) diff(other *cluster) (action string, nodeID string, err error) { lenFrom := len(c.Nodes) lenTo := len(other.Nodes) // Determine if a node is being added or removed. @@ -686,7 +686,7 @@ func (c *Cluster) diff(other *Cluster) (action string, nodeID string, err error) // fragSources returns a list of ResizeSources - for each node in the `to` cluster - // required to move from cluster `c` to cluster `to`. -func (c *Cluster) fragSources(to *Cluster, idx *Index) (map[string][]*internal.ResizeSource, error) { +func (c *cluster) fragSources(to *cluster, idx *Index) (map[string][]*internal.ResizeSource, error) { m := make(map[string][]*internal.ResizeSource) // Determine if a node is being added or removed. @@ -773,7 +773,7 @@ func (c *Cluster) fragSources(to *Cluster, idx *Index) (map[string][]*internal.R } // partition returns the partition that a shard belongs to. -func (c *Cluster) partition(index string, shard uint64) int { +func (c *cluster) partition(index string, shard uint64) int { var buf [8]byte binary.BigEndian.PutUint64(buf[:], shard) @@ -785,17 +785,17 @@ func (c *Cluster) partition(index string, shard uint64) int { } // shardNodes returns a list of nodes that own a fragment. -func (c *Cluster) shardNodes(index string, shard uint64) []*Node { +func (c *cluster) shardNodes(index string, shard uint64) []*Node { return c.partitionNodes(c.partition(index, shard)) } // ownsShard returns true if a host owns a fragment. -func (c *Cluster) ownsShard(nodeID string, index string, shard uint64) bool { +func (c *cluster) ownsShard(nodeID string, index string, shard uint64) bool { return Nodes(c.shardNodes(index, shard)).ContainsID(nodeID) } // partitionNodes returns a list of nodes that own a partition. -func (c *Cluster) partitionNodes(partitionID int) []*Node { +func (c *cluster) partitionNodes(partitionID int) []*Node { // Default replica count to between one and the number of nodes. // The replica count can be zero if there are no nodes. replicaN := c.ReplicaN @@ -818,7 +818,7 @@ func (c *Cluster) partitionNodes(partitionID int) []*Node { } // containsShards is like OwnsShards, but it includes replicas. -func (c *Cluster) containsShards(index string, maxShard uint64, node *Node) []uint64 { +func (c *cluster) containsShards(index string, maxShard uint64, node *Node) []uint64 { var shards []uint64 for i := uint64(0); i <= maxShard; i++ { p := c.partition(index, i) @@ -856,7 +856,7 @@ func (h *jmphasher) Hash(key uint64, n int) int { return int(b) } -func (c *Cluster) setup() error { +func (c *cluster) setup() error { // Cluster always comes up in state STARTING until cluster membership is determined. c.state = ClusterStateStarting @@ -883,7 +883,7 @@ func (c *Cluster) setup() error { return nil } -func (c *Cluster) open() error { +func (c *cluster) open() error { err := c.setup() if err != nil { return errors.Wrap(err, "setting up cluster") @@ -891,7 +891,7 @@ func (c *Cluster) open() error { return c.waitForStarted() } -func (c *Cluster) waitForStarted() error { +func (c *cluster) waitForStarted() error { // If not coordinator then wait for ClusterStatus from coordinator. if !c.isCoordinator() { // In the case where a node has been restarted and memberlist has @@ -918,7 +918,7 @@ func (c *Cluster) waitForStarted() error { return nil } -func (c *Cluster) close() error { +func (c *cluster) close() error { // Notify goroutines of closing and wait for completion. close(c.closing) c.wg.Wait() @@ -926,7 +926,7 @@ func (c *Cluster) close() error { return nil } -func (c *Cluster) markAsJoined() { +func (c *cluster) markAsJoined() { c.logger.Printf("mark node as joined (received coordinator update)") if !c.joined { c.joined = true @@ -934,18 +934,18 @@ func (c *Cluster) markAsJoined() { } } -func (c *Cluster) needTopologyAgreement() bool { +func (c *cluster) needTopologyAgreement() bool { return c.State() == ClusterStateStarting && !stringSlicesAreEqual(c.Topology.NodeIDs, c.nodeIDs()) } -func (c *Cluster) haveTopologyAgreement() bool { +func (c *cluster) haveTopologyAgreement() bool { if c.Static { return true } return stringSlicesAreEqual(c.Topology.NodeIDs, c.nodeIDs()) } -func (c *Cluster) allNodesReady() bool { +func (c *cluster) allNodesReady() bool { if c.Static { return true } @@ -957,7 +957,7 @@ func (c *Cluster) allNodesReady() bool { return true } -func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { +func (c *cluster) handleNodeAction(nodeAction nodeAction) error { j, err := c.generateResizeJob(nodeAction) if err != nil { c.logger.Printf("generateResizeJob error: err=%s", err) @@ -1004,7 +1004,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { return nil } -func (c *Cluster) setStateAndBroadcast(state string) error { +func (c *cluster) setStateAndBroadcast(state string) error { c.SetState(state) if c.Static { return nil @@ -1014,7 +1014,7 @@ func (c *Cluster) setStateAndBroadcast(state string) error { return c.broadcaster.SendSync(c.Status()) } -func (c *Cluster) sendTo(node *Node, msg proto.Message) error { +func (c *cluster) sendTo(node *Node, msg proto.Message) error { if err := c.broadcaster.SendTo(node, msg); err != nil { return errors.Wrap(err, "sending") } @@ -1022,7 +1022,7 @@ func (c *Cluster) sendTo(node *Node, msg proto.Message) error { } // listenForJoins handles cluster-resize events. -func (c *Cluster) listenForJoins() { +func (c *cluster) listenForJoins() { c.wg.Add(1) go func() { defer c.wg.Done() @@ -1077,7 +1077,7 @@ func (c *Cluster) listenForJoins() { // generateResizeJob creates a new resizeJob based on the new node being // 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) { +func (c *cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) { c.logger.Printf("generateResizeJob: %v", nodeAction) c.mu.Lock() defer c.mu.Unlock() @@ -1104,7 +1104,7 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) { // the difference between Cluster and a new Cluster with/without uri. // Broadcaster is associated to the resizeJob here for use in broadcasting // the resize instructions to other nodes in the cluster. -func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, error) { +func (c *cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, error) { j := newResizeJob(c.Nodes, nodeAction.node, nodeAction.action) j.Broadcaster = c.broadcaster @@ -1161,7 +1161,7 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, // completeCurrentJob sets the state of the current resizeJob // then removes the pointer to currentJob. -func (c *Cluster) completeCurrentJob(state string) error { +func (c *cluster) completeCurrentJob(state string) error { c.mu.Lock() defer c.mu.Unlock() if !c.unprotectedIsCoordinator() { @@ -1176,7 +1176,7 @@ 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 { +func (c *cluster) followResizeInstruction(instr *internal.ResizeInstruction) error { 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. @@ -1272,7 +1272,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err return nil } -func (c *Cluster) markResizeInstructionComplete(complete *internal.ResizeInstructionComplete) error { +func (c *cluster) markResizeInstructionComplete(complete *internal.ResizeInstructionComplete) error { j := c.job(complete.JobID) @@ -1300,7 +1300,7 @@ func (c *Cluster) markResizeInstructionComplete(complete *internal.ResizeInstruc } // job returns a resizeJob by id. -func (c *Cluster) job(id int64) *resizeJob { +func (c *cluster) job(id int64) *resizeJob { c.mu.RLock() defer c.mu.RUnlock() return c.jobs[id] @@ -1516,7 +1516,7 @@ func (t *Topology) Encode() *internal.Topology { } // loadTopology reads the topology for the node. -func (c *Cluster) loadTopology() error { +func (c *cluster) loadTopology() error { buf, err := ioutil.ReadFile(filepath.Join(c.Path, ".topology")) if os.IsNotExist(err) { c.Topology = NewTopology() @@ -1539,7 +1539,7 @@ func (c *Cluster) loadTopology() error { } // saveTopology writes the current topology to disk. -func (c *Cluster) saveTopology() error { +func (c *cluster) saveTopology() error { if err := os.MkdirAll(c.Path, 0777); err != nil { return errors.Wrap(err, "creating directory") @@ -1579,7 +1579,7 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { return t, nil } -func (c *Cluster) considerTopology() error { +func (c *cluster) considerTopology() error { // Create ClusterID if one does not already exist. if c.id == "" { u := uuid.NewV4() @@ -1612,7 +1612,7 @@ func (c *Cluster) considerTopology() error { } // ReceiveEvent represents an implementation of EventHandler. -func (c *Cluster) ReceiveEvent(e *nodeEvent) error { +func (c *cluster) ReceiveEvent(e *nodeEvent) error { // Ignore events sent from this node. if e.Node.ID == c.Node.ID { return nil @@ -1635,7 +1635,7 @@ func (c *Cluster) ReceiveEvent(e *nodeEvent) error { return nil } -func (c *Cluster) nodeJoin(node *Node) error { +func (c *cluster) nodeJoin(node *Node) error { if c.needTopologyAgreement() { // A host that is not part of the topology can't be added to the STARTING cluster. if !c.Topology.ContainsID(node.ID) { @@ -1699,7 +1699,7 @@ func (c *Cluster) nodeJoin(node *Node) error { } // nodeLeave initiates the removal of a node from the cluster. -func (c *Cluster) nodeLeave(node *Node) error { +func (c *cluster) nodeLeave(node *Node) error { // Refuse the request if this is not the coordinator. if !c.isCoordinator() { return fmt.Errorf("node removal requests are only valid on the coordinator node: %s", c.coordinatorNode().ID) @@ -1752,7 +1752,7 @@ func (c *Cluster) nodeLeave(node *Node) error { return nil } -func (c *Cluster) mergeClusterStatus(cs *internal.ClusterStatus) error { +func (c *cluster) mergeClusterStatus(cs *internal.ClusterStatus) error { c.mu.Lock() defer c.mu.Unlock() c.logger.Printf("merge cluster status: %v", cs) @@ -1801,7 +1801,7 @@ func (c *Cluster) mergeClusterStatus(cs *internal.ClusterStatus) error { return nil } -func (c *Cluster) setStatic(hosts []string) error { +func (c *cluster) setStatic(hosts []string) error { c.Static = true c.Coordinator = c.Node.ID for _, address := range hosts { diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 8e4621653..d607cd883 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -172,8 +172,8 @@ func TestFragSources(t *testing.T) { } tests := []struct { - from *Cluster - to *Cluster + from *cluster + to *cluster idx *Index expected map[string][]*internal.ResizeSource err string @@ -316,7 +316,7 @@ func TestResizeJob(t *testing.T) { // Ensure the cluster can fairly distribute partitions across the nodes. func TestCluster_Owners(t *testing.T) { - c := Cluster{ + c := cluster{ Nodes: []*Node{ {URI: NewTestURIFromHostPort("serverA", 1000)}, {URI: NewTestURIFromHostPort("serverB", 1000)}, diff --git a/executor.go b/executor.go index c3e9ec77c..1977665cb 100644 --- a/executor.go +++ b/executor.go @@ -43,7 +43,7 @@ type Executor struct { // Local hostname & cluster configuration. Node *Node - Cluster *Cluster + Cluster *cluster // Client used for remote requests. client InternalQueryClient diff --git a/fragment.go b/fragment.go index a5a40d68f..30fb67fca 100644 --- a/fragment.go +++ b/fragment.go @@ -1717,7 +1717,7 @@ type FragmentSyncer struct { Fragment *Fragment Node *Node - Cluster *Cluster + Cluster *cluster Closing <-chan struct{} } diff --git a/holder.go b/holder.go index eb6eacd57..0141d365e 100644 --- a/holder.go +++ b/holder.go @@ -569,7 +569,7 @@ type HolderSyncer struct { Holder *Holder Node *Node - Cluster *Cluster + Cluster *cluster // Stats Stats StatsClient @@ -778,7 +778,7 @@ type HolderCleaner struct { Node *Node Holder *Holder - Cluster *Cluster + Cluster *cluster // Signals that the sync should stop. Closing <-chan struct{} diff --git a/server.go b/server.go index d64b9efa7..47a97bd24 100644 --- a/server.go +++ b/server.go @@ -51,7 +51,7 @@ type Server struct { // Internal holder *Holder - cluster *Cluster + cluster *cluster translateFile *TranslateFile diagnostics *DiagnosticsCollector executor *Executor diff --git a/utils_internal_test.go b/utils_internal_test.go index 6a84d168d..171955b4e 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -28,7 +28,7 @@ import ( ) // NewTestCluster returns a cluster with n nodes and uses a mod-based hasher. -func NewTestCluster(n int) *Cluster { +func NewTestCluster(n int) *cluster { path, err := ioutil.TempDir("", "pilosa-cluster-") if err != nil { panic(err) @@ -82,7 +82,7 @@ func (*TestModHasher) Hash(key uint64, n int) int { return int(key) % n } // has a Cluster. // ClusterCluster implements Broadcaster interface. type ClusterCluster struct { - Clusters []*Cluster + Clusters []*cluster common *commonClusterSettings @@ -141,7 +141,7 @@ func (t *ClusterCluster) SetBit(index, field string, rowID, colID uint64, x *tim return nil } -func (t *ClusterCluster) clusterByID(id string) *Cluster { +func (t *ClusterCluster) clusterByID(id string) *cluster { for _, c := range t.Clusters { if c.Node.ID == id { return c @@ -194,7 +194,7 @@ func (t *ClusterCluster) WriteTopology(path string, top *Topology) error { return nil } -func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*Cluster, error) { +func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*cluster, error) { id := fmt.Sprintf("node%d", i) uri := NewTestURI("http", fmt.Sprintf("host%d", i), uint16(0))