From a6297dc48e97a7aa62c9032bf71649d2f58aab0e Mon Sep 17 00:00:00 2001 From: Travis Date: Wed, 10 Feb 2021 23:46:44 -0600 Subject: [PATCH] WIP: load schema from etcd on holder open; validate indexes, fields, views --- disco/disco.go | 180 ++++++++++++++++++++++++++++++++++++++-- etcd/embed.go | 4 +- holder.go | 60 +++++++------- holder_internal_test.go | 2 + index.go | 55 ++++++++++-- translate.go | 1 - 6 files changed, 252 insertions(+), 50 deletions(-) diff --git a/disco/disco.go b/disco/disco.go index 0725d2d82..a4029c97d 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -18,16 +18,21 @@ import ( "context" "fmt" "io" + "sync" "github.com/pilosa/pilosa/v2/roaring" ) var ( - ErrTooManyResults error = fmt.Errorf("too many results") - ErrNoResults error = fmt.Errorf("no results") - ErrKeyDeleted error = fmt.Errorf("key deleted") - ErrIndexExists error = fmt.Errorf("index already exists") - ErrFieldExists error = fmt.Errorf("field already exists") + ErrTooManyResults error = fmt.Errorf("too many results") + ErrNoResults error = fmt.Errorf("no results") + ErrKeyDeleted error = fmt.Errorf("key deleted") + ErrIndexExists error = fmt.Errorf("index already exists") + ErrIndexDoesNotExist error = fmt.Errorf("index does not exist") + ErrFieldExists error = fmt.Errorf("field already exists") + ErrFieldDoesNotExist error = fmt.Errorf("field does not exist") + ErrViewExists error = fmt.Errorf("view already exists") + ErrViewDoesNotExist error = fmt.Errorf("view does not exist") ) type Peer struct { @@ -84,6 +89,10 @@ type Stator interface { NodeStates(context.Context) (map[string]NodeState, error) } +// Schema is a map of all indexes, each of those being a map of fields, then +// views. +type Schema map[string]*Index + // Index is a struct which contains the data encoded for the index as well as // for each of its fields. type Index struct { @@ -99,7 +108,7 @@ type Field struct { } type Schemator interface { - Schema(ctx context.Context) (map[string]*Index, error) + Schema(ctx context.Context) (Schema, error) Index(ctx context.Context, name string) ([]byte, error) CreateIndex(ctx context.Context, name string, val []byte) error DeleteIndex(ctx context.Context, name string) error @@ -252,7 +261,7 @@ var NopSchemator Schemator = &nopSchemator{} type nopSchemator struct{} // Schema is a no-op implementation of the Schemator Schema method. -func (*nopSchemator) Schema(ctx context.Context) (map[string]*Index, error) { return nil, nil } +func (*nopSchemator) Schema(ctx context.Context) (Schema, error) { return nil, nil } // Index is a no-op implementation of the Schemator Index method. func (*nopSchemator) Index(ctx context.Context, name string) ([]byte, error) { return nil, nil } @@ -286,3 +295,160 @@ func (*nopSchemator) CreateView(ctx context.Context, index, field, view string, // DeleteView is a no-op implementation of the Schemator DeleteView method. func (*nopSchemator) DeleteView(ctx context.Context, index, field, view string) error { return nil } + +// InMemSchemator represents a Schemator that manages the schema in memory. The +// intention is that this would be used for testing. +var InMemSchemator Schemator = &inMemSchemator{ + schema: make(Schema), +} + +type inMemSchemator struct { + mu sync.RWMutex + schema Schema +} + +// Schema is an in-memory implementation of the Schemator Schema method. +func (s *inMemSchemator) Schema(ctx context.Context) (Schema, error) { + s.mu.RLock() + defer s.mu.RUnlock() + return s.schema, nil +} + +// Index is an in-memory implementation of the Schemator Index method. +func (s *inMemSchemator) Index(ctx context.Context, name string) ([]byte, error) { + s.mu.RLock() + defer s.mu.RUnlock() + idx, ok := s.schema[name] + if !ok { + return nil, ErrIndexDoesNotExist + } + return idx.Data, nil +} + +// CreateIndex is an in-memory implementation of the Schemator CreateIndex method. +func (s *inMemSchemator) CreateIndex(ctx context.Context, name string, val []byte) error { + s.mu.Lock() + defer s.mu.Unlock() + if idx, ok := s.schema[name]; ok { + // The current logic in pilosa doesn't allow us to return ErrIndexExists + // here, so for now we just update the Data value if the index already + // exists. + idx.Data = val + return nil + } + s.schema[name] = &Index{ + Data: val, + Fields: make(map[string]*Field), + } + return nil +} + +// DeleteIndex is an in-memory implementation of the Schemator DeleteIndex method. +func (s *inMemSchemator) DeleteIndex(ctx context.Context, name string) error { + s.mu.Lock() + defer s.mu.Unlock() + delete(s.schema, name) + return nil +} + +// Field is an in-memory implementation of the Schemator Field method. +func (s *inMemSchemator) Field(ctx context.Context, index, field string) ([]byte, error) { + s.mu.RLock() + defer s.mu.RUnlock() + idx, ok := s.schema[index] + if !ok { + return nil, ErrIndexDoesNotExist + } + fld, ok := idx.Fields[field] + if !ok { + return nil, ErrFieldDoesNotExist + } + return fld.Data, nil +} + +// CreateField is an in-memory implementation of the Schemator CreateField method. +func (s *inMemSchemator) CreateField(ctx context.Context, index, field string, val []byte) error { + s.mu.Lock() + defer s.mu.Unlock() + idx, ok := s.schema[index] + if !ok { + return ErrIndexDoesNotExist + } + if fld, ok := idx.Fields[field]; ok { + // The current logic in pilosa doesn't allow us to return ErrFieldExists + // here, so for now we just update the Data value if the field already + // exists. + fld.Data = val + return nil + } + idx.Fields[field] = &Field{ + Data: val, + Views: make(map[string][]byte), + } + return nil +} + +// DeleteField is an in-memory implementation of the Schemator DeleteField method. +func (s *inMemSchemator) DeleteField(ctx context.Context, index, field string) error { + s.mu.Lock() + defer s.mu.Unlock() + idx, ok := s.schema[index] + if !ok { + return ErrIndexDoesNotExist + } + delete(idx.Fields, field) + return nil +} + +// View is an in-memory implementation of the Schemator View method. +func (s *inMemSchemator) View(ctx context.Context, index, field, view string) ([]byte, error) { + s.mu.RLock() + defer s.mu.RUnlock() + idx, ok := s.schema[index] + if !ok { + return nil, ErrIndexDoesNotExist + } + fld, ok := idx.Fields[field] + if !ok { + return nil, ErrFieldDoesNotExist + } + data, ok := fld.Views[view] + if !ok { + return nil, ErrViewDoesNotExist + } + return data, nil +} + +// CreateView is an in-memory implementation of the Schemator CreateView method. +func (s *inMemSchemator) CreateView(ctx context.Context, index, field, view string, val []byte) error { + s.mu.Lock() + defer s.mu.Unlock() + idx, ok := s.schema[index] + if !ok { + return ErrIndexDoesNotExist + } + fld, ok := idx.Fields[field] + if !ok { + return ErrFieldDoesNotExist + } + // The current logic in pilosa doesn't allow us to return ErrViewExists + // here, so for now we just update the value if the view already exists. + fld.Views[view] = val + return nil +} + +// DeleteView is an in-memory implementation of the Schemator DeleteView method. +func (s *inMemSchemator) DeleteView(ctx context.Context, index, field, view string) error { + s.mu.Lock() + defer s.mu.Unlock() + idx, ok := s.schema[index] + if !ok { + return ErrIndexDoesNotExist + } + fld, ok := idx.Fields[field] + if !ok { + return ErrFieldDoesNotExist + } + delete(fld.Views, view) + return nil +} diff --git a/etcd/embed.go b/etcd/embed.go index b4d08dc04..09bd82608 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -477,7 +477,7 @@ func (e *Etcd) DeleteNode(ctx context.Context, nodeID string) error { return nil } -func (e *Etcd) Schema(ctx context.Context) (map[string]*disco.Index, error) { +func (e *Etcd) Schema(ctx context.Context) (disco.Schema, error) { keys, vals, err := e.getKey(ctx, schemaPrefix) if err != nil { return nil, err @@ -494,7 +494,7 @@ func (e *Etcd) Schema(ctx context.Context) (map[string]*disco.Index, error) { // /index2 // /index2/field1 // - m := make(map[string]*disco.Index) + m := make(disco.Schema) for i, k := range keys { tokens := strings.Split(strings.Trim(k, "/"), "/") // token[0] contains the schemaPrefix diff --git a/holder.go b/holder.go index f3c7999e6..a60b571ff 100644 --- a/holder.go +++ b/holder.go @@ -603,13 +603,6 @@ func (h *Holder) Open() error { return errors.Wrap(err, "creating directory") } - // Verify that we are not trying to open with v1 translation data. - if ok, err := h.hasV1TranslateKeysFile(); err != nil { - return errors.Wrap(err, "verify v1 translation file") - } else if !ok { - return ErrCannotOpenV1TranslateFile - } - tstore, err := h.OpenTransactionStore(h.path) if err != nil { return errors.Wrap(err, "opening transaction store") @@ -623,6 +616,12 @@ func (h *Holder) Open() error { return errors.Wrap(err, "opening ID allocator") } + // Load schema from etcd. + schema, err := h.schemator.Schema(context.Background()) + if err != nil { + return errors.Wrap(err, "getting schema") + } + // Open path to read all index directories. f, err := os.Open(h.path) if err != nil { @@ -645,6 +644,19 @@ func (h *Holder) Open() error { continue } + // Only continue with indexes which are present in schema. + idx, ok := schema[fi.Name()] + if !ok { + continue + } + + // decode the CreateIndexMessage from the schema data in order to + // get its metadata, such as CreateAt. + cim, err := h.decodeCreateIndexMessage(idx.Data) + if err != nil { + return errors.Wrap(err, "decoding create index message") + } + h.Logger.Printf("opening index: %s", filepath.Base(fi.Name())) index, err := h.newIndex(h.IndexPath(filepath.Base(fi.Name())), filepath.Base(fi.Name())) @@ -655,12 +667,16 @@ func (h *Holder) Open() error { return errors.Wrap(err, "opening index") } - if h.isPrimary() { - index.createdAt = timestamp() - err = index.OpenWithTimestamp() - } else { - err = index.Open() - } + // Since we don't have createAt stored on disk within the data + // directory, we need to populate it from the etcd schema data. + // TODO: we may no longer need the createdAt value stored in memory on + // the index struct; it may only be needed in the schema return value + // from the API, which already comes from etcd. In that case, this logic + // could be removed, and the createdAt on the index struct could be + // removed. + index.createdAt = cim.CreatedAt + + err = index.OpenWithSchema(idx) if err != nil { _ = h.txf.Close() if err == ErrName { @@ -841,16 +857,6 @@ func (h *Holder) HasData() (bool, error) { return false, nil } -// hasV1TranslateKeysFile returns true if a v1 translation data file exists on disk. -func (h *Holder) hasV1TranslateKeysFile() (bool, error) { - if _, err := os.Stat(filepath.Join(h.path, ".keys")); os.IsNotExist(err) { - return true, nil - } else if err != nil { - return false, err - } - return false, nil -} - // availableShardsByIndex returns a bitmap of all shards by indexes. func (h *Holder) availableShardsByIndex() map[string]*roaring.Bitmap { m := make(map[string]*roaring.Bitmap) @@ -1415,14 +1421,6 @@ func (h *Holder) recalculateCaches() { } } -// TODO: this needs to be removed -func (h *Holder) isPrimary() bool { - if s, ok := h.broadcaster.(*Server); ok { - return s.IsPrimary() - } - return false -} - // setFileLimit attempts to set the open file limit to the FileLimit constant defined above. func (h *Holder) setFileLimit() { oldLimit := &syscall.Rlimit{} diff --git a/holder_internal_test.go b/holder_internal_test.go index 6c3c02cb3..04ff81c58 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -20,6 +20,7 @@ import ( "os" "testing" + "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/testhook" ) @@ -263,5 +264,6 @@ func mustHolderConfig() *HolderConfig { _ = MustBackendToTxtype(backend) cfg.StorageConfig.Backend = backend } + cfg.Schemator = disco.InMemSchemator return cfg } diff --git a/index.go b/index.go index 83d5d8c44..18ba85385 100644 --- a/index.go +++ b/index.go @@ -177,13 +177,16 @@ func (i *Index) options() IndexOptions { // Open opens and initializes the index. func (i *Index) Open() error { - return i.open(false) + return i.open(nil) } -// OpenWithTimestamp opens and initializes the index and set a new CreatedAt timestamp for fields. -func (i *Index) OpenWithTimestamp() error { return i.open(true) } +// OpenWithSchema opens the index and uses the provided schema to verify that +// the index's fields are expected. +func (i *Index) OpenWithSchema(idx *disco.Index) error { + return i.open(idx) +} -func (i *Index) open(withTimestamp bool) (err error) { +func (i *Index) open(idx *disco.Index) (err error) { // Ensure the path exists. i.holder.Logger.Debugf("ensure index path exists: %s", i.path) if err := os.MkdirAll(i.path, 0777); err != nil { @@ -207,7 +210,7 @@ func (i *Index) open(withTimestamp bool) (err error) { i.fieldView2shard = fieldView2shard i.holder.Logger.Debugf("open fields for index: %s", i.name) - if err := i.openFields(withTimestamp); err != nil { + if err := i.openFields(idx); err != nil { return errors.Wrap(err, "opening fields") } @@ -256,7 +259,7 @@ func (i *Index) open(withTimestamp bool) (err error) { var indexQueue = make(chan struct{}, 8) // openFields opens and initializes the fields inside the index. -func (i *Index) openFields(withTimestamp bool) error { +func (i *Index) openFields(idx *disco.Index) error { f, err := os.Open(i.path) if err != nil { return errors.Wrap(err, "opening directory") @@ -270,6 +273,11 @@ func (i *Index) openFields(withTimestamp bool) error { eg, ctx := errgroup.WithContext(context.Background()) var mu sync.Mutex + // var flds map[string]*disco.Field + // if idx != nil { + // flds = idx.Fields + // } + fileLoop: for _, loopFi := range fis { select { @@ -285,6 +293,36 @@ fileLoop: continue } + var createdAt int64 + + // Only continue with indexes which are present in the provided, + // non-nil index schema. The reason we have to check for idx != nil + // here is because there are tests which call index.Open on an index + // with a NopSchemator. A better approach might be for those tests + // to use a mock Schemator which returns a schema containing the + // index. For an example, see TestField_SetTimeQuantum which + // re-opens a field and curiously has to re-open that field's index + // because at some point we introduced a pointer from the field back + // to its index (possibly related to transactions?). + if idx != nil { + fld, ok := idx.Fields[fi.Name()] + //fld, ok := flds[fi.Name()] + if !ok { + continue + } + + // decode the CreateIndexMessage from the schema data in order to + // get its metadata, such as CreateAt. + // TODO: similar to the createdAt TODO in holder, it may no + // longer be necessary to keep createdAt on the in-memory field + // struct. + cfm, err := i.holder.decodeCreateFieldMessage(fld.Data) + if err != nil { + return errors.Wrap(err, "decoding create field message") + } + createdAt = cfm.CreatedAt + } + indexQueue <- struct{}{} eg.Go(func() error { defer func() { @@ -298,9 +336,8 @@ fileLoop: i.holder.addIndex(i) fld, err := i.newField(i.fieldPath(filepath.Base(fi.Name())), filepath.Base(fi.Name())) - if withTimestamp { - fld.createdAt = timestamp() - } + fld.createdAt = createdAt + mu.Unlock() if err != nil { return errors.Wrapf(ErrName, "'%s'", fi.Name()) diff --git a/translate.go b/translate.go index dced91507..96011d24b 100644 --- a/translate.go +++ b/translate.go @@ -40,7 +40,6 @@ var ( ErrTranslateStoreReadOnly = errors.New("translate store could not find or create key, translate store read only") ErrTranslateStoreNotFound = errors.New("translate store not found") ErrTranslatingKeyNotFound = errors.New("translating key not found") - ErrCannotOpenV1TranslateFile = errors.New("cannot open v1 translate .keys file") ) // TranslateStore is the storage for translation string-to-uint64 values.