From 07ca6d57ad19f87412211efe760d81e1ed66f21a Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Wed, 11 Apr 2018 09:42:05 -0500 Subject: [PATCH 1/5] variety of cleanup in api and handler don't export QueryValidationSpecRequired add docs to some methods simplify /debug/vars handling using expvar.Handler() implement Version in API and simplify --- api.go | 22 +++++++++++++++++++--- handler.go | 52 ++++++++++++++++------------------------------------ 2 files changed, 35 insertions(+), 39 deletions(-) diff --git a/api.go b/api.go index 82ab3a4ce..35d478d30 100644 --- a/api.go +++ b/api.go @@ -31,6 +31,8 @@ import ( "github.com/pkg/errors" ) +// 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 // The execution engine for running queries. @@ -46,6 +48,7 @@ type API struct { Logger Logger } +// NewAPI returns a new API instance. func NewAPI() *API { return &API{ Broadcaster: NopBroadcaster, @@ -55,7 +58,8 @@ func NewAPI() *API { } } -func (a *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryResponse, error) { +// ExecuteQuery parses a PQL query out of the request and executes it. +func (api *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryResponse, error) { resp := QueryResponse{} q, err := pql.NewParser(strings.NewReader(req.Query)).Parse() @@ -67,7 +71,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 := api.Executor.Execute(ctx, req.Index, q, req.Slices, execOpts) if err != nil { return resp, err } @@ -86,7 +90,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 := api.readColumnAttrSets(api.Holder.Index(req.Index), columnIDs) if err != nil { return resp, err } @@ -118,6 +122,7 @@ func (api *API) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet return ax, nil } +// CreateIndex makes a new Pilosa index. func (api *API) CreateIndex(ctx context.Context, indexName string, options IndexOptions) (*Index, error) { // Create index. index, err := api.Holder.CreateIndex(indexName, options) @@ -845,6 +850,7 @@ func (api *API) inputJSONDataParser(req map[string]interface{}, index *Index, na return setBits, nil } +// SetCoordinator makes a new Node the cluster coordinator. func (api *API) SetCoordinator(ctx context.Context, id string) (oldNode, newNode *Node, err error) { oldNode = api.Cluster.nodeByID(api.Cluster.Coordinator) newNode = api.Cluster.nodeByID(id) @@ -869,6 +875,8 @@ func (api *API) SetCoordinator(ctx context.Context, id string) (oldNode, newNode return oldNode, newNode, nil } +// RemoveNode puts the cluster into the "RESIZING" state and begins the job of +// removing the given node. func (api *API) RemoveNode(id string) (*Node, error) { removeNode := api.Cluster.nodeByID(id) if removeNode == nil { @@ -891,6 +899,14 @@ func (api *API) ResizeAbort() error { return errors.Wrap(err, "complete current job") } +// State returns the cluster state which is usually "NORMAL", but could be +// "STARTING", "RESIZING", or potentially others. See cluster.go for more +// details. func (api *API) State() string { return api.Cluster.State() } + +// Version returns the Pilosa version. +func (api *API) Version() string { + return strings.TrimPrefix(Version, "v") +} diff --git a/handler.go b/handler.go index d7536f738..4cedeb866 100644 --- a/handler.go +++ b/handler.go @@ -108,14 +108,14 @@ func BuildRouters(handler *Handler) { func (h *Handler) populateValidators() { h.validators = map[string]*queryValidationSpec{} - h.validators["GetFragmentNodes"] = QueryValidationSpecRequired("slice").Optional("index") - h.validators["GetSliceMax"] = QueryValidationSpecRequired().Optional("inverse") - h.validators["PostQuery"] = QueryValidationSpecRequired().Optional("slices", "columnAttrs", "excludeAttrs", "excludeBits") - h.validators["GetExport"] = QueryValidationSpecRequired("index", "frame", "view", "slice") - h.validators["GetFragmentData"] = QueryValidationSpecRequired("index", "frame", "view", "slice") - h.validators["PostFragmentData"] = QueryValidationSpecRequired("index", "frame", "view", "slice") - h.validators["GetFragmentBlocks"] = QueryValidationSpecRequired("index", "frame", "view", "slice") - h.validators["PostFrameRestore"] = QueryValidationSpecRequired("host") + h.validators["GetFragmentNodes"] = queryValidationSpecRequired("slice").Optional("index") + h.validators["GetSliceMax"] = queryValidationSpecRequired().Optional("inverse") + h.validators["PostQuery"] = queryValidationSpecRequired().Optional("slices", "columnAttrs", "excludeAttrs", "excludeBits") + h.validators["GetExport"] = queryValidationSpecRequired("index", "frame", "view", "slice") + h.validators["GetFragmentData"] = queryValidationSpecRequired("index", "frame", "view", "slice") + h.validators["PostFragmentData"] = queryValidationSpecRequired("index", "frame", "view", "slice") + h.validators["GetFragmentBlocks"] = queryValidationSpecRequired("index", "frame", "view", "slice") + h.validators["PostFrameRestore"] = queryValidationSpecRequired("host") } func (h *Handler) queryArgValidator(next http.Handler) http.Handler { @@ -154,7 +154,7 @@ func loadCommon(router *mux.Router, handler *Handler) { router.HandleFunc("/cluster/message", handler.handlePostClusterMessage).Methods("POST") router.HandleFunc("/cluster/resize/set-coordinator", handler.handlePostClusterResizeSetCoordinator).Methods("POST") router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET") - router.HandleFunc("/debug/vars", handler.handleExpvar).Methods("GET") + router.Handle("/debug/vars", expvar.Handler()).Methods("GET") router.HandleFunc("/fragment/data", handler.handleGetFragmentData).Methods("GET").Name("GetFragmentData") router.HandleFunc("/hosts", handler.handleGetHosts).Methods("GET") router.HandleFunc("/id", handler.handleGetID).Methods("GET") @@ -174,7 +174,7 @@ func loadRestricted(router *mux.Router, handler *Handler) { func loadNormal(router *mux.Router, handler *Handler) { router.HandleFunc("/cluster/resize/remove-node", handler.handlePostClusterResizeRemoveNode).Methods("POST") router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET") - router.HandleFunc("/debug/vars", handler.handleExpvar).Methods("GET") + router.Handle("/debug/vars", expvar.Handler()).Methods("GET") router.HandleFunc("/export", handler.handleGetExport).Methods("GET").Name("GetExport") router.HandleFunc("/fragment/block/data", handler.handleGetFragmentBlockData).Methods("GET") router.HandleFunc("/fragment/blocks", handler.handleGetFragmentBlocks).Methods("GET").Name("GetFragmentBlocks") @@ -270,7 +270,7 @@ func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) { } filesystem, err := h.FileSystem.New() if err != nil { - h.writeQueryResponse(w, r, &QueryResponse{Err: err}) + _ = h.writeQueryResponse(w, r, &QueryResponse{Err: err}) h.Logger.Printf("Pilosa WebUI is not available. Please run `make generate-statik` before building Pilosa with `make install`.") return } @@ -1337,36 +1337,16 @@ func (h *Handler) handleGetHosts(w http.ResponseWriter, r *http.Request) { // handleGetVersion handles /version requests. func (h *Handler) handleGetVersion(w http.ResponseWriter, r *http.Request) { - version := Version - if strings.HasPrefix(version, "v") { - // make the version string semver-compatible - version = version[1:] - } - if err := json.NewEncoder(w).Encode(struct { + err := json.NewEncoder(w).Encode(struct { Version string `json:"version"` }{ - Version: version, - }); err != nil { + Version: h.API.Version(), + }) + if err != nil { h.Logger.Printf("write version response error: %s", err) } } -// handleExpvar handles /debug/vars requests. -func (h *Handler) handleExpvar(w http.ResponseWriter, r *http.Request) { - // Copied from $GOROOT/src/expvar/expvar.go - w.Header().Set("Content-Type", "application/json; charset=utf-8") - fmt.Fprintf(w, "{\n") - first := true - expvar.Do(func(kv expvar.KeyValue) { - if !first { - fmt.Fprintf(w, ",\n") - } - first = false - fmt.Fprintf(w, "%q: %s", kv.Key, kv.Value) - }) - fmt.Fprintf(w, "\n}\n") -} - // QueryResult types. const ( QueryResultTypeNil uint32 = iota @@ -1809,7 +1789,7 @@ type queryValidationSpec struct { args map[string]struct{} } -func QueryValidationSpecRequired(requiredArgs ...string) *queryValidationSpec { +func queryValidationSpecRequired(requiredArgs ...string) *queryValidationSpec { args := map[string]struct{}{} for _, arg := range requiredArgs { args[arg] = struct{}{} From feec5f07e1487b74a76a928f2ef45ae50dca9dd6 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Wed, 11 Apr 2018 13:30:35 -0500 Subject: [PATCH 2/5] add more doc comments to api --- api.go | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/api.go b/api.go index 35d478d30..be20e5143 100644 --- a/api.go +++ b/api.go @@ -151,6 +151,8 @@ func (api *API) ReadIndex(ctx context.Context, indexName string) (*Index, error) return index, nil } +// DeleteIndex removes the named index. If the index is not found it does +// nothing and returns no error. func (api *API) DeleteIndex(ctx context.Context, indexName string) error { // Delete index from the holder. err := api.Holder.DeleteIndex(indexName) @@ -170,6 +172,7 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error { return nil } +// CreateFrame makes the named frame in the named index with the given options. func (api *API) CreateFrame(ctx context.Context, indexName string, frameName string, options FrameOptions) (*Frame, error) { // Find index. index := api.Holder.Index(indexName) @@ -198,6 +201,9 @@ func (api *API) CreateFrame(ctx context.Context, indexName string, frameName str return frame, nil } +// DeleteFrame removes the named frame from the named index. If the index is not +// found, an error is returned. If the frame is not found, it is ignored and no +// action is taken. func (api *API) DeleteFrame(ctx context.Context, indexName string, frameName string) error { // Find index. index := api.Holder.Index(indexName) @@ -387,6 +393,8 @@ func (api *API) RestoreFrame(ctx context.Context, indexName string, frameName st return nil } +// ClusterHosts returns a list of the hosts in the cluster including their ID, +// URL, and which is the coordinator. func (api *API) ClusterHosts(ctx context.Context) []*Node { return api.Cluster.Nodes } @@ -481,6 +489,7 @@ func (api *API) WriteInput(ctx context.Context, indexName string, inputDefName s return nil } +// RecalculateCaches forces all TopN caches to be updated. Used mainly for integration tests. func (api *API) RecalculateCaches(ctx context.Context) error { err := api.Broadcaster.SendSync(&internal.RecalculateCaches{}) if err != nil { @@ -498,10 +507,13 @@ func (api *API) PostClusterMessage(ctx context.Context, pb proto.Message) error return nil } +// LocalID returns the current node's ID. func (api *API) LocalID() string { return api.Cluster.Node.ID } +// Schema returns information about each index in Pilosa including which frames +// and views they contain. func (api *API) Schema(ctx context.Context) []*IndexInfo { return api.Holder.Schema() } @@ -595,7 +607,7 @@ func (api *API) DeleteView(ctx context.Context, indexName string, frameName stri // Delete the view. if err := f.DeleteView(viewName); err != nil { - // Ingore this error becuase views do not exist on all nodes due to slice distribution. + // Ignore this error becuase views do not exist on all nodes due to slice distribution. if err != ErrInvalidView { return err } @@ -675,6 +687,7 @@ func (api *API) FrameAttrDiff(ctx context.Context, indexName string, frameName s return attrs, nil } +// Import bulk imports data into a particular index,frame,slice. func (api *API) Import(ctx context.Context, req internal.ImportRequest) error { _, frame, err := api.indexFrame(req.Index, req.Frame, req.Slice) if err != nil { @@ -699,6 +712,7 @@ func (api *API) Import(ctx context.Context, req internal.ImportRequest) error { return err } +// ImportValue bulk imports values into a particular field. func (api *API) ImportValue(ctx context.Context, req internal.ImportValueRequest) error { _, frame, err := api.indexFrame(req.Index, req.Frame, req.Slice) if err != nil { @@ -750,6 +764,8 @@ func (api *API) StatsWithTags(tags []string) StatsClient { return api.Holder.Stats.WithTags(tags...) } +// ClusterLongQueryTime returns the configured threshold for logging/statting +// long running queries. func (api *API) ClusterLongQueryTime() time.Duration { if api.Cluster == nil { return 0 From 9daef2180cd84bebdf827c37440b5b4df0814f7e Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Wed, 11 Apr 2018 13:49:51 -0500 Subject: [PATCH 3/5] rename a number of API methods --- api.go | 22 +++++++++++----------- handler.go | 16 ++++++++-------- 2 files changed, 19 insertions(+), 19 deletions(-) diff --git a/api.go b/api.go index be20e5143..7c570b59d 100644 --- a/api.go +++ b/api.go @@ -58,8 +58,8 @@ func NewAPI() *API { } } -// ExecuteQuery parses a PQL query out of the request and executes it. -func (api *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryResponse, error) { +// Query parses a PQL query out of the request and executes it. +func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, error) { resp := QueryResponse{} q, err := pql.NewParser(strings.NewReader(req.Query)).Parse() @@ -143,7 +143,7 @@ func (api *API) CreateIndex(ctx context.Context, indexName string, options Index return index, nil } -func (api *API) ReadIndex(ctx context.Context, indexName string) (*Index, error) { +func (api *API) Index(ctx context.Context, indexName string) (*Index, error) { index := api.Holder.Index(indexName) if index == nil { return nil, ErrIndexNotFound @@ -393,9 +393,9 @@ func (api *API) RestoreFrame(ctx context.Context, indexName string, frameName st return nil } -// ClusterHosts returns a list of the hosts in the cluster including their ID, +// Hosts returns a list of the hosts in the cluster including their ID, // URL, and which is the coordinator. -func (api *API) ClusterHosts(ctx context.Context) []*Node { +func (api *API) Hosts(ctx context.Context) []*Node { return api.Cluster.Nodes } @@ -522,7 +522,7 @@ func (api *API) Status(ctx context.Context) (proto.Message, error) { return api.StatusHandler.ClusterStatus() } -func (api *API) CreateFrameField(ctx context.Context, indexName string, frameName string, field *Field) error { +func (api *API) CreateField(ctx context.Context, indexName string, frameName string, field *Field) error { // Retrieve frame by name. f := api.Holder.Frame(indexName, frameName) if f == nil { @@ -547,7 +547,7 @@ func (api *API) CreateFrameField(ctx context.Context, indexName string, frameNam return err } -func (api *API) DeleteFrameField(ctx context.Context, indexName string, frameName string, fieldName string) error { +func (api *API) DeleteField(ctx context.Context, indexName string, frameName string, fieldName string) error { // Retrieve frame by name. f := api.Holder.Frame(indexName, frameName) if f == nil { @@ -572,7 +572,7 @@ func (api *API) DeleteFrameField(ctx context.Context, indexName string, frameNam return err } -func (api *API) FrameFields(ctx context.Context, indexName string, frameName string) ([]*Field, error) { +func (api *API) Fields(ctx context.Context, indexName string, frameName string) ([]*Field, error) { index := api.Holder.index(indexName) if index == nil { return nil, ErrIndexNotFound @@ -586,7 +586,7 @@ func (api *API) FrameFields(ctx context.Context, indexName string, frameName str return frame.GetFields() } -func (api *API) FrameViews(ctx context.Context, indexName string, frameName string) ([]*View, error) { +func (api *API) Views(ctx context.Context, indexName string, frameName string) ([]*View, error) { // Retrieve views. f := api.Holder.Frame(indexName, frameName) if f == nil { @@ -764,9 +764,9 @@ func (api *API) StatsWithTags(tags []string) StatsClient { return api.Holder.Stats.WithTags(tags...) } -// ClusterLongQueryTime returns the configured threshold for logging/statting +// LongQueryTime returns the configured threshold for logging/statting // long running queries. -func (api *API) ClusterLongQueryTime() time.Duration { +func (api *API) LongQueryTime() time.Duration { if api.Cluster == nil { return 0 } diff --git a/handler.go b/handler.go index 4cedeb866..208c5922f 100644 --- a/handler.go +++ b/handler.go @@ -241,7 +241,7 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { // Calculate per request StatsD metrics when the handler is fully configured. statsTags := make([]string, 0, 3) - longQueryTime := h.API.ClusterLongQueryTime() + longQueryTime := h.API.LongQueryTime() if longQueryTime > 0 && dif > longQueryTime { h.Logger.Printf("%s %s %v", r.Method, r.URL.String(), dif) statsTags = append(statsTags, "slow_query") @@ -328,7 +328,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.Query(r.Context(), req) if err != nil { w.WriteHeader(http.StatusBadRequest) h.writeQueryResponse(w, r, &QueryResponse{Err: err}) @@ -374,7 +374,7 @@ 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, err := h.API.ReadIndex(r.Context(), indexName) + index, err := h.API.Index(r.Context(), indexName) if err != nil { http.Error(w, err.Error(), http.StatusNotFound) return @@ -742,7 +742,7 @@ func (h *Handler) handlePostFrameField(w http.ResponseWriter, r *http.Request) { Max: req.Max, } - if err := h.API.CreateFrameField(r.Context(), indexName, frameName, field); err != nil { + if err := h.API.CreateField(r.Context(), indexName, frameName, field); err != nil { if err == ErrFrameNotFound { http.Error(w, err.Error(), http.StatusNotFound) } else { @@ -771,7 +771,7 @@ func (h *Handler) handleDeleteFrameField(w http.ResponseWriter, r *http.Request) frameName := mux.Vars(r)["frame"] fieldName := mux.Vars(r)["field"] - if err := h.API.DeleteFrameField(r.Context(), indexName, frameName, fieldName); err != nil { + if err := h.API.DeleteField(r.Context(), indexName, frameName, fieldName); err != nil { if err == ErrFrameNotFound { http.Error(w, err.Error(), http.StatusNotFound) } else { @@ -790,7 +790,7 @@ func (h *Handler) handleGetFrameFields(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] frameName := mux.Vars(r)["frame"] - fields, err := h.API.FrameFields(r.Context(), indexName, frameName) + fields, err := h.API.Fields(r.Context(), indexName, frameName) if err != nil { switch err { case ErrIndexNotFound: @@ -824,7 +824,7 @@ func (h *Handler) handleGetFrameViews(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] frameName := mux.Vars(r)["frame"] - views, err := h.API.FrameViews(r.Context(), indexName, frameName) + views, err := h.API.Views(r.Context(), indexName, frameName) if err != nil { if err == ErrFrameNotFound { http.Error(w, err.Error(), http.StatusNotFound) @@ -1329,7 +1329,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request) // handleGetHosts handles /hosts requests. func (h *Handler) handleGetHosts(w http.ResponseWriter, r *http.Request) { - hosts := h.API.ClusterHosts(r.Context()) + hosts := h.API.Hosts(r.Context()) if err := json.NewEncoder(w).Encode(hosts); err != nil { h.Logger.Printf("write version response error: %s", err) } From 9257029da5fb00736d1ac43d5ec89940d0b3e8eb Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Wed, 11 Apr 2018 14:25:02 -0500 Subject: [PATCH 4/5] simplify Status endpoint had to update test which was relying on a fake ClusterStatus implementation. Now the status endpoint uses information directly from Cluster.Nodes and Cluster.state - which is what Server (the usual ClusterStatus impl) uses, so it should make no difference for real clusters. --- api.go | 4 ---- handler.go | 17 ++++------------- handler_test.go | 3 ++- 3 files changed, 6 insertions(+), 18 deletions(-) diff --git a/api.go b/api.go index 7c570b59d..27827fcec 100644 --- a/api.go +++ b/api.go @@ -518,10 +518,6 @@ func (api *API) Schema(ctx context.Context) []*IndexInfo { return api.Holder.Schema() } -func (api *API) Status(ctx context.Context) (proto.Message, error) { - return api.StatusHandler.ClusterStatus() -} - func (api *API) CreateField(ctx context.Context, indexName string, frameName string, field *Field) error { // Retrieve frame by name. f := api.Holder.Frame(indexName, frameName) diff --git a/handler.go b/handler.go index 208c5922f..4ce8560b5 100644 --- a/handler.go +++ b/handler.go @@ -289,20 +289,11 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { // handleGetStatus handles GET /status requests. func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) { - pb, err := h.API.Status(r.Context()) - if err != nil { - h.Logger.Printf("cluster status error: %s", err) - return + status := getStatusResponse{ + State: h.API.State(), + Nodes: h.API.Hosts(r.Context()), } - - cs, ok := pb.(*internal.ClusterStatus) - if !ok { - panic("status is not a status") - } - if err := json.NewEncoder(w).Encode(getStatusResponse{ - State: cs.State, - Nodes: DecodeNodes(cs.Nodes), - }); err != nil { + if err := json.NewEncoder(w).Encode(status); err != nil { h.Logger.Printf("write status response error: %s", err) } } diff --git a/handler_test.go b/handler_test.go index b605a3650..c7d2a80b7 100644 --- a/handler_test.go +++ b/handler_test.go @@ -141,6 +141,7 @@ func TestHandler_Status(t *testing.T) { h := test.NewHandler() h.API.Holder = hldr.Holder h.API.Cluster = test.NewCluster(1) + h.API.Cluster.SetState(pilosa.ClusterStateNormal) h.API.StatusHandler = s s.Handler = h @@ -148,7 +149,7 @@ func TestHandler_Status(t *testing.T) { h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/status", nil)) if w.Code != http.StatusOK { t.Fatalf("unexpected status code: %d", w.Code) - } else if body := w.Body.String(); body != `{"state":"NORMAL","nodes":[{"id":"test-node","uri":{"scheme":"http","host":"localhost","port":10101},"isCoordinator":false}]}`+"\n" { + } else if body := w.Body.String(); body != `{"state":"NORMAL","nodes":[{"id":"node0","uri":{"scheme":"http","host":"host0"},"isCoordinator":false}]}`+"\n" { t.Fatalf("unexpected body: %s", body) } } From 04f5d5875c9ab49c48321be08bf0c5d2c1e13195 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Thu, 12 Apr 2018 11:28:35 -0500 Subject: [PATCH 5/5] api docs, rename funcs, refactor usage of internal All exported funcs in api.go are now documented Several poorly named methods of API and Cluster were renamed. Particularly, the word Fragment was often changed to Slice in cases where it was really a slice being specified and not a fragment. several methods which received or returned internal data structures have been refactored to be more opaque. Deprecation logging was added to input definition methods. --- api.go | 105 +++++++++++++++++++++++++++++++++++++++++------- cluster.go | 12 +++--- executor.go | 6 +-- fragment.go | 4 +- handler.go | 46 +++++---------------- holder.go | 2 +- pilosa.go | 7 ++++ test/cluster.go | 4 +- 8 files changed, 120 insertions(+), 66 deletions(-) diff --git a/api.go b/api.go index 27827fcec..f57d159f2 100644 --- a/api.go +++ b/api.go @@ -19,6 +19,7 @@ import ( "encoding/csv" "fmt" "io" + "io/ioutil" "net/http" "reflect" "strconv" @@ -143,6 +144,7 @@ func (api *API) CreateIndex(ctx context.Context, indexName string, options Index return index, nil } +// Index retrieves the named index. func (api *API) Index(ctx context.Context, indexName string) (*Index, error) { index := api.Holder.Index(indexName) if index == nil { @@ -230,9 +232,11 @@ func (api *API) DeleteFrame(ctx context.Context, indexName string, frameName str return nil } +// ExportCSV encodes the fragment designated by the index,frame,view,slice as +// CSV of the form , func (api *API) ExportCSV(ctx context.Context, indexName string, frameName string, viewName string, slice uint64, w io.Writer) error { // Validate that this handler owns the slice. - if !api.Cluster.OwnsFragment(api.LocalID(), indexName, slice) { + if !api.Cluster.OwnsSlice(api.LocalID(), indexName, slice) { api.Logger.Printf("host does not own slice %s-%s slice:%d", api.URI, indexName, slice) return ErrClusterDoesNotOwnSlice } @@ -262,11 +266,20 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, frameName strin return nil } -func (api *API) FragmentNodes(ctx context.Context, indexName string, slice uint64) []*Node { - return api.Cluster.FragmentNodes(indexName, slice) +// SliceNodes returns the node and all replicas which should contain a slice's data. +func (api *API) SliceNodes(ctx context.Context, indexName string, slice uint64) []*Node { + return api.Cluster.SliceNodes(indexName, slice) } -func (api *API) FragmentData(ctx context.Context, indexName string, frameName string, viewName string, slice uint64) (*Fragment, error) { +// WriterTo is an interface for any Object which knows how to serialize itself to an io.Writer +type WriterTo interface { + WriteTo(w io.Writer) (n int64, err error) +} + +// MarshalFragment returns an object which can write the specified fragment's data +// to an io.Writer. The serialized data can be read back into a fragment with +// the UnmarshalFragment API call. +func (api *API) MarshalFragment(ctx context.Context, indexName string, frameName string, viewName string, slice uint64) (WriterTo, error) { // Retrieve fragment from holder. f := api.Holder.Fragment(indexName, frameName, viewName, slice) if f == nil { @@ -275,7 +288,10 @@ func (api *API) FragmentData(ctx context.Context, indexName string, frameName st return f, nil } -func (api *API) WriteFragmentData(ctx context.Context, indexName string, frameName string, viewName string, slice uint64, reader io.ReadCloser) error { +// UnmarshalFragment creates a new fragment (if necessary) and reads data from a +// Reader which was previously written by MarshalFragment to populate the +// fragment's data. +func (api *API) UnmarshalFragment(ctx context.Context, indexName string, frameName string, viewName string, slice uint64, reader io.ReadCloser) error { // Retrieve frame. f := api.Holder.Frame(indexName, frameName) if f == nil { @@ -301,19 +317,38 @@ func (api *API) WriteFragmentData(ctx context.Context, indexName string, frameNa return nil } -func (api *API) FragmentBlockData(ctx context.Context, req internal.BlockDataRequest) (internal.BlockDataResponse, error) { +// FragmentBlockData is an endpoint for internal usage. It is not guaranteed to +// return anything useful. Currently it returns protobuf encoded row and column +// ids from a "block" which is a subdivision of a fragment. +func (api *API) FragmentBlockData(ctx context.Context, body io.Reader) ([]byte, error) { + reqBytes, err := ioutil.ReadAll(body) + if err != nil { + return nil, BadRequestError{errors.Wrap(err, "read body error")} + } + var req internal.BlockDataRequest + if err := proto.Unmarshal(reqBytes, &req); err != nil { + return nil, BadRequestError{errors.Wrap(err, "unmarshal body error")} + } + // Retrieve fragment from holder. f := api.Holder.Fragment(req.Index, req.Frame, req.View, req.Slice) if f == nil { - return internal.BlockDataResponse{}, ErrFragmentNotFound + return nil, ErrFragmentNotFound } - // Read data - var resp internal.BlockDataResponse + var resp = internal.BlockDataResponse{} resp.RowIDs, resp.ColumnIDs = f.BlockData(int(req.Block)) - return resp, nil + + // Encode response. + buf, err := proto.Marshal(&resp) + if err != nil { + return nil, errors.Wrap(err, "merge block response encoding error: %s") + } + + return buf, nil } +// FragmentBlocks returns the checksums and block ids for all blocks in the specified fragment. func (api *API) FragmentBlocks(ctx context.Context, indexName string, frameName string, viewName string, slice uint64) ([]FragmentBlock, error) { // Retrieve fragment from holder. f := api.Holder.Fragment(indexName, frameName, viewName, slice) @@ -326,6 +361,8 @@ func (api *API) FragmentBlocks(ctx context.Context, indexName string, frameName return blocks, nil } +// RestoreFrame reads all the data that this host should have for a given frame +// from replicas in the cluster and restores that data to it. func (api *API) RestoreFrame(ctx context.Context, indexName string, frameName string, host *URI) error { // Create a client for the remote cluster. client := NewInternalHTTPClientFromURI(host, api.RemoteClient) @@ -351,7 +388,7 @@ func (api *API) RestoreFrame(ctx context.Context, indexName string, frameName st // Loop over each slice and import it if this node owns it. for slice := uint64(0); slice <= maxSlices[indexName]; slice++ { // Ignore this slice if we don't own it. - if !api.Cluster.OwnsFragment(api.LocalID(), indexName, slice) { + if !api.Cluster.OwnsSlice(api.LocalID(), indexName, slice) { continue } @@ -399,7 +436,10 @@ func (api *API) Hosts(ctx context.Context) []*Node { return api.Cluster.Nodes } +// CreateInputDefinition is deprecated and will be removed. Do not use it. func (api *API) CreateInputDefinition(ctx context.Context, indexName string, inputDefName string, inputDef InputDefinitionInfo) error { + api.Logger.Printf(`CreateInputDefinition is deprecated and will be removed. +Please open an issue if you need to continue using it.`) // Find index. index := api.Holder.Index(indexName) if index == nil { @@ -430,7 +470,9 @@ func (api *API) CreateInputDefinition(ctx context.Context, indexName string, inp return nil } +// InputDefinition is deprecated and will be removed. func (api *API) InputDefinition(ctx context.Context, indexName string, inputDefName string) (*InputDefinition, error) { + api.Logger.Printf(`InputDefinition is deprecated and will be removed.`) // Find index. index := api.Holder.Index(indexName) if index == nil { @@ -444,7 +486,9 @@ func (api *API) InputDefinition(ctx context.Context, indexName string, inputDefN return inputDef, nil } +// DeleteInputDefinition is deprecated and will be removed. func (api *API) DeleteInputDefinition(ctx context.Context, indexName string, inputDefName string) error { + api.Logger.Printf("DeleteInputDefinition is deprecated and will be removed.") // Find index. index := api.Holder.Index(indexName) if index == nil { @@ -467,7 +511,9 @@ func (api *API) DeleteInputDefinition(ctx context.Context, indexName string, inp return nil } +// WriteInput is deprecated and will be removed. func (api *API) WriteInput(ctx context.Context, indexName string, inputDefName string, reqs []interface{}) error { + api.Logger.Printf("WriteInput is deprecated and will be removed.") // Find index. index := api.Holder.Index(indexName) if index == nil { @@ -499,10 +545,24 @@ func (api *API) RecalculateCaches(ctx context.Context) error { return nil } -func (api *API) PostClusterMessage(ctx context.Context, pb proto.Message) error { +// PostClusterMessage is for internal use. It decodes a protobuf message out of +// the body and forwards it to the BroadcastHandler. +func (api *API) PostClusterMessage(ctx context.Context, reqBody io.Reader) error { + // Read entire body. + body, err := ioutil.ReadAll(reqBody) + if err != nil { + return errors.Wrap(err, "reading body") + } + + // Marshal into request object. + pb, err := UnmarshalMessage(body) + if err != nil { + return errors.Wrap(err, "unmarshaling message") + } + // Forward the error message. if err := api.BroadcastHandler.ReceiveMessage(pb); err != nil { - return err + return errors.Wrap(err, "receiving message") } return nil } @@ -518,6 +578,7 @@ func (api *API) Schema(ctx context.Context) []*IndexInfo { return api.Holder.Schema() } +// CreateField creates a new BSI field in the given index and frame. func (api *API) CreateField(ctx context.Context, indexName string, frameName string, field *Field) error { // Retrieve frame by name. f := api.Holder.Frame(indexName, frameName) @@ -543,6 +604,7 @@ func (api *API) CreateField(ctx context.Context, indexName string, frameName str return err } +// DeleteField deletes the given field. func (api *API) DeleteField(ctx context.Context, indexName string, frameName string, fieldName string) error { // Retrieve frame by name. f := api.Holder.Frame(indexName, frameName) @@ -568,6 +630,7 @@ func (api *API) DeleteField(ctx context.Context, indexName string, frameName str return err } +// Fields returns the fields in the given frame. func (api *API) Fields(ctx context.Context, indexName string, frameName string) ([]*Field, error) { index := api.Holder.index(indexName) if index == nil { @@ -582,6 +645,7 @@ func (api *API) Fields(ctx context.Context, indexName string, frameName string) return frame.GetFields() } +// Views returns the views in the given frame. func (api *API) Views(ctx context.Context, indexName string, frameName string) ([]*View, error) { // Retrieve views. f := api.Holder.Frame(indexName, frameName) @@ -594,6 +658,7 @@ func (api *API) Views(ctx context.Context, indexName string, frameName string) ( return views, nil } +// DeleteView removes the given view. func (api *API) DeleteView(ctx context.Context, indexName string, frameName string, viewName string) error { // Retrieve frame. f := api.Holder.Frame(indexName, frameName) @@ -623,6 +688,7 @@ func (api *API) DeleteView(ctx context.Context, indexName string, frameName stri return err } +// IndexAttrDiff func (api *API) IndexAttrDiff(ctx context.Context, indexName string, blocks []AttrBlock) (map[uint64]map[string]interface{}, error) { // Retrieve index from holder. index := api.Holder.Index(indexName) @@ -723,6 +789,7 @@ func (api *API) ImportValue(ctx context.Context, req internal.ImportValueRequest return err } +// ModifyIndexTimeQuantum changes the default time quantum on the given index. func (api *API) ModifyIndexTimeQuantum(ctx context.Context, indexName string, timeQuantum TimeQuantum) error { // Retrieve index by name. index := api.Holder.Index(indexName) @@ -734,6 +801,8 @@ func (api *API) ModifyIndexTimeQuantum(ctx context.Context, indexName string, ti return index.SetTimeQuantum(timeQuantum) } +// ModifyFrameTimeQuantum changes the time quantum on the given frame. TODO: +// what happens if there is already data in the frame? func (api *API) ModifyFrameTimeQuantum(ctx context.Context, indexName string, frameName string, timeQuantum TimeQuantum) error { // Retrieve index by name. frame := api.Holder.Frame(indexName, frameName) @@ -745,14 +814,19 @@ func (api *API) ModifyFrameTimeQuantum(ctx context.Context, indexName string, fr return frame.SetTimeQuantum(timeQuantum) } +// MaxSlices returns the maximum slice number for each index in a map. func (api *API) MaxSlices(ctx context.Context) map[string]uint64 { return api.Holder.MaxSlices() } +// MaxInverseSlices returns the maximum inverse slice number for each index in a +// map. func (api *API) MaxInverseSlices(ctx context.Context) map[string]uint64 { return api.Holder.MaxInverseSlices() } +// 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 { return nil @@ -771,7 +845,7 @@ func (api *API) LongQueryTime() time.Duration { func (api *API) indexFrame(indexName string, frameName string, slice uint64) (*Index, *Frame, error) { // Validate that this handler owns the slice. - if !api.Cluster.OwnsFragment(api.LocalID(), indexName, slice) { + if !api.Cluster.OwnsSlice(api.LocalID(), indexName, slice) { api.Logger.Printf("host does not own slice %s-%s slice:%d", api.URI, indexName, slice) return nil, nil, ErrClusterDoesNotOwnSlice } @@ -793,7 +867,7 @@ func (api *API) indexFrame(indexName string, frameName string, slice uint64) (*I return index, frame, nil } -// inputJSONDataParser validates input json file and executes SetBit. +// inputJSONDataParser validates input json file and executes SetBit. Deprecated - remove with input definition stuff. func (api *API) inputJSONDataParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) { inputDef, err := index.InputDefinition(name) if err != nil { @@ -903,6 +977,7 @@ func (api *API) RemoveNode(id string) (*Node, error) { return removeNode, nil } +// ResizeAbort stops the current resize job. func (api *API) ResizeAbort() error { if !api.Cluster.IsCoordinator() { return ErrNodeNotCoordinator diff --git a/cluster.go b/cluster.go index c490c58cb..abf753892 100644 --- a/cluster.go +++ b/cluster.go @@ -654,7 +654,7 @@ func (c *Cluster) fragsByHost(idx *Index) fragsByHost { func (c *Cluster) fragCombos(idx string, maxSlice uint64, frameViews viewsByFrame) fragsByHost { t := make(fragsByHost) for i := uint64(0); i <= maxSlice; i++ { - nodes := c.FragmentNodes(idx, i) + nodes := c.SliceNodes(idx, i) for _, n := range nodes { // for each frame/view combination: for frame, views := range frameViews { @@ -807,14 +807,14 @@ func (c *Cluster) Partition(index string, slice uint64) int { return int(h.Sum64() % uint64(c.PartitionN)) } -// FragmentNodes returns a list of nodes that own a fragment. -func (c *Cluster) FragmentNodes(index string, slice uint64) []*Node { +// SliceNodes returns a list of nodes that own a fragment. +func (c *Cluster) SliceNodes(index string, slice uint64) []*Node { return c.PartitionNodes(c.Partition(index, slice)) } -// OwnsFragment returns true if a host owns a fragment. -func (c *Cluster) OwnsFragment(nodeID string, index string, slice uint64) bool { - return Nodes(c.FragmentNodes(index, slice)).ContainsID(nodeID) +// OwnsSlice returns true if a host owns a fragment. +func (c *Cluster) OwnsSlice(nodeID string, index string, slice uint64) bool { + return Nodes(c.SliceNodes(index, slice)).ContainsID(nodeID) } // PartitionNodes returns a list of nodes that own a partition. diff --git a/executor.go b/executor.go index 2154b60ff..c2cd860c4 100644 --- a/executor.go +++ b/executor.go @@ -941,7 +941,7 @@ func (e *Executor) executeClearBit(ctx context.Context, index string, c *pql.Cal func (e *Executor) executeClearBitView(ctx context.Context, index string, c *pql.Call, f *Frame, view string, colID, rowID uint64, opt *ExecOptions) (bool, error) { slice := colID / SliceWidth ret := false - for _, node := range e.Cluster.FragmentNodes(index, slice) { + for _, node := range e.Cluster.SliceNodes(index, slice) { // Update locally if host matches. if node.ID == e.Node.ID { val, err := f.ClearBit(view, rowID, colID, nil) @@ -1042,7 +1042,7 @@ func (e *Executor) executeSetBitView(ctx context.Context, index string, c *pql.C slice := colID / SliceWidth ret := false - for _, node := range e.Cluster.FragmentNodes(index, slice) { + for _, node := range e.Cluster.SliceNodes(index, slice) { // Update locally if host matches. if node.ID == e.Node.ID { val, err := f.SetBit(view, rowID, colID, timestamp) @@ -1385,7 +1385,7 @@ func (e *Executor) slicesByNode(nodes []*Node, index string, slices []uint64) (m loop: for _, slice := range slices { - for _, node := range e.Cluster.FragmentNodes(index, slice) { + for _, node := range e.Cluster.SliceNodes(index, slice) { if Nodes(nodes).Contains(node) { m[node] = append(m[node], slice) continue loop diff --git a/fragment.go b/fragment.go index 438a80ae7..c44a07380 100644 --- a/fragment.go +++ b/fragment.go @@ -1702,7 +1702,7 @@ func (s *FragmentSyncer) isClosing() bool { // then merges any blocks which have differences. func (s *FragmentSyncer) SyncFragment() error { // Determine replica set. - nodes := s.Cluster.FragmentNodes(s.Fragment.Index(), s.Fragment.Slice()) + nodes := s.Cluster.SliceNodes(s.Fragment.Index(), s.Fragment.Slice()) if len(nodes) == 1 { return nil } @@ -1784,7 +1784,7 @@ func (s *FragmentSyncer) syncBlock(id int) error { // Read pairs from each remote block. var pairSets []PairSet var clients []InternalClient - for _, node := range s.Cluster.FragmentNodes(f.Index(), f.Slice()) { + for _, node := range s.Cluster.SliceNodes(f.Index(), f.Slice()) { if s.Node.ID == node.ID { continue } diff --git a/handler.go b/handler.go index 4ce8560b5..33d2e5846 100644 --- a/handler.go +++ b/handler.go @@ -1169,7 +1169,7 @@ func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request) } // Retrieve fragment owner nodes. - nodes := h.API.FragmentNodes(r.Context(), index, slice) + nodes := h.API.SliceNodes(r.Context(), index, slice) // Write to response. if err := json.NewEncoder(w).Encode(nodes); err != nil { @@ -1188,7 +1188,7 @@ func (h *Handler) handleGetFragmentData(w http.ResponseWriter, r *http.Request) } // Retrieve fragment from holder. - f, err := h.API.FragmentData(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice) + f, err := h.API.MarshalFragment(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice) if err != nil { http.Error(w, err.Error(), http.StatusNotFound) return @@ -1210,7 +1210,7 @@ func (h *Handler) handlePostFragmentData(w http.ResponseWriter, r *http.Request) return } - if err = h.API.WriteFragmentData(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice, r.Body); err != nil { + if err = h.API.UnmarshalFragment(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice, r.Body); err != nil { if err == ErrFrameNotFound { http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) } else { @@ -1221,19 +1221,11 @@ func (h *Handler) handlePostFragmentData(w http.ResponseWriter, r *http.Request) // handleGetFragmentBlockData handles GET /fragment/block/data requests. func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Request) { - // Read request object. - var req internal.BlockDataRequest - if body, err := ioutil.ReadAll(r.Body); err != nil { - http.Error(w, "ready body error", http.StatusBadRequest) - return - } else if err := proto.Unmarshal(body, &req); err != nil { - http.Error(w, "unmarshal body error", http.StatusBadRequest) - return - } - - resp, err := h.API.FragmentBlockData(r.Context(), req) + buf, err := h.API.FragmentBlockData(r.Context(), r.Body) if err != nil { - if err == ErrFragmentNotFound { + if _, ok := err.(BadRequestError); ok { + http.Error(w, err.Error(), http.StatusBadRequest) + } else if err == ErrFragmentNotFound { http.Error(w, err.Error(), http.StatusNotFound) } else { http.Error(w, err.Error(), http.StatusInternalServerError) @@ -1241,13 +1233,6 @@ func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Requ return } - // Encode response. - buf, err := proto.Marshal(&resp) - if err != nil { - h.Logger.Printf("merge block response encoding error: %s", err) - return - } - // Write response. w.Header().Set("Content-Type", "application/protobuf") w.Header().Set("Content-Length", strconv.Itoa(len(buf))) @@ -1742,23 +1727,10 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques return } - // Read entire body. - body, err := ioutil.ReadAll(r.Body) + err := h.API.PostClusterMessage(r.Context(), r.Body) if err != nil { + // TODO this was the previous behavior, but perhaps not everything is a bad request http.Error(w, err.Error(), http.StatusBadRequest) - return - } - - // Marshal into request object. - pb, err := UnmarshalMessage(body) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } - - if err := h.API.PostClusterMessage(r.Context(), pb); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return } if err := json.NewEncoder(w).Encode(defaultClusterMessageResponse{}); err != nil { diff --git a/holder.go b/holder.go index 22c53c7a4..570433ce4 100644 --- a/holder.go +++ b/holder.go @@ -614,7 +614,7 @@ func (s *HolderSyncer) SyncHolder() error { for slice := uint64(0); slice <= s.Holder.Index(di.Name).MaxSlice(); slice++ { // Ignore slices that this host doesn't own. - if !s.Cluster.OwnsFragment(s.Node.ID, di.Name, slice) { + if !s.Cluster.OwnsSlice(s.Node.ID, di.Name, slice) { continue } diff --git a/pilosa.go b/pilosa.go index 4b28daccb..ffb836ae8 100644 --- a/pilosa.go +++ b/pilosa.go @@ -82,6 +82,13 @@ var ( ErrResizeNotRunning = errors.New("no resize job currently running") ) +// BadRequestError wraps an error value to signify that a request could not be +// read, decoded, or parsed such that in an HTTP scenario, http.StatusBadRequest +// would be returned. +type BadRequestError struct { + error +} + // Regular expression to validate index and frame names. var nameRegexp = regexp.MustCompile(`^[a-z][a-z0-9_-]{0,63}$`) diff --git a/test/cluster.go b/test/cluster.go index 8362ca964..308a7e69c 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -129,7 +129,7 @@ func (t *TestCluster) SetBit(index, frame, view string, rowID, colID uint64, x * // Determine which node should receive the SetBit. c0 := t.Clusters[0] // use the first node's cluster to determine slice location. slice := colID / pilosa.SliceWidth - nodes := c0.FragmentNodes(index, slice) + nodes := c0.SliceNodes(index, slice) for _, node := range nodes { c := t.clusterByID(node.ID) @@ -153,7 +153,7 @@ func (t *TestCluster) SetFieldValue(index, frame string, columnID uint64, name s // Determine which node should receive the SetFieldValue. c0 := t.Clusters[0] // use the first node's cluster to determine slice location. slice := columnID / pilosa.SliceWidth - nodes := c0.FragmentNodes(index, slice) + nodes := c0.SliceNodes(index, slice) for _, node := range nodes { c := t.clusterByID(node.ID)