diff --git a/executor.go b/executor.go index 5079dedee..29d792253 100644 --- a/executor.go +++ b/executor.go @@ -227,8 +227,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) @@ -960,14 +965,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 { @@ -1968,8 +1966,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) @@ -1988,18 +1984,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 { @@ -2019,14 +2010,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. @@ -5553,16 +5537,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 { @@ -5821,8 +5796,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 { @@ -5837,18 +5810,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() { @@ -5912,15 +5876,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 { @@ -5958,16 +5914,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 { @@ -6005,14 +5952,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 719c7f3d5..68c3c00de 100644 --- a/executor_internal_test.go +++ b/executor_internal_test.go @@ -32,9 +32,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()) @@ -42,17 +41,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 }{ @@ -564,33 +559,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) + 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 b1d51b7d2..e71cc01e7 100644 --- a/executor_test.go +++ b/executor_test.go @@ -1633,12 +1633,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") @@ -1646,7 +1645,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") @@ -1706,12 +1705,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") @@ -1719,7 +1717,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 4d85801ce..0b671e9aa 100644 --- a/field.go +++ b/field.go @@ -1068,7 +1068,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 } @@ -1078,7 +1078,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. @@ -1219,14 +1219,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()) } @@ -1256,7 +1256,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. @@ -1296,7 +1298,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. @@ -1408,13 +1412,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)) } @@ -1422,7 +1426,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 @@ -1444,7 +1450,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 { @@ -1496,7 +1504,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 @@ -1516,7 +1526,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 @@ -1545,13 +1557,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 @@ -1580,7 +1595,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 1d775e702..7033a511a 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -299,15 +299,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) } @@ -359,46 +359,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()) @@ -628,10 +633,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) } @@ -639,14 +643,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() }) } } @@ -785,10 +788,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) } @@ -796,7 +797,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) } @@ -867,24 +868,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") @@ -892,18 +889,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 85299d4ea..aa6c2ca92 100644 --- a/field_test.go +++ b/field_test.go @@ -27,18 +27,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) @@ -47,7 +53,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") @@ -62,25 +68,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) @@ -97,11 +103,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) } }) @@ -113,12 +119,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) } }) @@ -131,11 +135,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) } }) @@ -203,13 +207,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) @@ -243,18 +247,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) @@ -262,14 +266,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 fc8ca98ea..68802e15e 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -29,22 +29,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) } @@ -96,14 +96,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 60e3c8765..8a22543b2 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -751,8 +751,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() { @@ -997,7 +997,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") } @@ -1325,7 +1325,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 ec981ddf8..9cbff1346 100644 --- a/test/holder.go +++ b/test/holder.go @@ -9,13 +9,12 @@ import ( pilosa "github.com/molecula/featurebase/v3" "github.com/molecula/featurebase/v3/pql" "github.com/molecula/featurebase/v3/testhook" - "github.com/molecula/featurebase/v3/vprint" - "github.com/pkg/errors" ) // 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. @@ -28,7 +27,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 } @@ -37,7 +36,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 } @@ -57,7 +56,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} } @@ -69,37 +68,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() } @@ -109,15 +110,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() } @@ -134,15 +136,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. @@ -153,15 +157,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. @@ -179,16 +185,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 } @@ -200,11 +207,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) } @@ -229,6 +236,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 b96111c57..874464d9b 100644 --- a/txfactory.go +++ b/txfactory.go @@ -168,9 +168,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(), @@ -186,6 +186,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 2adebe9fc..042f3734c 100644 --- a/view.go +++ b/view.go @@ -448,16 +448,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 @@ -526,7 +524,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) @@ -534,102 +532,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) }