From 34496e137c67c5f0d23bd05d5912d61cc08fcda5 Mon Sep 17 00:00:00 2001 From: Will Assis <35489495+gassiss@users.noreply.github.com> Date: Tue, 19 Aug 2025 09:41:35 -0400 Subject: [PATCH] unified-storage: Search after write save rv to index (#109641) * save rv to index after index is built --- pkg/storage/unified/resource/search.go | 6 +- pkg/storage/unified/search/bleve.go | 144 ++++++++++++++++------- pkg/storage/unified/search/bleve_test.go | 56 ++++++--- 3 files changed, 146 insertions(+), 60 deletions(-) diff --git a/pkg/storage/unified/resource/search.go b/pkg/storage/unified/resource/search.go index 614d0dafb50..f7b822df600 100644 --- a/pkg/storage/unified/resource/search.go +++ b/pkg/storage/unified/resource/search.go @@ -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 }) } diff --git a/pkg/storage/unified/search/bleve.go b/pkg/storage/unified/search/bleve.go index 6c0ca8a14ff..347ee8103aa 100644 --- a/pkg/storage/unified/search/bleve.go +++ b/pkg/storage/unified/search/bleve.go @@ -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") diff --git a/pkg/storage/unified/search/bleve_test.go b/pkg/storage/unified/search/bleve_test.go index e28d8a80cb8..fbec10dc9f2 100644 --- a/pkg/storage/unified/search/bleve_test.go +++ b/pkg/storage/unified/search/bleve_test.go @@ -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) }