From 2ebc1914be0865b074520cda3bfd0a21e897d8c7 Mon Sep 17 00:00:00 2001 From: Yuce Tekol Date: Fri, 2 Mar 2018 18:40:50 +0300 Subject: [PATCH] More API functions --- api.go | 133 ++++++++++++++++++++++++++++++++++++++++++++++++----- handler.go | 112 +++++++++++--------------------------------- server.go | 8 ++-- 3 files changed, 150 insertions(+), 103 deletions(-) diff --git a/api.go b/api.go index 5ab4a36aa..e004821fb 100644 --- a/api.go +++ b/api.go @@ -16,23 +16,31 @@ package pilosa import ( "context" + "fmt" + "io/ioutil" + "log" "strings" + "github.com/pilosa/pilosa/internal" "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) + } + Broadcaster Broadcaster + logger *log.Logger } -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 NewAPI(logger *log.Logger) *API { + if logger == nil { + logger = log.New(ioutil.Discard, "", 0) + } + return &API{ + logger: logger, } } @@ -49,7 +57,7 @@ func (a *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryRespons ExcludeAttrs: req.ExcludeAttrs, ExcludeBits: req.ExcludeBits, } - results, err := a.executor.Execute(ctx, req.Index, q, req.Slices, execOpts) + results, err := a.Executor.Execute(ctx, req.Index, q, req.Slices, execOpts) if err != nil { return resp, err } @@ -68,7 +76,7 @@ func (a *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryRespons } // Retrieve column attributes across all calls. - columnAttrSets, err := a.readColumnAttrSets(a.holder.Index(req.Index), columnIDs) + columnAttrSets, err := a.readColumnAttrSets(a.Holder.Index(req.Index), columnIDs) if err != nil { return resp, err } @@ -99,3 +107,104 @@ func (api *API) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet return ax, nil } + +func (api *API) CreateIndex(ctx context.Context, indexName string, options IndexOptions) (*Index, error) { + // Create index. + index, err := api.Holder.CreateIndex(indexName, options) + if err != nil { + return nil, err + } + // Send the create index message to all nodes. + err = api.Broadcaster.SendSync( + &internal.CreateIndexMessage{ + Index: indexName, + Meta: options.Encode(), + }) + if err != nil { + api.logger.Printf("problem sending CreateIndex message: %s", err) + return nil, err + } + api.Holder.Stats.Count("createIndex", 1, 1.0) + return index, nil +} + +func (api *API) ReadIndex(ctx context.Context, indexName string) (*Index, error) { + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + return index, nil +} + +func (api *API) DeleteIndex(ctx context.Context, indexName string) error { + // Delete index from the holder. + err := api.Holder.DeleteIndex(indexName) + if err != nil { + return err + } + // Send the delete index message to all nodes. + err = api.Broadcaster.SendSync( + &internal.DeleteIndexMessage{ + Index: indexName, + }) + if err != nil { + api.logger.Printf("problem sending DeleteIndex message: %s", err) + return err + } + api.Holder.Stats.Count("deleteIndex", 1, 1.0) + return nil +} + +func (api *API) CreateFrame(ctx context.Context, indexName string, frameName string, options FrameOptions) (*Frame, error) { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + // Create frame. + frame, err := index.CreateFrame(frameName, options) + if err != nil { + return nil, err + } + + // Send the create frame message to all nodes. + err = api.Broadcaster.SendSync( + &internal.CreateFrameMessage{ + Index: indexName, + Frame: frameName, + Meta: options.Encode(), + }) + if err != nil { + api.logger.Printf("problem sending CreateFrame message: %s", err) + return nil, err + } + api.Holder.Stats.CountWithCustomTags("createFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) + return frame, nil +} + +func (api *API) DeleteFrame(ctx context.Context, indexName string, frameName string) error { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + // Delete frame from the index. + if err := index.DeleteFrame(frameName); err != nil { + return err + } + + // Send the delete frame message to all nodes. + err := api.Broadcaster.SendSync( + &internal.DeleteFrameMessage{ + Index: indexName, + Frame: frameName, + }) + if err != nil { + api.logger.Printf("problem sending DeleteFrame message: %s", err) + return err + } + api.Holder.Stats.CountWithCustomTags("deleteFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) + return nil +} diff --git a/handler.go b/handler.go index c40148b43..bd47524ed 100644 --- a/handler.go +++ b/handler.go @@ -70,7 +70,7 @@ type Handler struct { // Keeps the query argument validators for each handler validators map[string]*queryValidationSpec - api *API + API *API } // externalPrefixFlag denotes endpoints that are intended to be exposed to clients. @@ -340,7 +340,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { // TODO: Remove req.Index = mux.Vars(r)["index"] - resp, err := h.api.ExecuteQuery(r.Context(), req) + resp, err := h.API.ExecuteQuery(r.Context(), req) if err != nil { w.WriteHeader(http.StatusBadRequest) h.writeQueryResponse(w, r, &resp) @@ -385,9 +385,9 @@ func (h *Handler) handleGetIndexes(w http.ResponseWriter, r *http.Request) { // handleGetIndex handles GET /index/ requests. func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) + index, err := h.API.ReadIndex(r.Context(), indexName) + if err != nil { + http.Error(w, err.Error(), http.StatusNotFound) return } @@ -469,28 +469,17 @@ type postIndexResponse struct{} // handleDeleteIndex handles DELETE /index request. func (h *Handler) handleDeleteIndex(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] - - // Delete index from the holder. - if err := h.Holder.DeleteIndex(indexName); err != nil { + err := h.API.DeleteIndex(r.Context(), indexName) + if err != nil { + h.Logger.Printf("problem deleting index: %s", err) http.Error(w, err.Error(), http.StatusInternalServerError) return } - // Send the delete index message to all nodes. - err := h.Broadcaster.SendSync( - &internal.DeleteIndexMessage{ - Index: indexName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteIndex message: %s", err) - } - // Encode response. if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.Count("deleteIndex", 1, 1.0) } type deleteIndexResponse struct{} @@ -510,8 +499,7 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) { return } - // Create index. - _, err = h.Holder.CreateIndex(indexName, req.Options) + _, err = h.API.CreateIndex(r.Context(), indexName, req.Options) if err == ErrIndexExists { http.Error(w, err.Error(), http.StatusConflict) return @@ -520,24 +508,10 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) { return } - // Send the create index message to all nodes. - err = h.Broadcaster.SendSync( - &internal.CreateIndexMessage{ - Index: indexName, - Meta: req.Options.Encode(), - }) - if err != nil { - h.Logger.Printf("problem sending CreateIndex message: %s", err) - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - // Encode response. if err := json.NewEncoder(w).Encode(postIndexResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.Count("createIndex", 1, 1.0) } // handlePatchIndexTimeQuantum handles PATCH /index/time_quantum request. @@ -655,41 +629,22 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } - - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Create frame. - _, err = index.CreateFrame(frameName, req.Options) - if err == ErrFrameExists { - http.Error(w, err.Error(), http.StatusConflict) - return - } else if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Send the create frame message to all nodes. - err = h.Broadcaster.SendSync( - &internal.CreateFrameMessage{ - Index: indexName, - Frame: frameName, - Meta: req.Options.Encode(), - }) + _, err = h.API.CreateFrame(r.Context(), indexName, frameName, req.Options) if err != nil { - h.Logger.Printf("problem sending CreateFrame message: %s", err) + switch err { + case ErrIndexNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrFrameExists: + http.Error(w, err.Error(), http.StatusConflict) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } + return } - // Encode response. if err := json.NewEncoder(w).Encode(postFrameResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.CountWithCustomTags("createFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) } type _postFrameRequest postFrameRequest @@ -742,37 +697,22 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] frameName := mux.Vars(r)["frame"] - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil { - h.Logger.Printf("response encoding error: %s", err) + err := h.API.DeleteFrame(r.Context(), indexName, frameName) + if err != nil { + if err == ErrIndexNotFound { + if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil { + h.Logger.Printf("response encoding error: %s", err) + } + return } - return - } - - // Delete frame from the index. - if err := index.DeleteFrame(frameName); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } - // Send the delete frame message to all nodes. - err := h.Broadcaster.SendSync( - &internal.DeleteFrameMessage{ - Index: indexName, - Frame: frameName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteFrame message: %s", err) - } - // Encode response. if err := json.NewEncoder(w).Encode(deleteFrameResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.CountWithCustomTags("deleteFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) } type deleteFrameResponse struct{} diff --git a/server.go b/server.go index ede265e19..1973f746e 100644 --- a/server.go +++ b/server.go @@ -117,10 +117,7 @@ func NewServer() *Server { s.Handler.Holder = s.Holder s.diagnostics.server = s - s.Handler.api = &API{ - holder: s.Holder, - } - + s.Handler.API = NewAPI(s.logger) return s } @@ -169,6 +166,7 @@ func (s *Server) Open() error { // Initialize HTTP handler. s.Handler.Broadcaster = s.Broadcaster + s.Handler.API.Broadcaster = s.Broadcaster s.Handler.BroadcastHandler = s s.Handler.StatusHandler = s s.Handler.Node = node @@ -176,7 +174,7 @@ func (s *Server) Open() error { s.Handler.Executor = e s.Cluster.prefect = s.Handler - s.Handler.api.executor = e + s.Handler.API.Executor = e // Initialize Holder. s.Holder.Broadcaster = s.Broadcaster