From 6c3f3ac0f092340338604863d19388acf9bb8a98 Mon Sep 17 00:00:00 2001 From: Yuce Tekol Date: Thu, 1 Mar 2018 17:39:14 +0300 Subject: [PATCH] Added API struct; moved query logic to API --- api.go | 101 +++++++++++++++++++++++++++++++++++++++++++++++++++++ handler.go | 47 ++++--------------------- server.go | 5 +++ 3 files changed, 113 insertions(+), 40 deletions(-) create mode 100644 api.go diff --git a/api.go b/api.go new file mode 100644 index 000000000..5ab4a36aa --- /dev/null +++ b/api.go @@ -0,0 +1,101 @@ +// Copyright 2017 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package pilosa + +import ( + "context" + "strings" + + "github.com/pilosa/pilosa/pql" +) + +type QueryOptions struct { + Remote bool + ExcludeAttrs bool + ExcludeBits bool + ColumnAttrs bool +} + +type API struct { + holder *Holder + // The execution engine for running queries. + executor interface { + Execute(context context.Context, index string, query *pql.Query, slices []uint64, opt *ExecOptions) ([]interface{}, error) + } +} + +func (a *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryResponse, error) { + resp := QueryResponse{} + + q, err := pql.NewParser(strings.NewReader(req.Query)).Parse() + if err != nil { + // TODO: Wrap + return resp, err + } + execOpts := &ExecOptions{ + Remote: req.Remote, + ExcludeAttrs: req.ExcludeAttrs, + ExcludeBits: req.ExcludeBits, + } + results, err := a.executor.Execute(ctx, req.Index, q, req.Slices, execOpts) + if err != nil { + return resp, err + } + resp.Results = results + + // Fill column attributes if requested. + if req.ColumnAttrs && !req.ExcludeBits { + // Consolidate all column ids across all calls. + var columnIDs []uint64 + for _, result := range results { + bm, ok := result.(*Bitmap) + if !ok { + continue + } + columnIDs = uint64Slice(columnIDs).merge(bm.Bits()) + } + + // Retrieve column attributes across all calls. + columnAttrSets, err := a.readColumnAttrSets(a.holder.Index(req.Index), columnIDs) + if err != nil { + return resp, err + } + resp.ColumnAttrSets = columnAttrSets + } + return resp, nil +} + +// readColumnAttrSets returns a list of column attribute objects by id. +func (api *API) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet, error) { + if index == nil { + return nil, nil + } + + ax := make([]*ColumnAttrSet, 0, len(ids)) + for _, id := range ids { + // Read attributes for column. Skip column if empty. + attrs, err := index.ColumnAttrStore().Attrs(id) + if err != nil { + return nil, err + } else if len(attrs) == 0 { + continue + } + + // Append column with attributes. + ax = append(ax, &ColumnAttrSet{ID: id, Attrs: attrs}) + } + + return ax, nil +} diff --git a/handler.go b/handler.go index ffd28da46..c40148b43 100644 --- a/handler.go +++ b/handler.go @@ -69,6 +69,8 @@ type Handler struct { // Keeps the query argument validators for each handler validators map[string]*queryValidationSpec + + api *API } // externalPrefixFlag denotes endpoints that are intended to be exposed to clients. @@ -328,8 +330,6 @@ type getStatusResponse struct { // handlePostQuery handles /query requests. func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { - indexName := mux.Vars(r)["index"] - // Parse incoming request. req, err := h.readQueryRequest(r) if err != nil { @@ -337,46 +337,13 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { h.writeQueryResponse(w, r, &QueryResponse{Err: err}) return } + // TODO: Remove + req.Index = mux.Vars(r)["index"] - // Build execution options. - opt := &ExecOptions{ - Remote: req.Remote, - ExcludeAttrs: req.ExcludeAttrs, - ExcludeBits: req.ExcludeBits, - } - - // Parse query string. - q, err := pql.NewParser(strings.NewReader(req.Query)).Parse() + resp, err := h.api.ExecuteQuery(r.Context(), req) if err != nil { w.WriteHeader(http.StatusBadRequest) - h.writeQueryResponse(w, r, &QueryResponse{Err: err}) - return - } - - // Execute the query. - results, err := h.Executor.Execute(r.Context(), indexName, q, req.Slices, opt) - resp := &QueryResponse{Results: results, Err: err} - - // Fill column attributes if requested. - if req.ColumnAttrs && !req.ExcludeBits { - // Consolidate all column ids across all calls. - var columnIDs []uint64 - for _, result := range results { - bm, ok := result.(*Bitmap) - if !ok { - continue - } - columnIDs = uint64Slice(columnIDs).merge(bm.Bits()) - } - - // Retrieve column attributes across all calls. - columnAttrSets, err := h.readColumnAttrSets(h.Holder.Index(indexName), columnIDs) - if err != nil { - w.WriteHeader(http.StatusInternalServerError) - h.writeQueryResponse(w, r, &QueryResponse{Err: err}) - return - } - resp.ColumnAttrSets = columnAttrSets + h.writeQueryResponse(w, r, &resp) } // Set appropriate status code, if there is an error. @@ -390,7 +357,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } // Write response back to client. - if err := h.writeQueryResponse(w, r, resp); err != nil { + if err := h.writeQueryResponse(w, r, &resp); err != nil { h.Logger.Printf("write query response error: %s", err) } } diff --git a/server.go b/server.go index 08923f13a..ede265e19 100644 --- a/server.go +++ b/server.go @@ -117,6 +117,10 @@ func NewServer() *Server { s.Handler.Holder = s.Holder s.diagnostics.server = s + s.Handler.api = &API{ + holder: s.Holder, + } + return s } @@ -172,6 +176,7 @@ func (s *Server) Open() error { s.Handler.Executor = e s.Cluster.prefect = s.Handler + s.Handler.api.executor = e // Initialize Holder. s.Holder.Broadcaster = s.Broadcaster