Merge pull request #1416 from travisturner/disco-load-schema-on-open

Load schema from etcd on holder open; validate indexes, fields, views
This commit is contained in:
Travis Turner 2021-02-12 21:35:13 -06:00 • committed by GitHub
commit b7db27b7a7
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
19 changed files with 451 additions and 194 deletions

View file

@ -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)
}

View file

@ -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
}

View file

@ -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

View file

@ -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

View file

@ -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")
@ -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
}

View file

@ -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
@ -844,7 +843,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
@ -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 {
@ -1871,9 +1873,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 +1917,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 +2011,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 +2030,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 +2111,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 +2211,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))
}

View file

@ -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<<i) != 0 {
if c, err := f.unprotectedSetBit(tx, uint64(bsiOffsetBit+i), columnID); err != nil {
return err
@ -1125,14 +1125,14 @@ func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint, value
}
// importSetValue is a more efficient SetValue just for imports.
func (f *fragment) importSetValue(txb *TxBitmap, columnID uint64, bitDepth uint, value int64, clear bool) (changed int, err error) { // nolint: unparam
func (f *fragment) importSetValue(txb *TxBitmap, columnID uint64, bitDepth uint64, value int64, clear bool) (changed int, err error) { // nolint: unparam
// Convert value to an unsigned representation.
uvalue := uint64(value)
if value < 0 {
uvalue = uint64(-value)
}
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 changed, errors.Wrap(err, "getting pos")
@ -1194,7 +1194,7 @@ func (f *fragment) importSetValue(txb *TxBitmap, columnID uint64, bitDepth uint,
// sum returns the sum 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) sum(tx Tx, filter *Row, bitDepth uint) (sum int64, count uint64, err error) {
func (f *fragment) sum(tx Tx, filter *Row, bitDepth uint64) (sum int64, count uint64, err error) {
// Compute count based on the existence row.
consider, err := f.row(tx, bsiExistsBit)
if err != nil {
@ -1225,7 +1225,7 @@ func (f *fragment) sum(tx Tx, filter *Row, bitDepth uint) (sum int64, count uint
//
// Execute once for positive numbers and once for negative. Subtract the
// negative sum from the positive sum.
for i := uint(0); i < bitDepth; i++ {
for i := uint64(0); i < bitDepth; i++ {
row, err := f.row(tx, uint64(bsiOffsetBit+i))
if err != nil {
return sum, count, err
@ -1243,7 +1243,7 @@ func (f *fragment) sum(tx Tx, filter *Row, bitDepth uint) (sum int64, count uint
// min returns the min 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) min(tx Tx, filter *Row, bitDepth uint) (min int64, count uint64, err error) {
func (f *fragment) min(tx Tx, filter *Row, bitDepth uint64) (min int64, count uint64, err error) {
consider, err := f.row(tx, bsiExistsBit)
if err != nil {
return min, count, err
@ -1272,7 +1272,7 @@ func (f *fragment) min(tx Tx, filter *Row, bitDepth uint) (min int64, count uint
}
// minUnsigned the lowest value without considering the sign bit. Filter is required.
func (f *fragment) minUnsigned(tx Tx, filter *Row, bitDepth uint) (min int64, count uint64, err error) {
func (f *fragment) minUnsigned(tx Tx, filter *Row, bitDepth uint64) (min int64, count uint64, err error) {
for i := int(bitDepth - 1); i >= 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)-1 && allowEquality:
// This query matches all possible values.
@ -1567,7 +1567,7 @@ func (f *fragment) rangeLTUnsigned(tx Tx, filter *Row, bitDepth uint, predicate
case predicate == (1<<bitDepth)-1 && !allowEquality:
// This query matches everything that is not (1<<bitDepth)-1.
matches := NewRow()
for i := uint(0); i < bitDepth; i++ {
for i := uint64(0); i < bitDepth; i++ {
row, err := f.row(tx, uint64(bsiOffsetBit+i))
if err != nil {
return nil, err
@ -1602,7 +1602,7 @@ func (f *fragment) rangeLTUnsigned(tx Tx, filter *Row, bitDepth uint, predicate
return matched, nil
}
func (f *fragment) rangeGT(tx Tx, bitDepth uint, predicate int64, allowEquality bool) (*Row, error) {
func (f *fragment) rangeGT(tx Tx, bitDepth uint64, predicate int64, allowEquality bool) (*Row, error) {
if predicate == -1 && !allowEquality {
predicate, allowEquality = 0, true
}
@ -1644,7 +1644,7 @@ func (f *fragment) rangeGT(tx Tx, bitDepth uint, predicate int64, allowEquality
}
}
func (f *fragment) rangeGTUnsigned(tx Tx, filter *Row, bitDepth uint, predicate uint64, allowEquality bool) (*Row, error) {
func (f *fragment) rangeGTUnsigned(tx Tx, filter *Row, bitDepth uint64, predicate uint64, allowEquality bool) (*Row, error) {
prep:
switch {
case predicate == 0 && allowEquality:
@ -1653,7 +1653,7 @@ prep:
case predicate == 0 && !allowEquality:
// This query matches everything that is not 0.
matches := NewRow()
for i := uint(0); i < bitDepth; i++ {
for i := uint64(0); i < bitDepth; i++ {
row, err := f.row(tx, uint64(bsiOffsetBit+i))
if err != nil {
return nil, err
@ -1661,7 +1661,7 @@ prep:
matches = matches.Union(filter.Intersect(row))
}
return matches, nil
case !allowEquality && uint(bits.Len64(predicate)) > 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<<bitDepth)-1:
// The upper bound cannot be violated.
@ -1781,11 +1781,11 @@ func (f *fragment) rangeBetweenUnsigned(tx Tx, filter *Row, bitDepth uint, predi
predicateMax &^= equalMask
var err error
remaining, err = f.rangeGTUnsigned(tx, remaining, uint(diffLen), predicateMin, true)
remaining, err = f.rangeGTUnsigned(tx, remaining, uint64(diffLen), predicateMin, true)
if err != nil {
return nil, err
}
remaining, err = f.rangeLTUnsigned(tx, remaining, uint(diffLen), predicateMax, true)
remaining, err = f.rangeLTUnsigned(tx, remaining, uint64(diffLen), predicateMax, true)
if err != nil {
return nil, err
}
@ -2628,7 +2628,7 @@ func (f *fragment) bulkImportMutex(tx Tx, rowIDs, columnIDs []uint64) error {
return errors.Wrap(f.importPositions(tx, toSet, toClear, rowSet), "importing positions")
}
func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int64, bitDepth uint, clear bool) error {
func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int64, bitDepth uint64, clear bool) error {
// TODO figure out how to avoid re-allocating these each time. Probably
// possible to store them on the fragment with a capacity based on
// MaxOpN. For now, we know that the total number of bits to be
@ -2663,7 +2663,7 @@ func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int
return err
}
rowSet := make(map[uint64]struct{}, bitDepth+1)
for i := uint(0); i < bitDepth+1; i++ {
for i := uint64(0); i < bitDepth+1; i++ {
rowSet[uint64(i)] = struct{}{}
}
err := f.importPositions(tx, toSet, toClear, rowSet)
@ -2680,7 +2680,7 @@ func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int
}
// importValue bulk imports a set of range-encoded values.
func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDepth uint, clear bool) error {
func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDepth uint64, clear bool) error {
f.mu.Lock()
defer f.mu.Unlock()
@ -3282,7 +3282,7 @@ func (f *fragment) blockToRoaringData(block int) ([]byte, error) {
// upgradeRoaringBSIv2 upgrades a fragment that contains old BSI formatting
// to a new BSI format (v2). The new format moves the "exists" bit to the
// beginning & adds a negative sign bit.
func upgradeRoaringBSIv2(f *fragment, bitDepth uint) (string, error) {
func upgradeRoaringBSIv2(f *fragment, bitDepth uint64) (string, error) {
// If flag set, already upgraded. Exit.
if f.storage.Flags&roaringFlagBSIv2 == 1 {
return "", nil

View file

@ -456,7 +456,7 @@ func TestFragment_SetValue(t *testing.T) {
})
t.Run("QuickCheck", func(t *testing.T) {
if err := quick.Check(func(bitDepth uint, bitN uint64, values []uint64) bool {
if err := quick.Check(func(bitDepth uint64, bitN uint64, values []uint64) bool {
// Limit bit depth & maximum values.
bitDepth = (bitDepth % 62) + 1
bitN = (bitN % 99) + 1
@ -988,7 +988,7 @@ func TestFragment_Range(t *testing.T) {
// benchmarkSetValues is a helper function to explore, very roughly, the cost
// of setting values.
func benchmarkSetValues(b *testing.B, tx Tx, bitDepth uint, f *fragment, cfunc func(uint64) uint64) {
func benchmarkSetValues(b *testing.B, tx Tx, bitDepth uint64, f *fragment, cfunc func(uint64) uint64) {
column := uint64(0)
for i := 0; i < b.N; i++ {
// We're not checking the error because this is a benchmark.
@ -1000,7 +1000,7 @@ func benchmarkSetValues(b *testing.B, tx Tx, bitDepth uint, f *fragment, cfunc f
// Benchmark performance of setValue for BSI ranges.
func BenchmarkFragment_SetValue(b *testing.B) {
depths := []uint{4, 8, 16}
depths := []uint64{4, 8, 16}
for _, bitDepth := range depths {
name := fmt.Sprintf("Depth%d", bitDepth)
f, idx, tx := mustOpenFragment(b, "i", "f", viewBSIGroupPrefix+"foo", 0, "none")
@ -1021,7 +1021,7 @@ func BenchmarkFragment_SetValue(b *testing.B) {
// benchmarkImportValues is a helper function to explore, very roughly, the cost
// of setting values using the special setter used for imports.
func benchmarkImportValues(b *testing.B, tx Tx, bitDepth uint, f *fragment, cfunc func(uint64) uint64) {
func benchmarkImportValues(b *testing.B, tx Tx, bitDepth uint64, f *fragment, cfunc func(uint64) uint64) {
column := uint64(0)
b.StopTimer()
columns := make([]uint64, b.N)
@ -1040,7 +1040,7 @@ func benchmarkImportValues(b *testing.B, tx Tx, bitDepth uint, f *fragment, cfun
// Benchmark performance of setValue for BSI ranges.
func BenchmarkFragment_ImportValue(b *testing.B) {
depths := []uint{4, 8, 16}
depths := []uint64{4, 8, 16}
for _, bitDepth := range depths {
name := fmt.Sprintf("Depth%d", bitDepth)
f, idx, tx := mustOpenBSIFragment(b, "i", "f", viewBSIGroupPrefix+"foo", 0)
@ -4405,7 +4405,7 @@ func TestFragmentPositionsForValue(t *testing.T) {
tests := []struct {
columnID uint64
bitDepth uint
bitDepth uint64
value int64
clear bool
toSet []uint64
@ -5260,7 +5260,7 @@ func TestImportMultipleValues(t *testing.T) {
vals []int64
checkCols []uint64
checkVals []int64
depth uint
depth uint64
}{
{
cols: []uint64{0, 0},
@ -5312,7 +5312,7 @@ func TestImportValueRowCache(t *testing.T) {
cols []uint64
vals []int64
checkCols []uint64
depth uint
depth uint64
}
tests := []struct {
tc1 testCase

View file

@ -226,11 +226,11 @@ type ImportRequest struct {
// ValidateWithTimestamp ensures that the payload of the request is valid.
func (ir *ImportRequest) ValidateWithTimestamp(indexCreatedAt, fieldCreatedAt int64) error {
if ir.IndexCreatedAt != 0 && ir.FieldCreatedAt != 0 {
if ir.IndexCreatedAt != indexCreatedAt || ir.FieldCreatedAt != fieldCreatedAt {
return ErrPreconditionFailed
}
if (ir.IndexCreatedAt != 0 && ir.IndexCreatedAt != indexCreatedAt) ||
(ir.FieldCreatedAt != 0 && ir.FieldCreatedAt != fieldCreatedAt) {
return ErrPreconditionFailed
}
return nil
}

View file

@ -235,7 +235,7 @@ func DefaultHolderConfig() *HolderConfig {
OpenIDAllocator: func(string) (*idAllocator, error) { return &idAllocator{}, nil },
TranslationSyncer: NopTranslationSyncer,
Serializer: NopSerializer,
Schemator: disco.NopSchemator,
Schemator: disco.InMemSchemator,
CacheFlushInterval: defaultCacheFlushInterval,
StatsClient: stats.NopStatsClient,
NewAttrStore: newNopAttrStore,
@ -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)
@ -1275,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) {
@ -1415,14 +1429,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{}

View file

@ -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
}

View file

@ -15,7 +15,6 @@
package pilosa_test
import (
"bytes"
"context"
"math"
"os"
@ -32,31 +31,8 @@ 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) {
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.")
}
@ -75,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()
@ -96,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.")
}
@ -119,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()
@ -142,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()
@ -165,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 {
@ -202,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)

136
index.go
View file

@ -102,7 +102,7 @@ func NewIndex(holder *Holder, path, name string) (*Index, error) {
holder: holder,
trackExistence: true,
schemator: disco.NopSchemator,
schemator: disco.InMemSchemator,
serializer: NopSerializer,
translateStores: make(map[int]TranslateStore),
@ -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() {
@ -292,32 +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()))
if withTimestamp {
fld.createdAt = timestamp()
}
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
})
}
@ -334,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")
@ -465,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 {
@ -581,7 +646,7 @@ func (i *Index) CreateFieldIfNotExists(name string, opts ...FieldOption) (*Field
cfm := &CreateFieldMessage{
Index: i.name,
Field: name,
CreatedAt: 0,
CreatedAt: timestamp(),
Meta: fo,
}
@ -655,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()
@ -674,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 {

View file

@ -30,7 +30,7 @@ type cv struct {
}
func forceSnapshotsCheckMapping(t *testing.T) {
depth := uint(6)
depth := uint64(6)
f, idx, tx := mustOpenBSIFragment(t, "i", "f", viewStandard, 0)
tx.Rollback()
f.Logger = logger.NewLogfLogger(t)

View file

@ -20,6 +20,7 @@ import (
"regexp"
"time"
"github.com/pilosa/pilosa/v2/disco"
pnet "github.com/pilosa/pilosa/v2/net"
"github.com/pilosa/pilosa/v2/storage"
"github.com/pkg/errors"
@ -30,7 +31,7 @@ var (
ErrHostRequired = errors.New("host required")
ErrIndexRequired = errors.New("index required")
ErrIndexExists = errors.New("index already exists")
ErrIndexExists = disco.ErrIndexExists
ErrIndexNotFound = errors.New("index not found")
ErrForeignIndexNotFound = errors.New("foreign index not found")
@ -38,7 +39,7 @@ var (
// ErrFieldRequired is returned when no field is specified.
ErrFieldRequired = errors.New("field required")
ErrColumnRequired = errors.New("column required")
ErrFieldExists = errors.New("field already exists")
ErrFieldExists = disco.ErrFieldExists
ErrFieldNotFound = errors.New("field not found")
ErrBSIGroupNotFound = errors.New("bsigroup not found")
@ -54,7 +55,7 @@ var (
ErrDecimalOutOfRange = errors.New("decimal value out of range")
ErrViewRequired = errors.New("view required")
ErrViewExists = errors.New("view already exists")
ErrViewExists = disco.ErrViewExists
ErrInvalidView = errors.New("invalid view")
ErrInvalidCacheType = errors.New("invalid cache type")

View file

@ -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()

View file

@ -231,10 +231,28 @@ 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}},{"name":"f1","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}}],"shardWidth":%[1]d},{"name":"i1","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}}],"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
if err := json.Unmarshal([]byte(fmt.Sprintf(`{"indexes":[{"name":"i0","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}},{"name":"f1","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}}],"shardWidth":%d},{"name":"i1","options":{"keys":false,"trackExistence":false},"fields":[{"name":"f0","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}}],"shardWidth":%[1]d}]}`, pilosa.ShardWidth)),
&targetSchema); err != nil {
t.Fatalf("unexpected unmarshalling error: %v", err)
}
if !reflect.DeepEqual(targetSchema, bodySchema) {
t.Fatalf("target: %+v\nbody: %+v\n", targetSchema, bodySchema)
}
})
@ -300,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)
}
})
@ -313,11 +350,6 @@ func TestHandler_Endpoints(t *testing.T) {
t.Fatalf("getting schema: %v", err)
}
err = cmd.API.ApplySchema(context.Background(), &pilosa.Schema{Indexes: indexInfo}, false)
if err != nil {
t.Fatalf("applying schema: %v", err)
}
idx := indexInfo[0]
fld := indexInfo[0].Fields[0]
msg := pilosa.ImportRequest{
@ -357,7 +389,7 @@ func TestHandler_Endpoints(t *testing.T) {
h.ServeHTTP(w, httpReq)
if w.Code != 412 {
t.Fatal("expected: Precondition Failed, got:" + w.Body.String())
t.Fatalf("expected: Precondition Failed, got: %d", w.Code)
}
})

View file

@ -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.

10
view.go
View file

@ -494,7 +494,7 @@ func (v *view) clearBit(txOrig Tx, rowID, columnID uint64) (changed bool, err er
}
// value uses a column of bits to read a multi-bit value.
func (v *view) value(txOrig Tx, columnID uint64, bitDepth uint) (value int64, exists bool, err error) {
func (v *view) value(txOrig Tx, columnID uint64, bitDepth uint64) (value int64, exists bool, err error) {
shard := columnID / ShardWidth
frag, err := v.CreateFragmentIfNotExists(shard)
if err != nil {
@ -511,7 +511,7 @@ func (v *view) value(txOrig Tx, columnID uint64, bitDepth uint) (value int64, ex
}
// setValue uses a column of bits to set a multi-bit value.
func (v *view) setValue(txOrig Tx, columnID uint64, bitDepth uint, value int64) (changed bool, err error) {
func (v *view) setValue(txOrig Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) {
shard := columnID / ShardWidth
frag, err := v.CreateFragmentIfNotExists(shard)
if err != nil {
@ -534,7 +534,7 @@ func (v *view) setValue(txOrig Tx, columnID uint64, bitDepth uint, value int64)
}
// clearValue removes a specific value assigned to columnID
func (v *view) clearValue(txOrig Tx, columnID uint64, bitDepth uint, value int64) (changed bool, err error) {
func (v *view) clearValue(txOrig Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) {
shard := columnID / ShardWidth
frag := v.Fragment(shard)
if frag == nil {
@ -556,7 +556,7 @@ func (v *view) clearValue(txOrig Tx, columnID uint64, bitDepth uint, value int64
}
// rangeOp returns rows with a field value encoding matching the predicate.
func (v *view) rangeOp(qcx *Qcx, op pql.Token, bitDepth uint, predicate int64) (_ *Row, err0 error) {
func (v *view) rangeOp(qcx *Qcx, op pql.Token, bitDepth uint64, predicate int64) (_ *Row, err0 error) {
r := NewRow()
for _, frag := range v.allFragments() {
@ -576,7 +576,7 @@ func (v *view) rangeOp(qcx *Qcx, op pql.Token, bitDepth uint, predicate int64) (
}
// upgradeViewBSIv2 upgrades the fragments of v. Returns ok true if any fragment upgraded.
func upgradeViewBSIv2(v *view, bitDepth uint) (ok bool, _ error) {
func upgradeViewBSIv2(v *view, bitDepth uint64) (ok bool, _ error) {
// If reading from an old formatted BSI roaring bitmap, upgrade and reload.
for _, frag := range v.allFragments() {
if frag.storage.Flags&roaringFlagBSIv2 == 1 {