diff --git a/executor.go b/executor.go index cfe02d3d5..fb7039d8b 100644 --- a/executor.go +++ b/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 { diff --git a/executor_internal_test.go b/executor_internal_test.go index 3cfd9b7c7..96f049724 100644 --- a/executor_internal_test.go +++ b/executor_internal_test.go @@ -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") } diff --git a/executor_test.go b/executor_test.go index fe830a488..6f17dc1a6 100644 --- a/executor_test.go +++ b/executor_test.go @@ -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") diff --git a/field.go b/field.go index ef1910386..d87f0f3e1 100644 --- a/field.go +++ b/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 diff --git a/field_internal_test.go b/field_internal_test.go index 9636aafda..48b5bbb78 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -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") diff --git a/field_test.go b/field_test.go index 7018c71b9..d2992a9ae 100644 --- a/field_test.go +++ b/field_test.go @@ -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") diff --git a/holder_internal_test.go b/holder_internal_test.go index 32be7b568..08634ed90 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -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) } diff --git a/server/grpc.go b/server/grpc.go index a475a1bdd..7cb714c34 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -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") } diff --git a/test/holder.go b/test/holder.go index c755ac18c..8b25555cd 100644 --- a/test/holder.go +++ b/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() } diff --git a/txfactory.go b/txfactory.go index 4d4879a8a..1aee2bec8 100644 --- a/txfactory.go +++ b/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") diff --git a/view.go b/view.go index 3885709ac..edc6acd33 100644 --- a/view.go +++ b/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) }