From a6297dc48e97a7aa62c9032bf71649d2f58aab0e Mon Sep 17 00:00:00 2001 From: Travis Date: Wed, 10 Feb 2021 23:46:44 -0600 Subject: [PATCH 1/6] WIP: load schema from etcd on holder open; validate indexes, fields, views --- disco/disco.go | 180 ++++++++++++++++++++++++++++++++++++++-- etcd/embed.go | 4 +- holder.go | 60 +++++++------- holder_internal_test.go | 2 + index.go | 55 ++++++++++-- translate.go | 1 - 6 files changed, 252 insertions(+), 50 deletions(-) diff --git a/disco/disco.go b/disco/disco.go index 0725d2d82..a4029c97d 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -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 +} diff --git a/etcd/embed.go b/etcd/embed.go index b4d08dc04..09bd82608 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -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 diff --git a/holder.go b/holder.go index f3c7999e6..a60b571ff 100644 --- a/holder.go +++ b/holder.go @@ -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{} diff --git a/holder_internal_test.go b/holder_internal_test.go index 6c3c02cb3..04ff81c58 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -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 } diff --git a/index.go b/index.go index 83d5d8c44..18ba85385 100644 --- a/index.go +++ b/index.go @@ -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()) diff --git a/translate.go b/translate.go index dced91507..96011d24b 100644 --- a/translate.go +++ b/translate.go @@ -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. From 2f35b51db87cd60c0aa48ecb75fac686158ea759 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 11 Feb 2021 18:05:18 +0100 Subject: [PATCH 2/6] Fix endpoint tests + change BitDepth type to uint64 --- encoding/proto/proto.go | 2 +- executor.go | 2 +- field.go | 20 +++++------ fragment.go | 72 +++++++++++++++++++-------------------- fragment_internal_test.go | 16 ++++----- handler.go | 8 ++--- holder.go | 2 +- index.go | 4 +-- mmap_test.go | 2 +- pilosa.go | 7 ++-- server/handler_test.go | 33 ++++++++++++------ view.go | 10 +++--- 12 files changed, 96 insertions(+), 82 deletions(-) diff --git a/encoding/proto/proto.go b/encoding/proto/proto.go index 7226dce7f..f091c77ae 100644 --- a/encoding/proto/proto.go +++ b/encoding/proto/proto.go @@ -1046,7 +1046,7 @@ func (s Serializer) decodeFieldOptions(options *internal.FieldOptions, m *pilosa s.decodeDecimal(options.Max, &m.Max) m.Base = options.Base m.Scale = options.Scale - m.BitDepth = uint(options.BitDepth) + m.BitDepth = uint64(options.BitDepth) m.TimeQuantum = pilosa.TimeQuantum(options.TimeQuantum) m.Keys = options.Keys m.ForeignIndex = options.ForeignIndex diff --git a/executor.go b/executor.go index b1f7d2188..76974f0e0 100644 --- a/executor.go +++ b/executor.go @@ -4086,7 +4086,7 @@ func (e *executor) executeExtractShard(ctx context.Context, qcx *Qcx, index stri mergeBits(sign, 1<<63, data) // Copy in the significand. - for i := uint(0); i < bsig.BitDepth; i++ { + for i := uint64(0); i < bsig.BitDepth; i++ { bits, err := fragment.row(tx, bsiOffsetBit+uint64(i)) if err != nil { return ExtractedIDMatrix{}, errors.Wrap(err, "loading BSI significand bit from fragment") diff --git a/field.go b/field.go index e6ca53afb..86479f72d 100644 --- a/field.go +++ b/field.go @@ -844,7 +844,7 @@ func (f *Field) loadMeta() error { f.options.Max = max f.options.Base = pb.Base f.options.Scale = pb.Scale - f.options.BitDepth = uint(pb.BitDepth) + f.options.BitDepth = pb.BitDepth f.options.TimeQuantum = TimeQuantum(pb.TimeQuantum) f.options.Keys = pb.Keys f.options.NoStandardView = pb.NoStandardView @@ -1871,9 +1871,9 @@ func (f *Field) importRoaringOverwrite(ctx context.Context, tx Tx, data []byte, return err } - var bitDepth uint + var bitDepth uint64 if maxRowID+1 > bsiOffsetBit { - bitDepth = uint(maxRowID + 1 - bsiOffsetBit) + bitDepth = uint64(maxRowID + 1 - bsiOffsetBit) } bsig := f.bsiGroup(f.name) @@ -1915,7 +1915,7 @@ func (p fieldInfoSlice) Less(i, j int) bool { return p[i].Name < p[j].Name } // FieldOptions represents options to set when initializing a field. type FieldOptions struct { Base int64 `json:"base,omitempty"` - BitDepth uint `json:"bitDepth,omitempty"` + BitDepth uint64 `json:"bitDepth,omitempty"` Min pql.Decimal `json:"min,omitempty"` Max pql.Decimal `json:"max,omitempty"` Scale int64 `json:"scale,omitempty"` @@ -2009,7 +2009,7 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) { return json.Marshal(struct { Type string `json:"type"` Base int64 `json:"base"` - BitDepth uint `json:"bitDepth"` + BitDepth uint64 `json:"bitDepth"` Min pql.Decimal `json:"min"` Max pql.Decimal `json:"max"` Keys bool `json:"keys"` @@ -2028,7 +2028,7 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) { Type string `json:"type"` Base int64 `json:"base"` Scale int64 `json:"scale"` - BitDepth uint `json:"bitDepth"` + BitDepth uint64 `json:"bitDepth"` Min pql.Decimal `json:"min"` Max pql.Decimal `json:"max"` Keys bool `json:"keys"` @@ -2109,7 +2109,7 @@ type bsiGroup struct { Max int64 `json:"max,omitempty"` Base int64 `json:"base,omitempty"` Scale int64 `json:"scale,omitempty"` - BitDepth uint `json:"bitDepth,omitempty"` + BitDepth uint64 `json:"bitDepth,omitempty"` } // baseValue adjusts the value to align with the range for Field for a certain @@ -2209,12 +2209,12 @@ func isValidCacheType(v string) bool { } // bitDepth returns the number of bits required to store a value. -func bitDepth(v uint64) uint { - return uint(bits.Len64(v)) +func bitDepth(v uint64) uint64 { + return uint64(bits.Len64(v)) } // bitDepthInt64 returns the required bit depth for abs(v). -func bitDepthInt64(v int64) uint { +func bitDepthInt64(v int64) uint64 { if v < 0 { return bitDepth(uint64(-v)) } diff --git a/fragment.go b/fragment.go index 2d6f202b9..3fb39c86c 100644 --- a/fragment.go +++ b/fragment.go @@ -965,7 +965,7 @@ func (f *fragment) bit(tx Tx, rowID, columnID uint64) (bool, error) { } // value uses a column of bits to read a multi-bit value. -func (f *fragment) value(tx Tx, columnID uint64, bitDepth uint) (value int64, exists bool, err error) { +func (f *fragment) value(tx Tx, columnID uint64, bitDepth uint64) (value int64, exists bool, err error) { f.mu.Lock() defer f.mu.Unlock() @@ -977,7 +977,7 @@ func (f *fragment) value(tx Tx, columnID uint64, bitDepth uint) (value int64, ex } // Compute other bits into a value. - for i := uint(0); i < bitDepth; i++ { + for i := uint64(0); i < bitDepth; i++ { if v, err := f.bit(tx, uint64(bsiOffsetBit+i), columnID); err != nil { return 0, false, errors.Wrapf(err, "getting value bit %d", i) } else if v { @@ -996,16 +996,16 @@ func (f *fragment) value(tx Tx, columnID uint64, bitDepth uint) (value int64, ex } // clearValue uses a column of bits to clear a multi-bit value. -func (f *fragment) clearValue(tx Tx, columnID uint64, bitDepth uint, value int64) (changed bool, err error) { +func (f *fragment) clearValue(tx Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) { return f.setValueBase(tx, columnID, bitDepth, value, true) } // setValue uses a column of bits to set a multi-bit value. -func (f *fragment) setValue(tx Tx, columnID uint64, bitDepth uint, value int64) (changed bool, err error) { +func (f *fragment) setValue(tx Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) { return f.setValueBase(tx, columnID, bitDepth, value, false) } -func (f *fragment) positionsForValue(columnID uint64, bitDepth uint, value int64, clear bool, toSet, toClear []uint64) ([]uint64, []uint64, error) { +func (f *fragment) positionsForValue(columnID uint64, bitDepth uint64, value int64, clear bool, toSet, toClear []uint64) ([]uint64, []uint64, error) { // Convert value to an unsigned representation. uvalue := uint64(value) if value < 0 { @@ -1030,7 +1030,7 @@ func (f *fragment) positionsForValue(columnID uint64, bitDepth uint, value int64 toSet = append(toSet, bit) } - for i := uint(0); i < bitDepth; i++ { + for i := uint64(0); i < bitDepth; i++ { bit, err := f.pos(uint64(bsiOffsetBit+i), columnID) if err != nil { return toSet, toClear, errors.Wrap(err, "getting pos") @@ -1046,7 +1046,7 @@ func (f *fragment) positionsForValue(columnID uint64, bitDepth uint, value int64 } // TODO get rid of this and use positionsForValue to generate a single write op, and set that with importPositions. -func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint, value int64, clear bool) (changed bool, err error) { +func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint64, value int64, clear bool) (changed bool, err error) { f.mu.Lock() defer f.mu.Unlock() @@ -1073,7 +1073,7 @@ func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint, value uvalue = uint64(-value) } - for i := uint(0); i < bitDepth; i++ { + for i := uint64(0); i < bitDepth; i++ { if uvalue&(1<= 0; i-- { row, err := f.row(tx, uint64(bsiOffsetBit+i)) if err != nil { @@ -1294,7 +1294,7 @@ func (f *fragment) minUnsigned(tx Tx, filter *Row, bitDepth uint) (min int64, co // max returns the max of a given bsiGroup as well as the number of columns involved. // A bitmap can be passed in to optionally filter the computed columns. -func (f *fragment) max(tx Tx, filter *Row, bitDepth uint) (max int64, count uint64, err error) { +func (f *fragment) max(tx Tx, filter *Row, bitDepth uint64) (max int64, count uint64, err error) { consider, err := f.row(tx, bsiExistsBit) if err != nil { return max, count, err @@ -1323,7 +1323,7 @@ func (f *fragment) max(tx Tx, filter *Row, bitDepth uint) (max int64, count uint } // maxUnsigned the highest value without considering the sign bit. Filter is required. -func (f *fragment) maxUnsigned(tx Tx, filter *Row, bitDepth uint) (max int64, count uint64, err error) { +func (f *fragment) maxUnsigned(tx Tx, filter *Row, bitDepth uint64) (max int64, count uint64, err error) { for i := int(bitDepth - 1); i >= 0; i-- { row, err := f.row(tx, uint64(bsiOffsetBit+i)) if err != nil { @@ -1424,7 +1424,7 @@ func (f *fragment) maxRowID(tx Tx) (_ uint64, err error) { } // rangeOp returns bitmaps with a bsiGroup value encoding matching the predicate. -func (f *fragment) rangeOp(tx Tx, op pql.Token, bitDepth uint, predicate int64) (*Row, error) { +func (f *fragment) rangeOp(tx Tx, op pql.Token, bitDepth uint64, predicate int64) (*Row, error) { switch op { case pql.EQ: return f.rangeEQ(tx, bitDepth, predicate) @@ -1450,7 +1450,7 @@ func absInt64(v int64) uint64 { } } -func (f *fragment) rangeEQ(tx Tx, bitDepth uint, predicate int64) (*Row, error) { +func (f *fragment) rangeEQ(tx Tx, bitDepth uint64, predicate int64) (*Row, error) { // Start with set of columns with values set. b, err := f.row(tx, bsiExistsBit) if err != nil { @@ -1458,7 +1458,7 @@ func (f *fragment) rangeEQ(tx Tx, bitDepth uint, predicate int64) (*Row, error) } upredicate := absInt64(predicate) - if uint(bits.Len64(upredicate)) > bitDepth { + if uint64(bits.Len64(upredicate)) > bitDepth { // Predicate is out of range. return NewRow(), nil } @@ -1492,7 +1492,7 @@ func (f *fragment) rangeEQ(tx Tx, bitDepth uint, predicate int64) (*Row, error) return b, nil } -func (f *fragment) rangeNEQ(tx Tx, bitDepth uint, predicate int64) (*Row, error) { +func (f *fragment) rangeNEQ(tx Tx, bitDepth uint64, predicate int64) (*Row, error) { // Start with set of columns with values set. b, err := f.row(tx, bsiExistsBit) if err != nil { @@ -1511,7 +1511,7 @@ func (f *fragment) rangeNEQ(tx Tx, bitDepth uint, predicate int64) (*Row, error) return b, nil } -func (f *fragment) rangeLT(tx Tx, bitDepth uint, predicate int64, allowEquality bool) (*Row, error) { +func (f *fragment) rangeLT(tx Tx, bitDepth uint64, predicate int64, allowEquality bool) (*Row, error) { if predicate == 1 && !allowEquality { predicate, allowEquality = 0, true } @@ -1557,9 +1557,9 @@ func (f *fragment) rangeLT(tx Tx, bitDepth uint, predicate int64, allowEquality } // rangeLTUnsigned returns all bits LT/LTE the predicate without considering the sign bit. -func (f *fragment) rangeLTUnsigned(tx Tx, filter *Row, bitDepth uint, predicate uint64, allowEquality bool) (*Row, error) { +func (f *fragment) rangeLTUnsigned(tx Tx, filter *Row, bitDepth uint64, predicate uint64, allowEquality bool) (*Row, error) { switch { - case uint(bits.Len64(predicate)) > bitDepth: + case uint64(bits.Len64(predicate)) > bitDepth: fallthrough case predicate == (1< bitDepth: + case !allowEquality && uint64(bits.Len64(predicate)) > bitDepth: // The predicate is bigger than the BSI width, so nothing can be bigger. return NewRow(), nil case allowEquality: @@ -1700,7 +1700,7 @@ func (f *fragment) notNull(tx Tx) (*Row, error) { } // rangeBetween returns bitmaps with a bsiGroup value encoding matching any value between predicateMin and predicateMax. -func (f *fragment) rangeBetween(tx Tx, bitDepth uint, predicateMin, predicateMax int64) (*Row, error) { +func (f *fragment) rangeBetween(tx Tx, bitDepth uint64, predicateMin, predicateMax int64) (*Row, error) { b, err := f.row(tx, bsiExistsBit) if err != nil { return nil, err @@ -1749,7 +1749,7 @@ func (f *fragment) rangeBetween(tx Tx, bitDepth uint, predicateMin, predicateMax } // rangeBetweenUnsigned returns BSI columns for a range of values. Disregards the sign bit. -func (f *fragment) rangeBetweenUnsigned(tx Tx, filter *Row, bitDepth uint, predicateMin, predicateMax uint64) (*Row, error) { +func (f *fragment) rangeBetweenUnsigned(tx Tx, filter *Row, bitDepth uint64, predicateMin, predicateMax uint64) (*Row, error) { switch { case predicateMax > (1< Date: Fri, 12 Feb 2021 15:14:42 +0100 Subject: [PATCH 3/6] Remove not needed holder test --- holder_test.go | 25 ------------------------- 1 file changed, 25 deletions(-) diff --git a/holder_test.go b/holder_test.go index dbde1b039..376cec189 100644 --- a/holder_test.go +++ b/holder_test.go @@ -15,7 +15,6 @@ package pilosa_test import ( - "bytes" "context" "math" "os" @@ -32,30 +31,6 @@ import ( ) func TestHolder_Open(t *testing.T) { - t.Run("ErrIndexName", func(t *testing.T) { - h := test.MustOpenHolder(t) - - bufLogger := test.NewBufferLogger() - h.Holder.Logger = bufLogger - - defer h.Close() - - if err := os.Mkdir(h.IndexPath("!"), 0777); err != nil { - t.Fatal(err) - } else if err := h.Holder.Close(); err != nil { - t.Fatal(err) - } - if err := h.Reopen(); err != nil { - t.Fatal(err) - } - - if bufbytes, err := bufLogger.ReadAll(); err != nil { - t.Fatal(err) - } else if !bytes.Contains(bufbytes, []byte("ERROR opening index: !")) { - t.Fatalf("expected log error:\n%s", bufbytes) - } - }) - t.Run("ErrIndexPermission", func(t *testing.T) { if os.Geteuid() == 0 { t.Skip("Skipping permissions test since user is root.") From 8f0270acda319a0a252ca7abe1fceda1eb47c72a Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 12 Feb 2021 18:51:40 -0600 Subject: [PATCH 4/6] adjust openExistenceField() to check on disk first --- executor.go | 2 +- field.go | 8 +++-- holder.go | 10 ++++++- holder_test.go | 9 +++++- index.go | 81 ++++++++++++++++++++++++++++++++++++-------------- server.go | 1 - 6 files changed, 81 insertions(+), 30 deletions(-) diff --git a/executor.go b/executor.go index 76974f0e0..93619c432 100644 --- a/executor.go +++ b/executor.go @@ -5067,7 +5067,7 @@ func (e *executor) executeSet(ctx context.Context, qcx *Qcx, index string, c *pq // Set column on existence field. if ef := idx.existenceField(); ef != nil { // we create tx here, rather than just above, to avoid creating an extra empty shard. - tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) + tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Field: ef, Shard: shard}) if err != nil { return false, err } diff --git a/field.go b/field.go index 86479f72d..34a6aa72e 100644 --- a/field.go +++ b/field.go @@ -752,7 +752,6 @@ func (f *Field) ForeignIndex() string { // openViews opens and initializes the views inside the field. func (f *Field) openViews() error { - view2shards := f.idx.fieldView2shard.getViewsForField(f.name) if view2shards == nil { // no data @@ -1189,8 +1188,11 @@ func (f *Field) createViewIfNotExistsBase(cvm *CreateViewMessage) (*view, bool, defer f.mu.Unlock() // Create the view in etcd as the system of record. - if err := f.persistView(context.Background(), cvm); err != nil { - return nil, false, errors.Wrap(err, "persisting view") + // Don't persist views related to the existence field. + if f.name != existenceFieldName { + if err := f.persistView(context.Background(), cvm); err != nil { + return nil, false, errors.Wrap(err, "persisting view") + } } if view := f.viewMap[cvm.View]; view != nil { diff --git a/holder.go b/holder.go index 1ee947c4b..91d0e1211 100644 --- a/holder.go +++ b/holder.go @@ -1281,7 +1281,15 @@ func (h *Holder) loadView(indexName, fieldName, viewName string) (*view, error) return nil, errors.Wrap(err, "decoding CreateFieldMessage") } - return fld.createViewIfNotExists(cvm.View) + // I think we eventually want to get rid of storing the serialized view in + // etcd because all it keeps is the view name. So in that case we would always + // just use the viewName argument here. + vName := cvm.View + if fieldName == existenceFieldName { + vName = viewName + } + + return fld.createViewIfNotExists(vName) } func (h *Holder) newIndex(path, name string) (*Index, error) { diff --git a/holder_test.go b/holder_test.go index 376cec189..6904cef88 100644 --- a/holder_test.go +++ b/holder_test.go @@ -32,6 +32,7 @@ import ( func TestHolder_Open(t *testing.T) { t.Run("ErrIndexPermission", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") if os.Geteuid() == 0 { t.Skip("Skipping permissions test since user is root.") } @@ -50,10 +51,11 @@ func TestHolder_Open(t *testing.T) { }() if err := h.Reopen(); err == nil || !strings.Contains(err.Error(), "permission denied") { - t.Fatalf("unexpected error: %s", err) + t.Fatalf("unexpected error: %v", err) } }) t.Run("ErrIndexAttrStoreCorrupt", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") h := test.MustOpenHolder(t) defer h.Close() @@ -71,6 +73,7 @@ func TestHolder_Open(t *testing.T) { }) t.Run("ErrFieldPermission", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") if os.Geteuid() == 0 { t.Skip("Skipping permissions test since user is root.") } @@ -94,6 +97,7 @@ func TestHolder_Open(t *testing.T) { } }) t.Run("ErrFieldOptionsCorrupt", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") h := test.MustOpenHolder(t) defer h.Close() @@ -117,6 +121,7 @@ func TestHolder_Open(t *testing.T) { } }) t.Run("ErrFieldAttrStoreCorrupt", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") h := test.MustOpenHolder(t) defer h.Close() @@ -140,6 +145,7 @@ func TestHolder_Open(t *testing.T) { }) t.Run("ErrFragmentStoragePermission", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") roaringOnlyTest(t) if os.Geteuid() == 0 { @@ -177,6 +183,7 @@ func TestHolder_Open(t *testing.T) { } }) t.Run("ErrFragmentStorageCorrupt", func(t *testing.T) { + t.Skip("we don't open the holder directly from disk anymore; we use the etcd schema") roaringOnlyTest(t) h := test.MustOpenHolder(t) diff --git a/index.go b/index.go index b166b9613..ea04dcf45 100644 --- a/index.go +++ b/index.go @@ -330,31 +330,11 @@ fileLoop: }() i.holder.Logger.Debugf("open field: %s", fi.Name()) - mu.Lock() - - // goroutine safe - i.holder.addIndex(i) - - fld, err := i.newField(i.fieldPath(filepath.Base(fi.Name())), filepath.Base(fi.Name())) - fld.createdAt = createdAt - - mu.Unlock() + _, err := i.openField(&mu, createdAt, fi.Name()) if err != nil { - return errors.Wrapf(ErrName, "'%s'", fi.Name()) + return errors.Wrap(err, "opening field") } - // Pass holder through to the field for use in looking - // up a foreign index. - fld.holder = i.holder - - // open the views we have data for. - if err := fld.Open(); err != nil { - return fmt.Errorf("open field: name=%s, err=%s", fld.Name(), err) - } - i.holder.Logger.Debugf("add field to index.fields: %s", fi.Name()) - i.mu.Lock() - i.fields[fld.Name()] = fld - i.mu.Unlock() return nil }) } @@ -371,8 +351,54 @@ fileLoop: return err } +// openField opens the field directory, initializes the field, and adds it to +// the in-memory map of fields maintained by Index. +func (i *Index) openField(mu *sync.Mutex, createdAt int64, file string) (*Field, error) { + mu.Lock() + + // goroutine safe + i.holder.addIndex(i) + + fld, err := i.newField(i.fieldPath(filepath.Base(file)), filepath.Base(file)) + mu.Unlock() + if err != nil { + return nil, errors.Wrapf(ErrName, "'%s'", file) + } + + // Pass holder through to the field for use in looking + // up a foreign index. + fld.holder = i.holder + + fld.createdAt = createdAt + + // open the views we have data for. + if err := fld.Open(); err != nil { + return nil, fmt.Errorf("open field: name=%s, err=%s", fld.Name(), err) + } + + i.holder.Logger.Debugf("add field to index.fields: %s", file) + i.mu.Lock() + i.fields[fld.Name()] = fld + i.mu.Unlock() + + return fld, nil +} + // openExistenceField gets or creates the existence field and associates it to the index. func (i *Index) openExistenceField() error { + // First try opening the existence field from disk. If it doesn't already + // exist on disk, then we fall through to the code path which creates it. + var mu sync.Mutex + fld, err := i.openField(&mu, 0, existenceFieldName) + if err == nil { + i.existenceFld = fld + return nil + } else if errors.Cause(err) != ErrName { + return errors.Wrap(err, "opening existence file") + } + + // If we have gotten here, it means that we couldn't successfully open the + // existence field from disk, so we need to create it. f, err := i.createFieldIfNotExists(existenceFieldName, &FieldOptions{CacheType: CacheTypeNone, CacheSize: 0}) if err != nil { return errors.Wrap(err, "creating existence field") @@ -502,7 +528,9 @@ func (i *Index) Field(name string) *Field { return i.field(name) } -func (i *Index) field(name string) *Field { return i.fields[name] } +func (i *Index) field(name string) *Field { + return i.fields[name] +} // Fields returns a list of all fields in the index. func (i *Index) Fields() []*Field { @@ -692,6 +720,9 @@ func (i *Index) persistField(ctx context.Context, cfm *CreateFieldMessage) error return nil } +// createFieldIfNotExists creates the field if it does not already exist in the +// in-memory index structure. This is not related to whether or not the field +// exists in etcd. func (i *Index) createFieldIfNotExists(name string, opt *FieldOptions) (*Field, error) { i.mu.Lock() defer i.mu.Unlock() @@ -711,6 +742,10 @@ func (i *Index) createFieldIfNotExists(name string, opt *FieldOptions) (*Field, return i.createField(cfm, false) } +// createField, in addition to creating a new Field, calls Field.Open which +// potentially aquires a lock on Index. So until/unless we refactor the +// Index.createField() function call path, we cannot call Index.createField +// while holding an Index lock. func (i *Index) createField(cfm *CreateFieldMessage, broadcast bool) (*Field, error) { opt := cfm.Meta if opt == nil { diff --git a/server.go b/server.go index 2b455085c..2803a2c5d 100644 --- a/server.go +++ b/server.go @@ -557,7 +557,6 @@ func (s *Server) Open() error { if err != nil { return errors.Wrap(err, "starting DisCo") } - _ = initState // Set node ID. s.nodeID = s.disCo.ID() From 2d17405e94c7dfdf150810c4a8f2775838394d2e Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 12 Feb 2021 21:11:45 -0600 Subject: [PATCH 5/6] fix SchemaDetails test --- server/handler_test.go | 27 +++++++++++++++++++++++---- 1 file changed, 23 insertions(+), 4 deletions(-) diff --git a/server/handler_test.go b/server/handler_test.go index c98b07cf0..51e2827a6 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -318,10 +318,29 @@ func TestHandler_Endpoints(t *testing.T) { t.Fatalf("unexpected status code: %d", w.Code) } - body := strings.TrimSpace(w.Body.String()) - target := fmt.Sprintf(`{"indexes":[{"name":"i0","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false},"cardinality":0},{"name":"f1","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false},"cardinality":1}],"shardWidth":%[1]d},{"name":"i1","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false},"cardinality":1}],"shardWidth":%[1]d},{"name":"i2","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":1000,"keys":false},"cardinality":1},{"name":"f1","options":{"type":"int","base":0,"bitDepth":2,"min":-100,"max":100,"keys":false,"foreignIndex":""},"cardinality":4},{"name":"f2","options":{"type":"decimal","base":0,"scale":1,"bitDepth":3,"min":-10,"max":10,"keys":false},"cardinality":5},{"name":"f3","options":{"type":"time","timeQuantum":"YMDH","keys":false,"noStandardView":false},"cardinality":1},{"name":"f4","options":{"type":"mutex","cacheType":"ranked","cacheSize":5000,"keys":false},"cardinality":1},{"name":"f5","options":{"type":"bool"},"cardinality":1}],"shardWidth":%[1]d}]}`, pilosa.ShardWidth) - if body != target { - t.Fatalf("\n%s\n!=\n%s", target, body) + var bodySchema pilosa.Schema + if err := json.Unmarshal(w.Body.Bytes(), + &bodySchema); err != nil { + t.Fatalf("unexpected unmarshalling error: %v", err) + } + // DO NOT COMPARE `CreatedAt` - reset to 0 + for _, i := range bodySchema.Indexes { + i.CreatedAt = 0 + for _, f := range i.Fields { + f.CreatedAt = 0 + } + } + // + + var targetSchema pilosa.Schema + target := fmt.Sprintf(`{"indexes":[{"name":"i0","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false},"cardinality":0},{"name":"f1","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false},"cardinality":1,"views":[{"name":"standard"}]}],"shardWidth":%[1]d},{"name":"i1","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false},"cardinality":1,"views":[{"name":"standard"}]}],"shardWidth":%[1]d},{"name":"i2","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":1000,"keys":false},"cardinality":1,"views":[{"name":"standard"}]},{"name":"f1","options":{"type":"int","base":0,"bitDepth":0,"min":-100,"max":100,"keys":false,"foreignIndex":""},"cardinality":4,"views":[{"name":"bsig_f1"}]},{"name":"f2","options":{"type":"decimal","base":0,"scale":1,"bitDepth":0,"min":-10,"max":10,"keys":false},"cardinality":5,"views":[{"name":"bsig_f2"}]},{"name":"f3","options":{"type":"time","timeQuantum":"YMDH","keys":false,"noStandardView":false},"cardinality":1,"views":[{"name":"standard"}]},{"name":"f4","options":{"type":"mutex","cacheType":"ranked","cacheSize":5000,"keys":false},"cardinality":1,"views":[{"name":"standard"}]},{"name":"f5","options":{"type":"bool"},"cardinality":1,"views":[{"name":"standard"}]}],"shardWidth":%[1]d}]}`, pilosa.ShardWidth) + if err := json.Unmarshal([]byte(target), + &targetSchema); err != nil { + t.Fatalf("unexpected unmarshalling error: %v", err) + } + + if !reflect.DeepEqual(targetSchema, bodySchema) { + t.Fatalf("target: %+v\nbody: %+v\n", targetSchema, bodySchema) } }) From a2a6e91f6d0a7a00c34a2fbef75993f64e699755 Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 12 Feb 2021 21:21:56 -0600 Subject: [PATCH 6/6] remove old, now conflicting test value --- cmd/root_test.go | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/cmd/root_test.go b/cmd/root_test.go index caf53d74b..ec998d346 100644 --- a/cmd/root_test.go +++ b/cmd/root_test.go @@ -189,11 +189,7 @@ bind = "127.0.0.1:10101" [cluster] replicas = 2 - partitions = 128 - hosts = [ - "127.0.0.1:10101", - "127.0.0.1:10111", - ]` + partitions = 128` if _, err := file.Write([]byte(config)); err != nil { t.Fatalf("writing config file: %v", err) }