unified-storage: Search after write save rv to index (#109641)
* save rv to index after index is built
This commit is contained in:
@@ -722,7 +722,7 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size
|
||||
span := trace.SpanFromContext(ctx)
|
||||
span.AddEvent("building index", trace.WithAttributes(attribute.Int64("size", size), attribute.Int64("rv", rv), attribute.String("reason", indexBuildReason)))
|
||||
|
||||
rv, err = s.storage.ListIterator(ctx, &resourcepb.ListRequest{
|
||||
listRV, err := s.storage.ListIterator(ctx, &resourcepb.ListRequest{
|
||||
Limit: 1000000000000, // big number
|
||||
Options: &resourcepb.ListOptions{
|
||||
Key: &resourcepb.ResourceKey{
|
||||
@@ -790,7 +790,7 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size
|
||||
}
|
||||
return iter.Error()
|
||||
})
|
||||
return rv, err
|
||||
return listRV, err
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
@@ -821,7 +821,7 @@ func (s *searchSupport) buildEmptyIndex(ctx context.Context, nsr NamespacedResou
|
||||
// Build an empty index by passing a builder function that doesn't add any documents
|
||||
return s.search.BuildIndex(ctx, nsr, 0, rv, fields, "empty", func(index ResourceIndex) (int64, error) {
|
||||
// Return the resource version without adding any documents to the index
|
||||
return rv, nil
|
||||
return 0, nil
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ package search
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -190,8 +191,13 @@ func (b *bleveBackend) updateIndexSizeMetric(indexPath string) {
|
||||
}
|
||||
}
|
||||
|
||||
// BuildIndex builds an index from scratch.
|
||||
// BuildIndex builds an index from scratch or retrieves it from the filesystem.
|
||||
// If built successfully, the new index replaces the old index in the cache (if there was any).
|
||||
// An index in the file system is considered to be valid if the requested resourceVersion is smaller than or equal to
|
||||
// the resourceVersion used to build the index and the number of indexed objects matches the expected size.
|
||||
// The return value of "builder" should be the RV returned from List. This will be stored as the index RV
|
||||
//
|
||||
//nolint:gocyclo
|
||||
func (b *bleveBackend) BuildIndex(
|
||||
ctx context.Context,
|
||||
key resource.NamespacedResource,
|
||||
@@ -249,6 +255,7 @@ func (b *bleveBackend) BuildIndex(
|
||||
resourceDir := b.getResourceDir(key)
|
||||
|
||||
var index bleve.Index
|
||||
var indexRV int64
|
||||
cachedIndex := b.getCachedIndex(key)
|
||||
fileIndexName := "" // Name of the file-based index, or empty for in-memory indexes.
|
||||
newIndexType := indexStorageMemory
|
||||
@@ -261,12 +268,12 @@ func (b *bleveBackend) BuildIndex(
|
||||
// This happens on startup, or when memory-based index has expired. (We don't expire file-based indexes)
|
||||
// If we do have an unexpired cached index already, we always build a new index from scratch.
|
||||
if cachedIndex == nil && resourceVersion > 0 {
|
||||
index, fileIndexName = b.findPreviousFileBasedIndex(resourceDir, resourceVersion, size)
|
||||
index, fileIndexName, indexRV = b.findPreviousFileBasedIndex(resourceDir, resourceVersion, size)
|
||||
}
|
||||
|
||||
if index != nil {
|
||||
build = false
|
||||
logWithDetails.Debug("Existing index found on filesystem", "directory", filepath.Join(resourceDir, fileIndexName))
|
||||
logWithDetails.Debug("Existing index found on filesystem", "indexRV", indexRV, "directory", filepath.Join(resourceDir, fileIndexName))
|
||||
defer closeIndexOnExit(index, "") // Close index, but don't delete directory.
|
||||
} else {
|
||||
// Building index from scratch. Index name has a time component in it to be unique, but if
|
||||
@@ -275,7 +282,7 @@ func (b *bleveBackend) BuildIndex(
|
||||
indexDir := ""
|
||||
now := time.Now()
|
||||
for index == nil {
|
||||
fileIndexName = formatIndexName(time.Now(), resourceVersion)
|
||||
fileIndexName = formatIndexName(now)
|
||||
indexDir = filepath.Join(resourceDir, fileIndexName)
|
||||
if !isPathWithinRoot(indexDir, b.opts.Root) {
|
||||
return nil, fmt.Errorf("invalid path %s", indexDir)
|
||||
@@ -322,7 +329,7 @@ func (b *bleveBackend) BuildIndex(
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
_, err = builder(idx)
|
||||
listRV, err := builder(idx)
|
||||
if err != nil {
|
||||
logWithDetails.Error("Failed to build index", "err", err)
|
||||
if b.indexMetrics != nil {
|
||||
@@ -330,13 +337,25 @@ func (b *bleveBackend) BuildIndex(
|
||||
}
|
||||
return nil, fmt.Errorf("failed to build index: %w", err)
|
||||
}
|
||||
err = idx.updateResourceVersion(listRV)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("fail to persist rv to index: %w", err)
|
||||
}
|
||||
|
||||
elapsed := time.Since(start)
|
||||
logWithDetails.Info("Finished building index", "elapsed", elapsed)
|
||||
|
||||
if b.indexMetrics != nil {
|
||||
b.indexMetrics.IndexCreationTime.WithLabelValues().Observe(elapsed.Seconds())
|
||||
}
|
||||
} else {
|
||||
logWithDetails.Info("Skipping index build, using existing index")
|
||||
|
||||
idx.resourceVersion, err = getRV(index)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to get RV from bleve index: %w", err)
|
||||
}
|
||||
|
||||
if b.indexMetrics != nil {
|
||||
b.indexMetrics.IndexBuildSkipped.Inc()
|
||||
}
|
||||
@@ -474,62 +493,61 @@ func (b *bleveBackend) TotalDocs() int64 {
|
||||
return totalDocs
|
||||
}
|
||||
|
||||
func formatIndexName(now time.Time, resourceVersion int64) string {
|
||||
timestamp := now.Format("20060102-150405")
|
||||
return fmt.Sprintf("%s-%d", timestamp, resourceVersion)
|
||||
func formatIndexName(now time.Time) string {
|
||||
return now.Format("20060102-150405")
|
||||
}
|
||||
|
||||
func (b *bleveBackend) findPreviousFileBasedIndex(resourceDir string, resourceVersion int64, size int64) (bleve.Index, string) {
|
||||
func (b *bleveBackend) findPreviousFileBasedIndex(resourceDir string, resourceVersion int64, size int64) (bleve.Index, string, int64) {
|
||||
entries, err := os.ReadDir(resourceDir)
|
||||
if err != nil {
|
||||
return nil, ""
|
||||
return nil, "", 0
|
||||
}
|
||||
|
||||
indexName := ""
|
||||
for _, ent := range entries {
|
||||
if !ent.IsDir() {
|
||||
continue
|
||||
}
|
||||
|
||||
parts := strings.Split(ent.Name(), "-")
|
||||
if len(parts) != 3 {
|
||||
continue
|
||||
}
|
||||
|
||||
// Last part is resourceVersion
|
||||
indexRv, err := strconv.ParseInt(parts[2], 10, 64)
|
||||
indexName := ent.Name()
|
||||
indexDir := filepath.Join(resourceDir, indexName)
|
||||
idx, err := bleve.Open(indexDir)
|
||||
if err != nil {
|
||||
b.log.Debug("error opening index", "indexDir", indexDir, "err", err)
|
||||
continue
|
||||
}
|
||||
if indexRv != resourceVersion {
|
||||
|
||||
cnt, err := idx.DocCount()
|
||||
if err != nil {
|
||||
b.log.Debug("error getting count from index", "indexDir", indexDir, "err", err)
|
||||
_ = idx.Close()
|
||||
continue
|
||||
}
|
||||
indexName = ent.Name()
|
||||
break
|
||||
|
||||
if uint64(size) != cnt {
|
||||
b.log.Debug("index count mismatch. ignoring index", "indexDir", indexDir, "size", size, "cnt", cnt)
|
||||
_ = idx.Close()
|
||||
continue
|
||||
}
|
||||
|
||||
indexRV, err := getRV(idx)
|
||||
if err != nil {
|
||||
b.log.Error("error getting rv from index", "indexDir", indexDir, "err", err)
|
||||
if !errors.Is(err, bleve.ErrorIndexClosed) {
|
||||
_ = idx.Close()
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
if indexRV < resourceVersion {
|
||||
b.log.Debug("indexRV is less than requested resourceVersion. ignoring index", "indexDir", indexDir, "rv", indexRV, "resourceVersion", resourceVersion)
|
||||
_ = idx.Close()
|
||||
continue
|
||||
}
|
||||
|
||||
return idx, indexName, indexRV
|
||||
}
|
||||
|
||||
if indexName == "" {
|
||||
return nil, ""
|
||||
}
|
||||
|
||||
indexDir := filepath.Join(resourceDir, indexName)
|
||||
idx, err := bleve.Open(indexDir)
|
||||
if err != nil {
|
||||
return nil, ""
|
||||
}
|
||||
|
||||
cnt, err := idx.DocCount()
|
||||
if err != nil {
|
||||
_ = idx.Close()
|
||||
return nil, ""
|
||||
}
|
||||
|
||||
if uint64(size) != cnt {
|
||||
_ = idx.Close()
|
||||
return nil, ""
|
||||
}
|
||||
|
||||
return idx, indexName
|
||||
return nil, "", 0
|
||||
}
|
||||
|
||||
func (b *bleveBackend) CloseAllIndexes() {
|
||||
@@ -550,6 +568,8 @@ type bleveIndex struct {
|
||||
key resource.NamespacedResource
|
||||
index bleve.Index
|
||||
|
||||
resourceVersion int64
|
||||
|
||||
standard resource.SearchableDocumentFields
|
||||
fields resource.SearchableDocumentFields
|
||||
|
||||
@@ -592,6 +612,44 @@ func (b *bleveIndex) BulkIndex(req *resource.BulkIndexRequest) error {
|
||||
return b.index.Batch(batch)
|
||||
}
|
||||
|
||||
var internalRVKey = []byte("rv")
|
||||
|
||||
func (b *bleveIndex) updateResourceVersion(rv int64) error {
|
||||
if rv == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
if err := setRV(b.index, rv); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
b.resourceVersion = rv
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func setRV(index bleve.Index, rv int64) error {
|
||||
buf := make([]byte, 8)
|
||||
binary.BigEndian.PutUint64(buf, uint64(rv))
|
||||
|
||||
return index.SetInternal(internalRVKey, buf)
|
||||
}
|
||||
|
||||
// getRV will call index.GetInternal to retrieve the RV saved in the index. If index is closed, it will return a
|
||||
// bleve.ErrorIndexClosed error. If there's no RV saved in the index, or it's invalid format, it will return 0
|
||||
func getRV(index bleve.Index) (int64, error) {
|
||||
raw, err := index.GetInternal(internalRVKey)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
if len(raw) < 8 {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
return int64(binary.BigEndian.Uint64(raw)), nil
|
||||
}
|
||||
|
||||
func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest) (*resourcepb.ListManagedObjectsResponse, error) {
|
||||
if req.NextPageToken != "" {
|
||||
return nil, fmt.Errorf("next page not implemented yet")
|
||||
|
||||
@@ -768,7 +768,7 @@ func TestBleveInMemoryIndexExpiration(t *testing.T) {
|
||||
Resource: "resource",
|
||||
}
|
||||
|
||||
builtIndex, err := backend.BuildIndex(context.Background(), ns, 1 /* below FileThreshold */, 100, nil, "test", indexTestDocs(ns, 1))
|
||||
builtIndex, err := backend.BuildIndex(context.Background(), ns, 1 /* below FileThreshold */, 100, nil, "test", indexTestDocs(ns, 1, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Wait for index expiration, which is 1ns
|
||||
@@ -800,7 +800,7 @@ func TestBleveFileIndexExpiration(t *testing.T) {
|
||||
}
|
||||
|
||||
// size=100 is above FileThreshold, this will be file-based index
|
||||
builtIndex, err := backend.BuildIndex(context.Background(), ns, 100, 100, nil, "test", indexTestDocs(ns, 1))
|
||||
builtIndex, err := backend.BuildIndex(context.Background(), ns, 100, 100, nil, "test", indexTestDocs(ns, 1, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Wait for index expiration, which is 1ns
|
||||
@@ -822,7 +822,7 @@ func TestBleveFileIndexExpiration(t *testing.T) {
|
||||
`), "index_server_open_indexes"))
|
||||
}
|
||||
|
||||
func TestFileIndexIsReusedOnSameSizeAndRV(t *testing.T) {
|
||||
func TestFileIndexIsReusedOnSameSizeAndRVLessThanIndexRV(t *testing.T) {
|
||||
ns := resource.NamespacedResource{
|
||||
Namespace: "test",
|
||||
Group: "group",
|
||||
@@ -832,7 +832,7 @@ func TestFileIndexIsReusedOnSameSizeAndRV(t *testing.T) {
|
||||
tmpDir := t.TempDir()
|
||||
|
||||
backend1, reg1 := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
_, err := backend1.BuildIndex(context.Background(), ns, 10 /* file based */, 100, nil, "test", indexTestDocs(ns, 10))
|
||||
_, err := backend1.BuildIndex(context.Background(), ns, 10 /* file based */, 100, nil, "test", indexTestDocs(ns, 10, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify one open index.
|
||||
@@ -855,7 +855,7 @@ func TestFileIndexIsReusedOnSameSizeAndRV(t *testing.T) {
|
||||
|
||||
// We open new backend using same directory, and run indexing with same size (10) and RV (100). This should reuse existing index, and skip indexing.
|
||||
backend2, reg2 := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
idx, err := backend2.BuildIndex(context.Background(), ns, 10 /* file based */, 100, nil, "test", indexTestDocs(ns, 1000))
|
||||
idx, err := backend2.BuildIndex(context.Background(), ns, 10 /* file based */, 100, nil, "test", indexTestDocs(ns, 1000, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify that we're reusing existing index and there is only 10 documents in it, not 1000.
|
||||
@@ -869,6 +869,34 @@ func TestFileIndexIsReusedOnSameSizeAndRV(t *testing.T) {
|
||||
index_server_open_indexes{index_storage="memory"} 0
|
||||
index_server_open_indexes{index_storage="file"} 1
|
||||
`), "index_server_open_indexes"))
|
||||
|
||||
backend2.CloseAllIndexes()
|
||||
// Verify that there are no open indexes after closeAllIndexes call.
|
||||
require.NoError(t, testutil.GatherAndCompare(reg2, bytes.NewBufferString(`
|
||||
# HELP index_server_open_indexes Number of open indexes per storage type. An open index corresponds to single resource group.
|
||||
# TYPE index_server_open_indexes gauge
|
||||
index_server_open_indexes{index_storage="memory"} 0
|
||||
index_server_open_indexes{index_storage="file"} 0
|
||||
`), "index_server_open_indexes"))
|
||||
|
||||
// We repeat with backend3 and RV 99. This should also reuse existing index and skip indexing
|
||||
backend3, reg3 := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
idx, err = backend3.BuildIndex(context.Background(), ns, 10 /* file based */, 99, nil, "test", indexTestDocs(ns, 1000, 99))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify that we're reusing existing index and there is only 10 documents in it, not 1000.
|
||||
cnt, err = idx.DocCount(context.Background(), "")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(10), cnt)
|
||||
|
||||
require.NoError(t, testutil.GatherAndCompare(reg3, bytes.NewBufferString(`
|
||||
# HELP index_server_open_indexes Number of open indexes per storage type. An open index corresponds to single resource group.
|
||||
# TYPE index_server_open_indexes gauge
|
||||
index_server_open_indexes{index_storage="memory"} 0
|
||||
index_server_open_indexes{index_storage="file"} 1
|
||||
`), "index_server_open_indexes"))
|
||||
|
||||
backend3.CloseAllIndexes()
|
||||
}
|
||||
|
||||
func TestFileIndexIsNotReusedOnDifferentSize(t *testing.T) {
|
||||
@@ -881,13 +909,13 @@ func TestFileIndexIsNotReusedOnDifferentSize(t *testing.T) {
|
||||
tmpDir := t.TempDir()
|
||||
|
||||
backend1, _ := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
_, err := backend1.BuildIndex(context.Background(), ns, 10, 100, nil, "test", indexTestDocs(ns, 10))
|
||||
_, err := backend1.BuildIndex(context.Background(), ns, 10, 100, nil, "test", indexTestDocs(ns, 10, 100))
|
||||
require.NoError(t, err)
|
||||
backend1.CloseAllIndexes()
|
||||
|
||||
// We open new backend using same directory, but with different size. Index should be rebuilt.
|
||||
backend2, _ := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
idx, err := backend2.BuildIndex(context.Background(), ns, 100, 100, nil, "test", indexTestDocs(ns, 100))
|
||||
idx, err := backend2.BuildIndex(context.Background(), ns, 100, 100, nil, "test", indexTestDocs(ns, 100, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify that index has updated number of documents.
|
||||
@@ -906,13 +934,13 @@ func TestFileIndexIsNotReusedOnDifferentRV(t *testing.T) {
|
||||
tmpDir := t.TempDir()
|
||||
|
||||
backend1, _ := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
_, err := backend1.BuildIndex(context.Background(), ns, 10, 100, nil, "test", indexTestDocs(ns, 10))
|
||||
_, err := backend1.BuildIndex(context.Background(), ns, 10, 100, nil, "test", indexTestDocs(ns, 10, 100))
|
||||
require.NoError(t, err)
|
||||
backend1.CloseAllIndexes()
|
||||
|
||||
// We open new backend using same directory, but with different RV. Index should be rebuilt.
|
||||
backend2, _ := setupBleveBackend(t, 5, time.Nanosecond, tmpDir)
|
||||
idx, err := backend2.BuildIndex(context.Background(), ns, 10 /* file based */, 999999, nil, "test", indexTestDocs(ns, 100))
|
||||
idx, err := backend2.BuildIndex(context.Background(), ns, 10 /* file based */, 999999, nil, "test", indexTestDocs(ns, 100, 999999))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify that index has updated number of documents.
|
||||
@@ -944,7 +972,7 @@ func TestRebuildingIndexClosesPreviousCachedIndex(t *testing.T) {
|
||||
if testCase.firstInMemory {
|
||||
firstSize = 1
|
||||
}
|
||||
firstIndex, err := backend.BuildIndex(context.Background(), ns, int64(firstSize), 100, nil, "test", indexTestDocs(ns, firstSize))
|
||||
firstIndex, err := backend.BuildIndex(context.Background(), ns, int64(firstSize), 100, nil, "test", indexTestDocs(ns, firstSize, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
if testCase.firstInMemory {
|
||||
@@ -960,7 +988,7 @@ func TestRebuildingIndexClosesPreviousCachedIndex(t *testing.T) {
|
||||
secondSize = 1
|
||||
openInMemoryIndexes = 1
|
||||
}
|
||||
secondIndex, err := backend.BuildIndex(context.Background(), ns, int64(secondSize), 100, nil, "test", indexTestDocs(ns, secondSize))
|
||||
secondIndex, err := backend.BuildIndex(context.Background(), ns, int64(secondSize), 100, nil, "test", indexTestDocs(ns, secondSize, 100))
|
||||
require.NoError(t, err)
|
||||
|
||||
if testCase.secondInMemory {
|
||||
@@ -1002,7 +1030,7 @@ func verifyDirEntriesCount(t *testing.T, dir string, count int) {
|
||||
require.Len(t, ents, count)
|
||||
}
|
||||
|
||||
func indexTestDocs(ns resource.NamespacedResource, docs int) func(index resource.ResourceIndex) (int64, error) {
|
||||
func indexTestDocs(ns resource.NamespacedResource, docs int, listRV int64) func(index resource.ResourceIndex) (int64, error) {
|
||||
return func(index resource.ResourceIndex) (int64, error) {
|
||||
var items []*resource.BulkIndexItem
|
||||
for i := 0; i < docs; i++ {
|
||||
@@ -1021,7 +1049,7 @@ func indexTestDocs(ns resource.NamespacedResource, docs int) func(index resource
|
||||
}
|
||||
|
||||
err := index.BulkIndex(&resource.BulkIndexRequest{Items: items})
|
||||
return int64(docs), err
|
||||
return listRV, err
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1083,6 +1111,6 @@ func testBleveIndexWithFailures(t *testing.T, fileBased bool) {
|
||||
require.Error(t, err)
|
||||
|
||||
// Even though previous build of the index failed, new building of the index should work.
|
||||
_, err = backend.BuildIndex(context.Background(), ns, size, 100, nil, "test", indexTestDocs(ns, int(size)))
|
||||
_, err = backend.BuildIndex(context.Background(), ns, size, 100, nil, "test", indexTestDocs(ns, int(size), 100))
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user