WIP: load schema from etcd on holder open; validate indexes, fields, views

This commit is contained in:
Travis 2021-02-10 23:46:44 -06:00
parent b89c699a8e
commit a6297dc48e
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
6 changed files with 252 additions and 50 deletions

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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