mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 03:47:51 +00:00
Added API struct; moved query logic to API
This commit is contained in:
parent
298047903a
commit
6c3f3ac0f0
3 changed files with 113 additions and 40 deletions
101
api.go
Normal file
101
api.go
Normal file
|
|
@ -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
|
||||
}
|
||||
47
handler.go
47
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)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue