From 8a95ac344bfc3d1ddf246bef47bc63bc420d42c8 Mon Sep 17 00:00:00 2001 From: Souhaila Noor Date: Fri, 25 Feb 2022 09:28:13 -0600 Subject: [PATCH] We check for incomplete deletion when server is started. When deletion is started, _exists field is updated with row+1. After deletion is completed, we delete _exists=row+1. If _exists>=1, then deletion was not completed. Updated go version in docker to match other requirements. Removed duplicate error check for grpc. --- executor.go | 24 ++++---- executor_internal_test.go | 50 ++++++++++++++++ holder.go | 47 +++++++++++++++ holder_internal_test.go | 74 ++++++++++++++++++++++++ internal/clustertests/Dockerfile-fakeIDP | 2 +- server/server.go | 4 -- 6 files changed, 185 insertions(+), 16 deletions(-) diff --git a/executor.go b/executor.go index cd0fe44b4..7dfc994f2 100644 --- a/executor.go +++ b/executor.go @@ -8265,10 +8265,6 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i if len(row.segments) == 0 { //nothing to remove return false, nil } - columns := row.segments[0].data //should only be one segment - if columns.Count() == 0 { - return false, nil - } // Fetch index. idx := e.Holder.Index(index) @@ -8276,14 +8272,20 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i return false, newNotFoundError(ErrIndexNotFound, index) } + return DeleteRows(row, idx, shard) +} + +func DeleteRows(row *Row, idx *Index, shard uint64) (bool, error) { + tx := idx.holder.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard}) + defer tx.Rollback() + + columns := row.segments[0].data //should only be one segment + if columns.Count() == 0 { + return false, nil + } columnIDs := make([]uint64, 0) none := make([]uint64, 0) // no bits will be set - tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) - if err != nil { - return false, err - } - defer finisher(&err) changed := false colCounts := make([]int, 0) toClear := columnIDs[:0] @@ -8301,7 +8303,7 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i toClear = columnIDs[:0] rowSet = make(map[uint64]struct{}) - err = tx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) + err := tx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) if err != nil { return false, err } @@ -8332,5 +8334,5 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i } } } - return changed, nil + return changed, tx.Commit() } diff --git a/executor_internal_test.go b/executor_internal_test.go index 171cb2ae3..628b63155 100644 --- a/executor_internal_test.go +++ b/executor_internal_test.go @@ -545,3 +545,53 @@ func TestDistinctTimestampUnion(t *testing.T) { }) } } + +func TestExecutor_DeleteRows(t *testing.T) { + path, _ := testhook.TempDir(t, "pilosa-executor-") + holder := NewHolder(path, mustHolderConfig()) + defer holder.Close() + + if err := holder.Open(); err != nil { + t.Fatalf("opening holder: %v", err) + } + + idx, err := holder.CreateIndex("i", IndexOptions{TrackExistence: true}) + if err != nil { + t.Fatalf("creating index: %v", err) + } + + f, err := idx.CreateField("f", OptFieldTypeDefault()) + if err != nil { + 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 { + 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) + if err != nil { + t.Fatalf("failed to read row: %v", err) + } + + changed, err := DeleteRows(row, idx, shard) + if !changed || err != nil { + t.Fatalf("failed to delete row: %v", err) + } + + changed, err = DeleteRows(row, idx, shard) + if changed { + t.Fatalf("expected delete to not clear bit but it did") + } +} diff --git a/holder.go b/holder.go index db38a6db7..2f24793ee 100644 --- a/holder.go +++ b/holder.go @@ -287,6 +287,50 @@ func (h *Holder) IndexesPath() string { return filepath.Join(h.path, IndexesDir) } +// processDeleteInflight checks if deletion was in progress when server shutdown +// the _exists field is set to row+1 when delete is started. Upon completion, the row is deleted. +// if _exists>=1, we finish deleting the rows +func (h *Holder) processDeleteInflight() error { + for _, index := range h.indexes { + if index.trackExistence { + shards := index.AvailableShards(includeRemote).Slice() + + for _, shard := range shards { + inprocessRowIDs := NewRow() + + frag := h.fragment(index.name, existenceFieldName, viewStandard, shard) + if frag == nil { + continue + } + + tx := index.Txf().NewTx(Txo{Write: !writable, Index: index, Shard: shard}) + defer tx.Rollback() + + // filter rows based on having _exists>=1, which is used to flag delete in-flight + rows, err := frag.rows(context.Background(), tx, 1) + if err != nil { + return err + } + + // check if any rows are found + if len(rows) == 0 { + return nil + } + + for _, rowID := range rows { + row, err2 := frag.row(tx, rowID) + if err2 != nil { + return err2 + } + inprocessRowIDs = inprocessRowIDs.Union(row) + } + DeleteRows(inprocessRowIDs, index, shard) + } + } + } + return nil +} + // Open initializes the root data directory for the holder. func (h *Holder) Open() error { h.opening = true @@ -380,6 +424,9 @@ func (h *Holder) Open() error { return errors.Wrap(err, "processing foreign index fields") } + // Check if deletion was in progress when server was shutdown + h.processDeleteInflight() + h.Stats.Open() h.opened.Close() diff --git a/holder_internal_test.go b/holder_internal_test.go index c9e164f27..2fec338af 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -2,7 +2,10 @@ package pilosa import ( + "testing" + "github.com/molecula/featurebase/v3/disco" + "github.com/molecula/featurebase/v3/testhook" ) // mustHolderConfig sets up a default holder config for tests. @@ -14,3 +17,74 @@ func mustHolderConfig() *HolderConfig { cfg.Sharder = disco.InMemSharder return cfg } + +func TestHolder_ProcessDeleteInflight(t *testing.T) { + path, _ := testhook.TempDir(t, "delete-inflight") + h := NewHolder(path, mustHolderConfig()) + defer h.Close() + + err := h.Open() + if err != nil { + t.Fatalf("failed to open holder: %v", err) + } + + idx, err := h.CreateIndexIfNotExists("i", IndexOptions{TrackExistence: true}) + if err != nil { + t.Fatalf("failed to create index: %v", err) + } + f, err := idx.CreateFieldIfNotExists("f", OptFieldTypeDefault()) + if err != nil { + t.Fatalf("failed to create field: %v", err) + } + + existencefield := idx.existenceFld + shard := uint64(0) + tx := idx.Txf().NewTx(Txo{Write: true, Index: idx, Shard: shard}) + defer tx.Rollback() + + rowCol := []struct { + row uint64 + col uint64 + }{ + {1, 1}, + {1, 2}, + {30, 33}, + {22, 2}, + } + for _, r := range rowCol { + _, err = f.SetBit(tx, r.row, r.col, nil) + if err != nil { + t.Fatalf("failed to set bit: %v", err) + } + + _, err = existencefield.SetBit(tx, r.row, r.col, nil) + if err != nil { + t.Fatalf("failed to set bit: %v", err) + } + } + + if err = tx.Commit(); err != nil { + t.Fatalf("failed to commit tx: %v", err) + } + + err = h.processDeleteInflight() + if err != nil { + t.Fatalf("failed to delete: %v", err) + } + + tx = idx.Txf().NewTx(Txo{Write: false, Index: idx, Shard: shard}) + defer tx.Rollback() + for _, r := range rowCol { + row, err := f.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + existenceRow, err := existencefield.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + if len(row.Columns()) != 0 || len(existenceRow.Columns()) != 0 { + t.Fatalf("expected columns for fields to be empty after delete") + } + } +} diff --git a/internal/clustertests/Dockerfile-fakeIDP b/internal/clustertests/Dockerfile-fakeIDP index b53d4d2c2..b46556b80 100644 --- a/internal/clustertests/Dockerfile-fakeIDP +++ b/internal/clustertests/Dockerfile-fakeIDP @@ -1,4 +1,4 @@ -FROM golang:latest +FROM golang:1.16 WORKDIR / COPY fakeidp ./ diff --git a/server/server.go b/server/server.go index 697816f16..c73cabda3 100644 --- a/server/server.go +++ b/server/server.go @@ -519,10 +519,6 @@ func (m *Command) SetupServer() error { // Tell server about its new API, which its client will need. m.Server.SetAPI(m.API) - if err != nil { - return errors.Wrap(err, "new grpc server") - } - var p authz.GroupPermissions if m.Config.Auth.Enable { m.Config.MustValidateAuth()