From f062e5a85f26792feeb0e825436e229a913c6e54 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Peter=20=C5=A0tibran=C3=BD?= Date: Tue, 23 Sep 2025 18:04:49 +0200 Subject: [PATCH] Stop background tasks and when stopping bleve backend. (#111249) Co-authored-by: Stephanie Hingtgen --- pkg/storage/unified/search/bleve.go | 40 +++++++++++++++---- .../unified/search/bleve_integration_test.go | 4 +- .../unified/search/bleve_search_test.go | 2 +- pkg/storage/unified/search/bleve_test.go | 13 +++--- .../unified/sql/test/integration_test.go | 2 +- 5 files changed, 42 insertions(+), 19 deletions(-) diff --git a/pkg/storage/unified/search/bleve.go b/pkg/storage/unified/search/bleve.go index 95bb77c2dc0..572bc8eb224 100644 --- a/pkg/storage/unified/search/bleve.go +++ b/pkg/storage/unified/search/bleve.go @@ -88,6 +88,9 @@ type bleveBackend struct { cache map[resource.NamespacedResource]*bleveIndex indexMetrics *resource.BleveIndexMetrics + + metricsUpdaterCancel func() + metricsUpdaterWg sync.WaitGroup } func NewBleveBackend(opts BleveOptions, tracer trace.Tracer, indexMetrics *resource.BleveIndexMetrics) (*bleveBackend, error) { @@ -129,7 +132,14 @@ func NewBleveBackend(opts BleveOptions, tracer trace.Tracer, indexMetrics *resou indexMetrics: indexMetrics, } - go be.updateIndexSizeMetric(opts.Root) + if be.indexMetrics != nil { + ctx, cancel := context.WithCancel(context.Background()) + be.metricsUpdaterCancel = cancel + be.metricsUpdaterWg.Add(1) + go be.updateIndexSizeMetric(ctx, opts.Root) + } else { + be.metricsUpdaterCancel = func() { /* empty */ } + } return be, nil } @@ -184,18 +194,19 @@ func (b *bleveBackend) getCachedIndex(key resource.NamespacedResource) *bleveInd } // updateIndexSizeMetric sets the total size of all file-based indices metric. -func (b *bleveBackend) updateIndexSizeMetric(indexPath string) { - if b.indexMetrics == nil { - return - } +func (b *bleveBackend) updateIndexSizeMetric(ctx context.Context, indexPath string) { + defer b.metricsUpdaterWg.Done() - for { + for ctx.Err() == nil { var totalSize int64 err := filepath.WalkDir(indexPath, func(path string, info os.DirEntry, err error) error { if err != nil { return err } + if err = ctx.Err(); err != nil { + return err + } if !info.IsDir() { fileInfo, err := info.Info() if err != nil { @@ -212,7 +223,12 @@ func (b *bleveBackend) updateIndexSizeMetric(indexPath string) { b.log.Error("got error while trying to calculate bleve file index size", "error", err) } - time.Sleep(60 * time.Second) + select { + case <-ctx.Done(): + return + case <-time.After(60 * time.Second): + continue + } } } @@ -626,7 +642,15 @@ func (b *bleveBackend) findPreviousFileBasedIndex(resourceDir string, minBuildTi return nil, "", 0 } -func (b *bleveBackend) CloseAllIndexes() { +// Stop closes all indexes and stops background tasks. +func (b *bleveBackend) Stop() { + b.closeAllIndexes() + + b.metricsUpdaterCancel() + b.metricsUpdaterWg.Wait() +} + +func (b *bleveBackend) closeAllIndexes() { b.cacheMx.Lock() defer b.cacheMx.Unlock() diff --git a/pkg/storage/unified/search/bleve_integration_test.go b/pkg/storage/unified/search/bleve_integration_test.go index 36f19a0bfd3..9ee16a1b559 100644 --- a/pkg/storage/unified/search/bleve_integration_test.go +++ b/pkg/storage/unified/search/bleve_integration_test.go @@ -24,7 +24,7 @@ func TestBleveSearchBackend(t *testing.T) { require.NoError(t, err) require.NotNil(t, backend) - t.Cleanup(backend.CloseAllIndexes) + t.Cleanup(backend.Stop) return backend }, &unitest.TestOptions{ @@ -49,7 +49,7 @@ func TestSearchBackendBenchmark(t *testing.T) { require.NoError(t, err) require.NotNil(t, backend) - t.Cleanup(backend.CloseAllIndexes) + t.Cleanup(backend.Stop) unitest.BenchmarkSearchBackend(t, backend, opts) } diff --git a/pkg/storage/unified/search/bleve_search_test.go b/pkg/storage/unified/search/bleve_search_test.go index 3b65ba62316..f3ed214deb5 100644 --- a/pkg/storage/unified/search/bleve_search_test.go +++ b/pkg/storage/unified/search/bleve_search_test.go @@ -250,7 +250,7 @@ func newTestDashboardsIndex(t testing.TB, threshold int64, size int64, batchSize }, tracing.NewNoopTracerService(), nil) require.NoError(t, err) - t.Cleanup(backend.CloseAllIndexes) + t.Cleanup(backend.Stop) ctx := identity.WithRequester(context.Background(), &user.SignedInUser{Namespace: "ns"}) diff --git a/pkg/storage/unified/search/bleve_test.go b/pkg/storage/unified/search/bleve_test.go index 324bba1f60f..18afe19abb2 100644 --- a/pkg/storage/unified/search/bleve_test.go +++ b/pkg/storage/unified/search/bleve_test.go @@ -41,8 +41,7 @@ func TestMain(m *testing.M) { goleak.VerifyTestMain(m, goleak.IgnoreTopFunction("github.com/open-feature/go-sdk/openfeature.(*eventExecutor).startEventListener.func1.1"), goleak.IgnoreTopFunction("go.opencensus.io/stats/view.(*worker).start"), - goleak.IgnoreTopFunction("github.com/blevesearch/bleve_index_api.AnalysisWorker"), // These don't stop when index is closed. - goleak.IgnoreAnyFunction("github.com/grafana/grafana/pkg/storage/unified/search.(*bleveBackend).updateIndexSizeMetric"), // We don't have a way to stop this one yet. + goleak.IgnoreTopFunction("github.com/blevesearch/bleve_index_api.AnalysisWorker"), // These don't stop when index is closed. ) } @@ -66,7 +65,7 @@ func TestBleveBackend(t *testing.T) { }, tracing.NewNoopTracerService(), nil) require.NoError(t, err) - t.Cleanup(backend.CloseAllIndexes) + t.Cleanup(backend.Stop) rv := int64(10) ctx := identity.WithRequester(context.Background(), &user.SignedInUser{Namespace: "ns"}) @@ -783,7 +782,7 @@ func setupBleveBackend(t *testing.T, options ...setupOption) (*bleveBackend, pro backend, err := NewBleveBackend(opts, tracing.NewNoopTracerService(), metrics) require.NoError(t, err) require.NotNil(t, backend) - t.Cleanup(backend.CloseAllIndexes) + t.Cleanup(backend.Stop) return backend, reg } @@ -895,9 +894,9 @@ func TestCloseAllIndexes(t *testing.T) { // Verify two open indexes. checkOpenIndexes(t, reg, 1, 1) - backend1.CloseAllIndexes() + backend1.closeAllIndexes() - // Verify that there are no open indexes after CloseAllIndexes call. + // Verify that there are no open indexes after closeAllIndexes call. checkOpenIndexes(t, reg, 0, 0) } @@ -962,7 +961,7 @@ func TestBuildIndex(t *testing.T) { backend, _ := setupBleveBackend(t, withFileThreshold(5), withRootDir(tmpDir), withBuildVersion(version)) _, err := backend.BuildIndex(context.Background(), ns, firstIndexDocsCount, nil, "test", indexTestDocs(ns, firstIndexDocsCount, 100), nil, rebuild) require.NoError(t, err) - backend.CloseAllIndexes() + backend.Stop() } // Make sure we pass at least 1 nanosecond (alwaysRebuildDueToAge) to ensure that the index needs to be rebuild. diff --git a/pkg/storage/unified/sql/test/integration_test.go b/pkg/storage/unified/sql/test/integration_test.go index 0166a4a96f7..935f469c4af 100644 --- a/pkg/storage/unified/sql/test/integration_test.go +++ b/pkg/storage/unified/sql/test/integration_test.go @@ -110,7 +110,7 @@ func TestIntegrationSearchAndStorage(t *testing.T) { require.NoError(t, err) require.NotNil(t, search) - t.Cleanup(search.CloseAllIndexes) + t.Cleanup(search.Stop) // Create a new resource backend storage := newTestBackend(t, false, 0)