From a6a2f84bd54bd2d40cc66b0013dc2737061e48f8 Mon Sep 17 00:00:00 2001 From: Travis Date: Thu, 9 Jan 2020 11:51:06 -0600 Subject: [PATCH] During Holder.Open, apply foreign index after all indexes open In the case where a field with a foreign index opens before the foreign index has opened (and is available as a reference in the holder), push the field into a queue to have its foreign index applied once all indexes have opened. --- executor.go | 11 ++-------- field.go | 43 +++++++++++++++++++++++++++--------- holder.go | 48 ++++++++++++++++++++++++++++++++++++++++ holder_test.go | 59 ++++++++++++++++++++++++++++++++++++++++++++++++++ pilosa.go | 2 ++ 5 files changed, 144 insertions(+), 19 deletions(-) diff --git a/executor.go b/executor.go index b6d62ca03..b8fa0b5d5 100644 --- a/executor.go +++ b/executor.go @@ -3802,15 +3802,8 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res if fieldName := callArgString(call, "field"); fieldName != "" { field := idx.Field(fieldName) - if field != nil { - // Get the foreign index. - if fidx := field.Options().ForeignIndex; fidx != "" { - foreignIndex := e.Holder.Index(fidx) - if foreignIndex == nil { - return nil, errors.Errorf("foreign index does not exist: %s", fidx) - } - store = foreignIndex.translateStore - } + if field != nil && field.Keys() { + store = field.TranslateStore() } } diff --git a/field.go b/field.go index 800c998a0..884b45a7a 100644 --- a/field.go +++ b/field.go @@ -516,19 +516,13 @@ func (f *Field) Open() error { // If the field has a foreign index, and that index uses keys, // then use that index's translateStore instead. if f.options.ForeignIndex != "" { - foreignIndex := f.holder.Index(f.options.ForeignIndex) - if foreignIndex == nil { - return errors.Errorf("foreign index does not exist: %s", f.options.ForeignIndex) - } else if foreignIndex.Keys() { - f.usesKeys = true - f.translateStore = foreignIndex.translateStore + if err := f.holder.checkForeignIndex(f); err != nil { + return errors.Wrap(err, "checking foreign index") } } else { - // Instantiate & open translation store. - if f.translateStore, err = f.OpenTranslateStore(filepath.Join(f.path, "keys"), f.index, f.name); err != nil { - return errors.Wrap(err, "opening translate store") + if err := f.applyTranslateStore(); err != nil { + return errors.Wrap(err, "applying translate store") } - f.usesKeys = f.options.Keys } return nil @@ -541,6 +535,35 @@ func (f *Field) Open() error { return nil } +// applyTranslateStore opens the configured translate store. +func (f *Field) applyTranslateStore() error { + // Instantiate & open translation store. + var err error + f.translateStore, err = f.OpenTranslateStore(filepath.Join(f.path, "keys"), f.index, f.name) + if err != nil { + return errors.Wrap(err, "opening translate store") + } + f.usesKeys = f.options.Keys + return nil +} + +// applyForeignIndex sets the field's translateStore +// to that of a foreign index in the case where the +// foreign index uses keys. If the foreign index does +// not use keys, it falls back to applying the field's +// default translate store. +func (f *Field) applyForeignIndex() error { + foreignIndex := f.holder.Index(f.options.ForeignIndex) + if foreignIndex == nil { + return errors.Wrapf(ErrForeignIndexNotFound, "%s", f.options.ForeignIndex) + } else if foreignIndex.Keys() { + f.usesKeys = true + f.translateStore = foreignIndex.translateStore + return nil + } + return f.applyTranslateStore() +} + var fieldQueue = make(chan struct{}, 16) // openViews opens and initializes the views inside the field. diff --git a/holder.go b/holder.go index 5045bff4d..9cd93bd3e 100644 --- a/holder.go +++ b/holder.go @@ -84,6 +84,16 @@ type Holder struct { // Instantiates new translation stores for indexes & fields. OpenTranslateStore OpenTranslateStoreFunc // local store OpenTranslateReader OpenTranslateReaderFunc // replication + + // Queue of fields (having a foreign index) which have + // opened before their foreign index has opened. + foreignIndexFields []*Field + + // opening is set to true while Holder is opening. + // It's used to determine if foreign index application + // needs to be queued and completed after all indexes + // have opened. + opening bool } // lockedChan looks a little ridiculous admittedly, but exists for good reason. @@ -135,6 +145,9 @@ func NewHolder() *Holder { // Open initializes the root data directory for the holder. func (h *Holder) Open() error { + h.opening = true + defer func() { h.opening = false }() + // Reset closing in case Holder is being reopened. h.closing = make(chan struct{}) @@ -196,6 +209,14 @@ func (h *Holder) Open() error { h.indexes[index.Name()] = index h.mu.Unlock() } + + // If any fields were opened before their foreign index + // was opened, it's safe to process those now since all index + // opens have completed by this point. + if err := h.processForeignIndexFields(); err != nil { + return errors.Wrap(err, "processing foreign index fields") + } + h.Logger.Printf("open holder: complete") // Periodically flush cache. @@ -209,6 +230,33 @@ func (h *Holder) Open() error { return nil } +// checkForeignIndex is a check before applying a foreign +// index to a field; if the index is not yet available, +// (because holder is still opening and may not have opened +// the index yet), this method queues it up to be processed +// once all indexes have been opened. +func (h *Holder) checkForeignIndex(f *Field) error { + if h.opening { + if fi := h.Index(f.options.ForeignIndex); fi == nil { + h.foreignIndexFields = append(h.foreignIndexFields, f) + return nil + } + } + return f.applyForeignIndex() +} + +// processForeignIndexFields applies a foreign index to any +// fields which were opened before their foreign index. +func (h *Holder) processForeignIndexFields() error { + for _, f := range h.foreignIndexFields { + if err := f.applyForeignIndex(); err != nil { + return errors.Wrap(err, "applying foreign index") + } + } + h.foreignIndexFields = h.foreignIndexFields[:0] // reset + return nil +} + // Close closes all open fragments. func (h *Holder) Close() error { h.Stats.Close() diff --git a/holder_test.go b/holder_test.go index 09ec03d97..6d2c5f0c7 100644 --- a/holder_test.go +++ b/holder_test.go @@ -26,6 +26,7 @@ import ( "github.com/pilosa/pilosa/v2" "github.com/pilosa/pilosa/v2/test" + "github.com/pkg/errors" ) func TestHolder_Open(t *testing.T) { @@ -217,6 +218,64 @@ func TestHolder_Open(t *testing.T) { t.Fatalf("unexpected error: %s", err) } }) + + t.Run("ForeignIndex", func(t *testing.T) { + t.Run("ErrForeignIndexNotFound", func(t *testing.T) { + h := test.MustOpenHolder() + defer h.Close() + + if idx, err := h.CreateIndex("foo", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } else { + _, err := idx.CreateField("bar", pilosa.OptFieldTypeInt(0, 100), pilosa.OptFieldForeignIndex("nonexistent")) + if err == nil { + t.Fatalf("expected error: %s", pilosa.ErrForeignIndexNotFound) + } else if errors.Cause(err) != pilosa.ErrForeignIndexNotFound { + t.Fatalf("expected error: %s, but got: %s", pilosa.ErrForeignIndexNotFound, err) + } + } + }) + + // Foreign index zzz is opened after foo/bar. + t.Run("ForeignIndexNotOpenYet", func(t *testing.T) { + h := test.MustOpenHolder() + defer h.Close() + + if _, err := h.CreateIndex("zzz", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } else if idx, err := h.CreateIndex("foo", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } else if _, err := idx.CreateField("bar", pilosa.OptFieldTypeInt(0, 100), pilosa.OptFieldForeignIndex("zzz")); err != nil { + t.Fatal(err) + } else if err := h.Holder.Close(); err != nil { + t.Fatal(err) + } + + if err := h.Reopen(); err != nil { + t.Fatalf("unexpected error: %s", err) + } + }) + + // Foreign index aaa is opened before foo/bar. + t.Run("ForeignIndexIsOpen", func(t *testing.T) { + h := test.MustOpenHolder() + defer h.Close() + + if _, err := h.CreateIndex("aaa", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } else if idx, err := h.CreateIndex("foo", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } else if _, err := idx.CreateField("bar", pilosa.OptFieldTypeInt(0, 100), pilosa.OptFieldForeignIndex("aaa")); err != nil { + t.Fatal(err) + } else if err := h.Holder.Close(); err != nil { + t.Fatal(err) + } + + if err := h.Reopen(); err != nil { + t.Fatalf("unexpected error: %s", err) + } + }) + }) } func TestHolder_HasData(t *testing.T) { diff --git a/pilosa.go b/pilosa.go index 61f3cde53..3f1e8545b 100644 --- a/pilosa.go +++ b/pilosa.go @@ -29,6 +29,8 @@ var ( ErrIndexExists = errors.New("index already exists") ErrIndexNotFound = errors.New("index not found") + ErrForeignIndexNotFound = errors.New("foreign index not found") + // ErrFieldRequired is returned when no field is specified. ErrFieldRequired = errors.New("field required") ErrFieldExists = errors.New("field already exists")