diff --git a/disco/disco.go b/disco/disco.go index 7f4519f13..8ad815c28 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -312,6 +312,14 @@ type inMemSchemator struct { schema Schema } +// NewInMemSchemator instantiates an InMemSchemator +// this allows new holders to have thier own, and not rely on a shared instance +func NewInMemSchemator() *inMemSchemator { + return &inMemSchemator{ + schema: make(Schema), + } +} + // Schema is an in-memory implementation of the Schemator Schema method. func (s *inMemSchemator) Schema(ctx context.Context) (Schema, error) { s.mu.RLock() diff --git a/holder.go b/holder.go index 4ecb5ec6e..db38a6db7 100644 --- a/holder.go +++ b/holder.go @@ -9,7 +9,6 @@ import ( "path/filepath" "runtime" "sort" - "strings" "sync" "time" @@ -219,7 +218,7 @@ func DefaultHolderConfig() *HolderConfig { OpenIDAllocator: func(string, bool) (*idAllocator, error) { return &idAllocator{}, nil }, TranslationSyncer: NopTranslationSyncer, Serializer: GobSerializer, - Schemator: disco.InMemSchemator, + Schemator: disco.NewInMemSchemator(), Sharder: disco.InMemSharder, CacheFlushInterval: defaultCacheFlushInterval, StatsClient: stats.NopStatsClient, @@ -335,23 +334,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) @@ -359,11 +342,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") diff --git a/holder_internal_test.go b/holder_internal_test.go index e18704ac8..c9e164f27 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -10,7 +10,7 @@ func mustHolderConfig() *HolderConfig { cfg := DefaultHolderConfig() cfg.StorageConfig.FsyncEnabled = false cfg.RBFConfig.FsyncEnabled = false - cfg.Schemator = disco.InMemSchemator + cfg.Schemator = disco.NewInMemSchemator() cfg.Sharder = disco.InMemSharder return cfg } diff --git a/index.go b/index.go index ba3901d67..0e900cb05 100644 --- a/index.go +++ b/index.go @@ -77,7 +77,7 @@ func NewIndex(holder *Holder, path, name string) (*Index, error) { holder: holder, trackExistence: true, - Schemator: disco.InMemSchemator, + Schemator: disco.NewInMemSchemator(), serializer: NopSerializer, translateStores: make(map[int]TranslateStore), @@ -264,43 +264,26 @@ 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 + if idx == nil { + return nil + } fileLoop: - for _, loopFi := range fis { + for fname, fld := range idx.Fields { 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 - } - - // Decode the CreateFieldMessage from the schema data in order to - // get its metadata. - cfm, err = decodeCreateFieldMessage(i.holder.serializer, fld.Data) - if err != nil { - return errors.Wrap(err, "decoding create field message") - } + // Decode the CreateFieldMessage from the schema data in order to + // get its metadata. + cfm, err = decodeCreateFieldMessage(i.holder.serializer, fld.Data) + if err != nil { + return errors.Wrap(err, "decoding create field message") } indexQueue <- struct{}{} @@ -308,9 +291,9 @@ fileLoop: defer func() { <-indexQueue }() - i.holder.Logger.Debugf("open field: %s", fi.Name()) + i.holder.Logger.Debugf("open field: %s", fname) - _, err := i.openField(&mu, cfm, fi.Name()) + _, err := i.openField(&mu, cfm, fname) if err != nil { return errors.Wrap(err, "opening field") } @@ -319,6 +302,7 @@ fileLoop: }) } } + err = eg.Wait() if err != nil { // Close any fields which got opened, since the overall