From 114f6a87514086af57115b0bcd10b4a0d043b93d Mon Sep 17 00:00:00 2001 From: Travis Date: Sun, 7 Feb 2021 13:43:02 -0600 Subject: [PATCH] add withViews argument to api.Schema() method --- api.go | 9 +++++++-- api_test.go | 2 +- cluster.go | 33 +++++++++++++++++++++------------ cmd/random-query/main.go | 2 +- holder.go | 15 ++++++++++----- http/handler.go | 12 +++++++++--- server/grpc.go | 8 ++++---- server/grpc_test.go | 8 ++++---- server/handler_test.go | 2 +- server/server_test.go | 2 +- sql/show.go | 2 +- 11 files changed, 60 insertions(+), 35 deletions(-) diff --git a/api.go b/api.go index 2ce6d2607..39bb6dbcc 100644 --- a/api.go +++ b/api.go @@ -1020,14 +1020,19 @@ func (err MessageProcessingError) Unwrap() error { // Schema returns information about each index in Pilosa including which fields // they contain. -func (api *API) Schema(ctx context.Context) ([]*IndexInfo, error) { +func (api *API) Schema(ctx context.Context, withViews bool) ([]*IndexInfo, error) { if err := api.validate(apiSchema); err != nil { return nil, errors.Wrap(err, "validating api method") } span, _ := tracing.StartSpanFromContext(ctx, "API.Schema") defer span.Finish() - return api.holder.limitedSchema(), nil + + if withViews { + return api.holder.Schema() + } + + return api.holder.limitedSchema() } // ApplySchema takes the given schema and applies it across the diff --git a/api_test.go b/api_test.go index 936e7b8c4..c512131f3 100644 --- a/api_test.go +++ b/api_test.go @@ -269,7 +269,7 @@ func TestAPI_Import(t *testing.T) { // Relies on the previous test creating an index with TrackExistence and // adding some data. t.Run("SchemaHasNoExists", func(t *testing.T) { - schema, err := m1.API.Schema(context.Background()) + schema, err := m1.API.Schema(context.Background(), false) if err != nil { t.Fatal(err) } diff --git a/cluster.go b/cluster.go index 3b46a3bb8..b4c891aa5 100644 --- a/cluster.go +++ b/cluster.go @@ -410,11 +410,15 @@ func (c *cluster) generateResizeInstructionOnAdd(addNodeID string) (*ResizeInstr } myid := c.disCo.ID() + nodeStatus, err := c.nodeStatus() + if err != nil { + return nil, errors.Wrap(err, "getting node status") + } return &ResizeInstruction{ Node: c.unprotectedNodeByID(myid), Sources: fragmentSourcesByNode[myid], TranslationSources: translationSourcesByNode[myid], - NodeStatus: c.nodeStatus(), // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node. + NodeStatus: nodeStatus, // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node. ClusterStatus: status, }, nil } @@ -594,11 +598,15 @@ func (c *cluster) generateResizeInstructionOnRemove(removeNodeID string) (*Resiz } myid := c.disCo.ID() + nodeStatus, err := c.nodeStatus() + if err != nil { + return nil, errors.Wrap(err, "getting node status") + } return &ResizeInstruction{ Node: toCluster.unprotectedNodeByID(myid), Sources: fragmentSourcesByNode[myid], TranslationSources: translationSourcesByNode[myid], - NodeStatus: c.nodeStatus(), // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node. + NodeStatus: nodeStatus, // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node. ClusterStatus: status, }, nil } @@ -610,13 +618,10 @@ func (c *cluster) unprotectedStatus() (*ClusterStatus, error) { return nil, err } - // TODO: replace following code by following code, - // after schemator is implemented - // indexes, err := c.holder.Schema() - // if err != nil { - // return nil, errors.Wrap(err, "getting schema") - // } - indexes := c.holder.Schema() + indexes, err := c.holder.Schema() + if err != nil { + return nil, errors.Wrap(err, "getting schema") + } return &ClusterStatus{ State: string(state), @@ -1272,10 +1277,14 @@ func (c *cluster) SetNodeState(nodeID string, state string) {} /////////////////////////////////////////// -func (c *cluster) nodeStatus() *NodeStatus { +func (c *cluster) nodeStatus() (*NodeStatus, error) { + indexes, err := c.holder.Schema() + if err != nil { + return nil, errors.Wrap(err, "getting schema") + } ns := &NodeStatus{ Node: c.Node, - Schema: &Schema{Indexes: c.holder.Schema()}, + Schema: &Schema{Indexes: indexes}, } var availableShards *roaring.Bitmap for _, idx := range ns.Schema.Indexes { @@ -1294,7 +1303,7 @@ func (c *cluster) nodeStatus() *NodeStatus { } ns.Indexes = append(ns.Indexes, is) } - return ns + return ns, nil } // unprotectedPreviousNode returns the node listed before the current node in c.Nodes. diff --git a/cmd/random-query/main.go b/cmd/random-query/main.go index c271f3e26..ead3ccd9d 100644 --- a/cmd/random-query/main.go +++ b/cmd/random-query/main.go @@ -73,7 +73,7 @@ type wrapper struct { } func (w *wrapper) Schema(ctx context.Context) ([]*pilosa.IndexInfo, error) { - return w.api.Schema(ctx) + return w.api.Schema(ctx, false) } func (w *wrapper) Query(ctx context.Context, index string, queryRequest *pilosa.QueryRequest) (*pilosa.QueryResponse, error) { diff --git a/holder.go b/holder.go index 607707968..4a4c8183a 100644 --- a/holder.go +++ b/holder.go @@ -853,7 +853,7 @@ func (h *Holder) availableShardsByIndex() map[string]*roaring.Bitmap { } // Schema returns schema information for all indexes, fields, and views. -func (h *Holder) Schema() []*IndexInfo { +func (h *Holder) Schema() ([]*IndexInfo, error) { var a []*IndexInfo for _, index := range h.Indexes() { di := &IndexInfo{ @@ -877,11 +877,11 @@ func (h *Holder) Schema() []*IndexInfo { a = append(a, di) } sort.Sort(indexInfoSlice(a)) - return a + return a, nil } // limitedSchema returns schema information for all indexes and fields. -func (h *Holder) limitedSchema() []*IndexInfo { +func (h *Holder) limitedSchema() ([]*IndexInfo, error) { var a []*IndexInfo for _, index := range h.Indexes() { di := &IndexInfo{ @@ -906,7 +906,7 @@ func (h *Holder) limitedSchema() []*IndexInfo { a = append(a, di) } sort.Sort(indexInfoSlice(a)) - return a + return a, nil } // applySchema applies an internal Schema to Holder. @@ -1320,8 +1320,13 @@ func (s *holderSyncer) SyncHolder() error { // Create a snapshot of the cluster to use for node/partition calculations. snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + schema, err := s.Holder.Schema() + if err != nil { + return errors.Wrap(err, "getting schema") + } + // Iterate over schema in sorted order. - for _, di := range s.Holder.Schema() { + for _, di := range schema { // Verify syncer has not closed. if s.IsClosing() { return nil diff --git a/http/handler.go b/http/handler.go index a7291c316..35fb64c4f 100644 --- a/http/handler.go +++ b/http/handler.go @@ -234,7 +234,7 @@ func (h *Handler) populateValidators() { h.validators["PostQuery"] = queryValidationSpecRequired().Optional("shards", "columnAttrs", "excludeRowAttrs", "excludeColumns", "profile") h.validators["GetInfo"] = queryValidationSpecRequired() h.validators["RecalculateCaches"] = queryValidationSpecRequired() - h.validators["GetSchema"] = queryValidationSpecRequired() + h.validators["GetSchema"] = queryValidationSpecRequired().Optional("views") h.validators["PostSchema"] = queryValidationSpecRequired().Optional("remote") h.validators["GetStatus"] = queryValidationSpecRequired() h.validators["GetVersion"] = queryValidationSpecRequired() @@ -667,8 +667,11 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { return } + q := r.URL.Query() + withViews := q.Get("views") == "true" + w.Header().Set("Content-Type", "application/json") - schema, err := h.api.Schema(r.Context()) + schema, err := h.api.Schema(r.Context(), withViews) if err != nil { h.logger.Printf("getting schema error: %s", err) } @@ -980,8 +983,11 @@ func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) { http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable) return } + q := r.URL.Query() + withViews := q.Get("views") == "true" + indexName := mux.Vars(r)["index"] - schema, err := h.api.Schema(r.Context()) + schema, err := h.api.Schema(r.Context(), withViews) if err != nil { h.logger.Printf("getting schema error: %s", err) } diff --git a/server/grpc.go b/server/grpc.go index 40bcc07d0..cdfe8108d 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -316,7 +316,7 @@ func (h *GRPCHandler) CreateIndex(ctx context.Context, req *pb.CreateIndexReques // GetIndex returns a single Index given a name func (h *GRPCHandler) GetIndex(ctx context.Context, req *pb.GetIndexRequest) (*pb.GetIndexResponse, error) { - schema, err := h.api.Schema(ctx) + schema, err := h.api.Schema(ctx, false) if err != nil { return nil, errToStatusError(err) } @@ -331,7 +331,7 @@ func (h *GRPCHandler) GetIndex(ctx context.Context, req *pb.GetIndexRequest) (*p // GetIndexes returns a list of all Indexes func (h *GRPCHandler) GetIndexes(ctx context.Context, req *pb.GetIndexesRequest) (*pb.GetIndexesResponse, error) { - schema, err := h.api.Schema(ctx) + schema, err := h.api.Schema(ctx, false) if err != nil { return nil, errToStatusError(err) } @@ -381,7 +381,7 @@ func (h *VDSMGRPCHandler) GetVDS(ctx context.Context, req *vdsm_pb.GetVDSRequest case *vdsm_pb.GetVDSRequest_Id: return nil, status.Error(codes.InvalidArgument, "VDS IDs are no longer supported") case *vdsm_pb.GetVDSRequest_Name: - schema, err := h.api.Schema(ctx) + schema, err := h.api.Schema(ctx, false) if err != nil { return nil, errToStatusError(err) } @@ -399,7 +399,7 @@ func (h *VDSMGRPCHandler) GetVDS(ctx context.Context, req *vdsm_pb.GetVDSRequest // GetVDSs returns a list of all VDSs func (h *VDSMGRPCHandler) GetVDSs(ctx context.Context, req *vdsm_pb.GetVDSsRequest) (*vdsm_pb.GetVDSsResponse, error) { - schema, err := h.api.Schema(ctx) + schema, err := h.api.Schema(ctx, false) if err != nil { return nil, errToStatusError(err) } diff --git a/server/grpc_test.go b/server/grpc_test.go index d24c2d6e4..0b6178c16 100644 --- a/server/grpc_test.go +++ b/server/grpc_test.go @@ -1035,7 +1035,7 @@ func TestCRUDIndexes(t *testing.T) { t.Fatal(err) } - schema, err := m.API.Schema(ctx) + schema, err := m.API.Schema(ctx, false) if err != nil { t.Fatal("Getting schema error", err) } @@ -1058,7 +1058,7 @@ func TestCRUDIndexes(t *testing.T) { t.Fatal(err) } - schema, err = m.API.Schema(ctx) + schema, err = m.API.Schema(ctx, false) if err != nil { t.Fatal("Getting schema error", err) } @@ -1069,7 +1069,7 @@ func TestCRUDIndexes(t *testing.T) { _ = m.API.DeleteIndex(ctx, "testindex1") - schema, err = m.API.Schema(ctx) + schema, err = m.API.Schema(ctx, false) if err != nil { t.Fatal("Getting schema error", err) } @@ -1183,7 +1183,7 @@ func TestCRUDIndexes(t *testing.T) { t.Fatal(err) } - schema, err := m.API.Schema(ctx) + schema, err := m.API.Schema(ctx, false) if err != nil { t.Fatal("Getting schema error", err) } diff --git a/server/handler_test.go b/server/handler_test.go index 346ddc1b1..a14a04952 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -226,7 +226,7 @@ func TestHandler_Endpoints(t *testing.T) { }) t.Run("Import", func(t *testing.T) { - indexInfo, err := cmd.API.Schema(context.Background()) + indexInfo, err := cmd.API.Schema(context.Background(), false) if err != nil { t.Fatalf("getting schema: %v", err) } diff --git a/server/server_test.go b/server/server_test.go index 04c2da30c..4fcc7f253 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -1250,7 +1250,7 @@ func TestClusterCreatedAtRace(t *testing.T) { schemas := make([]*pilosa.IndexInfo, len(cluster.Nodes)) for i, cmd := range cluster.Nodes { - s, err := cmd.API.Schema(context.Background()) + s, err := cmd.API.Schema(context.Background(), false) if err != nil { t.Fatalf("getting schema: %v", err) } diff --git a/sql/show.go b/sql/show.go index c574ae4a5..0bac67d46 100644 --- a/sql/show.go +++ b/sql/show.go @@ -54,7 +54,7 @@ func (s *ShowHandler) Handle(ctx context.Context, mapped *MappedSQL) (pproto.ToR } func (s *ShowHandler) execShowTables(ctx context.Context, showStmt *sqlparser.Show) (pproto.ToRowser, error) { - indexInfo, err := s.api.Schema(ctx) + indexInfo, err := s.api.Schema(ctx, false) if err != nil { return nil, errors.Wrap(err, "getting schema") }