mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 02:44:59 +00:00
kill everything that tries to pass nil tx to field/view things
The "field/view will just synthesize a tx" behavior is awful and also hides a number of fundamental flaws. We distinguish between "we really do mean to work on a single shard here" and "we intend to work on the whole field or view", and the latter now take Qcx instead of Tx. This eliminates a lot of very weird cases where we checked for nil Tx and synthesized them, and also gets us away from field and view taking Tx parameters when no possible Tx can be constructed which is valid, because Tx are inherently shard-specific at this time.
This commit is contained in:
parent
204558b46f
commit
964e7d86c8
11 changed files with 241 additions and 324 deletions
94
executor.go
94
executor.go
|
|
@ -228,8 +228,13 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
|
|||
|
||||
// Can't do NewTx() this high up, because we need a specific shard.
|
||||
// So start a qcx with a TxGroup and pass it down.
|
||||
qcx := idx.holder.txf.NewQcx()
|
||||
qcx.write = needWriteTxn
|
||||
var qcx *Qcx
|
||||
if needWriteTxn {
|
||||
qcx = idx.holder.txf.NewWritableQcx()
|
||||
} else {
|
||||
qcx = idx.holder.txf.NewQcx()
|
||||
|
||||
}
|
||||
defer qcx.Abort()
|
||||
|
||||
results, err := e.execute(ctx, qcx, index, q, shards, opt)
|
||||
|
|
@ -961,14 +966,7 @@ func (e *executor) executeFieldValueCall(ctx context.Context, qcx *Qcx, index st
|
|||
}
|
||||
|
||||
func (e *executor) executeFieldValueCallShard(ctx context.Context, qcx *Qcx, field *Field, col uint64, shard uint64) (_ ValCount, err0 error) {
|
||||
idx := e.Holder.Index(field.index)
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return ValCount{}, err
|
||||
}
|
||||
defer finisher(&err0)
|
||||
|
||||
value, exists, err := field.Value(tx, col)
|
||||
value, exists, err := field.Value(qcx, col)
|
||||
if err != nil {
|
||||
return ValCount{}, errors.Wrap(err, "getting field value")
|
||||
} else if !exists {
|
||||
|
|
@ -1969,8 +1967,6 @@ func (e *executor) executeMinShard(ctx context.Context, qcx *Qcx, index string,
|
|||
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeMinShard")
|
||||
defer span.Finish()
|
||||
|
||||
idx := e.Holder.Index(index)
|
||||
|
||||
var filter *Row
|
||||
if len(c.Children) == 1 {
|
||||
row, err := e.executeBitmapCallShard(ctx, qcx, index, c.Children[0], shard)
|
||||
|
|
@ -1989,18 +1985,13 @@ func (e *executor) executeMinShard(ctx context.Context, qcx *Qcx, index string,
|
|||
if field == nil {
|
||||
return ValCount{}, ErrFieldNotFound
|
||||
}
|
||||
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return ValCount{}, err
|
||||
}
|
||||
defer finisher(&err0)
|
||||
return field.MinForShard(tx, shard, filter)
|
||||
return field.MinForShard(qcx, shard, filter)
|
||||
}
|
||||
|
||||
// executeMaxShard calculates the max for bsiGroups on a shard.
|
||||
func (e *executor) executeMaxShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (_ ValCount, err0 error) {
|
||||
idx := e.Holder.Index(index)
|
||||
span, ctx := tracing.StartSpanFromContext(ctx, "Executor.executeMaxShard")
|
||||
defer span.Finish()
|
||||
|
||||
var filter *Row
|
||||
if len(c.Children) == 1 {
|
||||
|
|
@ -2020,14 +2011,7 @@ func (e *executor) executeMaxShard(ctx context.Context, qcx *Qcx, index string,
|
|||
if field == nil {
|
||||
return ValCount{}, ErrFieldNotFound
|
||||
}
|
||||
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return ValCount{}, err
|
||||
}
|
||||
|
||||
defer finisher(&err0)
|
||||
return field.MaxForShard(tx, shard, filter)
|
||||
return field.MaxForShard(qcx, shard, filter)
|
||||
}
|
||||
|
||||
// executeMinRowShard returns the minimum row ID for a shard.
|
||||
|
|
@ -5554,16 +5538,7 @@ func (e *executor) executeClearBitField(ctx context.Context, qcx *Qcx, index str
|
|||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
idx := e.Holder.Index(index)
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
defer finisher(&err0)
|
||||
|
||||
val, err := f.ClearBit(tx, rowID, colID)
|
||||
val, err := f.ClearBit(qcx, rowID, colID)
|
||||
if err != nil {
|
||||
return false, err
|
||||
} else if val {
|
||||
|
|
@ -5822,8 +5797,6 @@ func (e *executor) executeSet(ctx context.Context, qcx *Qcx, index string, c *pq
|
|||
return false, ErrIndexNotFound
|
||||
}
|
||||
|
||||
shard := colID / ShardWidth
|
||||
|
||||
// Read field name.
|
||||
fieldName, err := c.FieldArg()
|
||||
if err != nil {
|
||||
|
|
@ -5838,18 +5811,9 @@ 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, Field: ef, Shard: shard})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
defer finisher(&err0)
|
||||
|
||||
if _, err := ef.SetBit(tx, 0, colID, nil); err != nil {
|
||||
if _, err := ef.SetBit(qcx, 0, colID, nil); err != nil {
|
||||
return false, errors.Wrap(err, "setting existence column")
|
||||
}
|
||||
finisher(nil) // commit to free of the write lock needed inside executeSetBitField
|
||||
}
|
||||
|
||||
switch f.Type() {
|
||||
|
|
@ -5913,15 +5877,7 @@ func (e *executor) executeSetBitField(ctx context.Context, qcx *Qcx, index strin
|
|||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
idx := e.Holder.Index(index)
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
defer finisher(&err0)
|
||||
|
||||
val, err := f.SetBit(tx, rowID, colID, timestamp)
|
||||
val, err := f.SetBit(qcx, rowID, colID, timestamp)
|
||||
if err != nil {
|
||||
return false, err
|
||||
} else if val {
|
||||
|
|
@ -5959,16 +5915,7 @@ func (e *executor) executeSetValueField(ctx context.Context, qcx *Qcx, index str
|
|||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
idx := e.Holder.Index(index)
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
defer finisher(&err0)
|
||||
|
||||
val, err := f.SetValue(tx, colID, value)
|
||||
val, err := f.SetValue(qcx, colID, value)
|
||||
if err != nil {
|
||||
return false, err
|
||||
} else if val {
|
||||
|
|
@ -6006,14 +5953,7 @@ func (e *executor) executeClearValueField(ctx context.Context, qcx *Qcx, index s
|
|||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
idx := e.Holder.Index(index)
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
defer finisher(&err0)
|
||||
|
||||
val, err := f.ClearValue(tx, colID)
|
||||
val, err := f.ClearValue(qcx, colID)
|
||||
if err != nil {
|
||||
return false, err
|
||||
} else if val {
|
||||
|
|
|
|||
|
|
@ -34,9 +34,8 @@ func TestExecutor_TranslateRowsOnBool(t *testing.T) {
|
|||
t.Fatalf("creating index: %v", err)
|
||||
}
|
||||
|
||||
shard := uint64(0)
|
||||
tx := idx.holder.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
fb, errb := idx.CreateField("b", OptFieldTypeBool())
|
||||
_, errbk := idx.CreateField("bk", OptFieldTypeBool(), OptFieldKeys())
|
||||
|
|
@ -44,17 +43,13 @@ func TestExecutor_TranslateRowsOnBool(t *testing.T) {
|
|||
t.Fatalf("creating fields %v, %v", errb, errbk)
|
||||
}
|
||||
|
||||
_, err1 := fb.SetBit(tx, 1, 1, nil)
|
||||
_, err2 := fb.SetBit(tx, 2, 2, nil)
|
||||
_, err3 := fb.SetBit(tx, 3, 3, nil)
|
||||
_, err1 := fb.SetBit(qcx, 1, 1, nil)
|
||||
_, err2 := fb.SetBit(qcx, 2, 2, nil)
|
||||
_, err3 := fb.SetBit(qcx, 3, 3, nil)
|
||||
if err1 != nil || err2 != nil || err3 != nil {
|
||||
t.Fatalf("setting bit %v, %v, %v", err1, err2, err3)
|
||||
}
|
||||
|
||||
if err := tx.Commit(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
tests := []struct {
|
||||
pql string
|
||||
}{
|
||||
|
|
@ -566,34 +561,25 @@ func TestExecutor_DeleteRows(t *testing.T) {
|
|||
t.Fatalf("creating field: %v", err)
|
||||
}
|
||||
|
||||
shard := uint64(0)
|
||||
tx := idx.holder.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
|
||||
if _, err = f.SetBit(tx, 1, 1, nil); err != nil {
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
if _, err = f.SetBit(qcx, 1, 1, nil); err != nil {
|
||||
t.Fatalf("setting bit: %v", err)
|
||||
}
|
||||
|
||||
if err := tx.Commit(); err != nil {
|
||||
t.Fatalf("failed to commit transaction: %v", err)
|
||||
}
|
||||
|
||||
tx = idx.holder.txf.NewTx(Txo{Write: !writable, Index: idx, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
|
||||
row, err := f.Row(tx, 1)
|
||||
// We rely here on the fact that write Qcx autocommit constantly.
|
||||
row, err := f.Row(qcx, 1)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to read row: %v", err)
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
changed, err := DeleteRows(ctx, row, idx, shard)
|
||||
changed, err := DeleteRows(ctx, row, idx, 0)
|
||||
if !changed || err != nil {
|
||||
t.Fatalf("failed to delete row: %v", err)
|
||||
}
|
||||
|
||||
changed, err = DeleteRows(ctx, row, idx, shard)
|
||||
assert.NoError(t, err)
|
||||
changed, err = DeleteRows(ctx, row, idx, 0)
|
||||
if changed {
|
||||
t.Fatalf("expected delete to not clear bit but it did")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1634,12 +1634,11 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
|
|||
|
||||
// Obtain transaction.
|
||||
idx := index.Index
|
||||
shard := uint64(0)
|
||||
tx := idx.Txf().NewTx(pilosa.Txo{Write: !writable, Index: idx, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
f := hldr.Field("i", "f")
|
||||
if value, exists, err := f.Value(tx, 10); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 10); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !exists {
|
||||
t.Fatal("expected value to exist")
|
||||
|
|
@ -1647,7 +1646,7 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
|
|||
t.Fatalf("unexpected value: %v", value)
|
||||
}
|
||||
|
||||
if value, exists, err := f.Value(tx, 100); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !exists {
|
||||
t.Fatal("expected value to exist")
|
||||
|
|
@ -1707,12 +1706,11 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
|
|||
|
||||
// Obtain transaction.
|
||||
idx := index.Index
|
||||
shard := uint64(0)
|
||||
tx := idx.Txf().NewTx(pilosa.Txo{Write: !writable, Index: idx, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
f := hldr.Field("i", "f")
|
||||
if value, exists, err := f.Value(tx, 10); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 10); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !exists {
|
||||
t.Fatal("expected value to exist")
|
||||
|
|
@ -1720,7 +1718,7 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
|
|||
t.Fatalf("unexpected value: %v", value)
|
||||
}
|
||||
|
||||
if value, exists, err := f.Value(tx, 100); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !exists {
|
||||
t.Fatal("expected value to exist")
|
||||
|
|
|
|||
46
field.go
46
field.go
|
|
@ -1069,7 +1069,7 @@ func (f *Field) viewsByTimeRange(from, to time.Time) (views []string, err error)
|
|||
|
||||
// RowTime gets the row at the particular time with the granularity specified by
|
||||
// the quantum.
|
||||
func (f *Field) RowTime(tx Tx, rowID uint64, time time.Time, quantum string) (*Row, error) {
|
||||
func (f *Field) RowTime(qcx *Qcx, rowID uint64, time time.Time, quantum string) (*Row, error) {
|
||||
if !TimeQuantum(quantum).Valid() {
|
||||
return nil, ErrInvalidTimeQuantum
|
||||
}
|
||||
|
|
@ -1079,7 +1079,7 @@ func (f *Field) RowTime(tx Tx, rowID uint64, time time.Time, quantum string) (*R
|
|||
return nil, errors.Errorf("view with quantum %v not found.", quantum)
|
||||
}
|
||||
|
||||
return view.row(tx, rowID)
|
||||
return view.row(qcx, rowID)
|
||||
}
|
||||
|
||||
// viewPath returns the path to a view in the field.
|
||||
|
|
@ -1220,14 +1220,14 @@ func (f *Field) deleteView(name string) error {
|
|||
// package, and the fact that it's only allowed on
|
||||
// `set`,`mutex`, and `bool` fields is odd. This may
|
||||
// be considered for deprecation in a future version.
|
||||
func (f *Field) Row(tx Tx, rowID uint64) (*Row, error) {
|
||||
func (f *Field) Row(qcx *Qcx, rowID uint64) (*Row, error) {
|
||||
switch f.Type() {
|
||||
case FieldTypeSet, FieldTypeMutex, FieldTypeBool:
|
||||
view := f.view(viewStandard)
|
||||
if view == nil {
|
||||
return nil, ErrInvalidView
|
||||
}
|
||||
return view.row(tx, rowID)
|
||||
return view.row(qcx, rowID)
|
||||
default:
|
||||
return nil, errors.Errorf("row method unsupported for field type: %s", f.Type())
|
||||
}
|
||||
|
|
@ -1257,7 +1257,9 @@ func (f *Field) MutexCheck(ctx context.Context, qcx *Qcx, details bool, limit in
|
|||
}
|
||||
|
||||
// SetBit sets a bit on a view within the field.
|
||||
func (f *Field) SetBit(tx Tx, rowID, colID uint64, t *time.Time) (changed bool, err error) {
|
||||
func (f *Field) SetBit(qcx *Qcx, rowID, colID uint64, t *time.Time) (changed bool, err error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: colID / ShardWidth})
|
||||
defer finisher(&err)
|
||||
viewName := viewStandard
|
||||
if !f.options.NoStandardView {
|
||||
// Retrieve view. Exit if it doesn't exist.
|
||||
|
|
@ -1297,7 +1299,9 @@ func (f *Field) SetBit(tx Tx, rowID, colID uint64, t *time.Time) (changed bool,
|
|||
}
|
||||
|
||||
// ClearBit clears a bit within the field.
|
||||
func (f *Field) ClearBit(tx Tx, rowID, colID uint64) (changed bool, err error) {
|
||||
func (f *Field) ClearBit(qcx *Qcx, rowID, colID uint64) (changed bool, err error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: colID / ShardWidth})
|
||||
defer finisher(&err)
|
||||
viewName := viewStandard
|
||||
|
||||
// Retrieve view. Exit if it doesn't exist.
|
||||
|
|
@ -1409,13 +1413,13 @@ func (f *Field) allTimeViewsSortedByQuantum() (me []*view) {
|
|||
|
||||
// StringValue reads an integer field value for a column, and converts
|
||||
// it to a string based on a foreign index string key.
|
||||
func (f *Field) StringValue(tx Tx, columnID uint64) (value string, exists bool, err error) {
|
||||
func (f *Field) StringValue(qcx *Qcx, columnID uint64) (value string, exists bool, err error) {
|
||||
bsig := f.bsiGroup(f.name)
|
||||
if bsig == nil {
|
||||
return value, false, ErrBSIGroupNotFound
|
||||
}
|
||||
|
||||
val, exists, err := f.Value(tx, columnID)
|
||||
val, exists, err := f.Value(qcx, columnID)
|
||||
if exists {
|
||||
value, err = f.translateStore.TranslateID(uint64(val))
|
||||
}
|
||||
|
|
@ -1423,7 +1427,9 @@ func (f *Field) StringValue(tx Tx, columnID uint64) (value string, exists bool,
|
|||
}
|
||||
|
||||
// Value reads a field value for a column.
|
||||
func (f *Field) Value(tx Tx, columnID uint64) (value int64, exists bool, err error) {
|
||||
func (f *Field) Value(qcx *Qcx, columnID uint64) (value int64, exists bool, err error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: columnID / ShardWidth})
|
||||
defer finisher(&err)
|
||||
bsig := f.bsiGroup(f.name)
|
||||
if bsig == nil {
|
||||
return 0, false, ErrBSIGroupNotFound
|
||||
|
|
@ -1445,7 +1451,9 @@ func (f *Field) Value(tx Tx, columnID uint64) (value int64, exists bool, err err
|
|||
}
|
||||
|
||||
// SetValue sets a field value for a column.
|
||||
func (f *Field) SetValue(tx Tx, columnID uint64, value int64) (changed bool, err error) {
|
||||
func (f *Field) SetValue(qcx *Qcx, columnID uint64, value int64) (changed bool, err error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: columnID / ShardWidth})
|
||||
defer finisher(&err)
|
||||
// Fetch bsiGroup & validate min/max.
|
||||
bsig := f.bsiGroup(f.name)
|
||||
if bsig == nil {
|
||||
|
|
@ -1497,7 +1505,9 @@ func (f *Field) SetValue(tx Tx, columnID uint64, value int64) (changed bool, err
|
|||
}
|
||||
|
||||
// ClearValue removes a field value for a column.
|
||||
func (f *Field) ClearValue(tx Tx, columnID uint64) (changed bool, err error) {
|
||||
func (f *Field) ClearValue(qcx *Qcx, columnID uint64) (changed bool, err error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: columnID / ShardWidth})
|
||||
defer finisher(&err)
|
||||
bsig := f.bsiGroup(f.name)
|
||||
if bsig == nil {
|
||||
return false, ErrBSIGroupNotFound
|
||||
|
|
@ -1517,7 +1527,9 @@ func (f *Field) ClearValue(tx Tx, columnID uint64) (changed bool, err error) {
|
|||
return false, nil
|
||||
}
|
||||
|
||||
func (f *Field) MaxForShard(tx Tx, shard uint64, filter *Row) (ValCount, error) {
|
||||
func (f *Field) MaxForShard(qcx *Qcx, shard uint64, filter *Row) (ValCount, error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: shard})
|
||||
defer finisher(&err)
|
||||
bsig := f.bsiGroup(f.name)
|
||||
if bsig == nil {
|
||||
return ValCount{}, ErrBSIGroupNotFound
|
||||
|
|
@ -1546,13 +1558,16 @@ func (f *Field) MaxForShard(tx Tx, shard uint64, filter *Row) (ValCount, error)
|
|||
return ValCount{}, errors.Wrap(err, "calling fragment.max")
|
||||
}
|
||||
|
||||
return f.valCountize(max, cnt, bsig)
|
||||
v, err := f.valCountize(max, cnt, bsig)
|
||||
return v, err
|
||||
}
|
||||
|
||||
// MinForShard returns the minimum value which appears in this shard
|
||||
// (this field must be an Int or Decimal field). It also returns the
|
||||
// number of times the minimum value appears.
|
||||
func (f *Field) MinForShard(tx Tx, shard uint64, filter *Row) (ValCount, error) {
|
||||
func (f *Field) MinForShard(qcx *Qcx, shard uint64, filter *Row) (ValCount, error) {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: shard})
|
||||
defer finisher(&err)
|
||||
bsig := f.bsiGroup(f.name)
|
||||
if bsig == nil {
|
||||
return ValCount{}, ErrBSIGroupNotFound
|
||||
|
|
@ -1581,7 +1596,8 @@ func (f *Field) MinForShard(tx Tx, shard uint64, filter *Row) (ValCount, error)
|
|||
return ValCount{}, errors.Wrap(err, "calling fragment.min")
|
||||
}
|
||||
|
||||
return f.valCountize(min, cnt, bsig)
|
||||
v, err := f.valCountize(min, cnt, bsig)
|
||||
return v, err
|
||||
}
|
||||
|
||||
// valCountize takes the "raw" value and count we get from the
|
||||
|
|
|
|||
|
|
@ -300,15 +300,15 @@ func (f *TestField) Reopen() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (f *TestField) MustSetBit(tx Tx, row, col uint64, ts ...time.Time) {
|
||||
func (f *TestField) MustSetBit(qcx *Qcx, row, col uint64, ts ...time.Time) {
|
||||
if len(ts) == 0 {
|
||||
_, err := f.Field.SetBit(tx, row, col, nil)
|
||||
_, err := f.Field.SetBit(qcx, row, col, nil)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
for _, t := range ts {
|
||||
_, err := f.Field.SetBit(tx, row, col, &t)
|
||||
_, err := f.Field.SetBit(qcx, row, col, &t)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
@ -360,46 +360,51 @@ func TestField_RowTime(t *testing.T) {
|
|||
f := OpenField(t, OptFieldTypeTime(TimeQuantum("YMDH"), "0"))
|
||||
|
||||
// Obtain transaction.
|
||||
tx := f.idx.holder.txf.NewTx(Txo{Write: writable, Index: f.idx, Field: f.Field, Shard: 0})
|
||||
defer tx.Rollback()
|
||||
qcx := f.idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
f.MustSetBit(tx, 1, 1, time.Date(2010, time.January, 5, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(tx, 1, 2, time.Date(2011, time.January, 5, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(tx, 1, 3, time.Date(2010, time.February, 5, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(tx, 1, 4, time.Date(2010, time.January, 6, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(tx, 1, 5, time.Date(2010, time.January, 5, 13, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(qcx, 1, 1, time.Date(2010, time.January, 5, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(qcx, 1, 2, time.Date(2011, time.January, 5, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(qcx, 1, 3, time.Date(2010, time.February, 5, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(qcx, 1, 4, time.Date(2010, time.January, 6, 12, 0, 0, 0, time.UTC))
|
||||
f.MustSetBit(qcx, 1, 5, time.Date(2010, time.January, 5, 13, 0, 0, 0, time.UTC))
|
||||
|
||||
PanicOn(tx.Commit())
|
||||
// Warning: Right now this is misleading, and doesn't really do anything. We
|
||||
// already committed each change as we got there. SOME DAY we will fix this.
|
||||
PanicOn(qcx.Finish())
|
||||
|
||||
qcx = f.idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
// obtain 2nd transaction to read it back.
|
||||
tx = f.idx.holder.txf.NewTx(Txo{Write: !writable, Index: f.idx, Field: f.Field, Shard: 0})
|
||||
defer tx.Rollback()
|
||||
qcx = f.idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
if r, err := f.RowTime(tx, 1, time.Date(2010, time.November, 5, 12, 0, 0, 0, time.UTC), "Y"); err != nil {
|
||||
if r, err := f.RowTime(qcx, 1, time.Date(2010, time.November, 5, 12, 0, 0, 0, time.UTC), "Y"); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(r.Columns(), []uint64{1, 3, 4, 5}) {
|
||||
t.Fatalf("wrong columns: %#v", r.Columns())
|
||||
}
|
||||
|
||||
if r, err := f.RowTime(tx, 1, time.Date(2010, time.February, 7, 13, 0, 0, 0, time.UTC), "YM"); err != nil {
|
||||
if r, err := f.RowTime(qcx, 1, time.Date(2010, time.February, 7, 13, 0, 0, 0, time.UTC), "YM"); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(r.Columns(), []uint64{3}) {
|
||||
t.Fatalf("wrong columns: %#v", r.Columns())
|
||||
}
|
||||
|
||||
if r, err := f.RowTime(tx, 1, time.Date(2010, time.February, 7, 13, 0, 0, 0, time.UTC), "M"); err != nil {
|
||||
if r, err := f.RowTime(qcx, 1, time.Date(2010, time.February, 7, 13, 0, 0, 0, time.UTC), "M"); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(r.Columns(), []uint64{3}) {
|
||||
t.Fatalf("wrong columns: %#v", r.Columns())
|
||||
}
|
||||
|
||||
if r, err := f.RowTime(tx, 1, time.Date(2010, time.January, 5, 12, 0, 0, 0, time.UTC), "MD"); err != nil {
|
||||
if r, err := f.RowTime(qcx, 1, time.Date(2010, time.January, 5, 12, 0, 0, 0, time.UTC), "MD"); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(r.Columns(), []uint64{1, 5}) {
|
||||
t.Fatalf("wrong columns: %#v", r.Columns())
|
||||
}
|
||||
|
||||
if r, err := f.RowTime(tx, 1, time.Date(2010, time.January, 5, 13, 0, 0, 0, time.UTC), "MDH"); err != nil {
|
||||
if r, err := f.RowTime(qcx, 1, time.Date(2010, time.January, 5, 13, 0, 0, 0, time.UTC), "MDH"); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(r.Columns(), []uint64{5}) {
|
||||
t.Fatalf("wrong columns: %#v", r.Columns())
|
||||
|
|
@ -629,10 +634,9 @@ func TestIntField_MinMaxForShard(t *testing.T) {
|
|||
qcx.Reset()
|
||||
|
||||
shard := uint64(0)
|
||||
tx := f.idx.holder.txf.NewTx(Txo{Write: !writable, Index: f.idx, Field: f.Field, Shard: shard})
|
||||
// Rollback below manually, because we are in a loop.
|
||||
|
||||
maxvc, err := f.MaxForShard(tx, shard, nil)
|
||||
maxvc, err := f.MaxForShard(qcx, shard, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("getting max for shard: %v", err)
|
||||
}
|
||||
|
|
@ -640,14 +644,13 @@ func TestIntField_MinMaxForShard(t *testing.T) {
|
|||
t.Fatalf("max expected:\n%+v\ngot:\n%+v", test.expMax, maxvc)
|
||||
}
|
||||
|
||||
minvc, err := f.MinForShard(tx, shard, nil)
|
||||
minvc, err := f.MinForShard(qcx, shard, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("getting min for shard: %v", err)
|
||||
}
|
||||
if minvc != test.expMin {
|
||||
t.Fatalf("min expected:\n%+v\ngot:\n%+v", test.expMin, minvc)
|
||||
}
|
||||
tx.Rollback()
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
@ -786,10 +789,8 @@ func TestDecimalField_MinMaxForShard(t *testing.T) {
|
|||
}
|
||||
|
||||
shard := uint64(0)
|
||||
tx := f.idx.holder.txf.NewTx(Txo{Write: !writable, Index: f.idx, Field: f.Field, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
|
||||
maxvc, err := f.MaxForShard(tx, shard, nil)
|
||||
maxvc, err := f.MaxForShard(qcx, shard, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("getting max for shard: %v", err)
|
||||
}
|
||||
|
|
@ -797,7 +798,7 @@ func TestDecimalField_MinMaxForShard(t *testing.T) {
|
|||
t.Fatalf("max expected:\n%+v\ngot:\n%+v", test.expMax, maxvc)
|
||||
}
|
||||
|
||||
minvc, err := f.MinForShard(tx, shard, nil)
|
||||
minvc, err := f.MinForShard(qcx, shard, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("getting min for shard: %v", err)
|
||||
}
|
||||
|
|
@ -868,24 +869,20 @@ func TestField_SaveMeta(t *testing.T) {
|
|||
expBitDepth := uint64(7)
|
||||
|
||||
// Obtain transaction.
|
||||
tx := f.idx.holder.txf.NewTx(Txo{Write: writable, Index: f.idx, Field: f.Field, Shard: 0})
|
||||
defer tx.Rollback()
|
||||
qcx := f.idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
if changed, err := f.SetValue(tx, colID, val); err != nil {
|
||||
if changed, err := f.SetValue(qcx, colID, val); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected SetValue to return changed = true")
|
||||
} else if err := tx.Commit(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if f.options.BitDepth != expBitDepth {
|
||||
t.Fatalf("expected BitDepth after set to be: %d, got: %d", expBitDepth, f.options.BitDepth)
|
||||
}
|
||||
|
||||
tx2 := f.idx.holder.txf.NewTx(Txo{Index: f.idx, Field: f.Field, Shard: 0})
|
||||
defer tx2.Rollback()
|
||||
if rslt, ok, err := f.Value(tx2, colID); err != nil {
|
||||
if rslt, ok, err := f.Value(qcx, colID); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !ok {
|
||||
t.Fatal("expected Value() to return exists = true")
|
||||
|
|
@ -893,18 +890,23 @@ func TestField_SaveMeta(t *testing.T) {
|
|||
t.Fatalf("expected value to be: %d, got: %d", val, rslt)
|
||||
}
|
||||
|
||||
if err := qcx.Finish(); err != nil {
|
||||
t.Fatalf("error finishing qcx: %v", err)
|
||||
}
|
||||
|
||||
// Reload field and verify that it is persisted.
|
||||
if err := f.Reopen(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
qcx = f.idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
if f.options.BitDepth != expBitDepth {
|
||||
t.Fatalf("expected BitDepth after reopen to be: %d, got: %d", expBitDepth, f.options.BitDepth)
|
||||
}
|
||||
|
||||
tx3 := f.idx.holder.txf.NewTx(Txo{Index: f.idx, Field: f.Field, Shard: 0})
|
||||
defer tx3.Rollback()
|
||||
if rslt, ok, err := f.Value(tx3, colID); err != nil {
|
||||
if rslt, ok, err := f.Value(qcx, colID); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !ok {
|
||||
t.Fatal("expected Value() after reopen to return exists = true")
|
||||
|
|
|
|||
|
|
@ -28,18 +28,24 @@ func TestField_SetValue(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It is okay to pass in a nil tx. f.SetValue will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
// You're going to note the lack of any commits here. That's
|
||||
// because, when you have a writable Qcx, *every individual
|
||||
// sub-transaction commits immediately*. In theory, we ought
|
||||
// to be doing provisional writes and the entire set of writes
|
||||
// ought to be able to be reverted. Actually no. We're just committing
|
||||
// everything as we go anyway.
|
||||
|
||||
// Set value on field.
|
||||
if changed, err := f.SetValue(tx, 100, 21); err != nil {
|
||||
if changed, err := f.SetValue(qcx, 100, 21); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Read value.
|
||||
if value, exists, err := f.Value(tx, 100); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != 21 {
|
||||
t.Fatalf("unexpected value: %d", value)
|
||||
|
|
@ -48,7 +54,7 @@ func TestField_SetValue(t *testing.T) {
|
|||
}
|
||||
|
||||
// Setting value should return no change.
|
||||
if changed, err := f.SetValue(tx, 100, 21); err != nil {
|
||||
if changed, err := f.SetValue(qcx, 100, 21); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if changed {
|
||||
t.Fatal("expected no change")
|
||||
|
|
@ -63,25 +69,25 @@ func TestField_SetValue(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It is okay to pass in a nil tx. f.SetValue will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.SetValue(tx, 100, 21); err != nil {
|
||||
if changed, err := f.SetValue(qcx, 100, 21); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Set different value.
|
||||
if changed, err := f.SetValue(tx, 100, 23); err != nil {
|
||||
if changed, err := f.SetValue(qcx, 100, 23); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Read value.
|
||||
if value, exists, err := f.Value(tx, 100); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != 23 {
|
||||
t.Fatalf("unexpected value: %d", value)
|
||||
|
|
@ -98,11 +104,11 @@ func TestField_SetValue(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It is okay to pass in a nil tx. f.SetValue will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
// Set value.
|
||||
if _, err := f.SetValue(tx, 100, 21); err != pilosa.ErrBSIGroupNotFound {
|
||||
if _, err := f.SetValue(qcx, 100, 21); err != pilosa.ErrBSIGroupNotFound {
|
||||
t.Fatalf("unexpected error: %s", err)
|
||||
}
|
||||
})
|
||||
|
|
@ -114,12 +120,10 @@ func TestField_SetValue(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It is okay to pass in a nil tx. f.SetValue will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
// Set value.
|
||||
if _, err := f.SetValue(tx, 100, 15); !errors.Is(err, pilosa.ErrBSIGroupValueTooLow) {
|
||||
if _, err := f.SetValue(qcx, 100, 15); !errors.Is(err, pilosa.ErrBSIGroupValueTooLow) {
|
||||
t.Fatalf("unexpected error: %s", err)
|
||||
}
|
||||
})
|
||||
|
|
@ -132,11 +136,11 @@ func TestField_SetValue(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It is okay to pass in a nil tx. f.SetValue will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
// Set value.
|
||||
if _, err := f.SetValue(tx, 100, 31); !errors.Is(err, pilosa.ErrBSIGroupValueTooHigh) {
|
||||
if _, err := f.SetValue(qcx, 100, 31); !errors.Is(err, pilosa.ErrBSIGroupValueTooHigh) {
|
||||
t.Fatalf("unexpected error: %s", err)
|
||||
}
|
||||
})
|
||||
|
|
@ -204,13 +208,13 @@ func TestField_AvailableShards(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It is okay to pass in a nil tx. f.SetBit will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
// Set values on shards 0 & 2, and verify.
|
||||
if _, err := f.SetBit(tx, 0, 100, nil); err != nil {
|
||||
if _, err := f.SetBit(qcx, 0, 100, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetBit(tx, 0, ShardWidth*2, nil); err != nil {
|
||||
} else if _, err := f.SetBit(qcx, 0, ShardWidth*2, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if diff := cmp.Diff(f.AvailableShards(includeRemote).Slice(), []uint64{0, 2}); diff != "" {
|
||||
t.Fatal(diff)
|
||||
|
|
@ -244,18 +248,18 @@ func TestField_ClearValue(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// It is okay to pass in a nil tx. f.SetValue will lazily instantiate Tx.
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
// Set value on field.
|
||||
if changed, err := f.SetValue(tx, 100, 21); err != nil {
|
||||
if changed, err := f.SetValue(qcx, 100, 21); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Read value.
|
||||
if value, exists, err := f.Value(tx, 100); err != nil {
|
||||
if value, exists, err := f.Value(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != 21 {
|
||||
t.Fatalf("unexpected value: %d", value)
|
||||
|
|
@ -263,14 +267,14 @@ func TestField_ClearValue(t *testing.T) {
|
|||
t.Fatal("expected value to exist")
|
||||
}
|
||||
|
||||
if changed, err := f.ClearValue(tx, 100); err != nil {
|
||||
if changed, err := f.ClearValue(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Read value.
|
||||
if _, exists, err := f.Value(tx, 100); err != nil {
|
||||
if _, exists, err := f.Value(qcx, 100); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if exists {
|
||||
t.Fatal("expected value to not exist")
|
||||
|
|
|
|||
|
|
@ -30,22 +30,22 @@ func setupTest(t *testing.T, h *Holder, rowCol []rowCols, indexName string) (*In
|
|||
}
|
||||
existencefield := idx.existenceFld
|
||||
|
||||
shard := uint64(0)
|
||||
tx := idx.Txf().NewTx(Txo{Write: true, Index: idx, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
for _, r := range rowCol {
|
||||
_, err = f.SetBit(tx, r.row, r.col, nil)
|
||||
_, err = f.SetBit(qcx, r.row, r.col, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to set bit in index %v: %v", indexName, err)
|
||||
}
|
||||
|
||||
_, err = existencefield.SetBit(tx, r.row, r.col, nil)
|
||||
_, err = existencefield.SetBit(qcx, r.row, r.col, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to set bit in index %v: %v", indexName, err)
|
||||
}
|
||||
}
|
||||
|
||||
if err = tx.Commit(); err != nil {
|
||||
if err = qcx.Finish(); err != nil {
|
||||
t.Fatalf("failed to commit tx for index %v: %v", indexName, err)
|
||||
}
|
||||
|
||||
|
|
@ -97,14 +97,14 @@ func TestHolder_ProcessDeleteInflight(t *testing.T) {
|
|||
for _, test := range tests {
|
||||
func() {
|
||||
idx, f := test.idx, test.f
|
||||
tx := idx.Txf().NewTx(Txo{Write: false, Index: idx1, Shard: uint64(0)})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
for _, r := range rowCol {
|
||||
row, err := f.Row(tx, r.row)
|
||||
row, err := f.Row(qcx, r.row)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to get row: %v", err)
|
||||
}
|
||||
existenceRow, err := idx.existenceFld.Row(tx, r.row)
|
||||
existenceRow, err := idx.existenceFld.Row(qcx, r.row)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to get row: %v", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -752,8 +752,8 @@ func (h *GRPCHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSe
|
|||
return errToStatusError(err)
|
||||
}
|
||||
|
||||
// It is okay to pass a nil Tx to field.StringValue(). It will lazily create it.
|
||||
var tx pilosa.Tx
|
||||
qcx := index.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
var fields []*pilosa.Field
|
||||
for _, field := range index.Fields() {
|
||||
|
|
@ -998,7 +998,7 @@ func (h *GRPCHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSe
|
|||
}
|
||||
}
|
||||
} else {
|
||||
value, exists, err = field.StringValue(tx, col)
|
||||
value, exists, err = field.StringValue(qcx, col)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "getting string field value for column")
|
||||
}
|
||||
|
|
@ -1326,7 +1326,7 @@ func (h *GRPCHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSe
|
|||
&pb.ColumnResponse{ColumnVal: nil})
|
||||
}
|
||||
} else {
|
||||
value, exists, err = field.StringValue(tx, id)
|
||||
value, exists, err = field.StringValue(qcx, id)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "getting string field value for column")
|
||||
}
|
||||
|
|
|
|||
101
test/holder.go
101
test/holder.go
|
|
@ -7,16 +7,15 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/testhook"
|
||||
"github.com/featurebasedb/featurebase/v3/vprint"
|
||||
"github.com/pkg/errors"
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/pql"
|
||||
"github.com/molecula/featurebase/v3/testhook"
|
||||
)
|
||||
|
||||
// Holder is a test wrapper for pilosa.Holder.
|
||||
type Holder struct {
|
||||
*pilosa.Holder
|
||||
tb testing.TB
|
||||
}
|
||||
|
||||
// NewHolder returns a new instance of Holder with a temporary path.
|
||||
|
|
@ -29,7 +28,7 @@ func NewHolder(tb testing.TB) *Holder {
|
|||
cfg := pilosa.DefaultHolderConfig()
|
||||
cfg.StorageConfig.FsyncEnabled = false
|
||||
cfg.RBFConfig.FsyncEnabled = false
|
||||
h := &Holder{Holder: pilosa.NewHolder(path, cfg)}
|
||||
h := &Holder{Holder: pilosa.NewHolder(path, cfg), tb: tb}
|
||||
|
||||
return h
|
||||
}
|
||||
|
|
@ -38,7 +37,7 @@ func NewHolder(tb testing.TB) *Holder {
|
|||
func MustOpenHolder(tb testing.TB) *Holder {
|
||||
h := NewHolder(tb)
|
||||
if err := h.Open(); err != nil {
|
||||
panic(err)
|
||||
tb.Fatalf("opening holder: %v", err)
|
||||
}
|
||||
return h
|
||||
}
|
||||
|
|
@ -58,7 +57,7 @@ func (h *Holder) Reopen() error {
|
|||
func (h *Holder) MustCreateIndexIfNotExists(index string, opt pilosa.IndexOptions) *Index {
|
||||
idx, err := h.Holder.CreateIndexIfNotExists(index, opt)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
h.tb.Fatalf("creating index: %v", err)
|
||||
}
|
||||
return &Index{Index: idx}
|
||||
}
|
||||
|
|
@ -70,37 +69,39 @@ func (h *Holder) Row(index, field string, rowID uint64) *pilosa.Row {
|
|||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
row, err := f.Row(tx, rowID)
|
||||
row, err := f.Row(qcx, rowID)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
h.tb.Fatalf("retrieving row: %v", err)
|
||||
}
|
||||
// clone it so that mmapped storage doesn't disappear from under it
|
||||
// once the tx goes away.
|
||||
return row.Clone()
|
||||
// once the qcx goes away.
|
||||
return row
|
||||
}
|
||||
|
||||
// ReadRow returns a Row for a given field. If the field does not exist,
|
||||
// it panics rather than creating the field.
|
||||
// it fails the holder's test rather than creating the field.
|
||||
func (h *Holder) ReadRow(index, field string, rowID uint64) *pilosa.Row {
|
||||
idx := h.Holder.Index(index)
|
||||
if idx == nil {
|
||||
panic(errors.Wrap(pilosa.ErrIndexNotFound, index))
|
||||
h.tb.Fatalf("read row from index %q: index not found", index)
|
||||
}
|
||||
f := idx.Field(field)
|
||||
if f == nil {
|
||||
panic(errors.Wrap(pilosa.ErrFieldNotFound, field))
|
||||
h.tb.Fatalf("read row from field %q/%q: field not found", index, field)
|
||||
}
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
row, err := f.Row(tx, rowID)
|
||||
row, err := f.Row(qcx, rowID)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
h.tb.Fatalf("retrieving row: %v", err)
|
||||
}
|
||||
|
||||
// clone it so that mmapped storage doesn't disappear from under it
|
||||
// once the tx goes away.
|
||||
// once the qcx goes away.
|
||||
return row.Clone()
|
||||
}
|
||||
|
||||
|
|
@ -110,15 +111,16 @@ func (h *Holder) RowTime(index, field string, rowID uint64, t time.Time, quantum
|
|||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
var tx pilosa.Tx
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
row, err := f.RowTime(tx, rowID, t, quantum)
|
||||
row, err := f.RowTime(qcx, rowID, t, quantum)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
// clone it so that mmapped storage doesn't disappear from under it
|
||||
// once the tx goes away.
|
||||
// once the qcx goes away.
|
||||
return row.Clone()
|
||||
}
|
||||
|
||||
|
|
@ -135,15 +137,17 @@ func (h *Holder) SetBitTime(index, field string, rowID, columnID uint64, t *time
|
|||
panic(err)
|
||||
}
|
||||
|
||||
shard := columnID / pilosa.ShardWidth
|
||||
tx := idx.Index.Txf().NewTx(pilosa.Txo{Write: true, Index: idx.Index, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
_, err = f.SetBit(tx, rowID, columnID, t)
|
||||
_, err = f.SetBit(qcx, rowID, columnID, t)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
h.tb.Fatalf("setting bit: %v", err)
|
||||
}
|
||||
err = qcx.Finish()
|
||||
if err != nil {
|
||||
h.tb.Fatalf("finishing qcx: %v", err)
|
||||
}
|
||||
vprint.PanicOn(tx.Commit())
|
||||
}
|
||||
|
||||
// ClearBit clears a bit on the given field.
|
||||
|
|
@ -154,15 +158,17 @@ func (h *Holder) ClearBit(index, field string, rowID, columnID uint64) {
|
|||
panic(err)
|
||||
}
|
||||
|
||||
shard := columnID / pilosa.ShardWidth
|
||||
tx := idx.Index.Txf().NewTx(pilosa.Txo{Write: true, Index: idx.Index, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
_, err = f.ClearBit(tx, rowID, columnID)
|
||||
_, err = f.ClearBit(qcx, rowID, columnID)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
h.tb.Fatalf("clearing bit: %v", err)
|
||||
}
|
||||
err = qcx.Finish()
|
||||
if err != nil {
|
||||
h.tb.Fatalf("finishing qcx: %v", err)
|
||||
}
|
||||
vprint.PanicOn(tx.Commit())
|
||||
}
|
||||
|
||||
// MustSetBits sets columns on a row. Panic on error.
|
||||
|
|
@ -180,16 +186,17 @@ func (h *Holder) SetValue(index, field string, columnID uint64, value int64) *In
|
|||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
shard := columnID / pilosa.ShardWidth
|
||||
tx := idx.Index.Txf().NewTx(pilosa.Txo{Write: true, Index: idx.Index, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
|
||||
_, err = f.SetValue(tx, columnID, value)
|
||||
qcx := idx.Txf().NewWritableQcx()
|
||||
defer qcx.Abort()
|
||||
_, err = f.SetValue(qcx, columnID, value)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
h.tb.Fatalf("setting value: %v", err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
panic(err)
|
||||
|
||||
err = qcx.Finish()
|
||||
if err != nil {
|
||||
h.tb.Fatalf("finishing qcx: %v", err)
|
||||
}
|
||||
return idx
|
||||
}
|
||||
|
|
@ -201,11 +208,11 @@ func (h *Holder) Value(index, field string, columnID uint64) (int64, bool) {
|
|||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
shard := columnID / pilosa.ShardWidth
|
||||
tx := idx.Index.Txf().NewTx(pilosa.Txo{Write: false, Index: idx.Index, Shard: shard})
|
||||
defer tx.Rollback()
|
||||
|
||||
val, exists, err := f.Value(tx, columnID)
|
||||
qcx := idx.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
val, exists, err := f.Value(qcx, columnID)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
@ -230,6 +237,6 @@ func (h *Holder) Range(index, field string, op pql.Token, predicate int64) *pilo
|
|||
}
|
||||
|
||||
// clone it so that mmapped storage doesn't disappear from under it
|
||||
// once the tx goes away.
|
||||
// once the qcx goes away.
|
||||
return row.Clone()
|
||||
}
|
||||
|
|
|
|||
24
txfactory.go
24
txfactory.go
|
|
@ -167,9 +167,9 @@ func (q *Qcx) unprotected_reset() {
|
|||
q.done = false
|
||||
}
|
||||
|
||||
// NewQcxWithGroup allocates a freshly allocated and empty Grp.
|
||||
// The top-level executor will set qcx.write = true manually
|
||||
// if the overall query is a write.
|
||||
// NewQcx allocates a freshly allocated and empty Grp.
|
||||
// The top-level Qcx is not marked writable. Non-writable
|
||||
// Qcx should not be used to request write Tx.
|
||||
func (f *TxFactory) NewQcx() (qcx *Qcx) {
|
||||
qcx = &Qcx{
|
||||
Grp: f.NewTxGroup(),
|
||||
|
|
@ -185,6 +185,24 @@ func (f *TxFactory) NewQcx() (qcx *Qcx) {
|
|||
return
|
||||
}
|
||||
|
||||
// NewWritableQcx allocates a freshly allocated and empty Grp.
|
||||
// The resulting Qcx is marked writable.
|
||||
func (f *TxFactory) NewWritableQcx() (qcx *Qcx) {
|
||||
qcx = &Qcx{
|
||||
Grp: f.NewTxGroup(),
|
||||
Txf: f,
|
||||
}
|
||||
if f.holder != nil && f.holder.executor != nil {
|
||||
qcx.workers = f.holder.executor.workers
|
||||
}
|
||||
if f.typeOfTx == "roaring" {
|
||||
qcx.isRoaring = true
|
||||
}
|
||||
_ = testhook.Opened(f.holder.Auditor, qcx, nil)
|
||||
qcx.write = true
|
||||
return
|
||||
}
|
||||
|
||||
var NoopFinisher = func(perr *error) {}
|
||||
|
||||
var ErrQcxDone = fmt.Errorf("Qcx already Aborted or Finished, so must call reset before re-use")
|
||||
|
|
|
|||
74
view.go
74
view.go
|
|
@ -449,16 +449,14 @@ func (v *view) deleteFragment(shard uint64) error {
|
|||
}
|
||||
|
||||
// row returns a row for a shard of the view.
|
||||
func (v *view) row(txOrig Tx, rowID uint64) (*Row, error) {
|
||||
func (v *view) row(qcx *Qcx, rowID uint64) (*Row, error) {
|
||||
row := NewRow()
|
||||
for _, frag := range v.allFragments() {
|
||||
|
||||
tx := txOrig
|
||||
if NilInside(tx) {
|
||||
tx = v.idx.holder.txf.NewTx(Txo{Write: !writable, Index: v.idx, Fragment: frag, Shard: frag.shard})
|
||||
defer tx.Rollback()
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: v.idx, Fragment: frag, Shard: frag.shard})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
defer finisher(&err)
|
||||
fr, err := frag.row(tx, rowID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -527,7 +525,7 @@ func (v *view) mutexCheck(ctx context.Context, qcx *Qcx, details bool, limit int
|
|||
}
|
||||
|
||||
// setBit sets a bit within the view.
|
||||
func (v *view) setBit(txOrig Tx, rowID, columnID uint64) (changed bool, err error) {
|
||||
func (v *view) setBit(tx Tx, rowID, columnID uint64) (changed bool, err error) {
|
||||
shard := columnID / ShardWidth
|
||||
var frag *fragment
|
||||
frag, err = v.CreateFragmentIfNotExists(shard)
|
||||
|
|
@ -535,102 +533,50 @@ func (v *view) setBit(txOrig Tx, rowID, columnID uint64) (changed bool, err erro
|
|||
return changed, err
|
||||
}
|
||||
|
||||
tx := txOrig
|
||||
if NilInside(tx) {
|
||||
tx = v.idx.holder.txf.NewTx(Txo{Write: writable, Index: v.idx, Fragment: frag, Shard: shard})
|
||||
defer func() {
|
||||
if err == nil {
|
||||
vprint.PanicOn(tx.Commit())
|
||||
} else {
|
||||
tx.Rollback()
|
||||
}
|
||||
}()
|
||||
}
|
||||
return frag.setBit(tx, rowID, columnID)
|
||||
}
|
||||
|
||||
// clearBit clears a bit within the view.
|
||||
func (v *view) clearBit(txOrig Tx, rowID, columnID uint64) (changed bool, err error) {
|
||||
func (v *view) clearBit(tx Tx, rowID, columnID uint64) (changed bool, err error) {
|
||||
shard := columnID / ShardWidth
|
||||
frag := v.Fragment(shard)
|
||||
if frag == nil {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
tx := txOrig
|
||||
if NilInside(tx) {
|
||||
tx = v.idx.holder.txf.NewTx(Txo{Write: writable, Index: v.idx, Fragment: frag, Shard: shard})
|
||||
defer func() {
|
||||
if err == nil {
|
||||
vprint.PanicOn(tx.Commit())
|
||||
} else {
|
||||
tx.Rollback()
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
return frag.clearBit(tx, rowID, columnID)
|
||||
}
|
||||
|
||||
// value uses a column of bits to read a multi-bit value.
|
||||
func (v *view) value(txOrig Tx, columnID uint64, bitDepth uint64) (value int64, exists bool, err error) {
|
||||
func (v *view) value(tx Tx, columnID uint64, bitDepth uint64) (value int64, exists bool, err error) {
|
||||
shard := columnID / ShardWidth
|
||||
frag, err := v.CreateFragmentIfNotExists(shard)
|
||||
if err != nil {
|
||||
return value, exists, err
|
||||
}
|
||||
|
||||
tx := txOrig
|
||||
if NilInside(tx) {
|
||||
tx = frag.idx.holder.txf.NewTx(Txo{Write: !writable, Index: frag.idx, Fragment: frag, Shard: frag.shard})
|
||||
defer tx.Rollback()
|
||||
}
|
||||
|
||||
return frag.value(tx, columnID, bitDepth)
|
||||
}
|
||||
|
||||
// setValue uses a column of bits to set a multi-bit value.
|
||||
func (v *view) setValue(txOrig Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) {
|
||||
func (v *view) setValue(tx Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) {
|
||||
shard := columnID / ShardWidth
|
||||
frag, err := v.CreateFragmentIfNotExists(shard)
|
||||
if err != nil {
|
||||
return changed, err
|
||||
}
|
||||
|
||||
tx := txOrig
|
||||
if NilInside(tx) {
|
||||
tx = v.idx.holder.txf.NewTx(Txo{Write: writable, Index: v.idx, Fragment: frag, Shard: shard})
|
||||
defer func() {
|
||||
if err == nil {
|
||||
vprint.PanicOn(tx.Commit())
|
||||
} else {
|
||||
tx.Rollback()
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
return frag.setValue(tx, columnID, bitDepth, value)
|
||||
}
|
||||
|
||||
// clearValue removes a specific value assigned to columnID
|
||||
func (v *view) clearValue(txOrig Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) {
|
||||
func (v *view) clearValue(tx Tx, columnID uint64, bitDepth uint64, value int64) (changed bool, err error) {
|
||||
shard := columnID / ShardWidth
|
||||
frag := v.Fragment(shard)
|
||||
if frag == nil {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
tx := txOrig
|
||||
if NilInside(tx) {
|
||||
tx = v.idx.holder.txf.NewTx(Txo{Write: writable, Index: v.idx, Fragment: frag, Shard: shard})
|
||||
defer func() {
|
||||
if err == nil {
|
||||
vprint.PanicOn(tx.Commit())
|
||||
} else {
|
||||
tx.Rollback()
|
||||
}
|
||||
}()
|
||||
}
|
||||
return frag.clearValue(tx, columnID, bitDepth, value)
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue