add withViews argument to api.Schema() method

This commit is contained in:
Travis 2021-02-07 13:43:02 -06:00
parent ecd8e3fa10
commit 114f6a8751
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
11 changed files with 60 additions and 35 deletions

9
api.go
View file

@ -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

View file

@ -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)
}

View file

@ -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.

View file

@ -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) {

View file

@ -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

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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")
}