diff --git a/pkg/storage/unified/search/bleve.go b/pkg/storage/unified/search/bleve.go index 635bb09e04f..54d2de456fc 100644 --- a/pkg/storage/unified/search/bleve.go +++ b/pkg/storage/unified/search/bleve.go @@ -4,9 +4,11 @@ import ( "context" "fmt" "log/slog" + "os" "path/filepath" "strings" "sync" + "time" "github.com/blevesearch/bleve/v2" "github.com/blevesearch/bleve/v2/search" @@ -39,21 +41,32 @@ type bleveBackend struct { tracer trace.Tracer log *slog.Logger opts BleveOptions + start time.Time // cache info cache map[resource.NamespacedResource]*bleveIndex cacheMu sync.RWMutex } -func NewBleveBackend(opts BleveOptions, tracer trace.Tracer) *bleveBackend { - b := &bleveBackend{ +func NewBleveBackend(opts BleveOptions, tracer trace.Tracer) (*bleveBackend, error) { + if opts.Root == "" { + return nil, fmt.Errorf("bleve backend missing root folder configuration") + } + root, err := os.Stat(opts.Root) + if err != nil { + return nil, fmt.Errorf("error opening bleve root folder %w", err) + } + if !root.IsDir() { + return nil, fmt.Errorf("bleve root is configured against a file (not folder)") + } + + return &bleveBackend{ log: slog.Default().With("logger", "bleve-backend"), tracer: tracer, cache: make(map[resource.NamespacedResource]*bleveIndex), opts: opts, - } - - return b + start: time.Now(), + }, nil } // This will return nil if the key does not exist @@ -91,13 +104,38 @@ func (b *bleveBackend) BuildIndex(ctx context.Context, var err error var index bleve.Index + build := true mapper := getBleveMappings(fields) if size > b.opts.FileThreshold { - dir := filepath.Join(b.opts.Root, key.Namespace, fmt.Sprintf("%s.%s", key.Resource, key.Group)) - index, err = bleve.New(dir, mapper) + fname := fmt.Sprintf("rv%d", resourceVersion) + if resourceVersion == 0 { + fname = b.start.Format("tmp-20060102-150405") + } + dir := filepath.Join(b.opts.Root, key.Namespace, + fmt.Sprintf("%s.%s", key.Resource, key.Group), + fname, + ) + if resourceVersion > 0 { + info, _ := os.Stat(dir) + if info != nil && info.IsDir() { + index, err = bleve.Open(dir) // NOTE, will use the same mappings!!! + if err == nil { + found, err := index.DocCount() + if err != nil || int64(found) != size { + b.log.Info("this size changed since the last time the index opened") + _ = index.Close() + index = nil + } else { + build = false // no need to build the index + } + } + } + } - // TODO, check last RV so we can see if the numbers have changed + if index == nil { + index, err = bleve.New(dir, mapper) + } resource.IndexMetrics.IndexTenants.WithLabelValues(key.Namespace, "file").Inc() } else { @@ -123,15 +161,17 @@ func (b *bleveBackend) BuildIndex(ctx context.Context, return nil, err } - _, err = builder(idx) - if err != nil { - return nil, err - } + if build { + _, err = builder(idx) + if err != nil { + return nil, err + } - // Flush the batch - err = idx.Flush() - if err != nil { - return nil, err + // Flush the batch + err = idx.Flush() + if err != nil { + return nil, err + } } b.cacheMu.Lock() diff --git a/pkg/storage/unified/search/bleve_test.go b/pkg/storage/unified/search/bleve_test.go index c99dff15833..20571f82b81 100644 --- a/pkg/storage/unified/search/bleve_test.go +++ b/pkg/storage/unified/search/bleve_test.go @@ -27,13 +27,14 @@ func TestBleveBackend(t *testing.T) { Group: "folder.grafana.app", Resource: "folders", } - tmpdir, err := os.CreateTemp("", "bleve-test") + tmpdir, err := os.MkdirTemp("", "grafana-bleve-test") require.NoError(t, err) - backend := NewBleveBackend(BleveOptions{ - Root: tmpdir.Name(), + backend, err := NewBleveBackend(BleveOptions{ + Root: tmpdir, FileThreshold: 5, // with more than 5 items we create a file on disk }, tracing.NewNoopTracerService()) + require.NoError(t, err) // AVOID NPE in test resource.NewIndexMetrics(backend.opts.Root, backend) diff --git a/pkg/storage/unified/sql/server.go b/pkg/storage/unified/sql/server.go index 19dbfca1f12..3e0a93d966f 100644 --- a/pkg/storage/unified/sql/server.go +++ b/pkg/storage/unified/sql/server.go @@ -4,11 +4,13 @@ import ( "context" "log/slog" "os" + "path/filepath" "strings" - "github.com/grafana/grafana/pkg/storage/unified/search" "github.com/prometheus/client_golang/prometheus" + "github.com/grafana/grafana/pkg/storage/unified/search" + infraDB "github.com/grafana/grafana/pkg/infra/db" "github.com/grafana/grafana/pkg/infra/tracing" "github.com/grafana/grafana/pkg/services/authz" @@ -57,12 +59,25 @@ func NewResourceServer(ctx context.Context, db infraDB.DB, cfg *setting.Cfg, // Setup the search server if features.IsEnabledGlobally(featuremgmt.FlagUnifiedStorageSearch) { + root := cfg.IndexPath + if root == "" { + root = filepath.Join(cfg.DataPath, "unified-search", "bleve") + } + err = os.MkdirAll(root, 0750) + if err != nil { + return nil, err + } + bleve, err := search.NewBleveBackend(search.BleveOptions{ + Root: root, + FileThreshold: int64(cfg.IndexFileThreshold), // fewer than X items will use a memory index + BatchSize: cfg.IndexMaxBatchSize, // This is the batch size for how many objects to add to the index at once + }, tracer) + if err != nil { + return nil, err + } + opts.Search = resource.SearchOptions{ - Backend: search.NewBleveBackend(search.BleveOptions{ - Root: cfg.IndexPath, - FileThreshold: int64(cfg.IndexFileThreshold), // fewer than X items will use a memory index - BatchSize: cfg.IndexMaxBatchSize, // This is the batch size for how many objects to add to the index at once - }, tracer), + Backend: bleve, Resources: docs, WorkerThreads: cfg.IndexWorkers, InitMinCount: cfg.IndexMinCount,