From 91f531f2cdc9145600c38bed0ffd0abf4a41e3a6 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 2 Jul 2018 08:14:13 -0500 Subject: [PATCH] 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