From 964e7d86c88ba1db2982eacefd78ce3c8e722ead Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 27 Jul 2022 15:32:39 -0500 Subject: [PATCH] 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. --- executor.go | 94 +++++++---------------------------- executor_internal_test.go | 38 +++++--------- executor_test.go | 18 +++---- field.go | 46 +++++++++++------ field_internal_test.go | 76 ++++++++++++++-------------- field_test.go | 66 +++++++++++++------------ holder_internal_test.go | 20 ++++---- server/grpc.go | 8 +-- test/holder.go | 101 ++++++++++++++++++++------------------ txfactory.go | 24 +++++++-- view.go | 74 ++++------------------------ 11 files changed, 241 insertions(+), 324 deletions(-) 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) }