diff --git a/holder.go b/holder.go index 4e9030f40..e40edc3dc 100644 --- a/holder.go +++ b/holder.go @@ -287,6 +287,46 @@ func (h *Holder) IndexesPath() string { return filepath.Join(h.path, IndexesDir) } +func (h *Holder) deletePerShard(index *Index, shard uint64) error { + inprocessRecords := NewRow() + + frag := h.fragment(index.name, existenceFieldName, viewStandard, shard) + if frag == nil { + return nil + } + + 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 _, record := range rows { + row, err2 := frag.row(tx, record) + if err2 != nil { + return fmt.Errorf("getting row IDs: %v", err2) + } + inprocessRecords = inprocessRecords.Union(row) + } + h.Logger.Printf("retrying delete: index=%v shard=%v record count=%v", index.name, shard, inprocessRecords.Count()) + + tx.Rollback() // release the read tx in case a checksum is needed in DeleteRows + + _, err = DeleteRows(context.Background(), inprocessRecords, index, shard) + if err != nil { + return fmt.Errorf("deleting rows: %v", err) + } + return nil +} + // 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 @@ -294,37 +334,26 @@ func (h *Holder) processDeleteInflight() error { for _, index := range h.Indexes() { if index.trackExistence { shards := index.AvailableShards(includeRemote).Slice() - + index := index + ch := make(chan uint64, len(shards)) for _, shard := range shards { - inprocessRowIDs := NewRow() + ch <- shard + } + close(ch) - 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 + g := new(errgroup.Group) + for i := 0; i < runtime.NumCPU(); i++ { + g.Go(func() error { + for shard := range ch { + if err := h.deletePerShard(index, shard); err != nil { + return fmt.Errorf("delete shard %d: %w", shard, err) + } } - inprocessRowIDs = inprocessRowIDs.Union(row) - } - DeleteRows(context.Background(), inprocessRowIDs, index, shard) + return nil + }) + } + if err := g.Wait(); err != nil { + return err } } } diff --git a/holder_internal_test.go b/holder_internal_test.go index 2fec338af..fc8ca98ea 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -18,6 +18,48 @@ func mustHolderConfig() *HolderConfig { return cfg } +func setupTest(t *testing.T, h *Holder, rowCol []rowCols, indexName string) (*Index, *Field) { + idx, err := h.CreateIndexIfNotExists(indexName, IndexOptions{TrackExistence: true}) + if err != nil { + t.Fatalf("failed to create index %v: %v", indexName, err) + } + f, err := idx.CreateFieldIfNotExists("f", OptFieldTypeDefault()) + if err != nil { + t.Fatalf("failed to create field in index %v: %v", indexName, err) + } + existencefield := idx.existenceFld + + shard := uint64(0) + tx := idx.Txf().NewTx(Txo{Write: true, Index: idx, Shard: shard}) + defer tx.Rollback() + for _, r := range rowCol { + _, err = f.SetBit(tx, 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) + if err != nil { + t.Fatalf("failed to set bit in index %v: %v", indexName, err) + } + } + + if err = tx.Commit(); err != nil { + t.Fatalf("failed to commit tx for index %v: %v", indexName, err) + } + + shardsFound := idx.AvailableShards(includeRemote).Slice() + if len(shardsFound) != 3 { + t.Fatalf("expected 3 shards for index %v, got %v", indexName, len(shardsFound)) + } + return idx, f +} + +type rowCols struct { + row uint64 + col uint64 +} + func TestHolder_ProcessDeleteInflight(t *testing.T) { path, _ := testhook.TempDir(t, "delete-inflight") h := NewHolder(path, mustHolderConfig()) @@ -28,63 +70,48 @@ func TestHolder_ProcessDeleteInflight(t *testing.T) { 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 - }{ + rowCol := []rowCols{ {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) - } + {10, ShardWidth + 1}, + {1, ShardWidth * 2}, } - if err = tx.Commit(); err != nil { - t.Fatalf("failed to commit tx: %v", err) - } + idx1, f1 := setupTest(t, h, rowCol, "idxdelete1") + idx2, f2 := setupTest(t, h, rowCol, "idxdelete2") 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") - } + tests := []struct { + idx *Index + f *Field + }{ + {idx1, f1}, + {idx2, f2}, + } + + 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() + for _, r := range rowCol { + row, err := f.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + existenceRow, err := idx.existenceFld.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") + } + } + }() + } }