From 4e7c72cc00aaa70d077f0df0744c49ed66c10882 Mon Sep 17 00:00:00 2001 From: kcrodgers24 Date: Fri, 11 Feb 2022 11:47:49 -0800 Subject: [PATCH 1/3] make etcd schema primary source of truth for indexes and fields --- holder.go | 25 ++++----------------- index.go | 66 ++++++++++++++++++++----------------------------------- 2 files changed, 28 insertions(+), 63 deletions(-) diff --git a/holder.go b/holder.go index 18cd7154c..db9aa0cd2 100644 --- a/holder.go +++ b/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") diff --git a/index.go b/index.go index ba3901d67..11c924c08 100644 --- a/index.go +++ b/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() From 7013910158a9b2689afe051742ac19249d54436d Mon Sep 17 00:00:00 2001 From: kcrodgers24 Date: Mon, 14 Feb 2022 09:15:23 -0800 Subject: [PATCH 2/3] give each test its own InMemSchemator --- disco/disco.go | 6 ++++++ holder.go | 2 +- holder_internal_test.go | 2 +- index.go | 2 +- 4 files changed, 9 insertions(+), 3 deletions(-) diff --git a/disco/disco.go b/disco/disco.go index 7f4519f13..6ec6641a0 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -312,6 +312,12 @@ type inMemSchemator struct { schema Schema } +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 db9aa0cd2..e98c4e4bc 100644 --- a/holder.go +++ b/holder.go @@ -215,7 +215,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, 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 11c924c08..0901f1a36 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), From d57050d966cba7137e8b04e2725d7ead71b8abaa Mon Sep 17 00:00:00 2001 From: kcrodgers24 Date: Mon, 14 Feb 2022 10:09:33 -0800 Subject: [PATCH 3/3] requested idx == nil fix; add doc comment --- disco/disco.go | 2 ++ index.go | 58 ++++++++++++++++++++++++++------------------------ 2 files changed, 32 insertions(+), 28 deletions(-) diff --git a/disco/disco.go b/disco/disco.go index 6ec6641a0..8ad815c28 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -312,6 +312,8 @@ 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), diff --git a/index.go b/index.go index 0901f1a36..0e900cb05 100644 --- a/index.go +++ b/index.go @@ -267,40 +267,42 @@ func (i *Index) openFields(idx *disco.Index) error { eg, ctx := errgroup.WithContext(context.Background()) var mu sync.Mutex - if idx != nil { - fileLoop: - for fname, fld := range idx.Fields { - select { - case <-ctx.Done(): - break fileLoop - default: - var cfm *CreateFieldMessage = &CreateFieldMessage{} - var err error + if idx == nil { + return 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. - cfm, err = decodeCreateFieldMessage(i.holder.serializer, fld.Data) + // 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{}{} + 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, "decoding create field message") + return errors.Wrap(err, "opening field") } - 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 - }) - } + return nil + }) } } + err = eg.Wait() if err != nil { // Close any fields which got opened, since the overall