Stop background tasks and when stopping bleve backend. (#111249)
Co-authored-by: Stephanie Hingtgen <stephanie.hingtgen@grafana.com>
This commit is contained in:
co-authored by
Stephanie Hingtgen
parent
9494175984
commit
f062e5a85f
@@ -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()
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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"})
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user