mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 09:05:55 +00:00
make etcd schema primary source of truth for indexes and fields
This commit is contained in:
parent
7a2929d788
commit
4e7c72cc00
2 changed files with 28 additions and 63 deletions
25
holder.go
25
holder.go
|
|
@ -9,7 +9,6 @@ import (
|
|||
"path/filepath"
|
||||
"runtime"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
|
|
@ -332,23 +331,7 @@ func (h *Holder) Open() error {
|
|||
}
|
||||
defer f.Close()
|
||||
|
||||
fis, err := f.Readdir(0)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading directory")
|
||||
}
|
||||
|
||||
for _, fi := range fis {
|
||||
// Skip files or hidden directories.
|
||||
if !fi.IsDir() || strings.HasPrefix(fi.Name(), ".") {
|
||||
continue
|
||||
}
|
||||
|
||||
// Only continue with indexes which are present in schema.
|
||||
idx, ok := schema[fi.Name()]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
for idxKey, idx := range schema {
|
||||
// decode the CreateIndexMessage from the schema data in order to
|
||||
// get its metadata, such as CreateAt.
|
||||
cim, err := decodeCreateIndexMessage(h.serializer, idx.Data)
|
||||
|
|
@ -356,11 +339,11 @@ func (h *Holder) Open() error {
|
|||
return errors.Wrap(err, "decoding create index message")
|
||||
}
|
||||
|
||||
h.Logger.Printf("opening index: %s", filepath.Base(fi.Name()))
|
||||
h.Logger.Printf("opening index: %s", idxKey)
|
||||
|
||||
index, err := h.newIndex(h.IndexPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
|
||||
index, err := h.newIndex(h.IndexPath(idxKey), idxKey)
|
||||
if errors.Cause(err) == ErrName {
|
||||
h.Logger.Errorf("opening index: %s, err=%s", fi.Name(), err)
|
||||
h.Logger.Errorf("opening index: %s, err=%s", idxKey, err)
|
||||
continue
|
||||
} else if err != nil {
|
||||
return errors.Wrap(err, "opening index")
|
||||
|
|
|
|||
66
index.go
66
index.go
|
|
@ -264,36 +264,18 @@ func (i *Index) openFields(idx *disco.Index) error {
|
|||
}
|
||||
defer f.Close()
|
||||
|
||||
fis, err := f.Readdir(0)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading directory")
|
||||
}
|
||||
eg, ctx := errgroup.WithContext(context.Background())
|
||||
var mu sync.Mutex
|
||||
|
||||
fileLoop:
|
||||
for _, loopFi := range fis {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
break fileLoop
|
||||
default:
|
||||
fi := loopFi
|
||||
if !fi.IsDir() {
|
||||
continue
|
||||
}
|
||||
|
||||
var cfm *CreateFieldMessage = &CreateFieldMessage{}
|
||||
var err error
|
||||
|
||||
// Only continue with fields 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 without
|
||||
// having a disco.Index available.
|
||||
if idx != nil {
|
||||
fld, ok := idx.Fields[fi.Name()]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
if idx != nil {
|
||||
fileLoop:
|
||||
for fname, fld := range idx.Fields {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
break fileLoop
|
||||
default:
|
||||
var cfm *CreateFieldMessage = &CreateFieldMessage{}
|
||||
var err error
|
||||
|
||||
// Decode the CreateFieldMessage from the schema data in order to
|
||||
// get its metadata.
|
||||
|
|
@ -301,22 +283,22 @@ fileLoop:
|
|||
if err != nil {
|
||||
return errors.Wrap(err, "decoding create field message")
|
||||
}
|
||||
|
||||
indexQueue <- struct{}{}
|
||||
eg.Go(func() error {
|
||||
defer func() {
|
||||
<-indexQueue
|
||||
}()
|
||||
i.holder.Logger.Debugf("open field: %s", fname)
|
||||
|
||||
_, err := i.openField(&mu, cfm, fname)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "opening field")
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
indexQueue <- struct{}{}
|
||||
eg.Go(func() error {
|
||||
defer func() {
|
||||
<-indexQueue
|
||||
}()
|
||||
i.holder.Logger.Debugf("open field: %s", fi.Name())
|
||||
|
||||
_, err := i.openField(&mu, cfm, fi.Name())
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "opening field")
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
}
|
||||
}
|
||||
err = eg.Wait()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue