From 6ff792c164183289450cd12131d3bf3d6f9a0818 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Fri, 29 Jun 2018 08:12:04 -0500 Subject: [PATCH 1/4] remove dead code (deadcode) --- field.go | 37 ------------------------------------- translate.go | 3 --- view.go | 5 ----- 3 files changed, 45 deletions(-) diff --git a/field.go b/field.go index 9eb4f75b8..37cf17c2d 100644 --- a/field.go +++ b/field.go @@ -1298,43 +1298,6 @@ func (b *bsiGroup) validate() error { return nil } -func encodeBSIGroup(b *bsiGroup) *internal.BSIGroup { - if b == nil { - return nil - } - return &internal.BSIGroup{ - Name: b.Name, - Type: b.Type, - Min: int64(b.Min), - Max: int64(b.Max), - } -} - -func decodeBSIGroup(b *internal.BSIGroup) *bsiGroup { - if b == nil { - return nil - } - return &bsiGroup{ - Name: b.Name, - Type: b.Type, - Min: b.Min, - Max: b.Max, - } -} - -// importBitSet represents slices of row and column ids. -// This is used to sort data during import. -type importBitSet struct { - rowIDs, columnIDs []uint64 -} - -func (p importBitSet) Swap(i, j int) { - p.rowIDs[i], p.rowIDs[j] = p.rowIDs[j], p.rowIDs[i] - p.columnIDs[i], p.columnIDs[j] = p.columnIDs[j], p.columnIDs[i] -} -func (p importBitSet) Len() int { return len(p.rowIDs) } -func (p importBitSet) Less(i, j int) bool { return p.rowIDs[i] < p.rowIDs[j] } - // Cache types. const ( CacheTypeLRU = "lru" diff --git a/translate.go b/translate.go index 7c716640b..99a52acf5 100644 --- a/translate.go +++ b/translate.go @@ -5,7 +5,6 @@ import ( "bytes" "context" "encoding/binary" - "encoding/hex" "errors" "fmt" "io" @@ -1004,5 +1003,3 @@ func UvarintSize(x uint64) (i int) { } return i + 1 } - -func hexdump(b []byte) { os.Stderr.Write([]byte(hex.Dump(b))) } diff --git a/view.go b/view.go index 3256fbc1b..fab7d7e1b 100644 --- a/view.go +++ b/view.go @@ -34,11 +34,6 @@ const ( viewBSIGroupPrefix = "bsig_" ) -// isValidView returns true if name is valid. -func isValidView(name string) bool { - return name == ViewStandard -} - // View represents a container for field data. type View struct { mu sync.RWMutex From d73ff9bfe2434335feaaa6a9b4596cea3ba8af72 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Fri, 29 Jun 2018 08:12:20 -0500 Subject: [PATCH 2/4] collapse memAttrStore memAttrStore is not used, but I opted to leave it in order to encourage its use by future unit tests. Since it is unexported, however, and its fairly obvious what it does, I didn't think the docstrings were adding much value, and it's more readable in this compact form. --- attr.go | 37 ++++++------------------------------- 1 file changed, 6 insertions(+), 31 deletions(-) diff --git a/attr.go b/attr.go index e73f303af..34d14a280 100644 --- a/attr.go +++ b/attr.go @@ -233,44 +233,19 @@ type memAttrStore struct { store map[uint64]map[string]interface{} } -// Path is an in-memory implementation of AttrStore Path method. -func (s *memAttrStore) Path() string { return "" } - -// Open is an in-memory implementation of AttrStore Open method. -func (s *memAttrStore) Open() error { - return nil -} - -// Close is an in-memory implementation of AttrStore Close method. -func (s *memAttrStore) Close() error { - return nil -} - -// Attrs returns a set of attributes by ID. -func (s *memAttrStore) Attrs(id uint64) (m map[string]interface{}, err error) { - return s.store[id], nil -} - -// SetAttrs sets attribute values for a given ID. +func (s *memAttrStore) Path() string { return "" } +func (s *memAttrStore) Open() error { return nil } +func (s *memAttrStore) Close() error { return nil } +func (s *memAttrStore) Attrs(id uint64) (m map[string]interface{}, err error) { return s.store[id], nil } func (s *memAttrStore) SetAttrs(id uint64, m map[string]interface{}) error { s.store[id] = m return nil } - -// SetBulkAttrs sets attribute values for a set of ids. func (s *memAttrStore) SetBulkAttrs(m map[uint64]map[string]interface{}) error { for id, v := range m { s.store[id] = v } return nil } - -// Blocks is an in-memory implementation of AttrStore Blocks method. -func (s *memAttrStore) Blocks() ([]AttrBlock, error) { - return nil, nil -} - -// BlockData is an in-memory implementation of AttrStore BlockData method. -func (s *memAttrStore) BlockData(i uint64) (map[uint64]map[string]interface{}, error) { - return nil, nil -} +func (s *memAttrStore) Blocks() ([]AttrBlock, error) { return nil, nil } +func (s *memAttrStore) BlockData(i uint64) (map[uint64]map[string]interface{}, error) { return nil, nil } From 7801b81b10c2b3ef0fbc2c308c1f7fa9ab1e758e Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 2 Jul 2018 07:56:51 -0500 Subject: [PATCH 3/4] 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)) From 91f531f2cdc9145600c38bed0ffd0abf4a41e3a6 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 2 Jul 2018 08:14:13 -0500 Subject: [PATCH 4/4] more unexports - executor, api fields --- api.go | 94 ++++++++++++++++++++++++++--------------------------- executor.go | 94 ++++++++++++++++++++++++++--------------------------- server.go | 4 +-- 3 files changed, 96 insertions(+), 96 deletions(-) diff --git a/api.go b/api.go index bc3840ee5..0c1f85950 100644 --- a/api.go +++ b/api.go @@ -35,8 +35,8 @@ import ( // API provides the top level programmatic interface to Pilosa. It is usually // wrapped by a handler which provides an external interface (e.g. HTTP). type API struct { - Holder *Holder - Cluster *cluster + holder *Holder + cluster *cluster server *Server } @@ -46,8 +46,8 @@ type APIOption func(*API) error func OptAPIServer(s *Server) APIOption { return func(a *API) error { a.server = s - a.Holder = s.holder - a.Cluster = s.cluster + a.holder = s.holder + a.cluster = s.cluster return nil } } @@ -85,7 +85,7 @@ func appendMap(a, b map[apiMethod]struct{}) map[apiMethod]struct{} { } func (api *API) validate(f apiMethod) error { - state := api.Cluster.State() + state := api.cluster.State() if _, ok := validAPIMethods[state][f]; ok { return nil } @@ -128,7 +128,7 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er } // Retrieve column attributes across all calls. - columnAttrSets, err := api.readColumnAttrSets(api.Holder.Index(req.Index), columnIDs) + columnAttrSets, err := api.readColumnAttrSets(api.holder.Index(req.Index), columnIDs) if err != nil { return resp, errors.Wrap(err, "reading column attrs") } @@ -179,7 +179,7 @@ func (api *API) CreateIndex(ctx context.Context, indexName string, options Index } // Create index. - index, err := api.Holder.CreateIndex(indexName, options) + index, err := api.holder.CreateIndex(indexName, options) if err != nil { return nil, errors.Wrap(err, "creating index") } @@ -193,7 +193,7 @@ func (api *API) CreateIndex(ctx context.Context, indexName string, options Index api.server.logger.Printf("problem sending CreateIndex message: %s", err) return nil, errors.Wrap(err, "sending CreateIndex message") } - api.Holder.Stats.Count("createIndex", 1, 1.0) + api.holder.Stats.Count("createIndex", 1, 1.0) return index, nil } @@ -203,7 +203,7 @@ func (api *API) Index(ctx context.Context, indexName string) (*Index, error) { return nil, errors.Wrap(err, "validating api method") } - index := api.Holder.Index(indexName) + index := api.holder.Index(indexName) if index == nil { return nil, ErrIndexNotFound } @@ -218,7 +218,7 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error { } // Delete index from the holder. - err := api.Holder.DeleteIndex(indexName) + err := api.holder.DeleteIndex(indexName) if err != nil { return errors.Wrap(err, "deleting index") } @@ -231,7 +231,7 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error { api.server.logger.Printf("problem sending DeleteIndex message: %s", err) return errors.Wrap(err, "sending DeleteIndex message") } - api.Holder.Stats.Count("deleteIndex", 1, 1.0) + api.holder.Stats.Count("deleteIndex", 1, 1.0) return nil } @@ -251,7 +251,7 @@ func (api *API) CreateField(ctx context.Context, indexName string, fieldName str } // Find index. - index := api.Holder.Index(indexName) + index := api.holder.Index(indexName) if index == nil { return nil, ErrIndexNotFound } @@ -273,7 +273,7 @@ func (api *API) CreateField(ctx context.Context, indexName string, fieldName str api.server.logger.Printf("problem sending CreateField message: %s", err) return nil, errors.Wrap(err, "sending CreateField message") } - api.Holder.Stats.CountWithCustomTags("createField", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) + api.holder.Stats.CountWithCustomTags("createField", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) return field, nil } @@ -286,7 +286,7 @@ func (api *API) DeleteField(ctx context.Context, indexName string, fieldName str } // Find index. - index := api.Holder.Index(indexName) + index := api.holder.Index(indexName) if index == nil { return ErrIndexNotFound } @@ -306,7 +306,7 @@ func (api *API) DeleteField(ctx context.Context, indexName string, fieldName str api.server.logger.Printf("problem sending DeleteField message: %s", err) return errors.Wrap(err, "sending DeleteField message") } - api.Holder.Stats.CountWithCustomTags("deleteField", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) + api.holder.Stats.CountWithCustomTags("deleteField", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) return nil } @@ -318,13 +318,13 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin } // Validate that this handler owns the shard. - if !api.Cluster.ownsShard(api.LocalID(), indexName, shard) { + if !api.cluster.ownsShard(api.LocalID(), indexName, shard) { api.server.logger.Printf("node %s does not own shard %d of index %s", api.LocalID(), shard, indexName) return ErrClusterDoesNotOwnShard } // Find the fragment. - f := api.Holder.Fragment(indexName, fieldName, ViewStandard, shard) + f := api.holder.Fragment(indexName, fieldName, ViewStandard, shard) if f == nil { return ErrFragmentNotFound } @@ -354,7 +354,7 @@ func (api *API) ShardNodes(ctx context.Context, indexName string, shard uint64) return nil, errors.Wrap(err, "validating api method") } - return api.Cluster.shardNodes(indexName, shard), nil + return api.cluster.shardNodes(indexName, shard), nil } // MarshalFragment returns an object which can write the specified fragment's data @@ -366,7 +366,7 @@ func (api *API) MarshalFragment(ctx context.Context, indexName string, fieldName } // Retrieve fragment from holder. - f := api.Holder.Fragment(indexName, fieldName, ViewStandard, shard) + f := api.holder.Fragment(indexName, fieldName, ViewStandard, shard) if f == nil { return nil, ErrFragmentNotFound } @@ -382,7 +382,7 @@ func (api *API) UnmarshalFragment(ctx context.Context, indexName string, fieldNa } // Retrieve field. - f := api.Holder.Field(indexName, fieldName) + f := api.holder.Field(indexName, fieldName) if f == nil { return ErrFieldNotFound } @@ -424,7 +424,7 @@ func (api *API) FragmentBlockData(ctx context.Context, body io.Reader) ([]byte, } // Retrieve fragment from holder. - f := api.Holder.Fragment(req.Index, req.Field, ViewStandard, req.Shard) + f := api.holder.Fragment(req.Index, req.Field, ViewStandard, req.Shard) if f == nil { return nil, ErrFragmentNotFound } @@ -448,7 +448,7 @@ func (api *API) FragmentBlocks(ctx context.Context, indexName string, fieldName } // Retrieve fragment from holder. - f := api.Holder.Fragment(indexName, fieldName, ViewStandard, shard) + f := api.holder.Fragment(indexName, fieldName, ViewStandard, shard) if f == nil { return nil, ErrFragmentNotFound } @@ -461,7 +461,7 @@ func (api *API) FragmentBlocks(ctx context.Context, indexName string, fieldName // Hosts returns a list of the hosts in the cluster including their ID, // URL, and which is the coordinator. func (api *API) Hosts(ctx context.Context) []*Node { - return api.Cluster.Nodes + return api.cluster.Nodes } // RecalculateCaches forces all TopN caches to be updated. Used mainly for integration tests. @@ -474,7 +474,7 @@ func (api *API) RecalculateCaches(ctx context.Context) error { if err != nil { return errors.Wrap(err, "broacasting message") } - api.Holder.RecalculateCaches() + api.holder.RecalculateCaches() return nil } @@ -506,13 +506,13 @@ func (api *API) ClusterMessage(ctx context.Context, reqBody io.Reader) error { // LocalID returns the current node's ID. func (api *API) LocalID() string { - return api.Cluster.Node.ID + return api.cluster.Node.ID } // Schema returns information about each index in Pilosa including which fields // and views they contain. func (api *API) Schema(ctx context.Context) []*IndexInfo { - return api.Holder.Schema() + return api.holder.Schema() } // Views returns the views in the given field. @@ -522,7 +522,7 @@ func (api *API) Views(ctx context.Context, indexName string, fieldName string) ( } // Retrieve views. - f := api.Holder.Field(indexName, fieldName) + f := api.holder.Field(indexName, fieldName) if f == nil { return nil, ErrFieldNotFound } @@ -539,7 +539,7 @@ func (api *API) DeleteView(ctx context.Context, indexName string, fieldName stri } // Retrieve field. - f := api.Holder.Field(indexName, fieldName) + f := api.holder.Field(indexName, fieldName) if f == nil { return ErrFieldNotFound } @@ -573,7 +573,7 @@ func (api *API) IndexAttrDiff(ctx context.Context, indexName string, blocks []At } // Retrieve index from holder. - index := api.Holder.Index(indexName) + index := api.holder.Index(indexName) if index == nil { return nil, ErrIndexNotFound } @@ -607,7 +607,7 @@ func (api *API) FieldAttrDiff(ctx context.Context, indexName string, fieldName s } // Retrieve index from holder. - f := api.Holder.Field(indexName, fieldName) + f := api.holder.Field(indexName, fieldName) if f == nil { return nil, ErrFieldNotFound } @@ -684,37 +684,37 @@ func (api *API) ImportValue(ctx context.Context, req internal.ImportValueRequest // MaxShards returns the maximum shard number for each index in a map. func (api *API) MaxShards(ctx context.Context) map[string]uint64 { - return api.Holder.MaxShards() + return api.holder.MaxShards() } // StatsWithTags returns an instance of whatever implementation of StatsClient // pilosa is using with the given tags. func (api *API) StatsWithTags(tags []string) StatsClient { - if api.Holder == nil || api.Cluster == nil { + if api.holder == nil || api.cluster == nil { return nil } - return api.Holder.Stats.WithTags(tags...) + return api.holder.Stats.WithTags(tags...) } // LongQueryTime returns the configured threshold for logging/statting // long running queries. func (api *API) LongQueryTime() time.Duration { - if api.Cluster == nil { + if api.cluster == nil { return 0 } - return api.Cluster.longQueryTime + return api.cluster.longQueryTime } func (api *API) indexField(indexName string, fieldName string, shard uint64) (*Index, *Field, error) { // Validate that this handler owns the shard. - if !api.Cluster.ownsShard(api.LocalID(), indexName, shard) { + if !api.cluster.ownsShard(api.LocalID(), indexName, shard) { api.server.logger.Printf("node %s does not own shard %d of index %s", api.LocalID(), shard, indexName) return nil, nil, ErrClusterDoesNotOwnShard } // Find the Index. api.server.logger.Printf("importing: %v %v %v", indexName, fieldName, shard) - index := api.Holder.Index(indexName) + index := api.holder.Index(indexName) if index == nil { api.server.logger.Printf("fragment error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrIndexNotFound.Error()) return nil, nil, ErrIndexNotFound @@ -735,15 +735,15 @@ func (api *API) SetCoordinator(ctx context.Context, id string) (oldNode, newNode return nil, nil, errors.Wrap(err, "validating api method") } - oldNode = api.Cluster.nodeByID(api.Cluster.Coordinator) - newNode = api.Cluster.nodeByID(id) + oldNode = api.cluster.nodeByID(api.cluster.Coordinator) + newNode = api.cluster.nodeByID(id) if newNode == nil { return nil, nil, errors.Wrap(ErrNodeIDNotExists, "getting new node") } // If the new coordinator is this node, do the SetCoordinator directly. if newNode.ID == api.LocalID() { - return oldNode, newNode, api.Cluster.setCoordinator(newNode) + return oldNode, newNode, api.cluster.setCoordinator(newNode) } // Send the set-coordinator message to new node. @@ -765,13 +765,13 @@ func (api *API) RemoveNode(id string) (*Node, error) { return nil, errors.Wrap(err, "validating api method") } - removeNode := api.Cluster.unprotectedNodeByID(id) + removeNode := api.cluster.unprotectedNodeByID(id) if removeNode == nil { return nil, errors.Wrap(ErrNodeIDNotExists, "finding node to remove") } // Start the resize process (similar to NodeJoin) - err := api.Cluster.nodeLeave(removeNode) + err := api.cluster.nodeLeave(removeNode) if err != nil { return removeNode, errors.Wrap(err, "calling node leave") } @@ -784,7 +784,7 @@ func (api *API) ResizeAbort() error { return errors.Wrap(err, "validating api method") } - err := api.Cluster.completeCurrentJob(resizeJobStateAborted) + err := api.cluster.completeCurrentJob(resizeJobStateAborted) return errors.Wrap(err, "complete current job") } @@ -834,7 +834,7 @@ func (api *API) GetTranslateData(ctx context.Context, w io.WriteCloser, offset i // "STARTING", "RESIZING", or potentially others. See cluster.go for more // details. func (api *API) State() string { - return api.Cluster.State() + return api.cluster.State() } // Version returns the Pilosa version. @@ -843,13 +843,13 @@ func (api *API) Version() string { } // Info returns information about this server instance -func (api *API) Info() ServerInfo { - return ServerInfo{ +func (api *API) Info() serverInfo { + return serverInfo{ ShardWidth: ShardWidth, } } -type ServerInfo struct { +type serverInfo struct { ShardWidth uint64 `json:"shardWidth"` } diff --git a/executor.go b/executor.go index 1977665cb..032bf4225 100644 --- a/executor.go +++ b/executor.go @@ -37,8 +37,8 @@ const ( rowLabel = "row" ) -// Executor recursively executes calls in a PQL query across all shards. -type Executor struct { +// executor recursively executes calls in a PQL query across all shards. +type executor struct { Holder *Holder // Local hostname & cluster configuration. @@ -55,19 +55,19 @@ type Executor struct { TranslateStore TranslateStore } -// ExecutorOption is a functional option type for pilosa.Executor -type ExecutorOption func(e *Executor) error +// executorOption is a functional option type for pilosa.Executor +type executorOption func(e *executor) error -func OptExecutorInternalQueryClient(c InternalQueryClient) ExecutorOption { - return func(e *Executor) error { +func optExecutorInternalQueryClient(c InternalQueryClient) executorOption { + return func(e *executor) error { e.client = c return nil } } -// NewExecutor returns a new instance of Executor. -func NewExecutor(opts ...ExecutorOption) *Executor { - e := &Executor{ +// newExecutor returns a new instance of Executor. +func newExecutor(opts ...executorOption) *executor { + e := &executor{ client: NewNopInternalQueryClient(), } for _, opt := range opts { @@ -80,7 +80,7 @@ func NewExecutor(opts ...ExecutorOption) *Executor { } // Execute executes a PQL query. -func (e *Executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *ExecOptions) ([]interface{}, error) { +func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *ExecOptions) ([]interface{}, error) { // Verify that an index is set. if index == "" { return nil, ErrIndexRequired @@ -123,7 +123,7 @@ func (e *Executor) Execute(ctx context.Context, index string, q *pql.Query, shar return results, nil } -func (e *Executor) execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *ExecOptions) ([]interface{}, error) { +func (e *executor) execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *ExecOptions) ([]interface{}, error) { // Don't bother calculating shards for query types that don't require it. needsShards := needsShards(q.Calls) @@ -162,7 +162,7 @@ func (e *Executor) execute(ctx context.Context, index string, q *pql.Query, shar } // executeCall executes a call. -func (e *Executor) executeCall(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) { +func (e *executor) executeCall(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) { if err := e.validateCallArgs(c); err != nil { return nil, errors.Wrap(err, "validating args") } @@ -201,7 +201,7 @@ func (e *Executor) executeCall(ctx context.Context, index string, c *pql.Call, s } // validateCallArgs ensures that the value types in call.Args are expected. -func (e *Executor) validateCallArgs(c *pql.Call) error { +func (e *executor) validateCallArgs(c *pql.Call) error { if _, ok := c.Args["ids"]; ok { switch v := c.Args["ids"].(type) { case []int64, []uint64: @@ -220,7 +220,7 @@ func (e *Executor) validateCallArgs(c *pql.Call) error { } // executeSum executes a Sum() call. -func (e *Executor) executeSum(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (ValCount, error) { +func (e *executor) executeSum(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (ValCount, error) { if field := c.Args["field"]; field == "" { return ValCount{}, errors.New("Sum(): field required") } @@ -253,7 +253,7 @@ func (e *Executor) executeSum(ctx context.Context, index string, c *pql.Call, sh } // executeMin executes a Min() call. -func (e *Executor) executeMin(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (ValCount, error) { +func (e *executor) executeMin(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (ValCount, error) { if field := c.Args["field"]; field == "" { return ValCount{}, errors.New("Min(): field required") } @@ -286,7 +286,7 @@ func (e *Executor) executeMin(ctx context.Context, index string, c *pql.Call, sh } // executeMax executes a Max() call. -func (e *Executor) executeMax(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (ValCount, error) { +func (e *executor) executeMax(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (ValCount, error) { if field := c.Args["field"]; field == "" { return ValCount{}, errors.New("Max(): field required") } @@ -319,7 +319,7 @@ func (e *Executor) executeMax(ctx context.Context, index string, c *pql.Call, sh } // executeBitmapCall executes a call that returns a bitmap. -func (e *Executor) executeBitmapCall(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*Row, error) { +func (e *executor) executeBitmapCall(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*Row, error) { // Execute calls in bulk on each remote node and merge. mapFn := func(shard uint64) (interface{}, error) { return e.executeBitmapCallShard(ctx, index, c, shard) @@ -385,7 +385,7 @@ func (e *Executor) executeBitmapCall(ctx context.Context, index string, c *pql.C } // executeBitmapCallShard executes a bitmap call for a single shard. -func (e *Executor) executeBitmapCallShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeBitmapCallShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { switch c.Name { case "Row": return e.executeBitmapShard(ctx, index, c, shard) @@ -405,7 +405,7 @@ func (e *Executor) executeBitmapCallShard(ctx context.Context, index string, c * } // executeSumCountShard calculates the sum and count for bsiGroups on a shard. -func (e *Executor) executeSumCountShard(ctx context.Context, index string, c *pql.Call, shard uint64) (ValCount, error) { +func (e *executor) executeSumCountShard(ctx context.Context, index string, c *pql.Call, shard uint64) (ValCount, error) { var filter *Row if len(c.Children) == 1 { row, err := e.executeBitmapCallShard(ctx, index, c.Children[0], shard) @@ -443,7 +443,7 @@ func (e *Executor) executeSumCountShard(ctx context.Context, index string, c *pq } // executeMinShard calculates the min for bsiGroups on a shard. -func (e *Executor) executeMinShard(ctx context.Context, index string, c *pql.Call, shard uint64) (ValCount, error) { +func (e *executor) executeMinShard(ctx context.Context, index string, c *pql.Call, shard uint64) (ValCount, error) { var filter *Row if len(c.Children) == 1 { row, err := e.executeBitmapCallShard(ctx, index, c.Children[0], shard) @@ -481,7 +481,7 @@ func (e *Executor) executeMinShard(ctx context.Context, index string, c *pql.Cal } // executeMaxShard calculates the max for bsiGroups on a shard. -func (e *Executor) executeMaxShard(ctx context.Context, index string, c *pql.Call, shard uint64) (ValCount, error) { +func (e *executor) executeMaxShard(ctx context.Context, index string, c *pql.Call, shard uint64) (ValCount, error) { var filter *Row if len(c.Children) == 1 { row, err := e.executeBitmapCallShard(ctx, index, c.Children[0], shard) @@ -521,7 +521,7 @@ func (e *Executor) executeMaxShard(ctx context.Context, index string, c *pql.Cal // executeTopN executes a TopN() call. // This first performs the TopN() to determine the top results and then // requeries to retrieve the full counts for each of the top results. -func (e *Executor) executeTopN(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) ([]Pair, error) { +func (e *executor) executeTopN(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) ([]Pair, error) { idsArg, _, err := c.UintSliceArg("ids") if err != nil { return nil, fmt.Errorf("executeTopN: %v", err) @@ -560,7 +560,7 @@ func (e *Executor) executeTopN(ctx context.Context, index string, c *pql.Call, s return trimmedList, nil } -func (e *Executor) executeTopNShards(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) ([]Pair, error) { +func (e *executor) executeTopNShards(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) ([]Pair, error) { // Execute calls in bulk on each remote node and merge. mapFn := func(shard uint64) (interface{}, error) { return e.executeTopNShard(ctx, index, c, shard) @@ -585,7 +585,7 @@ func (e *Executor) executeTopNShards(ctx context.Context, index string, c *pql.C } // executeTopNShard executes a TopN call for a single shard. -func (e *Executor) executeTopNShard(ctx context.Context, index string, c *pql.Call, shard uint64) ([]Pair, error) { +func (e *executor) executeTopNShard(ctx context.Context, index string, c *pql.Call, shard uint64) ([]Pair, error) { field, _ := c.Args["_field"].(string) n, _, err := c.UintArg("n") if err != nil { @@ -647,7 +647,7 @@ func (e *Executor) executeTopNShard(ctx context.Context, index string, c *pql.Ca } // executeDifferenceShard executes a difference() call for a local shard. -func (e *Executor) executeDifferenceShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeDifferenceShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { var other *Row if len(c.Children) == 0 { return nil, fmt.Errorf("empty Difference query is currently not supported") @@ -668,7 +668,7 @@ func (e *Executor) executeDifferenceShard(ctx context.Context, index string, c * return other, nil } -func (e *Executor) executeBitmapShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeBitmapShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { // Fetch column label from index. idx := e.Holder.Index(index) if idx == nil { @@ -701,7 +701,7 @@ func (e *Executor) executeBitmapShard(ctx context.Context, index string, c *pql. } // executeIntersectShard executes a intersect() call for a local shard. -func (e *Executor) executeIntersectShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeIntersectShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { var other *Row if len(c.Children) == 0 { return nil, fmt.Errorf("empty Intersect query is currently not supported") @@ -723,7 +723,7 @@ func (e *Executor) executeIntersectShard(ctx context.Context, index string, c *p } // executeRangeShard executes a range() call for a local shard. -func (e *Executor) executeRangeShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeRangeShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { // Handle bsiGroup ranges differently. if c.HasConditionArg() { return e.executeBSIGroupRangeShard(ctx, index, c, shard) @@ -796,7 +796,7 @@ func (e *Executor) executeRangeShard(ctx context.Context, index string, c *pql.C } // executeBSIGroupRangeShard executes a range(bsiGroup) call for a local shard. -func (e *Executor) executeBSIGroupRangeShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeBSIGroupRangeShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { // Only one conditional should be present. if len(c.Args) == 0 { return nil, errors.New("Range(): condition required") @@ -926,7 +926,7 @@ func (e *Executor) executeBSIGroupRangeShard(ctx context.Context, index string, } // executeUnionShard executes a union() call for a local shard. -func (e *Executor) executeUnionShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeUnionShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { other := NewRow() for i, input := range c.Children { row, err := e.executeBitmapCallShard(ctx, index, input, shard) @@ -945,7 +945,7 @@ func (e *Executor) executeUnionShard(ctx context.Context, index string, c *pql.C } // executeXorShard executes a xor() call for a local shard. -func (e *Executor) executeXorShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { +func (e *executor) executeXorShard(ctx context.Context, index string, c *pql.Call, shard uint64) (*Row, error) { other := NewRow() for i, input := range c.Children { row, err := e.executeBitmapCallShard(ctx, index, input, shard) @@ -964,7 +964,7 @@ func (e *Executor) executeXorShard(ctx context.Context, index string, c *pql.Cal } // executeCount executes a count() call. -func (e *Executor) executeCount(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (uint64, error) { +func (e *executor) executeCount(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (uint64, error) { if len(c.Children) == 0 { return 0, errors.New("Count() requires an input bitmap") } else if len(c.Children) > 1 { @@ -996,7 +996,7 @@ func (e *Executor) executeCount(ctx context.Context, index string, c *pql.Call, } // executeClearBit executes a Clear() call. -func (e *Executor) executeClearBit(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) (bool, error) { +func (e *executor) executeClearBit(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) (bool, error) { fieldName, err := c.FieldArg() if err != nil { return false, errors.New("Clear() argument required: field") @@ -1031,7 +1031,7 @@ func (e *Executor) executeClearBit(ctx context.Context, index string, c *pql.Cal } // executeClearBitField executes a Clear() call for a single view. -func (e *Executor) executeClearBitField(ctx context.Context, index string, c *pql.Call, f *Field, colID, rowID uint64, opt *ExecOptions) (bool, error) { +func (e *executor) executeClearBitField(ctx context.Context, index string, c *pql.Call, f *Field, colID, rowID uint64, opt *ExecOptions) (bool, error) { shard := colID / ShardWidth ret := false for _, node := range e.Cluster.shardNodes(index, shard) { @@ -1061,7 +1061,7 @@ func (e *Executor) executeClearBitField(ctx context.Context, index string, c *pq } // executeSetBit executes a Set() call. -func (e *Executor) executeSetBit(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) (bool, error) { +func (e *executor) executeSetBit(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) (bool, error) { fieldName, err := c.FieldArg() if err != nil { return false, errors.New("Set() argument required: field") @@ -1106,7 +1106,7 @@ func (e *Executor) executeSetBit(ctx context.Context, index string, c *pql.Call, } // executeSetBitField executes a Set() call for a specific view. -func (e *Executor) executeSetBitField(ctx context.Context, index string, c *pql.Call, f *Field, colID, rowID uint64, timestamp *time.Time, opt *ExecOptions) (bool, error) { +func (e *executor) executeSetBitField(ctx context.Context, index string, c *pql.Call, f *Field, colID, rowID uint64, timestamp *time.Time, opt *ExecOptions) (bool, error) { shard := colID / ShardWidth ret := false @@ -1138,7 +1138,7 @@ func (e *Executor) executeSetBitField(ctx context.Context, index string, c *pql. } // executeSetValue executes a SetValue() call. -func (e *Executor) executeSetValue(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) error { +func (e *executor) executeSetValue(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) error { // Parse labels. columnID, ok, err := c.UintArg(columnLabel) if err != nil { @@ -1198,7 +1198,7 @@ func (e *Executor) executeSetValue(ctx context.Context, index string, c *pql.Cal } // executeSetRowAttrs executes a SetRowAttrs() call. -func (e *Executor) executeSetRowAttrs(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) error { +func (e *executor) executeSetRowAttrs(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) error { fieldName, ok := c.Args["_field"].(string) if !ok { return errors.New("SetRowAttrs() field required") @@ -1255,7 +1255,7 @@ func (e *Executor) executeSetRowAttrs(ctx context.Context, index string, c *pql. } // executeBulkSetRowAttrs executes a set of SetRowAttrs() calls. -func (e *Executor) executeBulkSetRowAttrs(ctx context.Context, index string, calls []*pql.Call, opt *ExecOptions) ([]interface{}, error) { +func (e *executor) executeBulkSetRowAttrs(ctx context.Context, index string, calls []*pql.Call, opt *ExecOptions) ([]interface{}, error) { // Collect attributes by field/id. m := make(map[string]map[uint64]map[string]interface{}) for _, c := range calls { @@ -1342,7 +1342,7 @@ func (e *Executor) executeBulkSetRowAttrs(ctx context.Context, index string, cal } // executeSetColumnAttrs executes a SetColumnAttrs() call. -func (e *Executor) executeSetColumnAttrs(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) error { +func (e *executor) executeSetColumnAttrs(ctx context.Context, index string, c *pql.Call, opt *ExecOptions) error { // Retrieve index. idx := e.Holder.Index(index) if idx == nil { @@ -1390,7 +1390,7 @@ func (e *Executor) executeSetColumnAttrs(ctx context.Context, index string, c *p } // exec executes a PQL query remotely for a set of shards on a node. -func (e *Executor) remoteExec(ctx context.Context, node *Node, index string, q *pql.Query, shards []uint64, opt *ExecOptions) (results []interface{}, err error) { +func (e *executor) remoteExec(ctx context.Context, node *Node, index string, q *pql.Query, shards []uint64, opt *ExecOptions) (results []interface{}, err error) { // Encode request object. pbreq := &internal.QueryRequest{ Query: q.String(), @@ -1441,7 +1441,7 @@ func (e *Executor) remoteExec(ctx context.Context, node *Node, index string, q * // shardsByNode returns a mapping of nodes to shards. // Returns errShardUnavailable if a shard cannot be allocated to a node. -func (e *Executor) shardsByNode(nodes []*Node, index string, shards []uint64) (map[*Node][]uint64, error) { +func (e *executor) shardsByNode(nodes []*Node, index string, shards []uint64) (map[*Node][]uint64, error) { m := make(map[*Node][]uint64) loop: @@ -1461,7 +1461,7 @@ loop: // // If a mapping of shards to a node fails then the shards are resplit across // secondary nodes and retried. This continues to occur until all nodes are exhausted. -func (e *Executor) mapReduce(ctx context.Context, index string, shards []uint64, c *pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) (interface{}, error) { +func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64, c *pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) (interface{}, error) { ch := make(chan mapResponse) // Wrap context with a cancel to kill goroutines on exit. @@ -1520,7 +1520,7 @@ func (e *Executor) mapReduce(ctx context.Context, index string, shards []uint64, } } -func (e *Executor) mapper(ctx context.Context, ch chan mapResponse, nodes []*Node, index string, shards []uint64, c *pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) error { +func (e *executor) mapper(ctx context.Context, ch chan mapResponse, nodes []*Node, index string, shards []uint64, c *pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) error { // Group shards together by nodes. m, err := e.shardsByNode(nodes, index, shards) if err != nil { @@ -1555,7 +1555,7 @@ func (e *Executor) mapper(ctx context.Context, ch chan mapResponse, nodes []*Nod } // mapperLocal performs map & reduce entirely on the local node. -func (e *Executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFunc, reduceFn reduceFunc) (interface{}, error) { +func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFunc, reduceFn reduceFunc) (interface{}, error) { ch := make(chan mapResponse, len(shards)) for _, shard := range shards { @@ -1592,7 +1592,7 @@ func (e *Executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu } } -func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error { +func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error { var colKey, rowKey, fieldName string if c.Name == "Set" || c.Name == "Clear" || c.Name == "Row" { // Positional args in new PQL syntax require special handling here. @@ -1656,7 +1656,7 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error { return nil } -func (e *Executor) translateResult(index string, idx *Index, call *pql.Call, result interface{}) (interface{}, error) { +func (e *executor) translateResult(index string, idx *Index, call *pql.Call, result interface{}) (interface{}, error) { switch result := result.(type) { case *Row: if idx.Keys() { diff --git a/server.go b/server.go index 47a97bd24..fc85f890f 100644 --- a/server.go +++ b/server.go @@ -54,7 +54,7 @@ type Server struct { cluster *cluster translateFile *TranslateFile diagnostics *DiagnosticsCollector - executor *Executor + executor *executor hosts []string clusterDisabled bool @@ -158,7 +158,7 @@ func OptServerGCNotifier(gcn GCNotifier) ServerOption { func OptServerInternalClient(c InternalClient) ServerOption { return func(s *Server) error { - s.executor = NewExecutor(OptExecutorInternalQueryClient(c)) + s.executor = newExecutor(optExecutorInternalQueryClient(c)) s.defaultClient = c s.cluster.InternalClient = c return nil