diff --git a/pkg/storage/unified/resource/search.go b/pkg/storage/unified/resource/search.go index 7b5a6544e5b..2752c59650d 100644 --- a/pkg/storage/unified/resource/search.go +++ b/pkg/storage/unified/resource/search.go @@ -75,16 +75,16 @@ type ResourceIndex interface { // Search within a namespaced resource // When working with federated queries, the additional indexes will be passed in explicitly - Search(ctx context.Context, access types.AccessClient, req *resourcepb.ResourceSearchRequest, federate []ResourceIndex) (*resourcepb.ResourceSearchResponse, error) + Search(ctx context.Context, access types.AccessClient, req *resourcepb.ResourceSearchRequest, federate []ResourceIndex, stats *SearchStats) (*resourcepb.ResourceSearchResponse, error) // List within an response - ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest) (*resourcepb.ListManagedObjectsResponse, error) + ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest, stats *SearchStats) (*resourcepb.ListManagedObjectsResponse, error) // Counts the values in a repo - CountManagedObjects(ctx context.Context) ([]*resourcepb.CountManagedObjectsResponse_ResourceCount, error) + CountManagedObjects(ctx context.Context, stats *SearchStats) ([]*resourcepb.CountManagedObjectsResponse_ResourceCount, error) // Get the number of documents in the index - DocCount(ctx context.Context, folder string) (int64, error) + DocCount(ctx context.Context, folder string, stats *SearchStats) (int64, error) // UpdateIndex updates the index with the latest data (using update function provided when index was built) to guarantee strong consistency during the search. // Returns RV to which index was updated. @@ -238,6 +238,9 @@ func combineRebuildRequests(a, b rebuildRequest) (c rebuildRequest, ok bool) { } func (s *searchSupport) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest) (*resourcepb.ListManagedObjectsResponse, error) { + ctx, span := tracer.Start(ctx, "resource.searchSupport.ListManagedObjects") + defer span.End() + if req.NextPageToken != "" { return &resourcepb.ListManagedObjectsResponse{ Error: NewBadRequestError("multiple pages not yet supported"), @@ -248,14 +251,17 @@ func (s *searchSupport) ListManagedObjects(ctx context.Context, req *resourcepb. nsr := NamespacedResource{ Namespace: req.Namespace, } - stats, err := s.storage.GetResourceStats(ctx, nsr, 0) + resourceStats, err := s.storage.GetResourceStats(ctx, nsr, 0) if err != nil { rsp.Error = AsErrorResult(err) return rsp, nil } - for _, info := range stats { - idx, err := s.getOrCreateIndex(ctx, NamespacedResource{ + stats := NewSearchStats("ListManagedObjects") + defer s.logStats(ctx, stats, span, "namespace", req.Namespace) + + for _, info := range resourceStats { + idx, err := s.getOrCreateIndex(ctx, stats, NamespacedResource{ Namespace: req.Namespace, Group: info.Group, Resource: info.Resource, @@ -265,7 +271,7 @@ func (s *searchSupport) ListManagedObjects(ctx context.Context, req *resourcepb. return rsp, nil } - kind, err := idx.ListManagedObjects(ctx, req) + kind, err := idx.ListManagedObjects(ctx, req, stats) if err != nil { rsp.Error = AsErrorResult(err) return rsp, nil @@ -280,26 +286,69 @@ func (s *searchSupport) ListManagedObjects(ctx context.Context, req *resourcepb. } // Sort based on path + start := time.Now() slices.SortFunc(rsp.Items, func(a, b *resourcepb.ListManagedObjectsResponse_Item) int { return cmp.Compare(a.Path, b.Path) }) + stats.AddResultsConversionTime(time.Since(start)) return rsp, nil } +func (s *searchSupport) logWithTraceID(ctx context.Context) *slog.Logger { + l := s.log + if traceID := tracing.TraceIDFromContext(ctx, false); traceID != "" { + l = l.With("traceID", traceID) + } + return l +} + +func (s *searchSupport) logStats(ctx context.Context, stats *SearchStats, span trace.Span, params ...any) { + elapsed := time.Since(stats.startTime) + + args := []any{ + "operation", stats.operation, + "elapsedTime", elapsed, + "indexBuildTime", stats.indexBuildTime, + "indexUpdateTime", stats.indexUpdateTime, + "requestConversionTime", stats.requestConversion, + "searchTime", stats.searchTime, + "totalHits", stats.totalHits, + "returnedDocuments", stats.returnedDocuments, + "resultsConversionTime", stats.resultsConversionTime, + } + args = append(args, params...) + + s.logWithTraceID(ctx).Debug("Search stats", args...) + + if span != nil { + attrs := make([]attribute.KeyValue, 0, len(args)/2) + for i := 0; i < len(args); i += 2 { + attrs = append(attrs, attribute.String(fmt.Sprint(args[i]), fmt.Sprint(args[i+1]))) + } + span.AddEvent("search stats", trace.WithAttributes(attrs...)) + } +} + func (s *searchSupport) CountManagedObjects(ctx context.Context, req *resourcepb.CountManagedObjectsRequest) (*resourcepb.CountManagedObjectsResponse, error) { + ctx, span := tracer.Start(ctx, "resource.searchSupport.CountManagedObjects") + defer span.End() + + stats := NewSearchStats("CountManagedObjects") + defer s.logStats(ctx, stats, span, "namespace", req.Namespace) + rsp := &resourcepb.CountManagedObjectsResponse{} nsr := NamespacedResource{ Namespace: req.Namespace, } - stats, err := s.storage.GetResourceStats(ctx, nsr, 0) + resourceStats, err := s.storage.GetResourceStats(ctx, nsr, 0) if err != nil { rsp.Error = AsErrorResult(err) return rsp, nil } - for _, info := range stats { - idx, err := s.getOrCreateIndex(ctx, NamespacedResource{ + for _, info := range resourceStats { + idx, err := s.getOrCreateIndex(ctx, stats, NamespacedResource{ Namespace: req.Namespace, Group: info.Group, Resource: info.Resource, @@ -309,7 +358,7 @@ func (s *searchSupport) CountManagedObjects(ctx context.Context, req *resourcepb return rsp, nil } - counts, err := idx.CountManagedObjects(ctx) + counts, err := idx.CountManagedObjects(ctx, stats) if err != nil { rsp.Error = AsErrorResult(err) return rsp, nil @@ -349,12 +398,15 @@ func (s *searchSupport) Search(ctx context.Context, req *resourcepb.ResourceSear }, nil } + stats := NewSearchStats("Search") + defer s.logStats(ctx, stats, span, "namespace", req.Options.Key.Namespace, "group", req.Options.Key.Group, "resource", req.Options.Key.Resource, "query", req.Query) + nsr := NamespacedResource{ Group: req.Options.Key.Group, Namespace: req.Options.Key.Namespace, Resource: req.Options.Key.Resource, } - idx, err := s.getOrCreateIndex(ctx, nsr, "search") + idx, err := s.getOrCreateIndex(ctx, stats, nsr, "search") if err != nil { return &resourcepb.ResourceSearchResponse{ Error: AsErrorResult(err), @@ -366,7 +418,7 @@ func (s *searchSupport) Search(ctx context.Context, req *resourcepb.ResourceSear for i, f := range req.Federated { nsr.Group = f.Group nsr.Resource = f.Resource - federate[i], err = s.getOrCreateIndex(ctx, nsr, "federatedSearch") + federate[i], err = s.getOrCreateIndex(ctx, stats, nsr, "federatedSearch") if err != nil { return &resourcepb.ResourceSearchResponse{ Error: AsErrorResult(err), @@ -374,16 +426,23 @@ func (s *searchSupport) Search(ctx context.Context, req *resourcepb.ResourceSear } } - return idx.Search(ctx, s.access, req, federate) + return idx.Search(ctx, s.access, req, federate, stats) } // GetStats implements ResourceServer. func (s *searchSupport) GetStats(ctx context.Context, req *resourcepb.ResourceStatsRequest) (*resourcepb.ResourceStatsResponse, error) { + ctx, span := tracer.Start(ctx, "resource.searchSupport.GetStats") + defer span.End() + if req.Namespace == "" { return &resourcepb.ResourceStatsResponse{ Error: NewBadRequestError("missing namespace"), }, nil } + + stats := NewSearchStats("GetStats") + defer s.logStats(ctx, stats, span, "namespace", req.Namespace, "group", strings.Join(req.Kinds, ","), "folder", req.Folder) + rsp := &resourcepb.ResourceStatsResponse{} // Explicit list of kinds @@ -391,7 +450,7 @@ func (s *searchSupport) GetStats(ctx context.Context, req *resourcepb.ResourceSt rsp.Stats = make([]*resourcepb.ResourceStatsResponse_Stats, len(req.Kinds)) for i, k := range req.Kinds { parts := strings.SplitN(k, "/", 2) - index, err := s.getOrCreateIndex(ctx, NamespacedResource{ + index, err := s.getOrCreateIndex(ctx, stats, NamespacedResource{ Namespace: req.Namespace, Group: parts[0], Resource: parts[1], @@ -400,7 +459,7 @@ func (s *searchSupport) GetStats(ctx context.Context, req *resourcepb.ResourceSt rsp.Error = AsErrorResult(err) return rsp, nil } - count, err := index.DocCount(ctx, req.Folder) + count, err := index.DocCount(ctx, req.Folder, stats) if err != nil { rsp.Error = AsErrorResult(err) return rsp, nil @@ -417,17 +476,17 @@ func (s *searchSupport) GetStats(ctx context.Context, req *resourcepb.ResourceSt nsr := NamespacedResource{ Namespace: req.Namespace, } - stats, err := s.storage.GetResourceStats(ctx, nsr, 0) + resourceStats, err := s.storage.GetResourceStats(ctx, nsr, 0) if err != nil { return &resourcepb.ResourceStatsResponse{ Error: AsErrorResult(err), }, nil } - rsp.Stats = make([]*resourcepb.ResourceStatsResponse_Stats, len(stats)) + rsp.Stats = make([]*resourcepb.ResourceStatsResponse_Stats, len(resourceStats)) // When not filtered by folder or repository, we can use the results directly if req.Folder == "" { - for i, stat := range stats { + for i, stat := range resourceStats { rsp.Stats[i] = &resourcepb.ResourceStatsResponse_Stats{ Group: stat.Group, Resource: stat.Resource, @@ -437,8 +496,8 @@ func (s *searchSupport) GetStats(ctx context.Context, req *resourcepb.ResourceSt return rsp, nil } - for i, stat := range stats { - index, err := s.getOrCreateIndex(ctx, NamespacedResource{ + for i, stat := range resourceStats { + index, err := s.getOrCreateIndex(ctx, stats, NamespacedResource{ Namespace: req.Namespace, Group: stat.Group, Resource: stat.Resource, @@ -447,7 +506,7 @@ func (s *searchSupport) GetStats(ctx context.Context, req *resourcepb.ResourceSt rsp.Error = AsErrorResult(err) return rsp, nil } - count, err := index.DocCount(ctx, req.Folder) + count, err := index.DocCount(ctx, req.Folder, stats) if err != nil { rsp.Error = AsErrorResult(err) return rsp, nil @@ -733,7 +792,7 @@ type rebuildRequest struct { minBuildVersion *semver.Version // if not nil, rebuild index with build version older than this. } -func (s *searchSupport) getOrCreateIndex(ctx context.Context, key NamespacedResource, reason string) (ResourceIndex, error) { +func (s *searchSupport) getOrCreateIndex(ctx context.Context, stats *SearchStats, key NamespacedResource, reason string) (ResourceIndex, error) { if s == nil || s.search == nil { return nil, fmt.Errorf("search is not configured properly (missing unifiedStorageSearch feature toggle?)") } @@ -750,6 +809,7 @@ func (s *searchSupport) getOrCreateIndex(ctx context.Context, key NamespacedReso idx := s.search.GetIndex(key) if idx == nil { span.AddEvent("Building index") + buildStartTime := time.Now() ch := s.buildIndex.DoChan(key.String(), func() (interface{}, error) { // We want to finish building of the index even if original context is canceled. // We reuse original context without cancel to keep the tracing spans correct. @@ -795,6 +855,7 @@ func (s *searchSupport) getOrCreateIndex(ctx context.Context, key NamespacedReso if res.Err != nil { return nil, tracing.Error(span, res.Err) } + stats.AddIndexBuildTime(time.Since(buildStartTime)) idx = res.Val.(ResourceIndex) case <-ctx.Done(): return nil, tracing.Error(span, fmt.Errorf("failed to get index: %w", ctx.Err())) @@ -808,10 +869,11 @@ func (s *searchSupport) getOrCreateIndex(ctx context.Context, key NamespacedReso return nil, tracing.Error(span, fmt.Errorf("failed to update index to guarantee strong consistency: %w", err)) } elapsed := time.Since(start) + stats.AddIndexUpdateTime(elapsed) if s.indexMetrics != nil { s.indexMetrics.SearchUpdateWaitTime.WithLabelValues(reason).Observe(elapsed.Seconds()) } - s.log.Debug("Index updated before search", "namespace", key.Namespace, "group", key.Group, "resource", key.Resource, "reason", reason, "duration", elapsed, "rv", rv) + s.logWithTraceID(ctx).Debug("Index updated before search", "namespace", key.Namespace, "group", key.Group, "resource", key.Resource, "reason", reason, "duration", elapsed, "rv", rv) span.AddEvent("Index updated") return idx, nil @@ -828,7 +890,7 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size attribute.Int64("size", size), ) - logger := s.log.With("namespace", nsr.Namespace, "group", nsr.Group, "resource", nsr.Resource) + logger := s.logWithTraceID(ctx).With("namespace", nsr.Namespace, "group", nsr.Group, "resource", nsr.Resource) builder, err := s.builders.get(ctx, nsr) if err != nil { @@ -991,7 +1053,9 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size } // Record the number of objects indexed for the kind/resource - docCount, err := index.DocCount(ctx, "") + // We don't pass searchStats to DocCount here, as it's not really user-initiated search. Time spent + // here will be recorded in the index build time instead. + docCount, err := index.DocCount(ctx, "", nil) if err != nil { logger.Warn("error getting doc count", "error", err) } diff --git a/pkg/storage/unified/resource/search_test.go b/pkg/storage/unified/resource/search_test.go index c3f852c837e..215a88cd2cd 100644 --- a/pkg/storage/unified/resource/search_test.go +++ b/pkg/storage/unified/resource/search_test.go @@ -41,22 +41,22 @@ func (m *MockResourceIndex) BulkIndex(req *BulkIndexRequest) error { return args.Error(0) } -func (m *MockResourceIndex) Search(ctx context.Context, access types.AccessClient, req *resourcepb.ResourceSearchRequest, federate []ResourceIndex) (*resourcepb.ResourceSearchResponse, error) { +func (m *MockResourceIndex) Search(ctx context.Context, access types.AccessClient, req *resourcepb.ResourceSearchRequest, federate []ResourceIndex, stats *SearchStats) (*resourcepb.ResourceSearchResponse, error) { args := m.Called(ctx, access, req, federate) return args.Get(0).(*resourcepb.ResourceSearchResponse), args.Error(1) } -func (m *MockResourceIndex) CountManagedObjects(ctx context.Context) ([]*resourcepb.CountManagedObjectsResponse_ResourceCount, error) { +func (m *MockResourceIndex) CountManagedObjects(ctx context.Context, stats *SearchStats) ([]*resourcepb.CountManagedObjectsResponse_ResourceCount, error) { args := m.Called(ctx) return args.Get(0).([]*resourcepb.CountManagedObjectsResponse_ResourceCount), args.Error(1) } -func (m *MockResourceIndex) DocCount(ctx context.Context, folder string) (int64, error) { +func (m *MockResourceIndex) DocCount(ctx context.Context, folder string, stats *SearchStats) (int64, error) { args := m.Called(ctx, folder) return args.Get(0).(int64), args.Error(1) } -func (m *MockResourceIndex) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest) (*resourcepb.ListManagedObjectsResponse, error) { +func (m *MockResourceIndex) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest, stats *SearchStats) (*resourcepb.ListManagedObjectsResponse, error) { args := m.Called(ctx, req) return args.Get(0).(*resourcepb.ListManagedObjectsResponse), args.Error(1) } @@ -223,7 +223,7 @@ func TestSearchGetOrCreateIndex(t *testing.T) { go func() { defer wg.Done() <-start - _, _ = support.getOrCreateIndex(context.Background(), NamespacedResource{Namespace: "ns", Group: "group", Resource: "resource"}, "test") + _, _ = support.getOrCreateIndex(context.Background(), nil, NamespacedResource{Namespace: "ns", Group: "group", Resource: "resource"}, "test") }() } @@ -270,17 +270,17 @@ func TestSearchGetOrCreateIndexWithIndexUpdate(t *testing.T) { require.NoError(t, err) require.NotNil(t, support) - idx, err := support.getOrCreateIndex(context.Background(), NamespacedResource{Namespace: "ns", Group: "group", Resource: "resource"}, "initial call") + idx, err := support.getOrCreateIndex(context.Background(), nil, NamespacedResource{Namespace: "ns", Group: "group", Resource: "resource"}, "initial call") require.NoError(t, err) require.NotNil(t, idx) checkMockIndexUpdateCalls(t, idx, 1) - idx, err = support.getOrCreateIndex(context.Background(), NamespacedResource{Namespace: "ns", Group: "group", Resource: "resource"}, "second call") + idx, err = support.getOrCreateIndex(context.Background(), nil, NamespacedResource{Namespace: "ns", Group: "group", Resource: "resource"}, "second call") require.NoError(t, err) require.NotNil(t, idx) checkMockIndexUpdateCalls(t, idx, 2) - idx, err = support.getOrCreateIndex(context.Background(), NamespacedResource{Namespace: "ns", Group: "group", Resource: "bad"}, "call to bad index") + idx, err = support.getOrCreateIndex(context.Background(), nil, NamespacedResource{Namespace: "ns", Group: "group", Resource: "bad"}, "call to bad index") require.ErrorIs(t, err, failedErr) require.Nil(t, idx) } @@ -324,7 +324,7 @@ func TestSearchGetOrCreateIndexWithCancellation(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 1*time.Millisecond) defer cancel() - _, err = support.getOrCreateIndex(ctx, key, "test") + _, err = support.getOrCreateIndex(ctx, nil, key, "test") // Make sure we get context deadline error require.ErrorIs(t, err, context.DeadlineExceeded) @@ -340,7 +340,7 @@ func TestSearchGetOrCreateIndexWithCancellation(t *testing.T) { }, 1*time.Second, 100*time.Millisecond, "Indexing finishes despite context cancellation") // Second call to getOrCreateIndex returns index immediately, even if context is canceled, as the index is now ready and cached. - _, err = support.getOrCreateIndex(ctx, key, "test") + _, err = support.getOrCreateIndex(ctx, nil, key, "test") require.NoError(t, err) } diff --git a/pkg/storage/unified/resource/stats.go b/pkg/storage/unified/resource/stats.go new file mode 100644 index 00000000000..f35bfd4f4ed --- /dev/null +++ b/pkg/storage/unified/resource/stats.go @@ -0,0 +1,65 @@ +package resource + +import "time" + +type SearchStats struct { + operation string + startTime time.Time // Time when the operation. + + indexBuildTime time.Duration // Time to build indexes if it wasn't open before. + indexUpdateTime time.Duration // Time to update indexes to the latest state. + + requestConversion time.Duration // How long does it take to convert search request to the bleve query + + searchTime time.Duration + totalHits int // Total hits across all pages + returnedDocuments int // Hits returned in this page + + resultsConversionTime time.Duration // How long does it take to convert search results to final results. +} + +func NewSearchStats(op string) *SearchStats { + return &SearchStats{operation: op, startTime: time.Now()} +} + +func (s *SearchStats) AddIndexBuildTime(d time.Duration) { + if s != nil { + s.indexBuildTime += d + } +} + +func (s *SearchStats) AddIndexUpdateTime(d time.Duration) { + if s != nil { + s.indexUpdateTime += d + } +} + +func (s *SearchStats) AddRequestConversionTime(d time.Duration) { + if s != nil { + s.requestConversion += d + } +} + +func (s *SearchStats) AddSearchTime(d time.Duration) { + if s != nil { + s.searchTime += d + } +} + +func (s *SearchStats) AddTotalHits(total int) { + if s != nil { + s.totalHits += total + } +} + +func (s *SearchStats) AddReturnedDocuments(docs int) { + if s != nil { + s.returnedDocuments += docs + } +} + +func (s *SearchStats) AddResultsConversionTime(d time.Duration) { + if s != nil { + s.resultsConversionTime += d + } +} diff --git a/pkg/storage/unified/search/bleve.go b/pkg/storage/unified/search/bleve.go index 863c349186f..836c491b3c7 100644 --- a/pkg/storage/unified/search/bleve.go +++ b/pkg/storage/unified/search/bleve.go @@ -866,7 +866,7 @@ func (b *bleveIndex) BuildInfo() (resource.IndexBuildInfo, error) { }, nil } -func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest) (*resourcepb.ListManagedObjectsResponse, error) { +func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest, stats *resource.SearchStats) (*resourcepb.ListManagedObjectsResponse, error) { if req.NextPageToken != "" { return nil, fmt.Errorf("next page not implemented yet") } @@ -881,6 +881,7 @@ func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.Lis }, nil } + start := time.Now() q := bleve.NewBooleanQuery() q.AddMust(&query.TermQuery{ Term: req.Kind, @@ -890,6 +891,7 @@ func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.Lis Term: req.Id, FieldVal: resource.SEARCH_FIELD_MANAGER_ID, }) + stats.AddResultsConversionTime(time.Since(start)) found, err := b.index.SearchInContext(ctx, &bleve.SearchRequest{ Query: q, @@ -916,6 +918,10 @@ func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.Lis return nil, err } + stats.AddTotalHits(int(found.Total)) + stats.AddSearchTime(found.Took) + stats.AddReturnedDocuments(len(found.Hits)) + asString := func(v any) string { if v == nil { return "" @@ -947,6 +953,7 @@ func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.Lis return 0 } + start = time.Now() rsp := &resourcepb.ListManagedObjectsResponse{} for _, hit := range found.Hits { item := &resourcepb.ListManagedObjectsResponse_Item{ @@ -963,10 +970,11 @@ func (b *bleveIndex) ListManagedObjects(ctx context.Context, req *resourcepb.Lis } rsp.Items = append(rsp.Items, item) } + stats.AddResultsConversionTime(time.Since(start)) return rsp, nil } -func (b *bleveIndex) CountManagedObjects(ctx context.Context) ([]*resourcepb.CountManagedObjectsResponse_ResourceCount, error) { +func (b *bleveIndex) CountManagedObjects(ctx context.Context, stats *resource.SearchStats) ([]*resourcepb.CountManagedObjectsResponse_ResourceCount, error) { found, err := b.index.SearchInContext(ctx, &bleve.SearchRequest{ Query: bleve.NewMatchAllQuery(), Size: 0, @@ -977,6 +985,11 @@ func (b *bleveIndex) CountManagedObjects(ctx context.Context) ([]*resourcepb.Cou if err != nil { return nil, err } + + stats.AddSearchTime(found.Took) + stats.AddTotalHits(int(found.Total)) + stats.AddReturnedDocuments(len(found.Hits)) + vals := make([]*resourcepb.CountManagedObjectsResponse_ResourceCount, 0) f, ok := found.Facets["count"] if ok && f.Terms != nil { @@ -1003,6 +1016,7 @@ func (b *bleveIndex) Search( access authlib.AccessClient, req *resourcepb.ResourceSearchRequest, federate []resource.ResourceIndex, // For federated queries, these will match the values in req.federate + stats *resource.SearchStats, ) (*resourcepb.ResourceSearchResponse, error) { ctx, span := b.tracing.Start(ctx, tracingPrexfixBleve+"Search") defer span.End() @@ -1026,6 +1040,7 @@ func (b *bleveIndex) Search( return nil, err } + conversionStarts := time.Now() // convert protobuf request to bleve request searchrequest, e := b.toBleveSearchRequest(ctx, req, access) if e != nil { @@ -1050,6 +1065,7 @@ func (b *bleveIndex) Search( } } } + stats.AddRequestConversionTime(time.Since(conversionStarts)) res, err := index.SearchInContext(ctx, searchrequest) if err != nil { @@ -1059,7 +1075,11 @@ func (b *bleveIndex) Search( response.TotalHits = int64(res.Total) response.QueryCost = float64(res.Cost) response.MaxScore = res.MaxScore + stats.AddSearchTime(res.Took) + stats.AddTotalHits(int(res.Total)) + stats.AddReturnedDocuments(len(res.Hits)) + resultsConversionStart := time.Now() response.Results, err = b.hitsToTable(ctx, searchrequest.Fields, res.Hits, req.Explain) if err != nil { return nil, err @@ -1073,10 +1093,11 @@ func (b *bleveIndex) Search( } response.Facet[k] = f } + stats.AddResultsConversionTime(time.Since(resultsConversionStart)) return response, nil } -func (b *bleveIndex) DocCount(ctx context.Context, folder string) (int64, error) { +func (b *bleveIndex) DocCount(ctx context.Context, folder string, stats *resource.SearchStats) (int64, error) { ctx, span := b.tracing.Start(ctx, tracingPrexfixBleve+"DocCount") defer span.End() @@ -1097,6 +1118,10 @@ func (b *bleveIndex) DocCount(ctx context.Context, folder string) (int64, error) if rsp == nil { return 0, err } + if stats != nil { + stats.AddTotalHits(int(rsp.Total)) + stats.AddSearchTime(rsp.Took) + } return int64(rsp.Total), err } diff --git a/pkg/storage/unified/search/bleve_performance_test.go b/pkg/storage/unified/search/bleve_performance_test.go index 00e23d78259..0dfad7ac8fe 100644 --- a/pkg/storage/unified/search/bleve_performance_test.go +++ b/pkg/storage/unified/search/bleve_performance_test.go @@ -57,7 +57,7 @@ func runBenchmark(b *testing.B, testIndex resource.ResourceIndex) { var memStatsBefore, memStatsAfter runtime.MemStats runtime.ReadMemStats(&memStatsBefore) - _, err := testIndex.Search(context.Background(), nil, searchRequest, nil) + _, err := testIndex.Search(context.Background(), nil, searchRequest, nil, nil) elapsed := time.Since(start) // Calculate elapsed time runtime.ReadMemStats(&memStatsAfter) diff --git a/pkg/storage/unified/search/bleve_search_test.go b/pkg/storage/unified/search/bleve_search_test.go index c1b18724d96..600dd0cfe04 100644 --- a/pkg/storage/unified/search/bleve_search_test.go +++ b/pkg/storage/unified/search/bleve_search_test.go @@ -43,7 +43,7 @@ func indexDocumentsWithTitles(t *testing.T, index resource.ResourceIndex, key re } func checkSearchQuery(t *testing.T, index resource.ResourceIndex, query *resourcepb.ResourceSearchRequest, orderedExpectedNames []string) { - res, err := index.Search(context.Background(), nil, query, nil) + res, err := index.Search(context.Background(), nil, query, nil, nil) require.NoError(t, err) require.Equal(t, int64(len(orderedExpectedNames)), res.TotalHits) for ix, name := range orderedExpectedNames { diff --git a/pkg/storage/unified/search/bleve_test.go b/pkg/storage/unified/search/bleve_test.go index 9bc867ba0d3..b303d11c6e8 100644 --- a/pkg/storage/unified/search/bleve_test.go +++ b/pkg/storage/unified/search/bleve_test.go @@ -213,7 +213,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { Limit: 100, }, }, - }, nil) + }, nil, nil) require.NoError(t, err) require.Nil(t, rsp.Error) require.NotNil(t, rsp.Results) @@ -242,10 +242,10 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { ] }`, string(disp)) - count, _ := index.DocCount(ctx, "") + count, _ := index.DocCount(ctx, "", nil) assert.Equal(t, int64(3), count) - count, _ = index.DocCount(ctx, "zzz") + count, _ = index.DocCount(ctx, "zzz", nil) assert.Equal(t, int64(1), count) rsp, err = index.Search(ctx, NewStubAccessClient(map[string]bool{"dashboards": true}), &resourcepb.ResourceSearchRequest{ @@ -258,7 +258,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { }}, }, Limit: 100000, - }, nil) + }, nil, nil) require.NoError(t, err) require.Equal(t, int64(2), rsp.TotalHits) require.Equal(t, []string{"aaa", "bbb"}, []string{ @@ -276,7 +276,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { SortBy: []*resourcepb.ResourceSearchRequest_Sort{ {Field: "fields." + DASHBOARD_VIEWS_LAST_1_DAYS, Desc: true}, }, - }, nil) + }, nil, nil) require.NoError(t, err) require.Equal(t, 2, len(rsp.Results.Columns)) require.Equal(t, DASHBOARD_ERRORS_TODAY, rsp.Results.Columns[0].Name) @@ -296,7 +296,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { SortBy: []*resourcepb.ResourceSearchRequest_Sort{ {Field: "fields." + DASHBOARD_VIEWS_LAST_1_DAYS, Desc: true}, }, - }, nil) + }, nil, nil) require.NoError(t, err) require.Equal(t, 0, len(rsp.Results.Rows)) @@ -304,7 +304,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { found, err := index.ListManagedObjects(ctx, &resourcepb.ListManagedObjectsRequest{ Kind: "repo", Id: "repo-1", - }) + }, nil) require.NoError(t, err) jj, err := json.MarshalIndent(found, "", " ") require.NoError(t, err) @@ -341,7 +341,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { ] }`, string(jj)) - counts, err := index.CountManagedObjects(ctx) + counts, err := index.CountManagedObjects(ctx, nil) require.NoError(t, err) jj, err = json.MarshalIndent(counts, "", " ") require.NoError(t, err) @@ -433,7 +433,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { Key: key, }, Limit: 100000, - }, nil) + }, nil, nil) require.NoError(t, err) require.Nil(t, rsp.Error) require.NotNil(t, rsp.Results) @@ -468,7 +468,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { Limit: 100, }, }, - }, []resource.ResourceIndex{foldersIndex}) // << note the folder index matches the federation request + }, []resource.ResourceIndex{foldersIndex}, nil) // << note the folder index matches the federation request require.NoError(t, err) require.Nil(t, rsp.Error) require.NotNil(t, rsp.Results) @@ -532,7 +532,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { Limit: 100, }, }, - }, []resource.ResourceIndex{foldersIndex}) // << note the folder index matches the federation request + }, []resource.ResourceIndex{foldersIndex}, nil) // << note the folder index matches the federation request require.NoError(t, err) require.Equal(t, 3, len(rsp.Results.Rows)) @@ -561,7 +561,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { Limit: 100, }, }, - }, []resource.ResourceIndex{foldersIndex}) // << note the folder index matches the federation request + }, []resource.ResourceIndex{foldersIndex}, nil) // << note the folder index matches the federation request require.NoError(t, err) require.Equal(t, 2, len(rsp.Results.Rows)) @@ -589,7 +589,7 @@ func testBleveBackend(t *testing.T, backend *bleveBackend) { Limit: 100, }, }, - }, []resource.ResourceIndex{foldersIndex}) // << note the folder index matches the federation request + }, []resource.ResourceIndex{foldersIndex}, nil) // << note the folder index matches the federation request require.NoError(t, err) require.Equal(t, 0, len(rsp.Results.Rows)) @@ -888,7 +888,7 @@ func TestBuildIndexExpiration(t *testing.T) { idx := backend.GetIndex(ns) require.Nil(t, idx) - _, err = builtIndex.DocCount(context.Background(), "") + _, err = builtIndex.DocCount(context.Background(), "", nil) require.ErrorIs(t, err, bleve.ErrorIndexClosed) // Verify that there are no open indexes. @@ -897,7 +897,7 @@ func TestBuildIndexExpiration(t *testing.T) { idx := backend.GetIndex(ns) require.NotNil(t, idx) - cnt, err := builtIndex.DocCount(context.Background(), "") + cnt, err := builtIndex.DocCount(context.Background(), "", nil) require.NoError(t, err) require.Equal(t, int64(1), cnt) @@ -971,7 +971,7 @@ func TestBuildIndex(t *testing.T) { idx, err := newBackend.BuildIndex(context.Background(), ns, secondIndexDocsCount, nil, "test", indexTestDocs(ns, secondIndexDocsCount, 100), nil, rebuild) require.NoError(t, err) - cnt, err := idx.DocCount(context.Background(), "") + cnt, err := idx.DocCount(context.Background(), "", nil) require.NoError(t, err) if rebuild { require.Equal(t, int64(secondIndexDocsCount), cnt, "Index has been not rebuilt") @@ -1033,10 +1033,10 @@ func TestRebuildingIndexClosesPreviousCachedIndex(t *testing.T) { // Verify that first and second index are different, and first one is now closed. require.NotEqual(t, firstIndex, secondIndex) - _, err = firstIndex.DocCount(context.Background(), "") + _, err = firstIndex.DocCount(context.Background(), "", nil) require.ErrorIs(t, err, bleve.ErrorIndexClosed) - cnt, err := secondIndex.DocCount(context.Background(), "") + cnt, err := secondIndex.DocCount(context.Background(), "", nil) require.NoError(t, err) require.Equal(t, int64(secondSize), cnt) @@ -1341,7 +1341,7 @@ func TestConcurrentIndexUpdateSearchAndRebuild(t *testing.T) { Fields: []string{"title"}, Query: "Document", Limit: 10, - }, nil) + }, nil, nil) if err != nil { if errors.Is(err, bleve.ErrorIndexClosed) || errors.Is(err, context.Canceled) { continue @@ -1574,13 +1574,13 @@ func searchTitle(t *testing.T, idx resource.ResourceIndex, query string, limit i Fields: []string{"title"}, Query: query, Limit: int64(limit), - }, nil) + }, nil, nil) require.NoError(t, err) return resp } func docCount(t *testing.T, idx resource.ResourceIndex) int { - cnt, err := idx.DocCount(context.Background(), "") + cnt, err := idx.DocCount(context.Background(), "", nil) require.NoError(t, err) return int(cnt) } diff --git a/pkg/storage/unified/search/user_test.go b/pkg/storage/unified/search/user_test.go index a58516842c6..67b6777f895 100644 --- a/pkg/storage/unified/search/user_test.go +++ b/pkg/storage/unified/search/user_test.go @@ -7,6 +7,9 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/selection" + iamv0 "github.com/grafana/grafana/apps/iam/pkg/apis/iam/v0alpha1" "github.com/grafana/grafana/pkg/apimachinery/identity" "github.com/grafana/grafana/pkg/infra/tracing" @@ -14,8 +17,6 @@ import ( "github.com/grafana/grafana/pkg/storage/unified/resource" "github.com/grafana/grafana/pkg/storage/unified/resourcepb" "github.com/grafana/grafana/pkg/storage/unified/search" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/selection" ) func TestUserDocumentBuilder(t *testing.T) { @@ -175,7 +176,7 @@ func indexUserDocuments(t *testing.T, index resource.ResourceIndex, key resource func checkUserSearchQuery(t *testing.T, index resource.ResourceIndex, query *resourcepb.ResourceSearchRequest, orderedExpectedNames []string) { t.Helper() - res, err := index.Search(context.Background(), nil, query, nil) + res, err := index.Search(context.Background(), nil, query, nil, nil) require.NoError(t, err) require.Equal(t, int64(len(orderedExpectedNames)), res.TotalHits) names := make([]string, len(res.Results.Rows)) diff --git a/pkg/storage/unified/testing/search_backend.go b/pkg/storage/unified/testing/search_backend.go index 4ef4b6d2753..0b3c5f97d25 100644 --- a/pkg/storage/unified/testing/search_backend.go +++ b/pkg/storage/unified/testing/search_backend.go @@ -170,7 +170,7 @@ func runTestResourceIndex(t *testing.T, backend resource.SearchBackend, nsPrefix }, Fields: []string{"title", "folder", "tags"}, Limit: 10, - }, nil) + }, nil, nil) require.NoError(t, err) require.NotNil(t, resp) require.Equal(t, int64(1), resp.TotalHits) // Only doc2 should have tag3 now @@ -187,7 +187,7 @@ func runTestResourceIndex(t *testing.T, backend resource.SearchBackend, nsPrefix Query: "Document", Fields: []string{"title", "folder", "tags"}, Limit: 10, - }, nil) + }, nil, nil) require.NoError(t, err) require.NotNil(t, resp) require.Equal(t, int64(2), resp.TotalHits) // Both doc1 and doc2 should have doc now @@ -229,7 +229,7 @@ func runTestResourceIndex(t *testing.T, backend resource.SearchBackend, nsPrefix Query: "Document", Fields: []string{"title", "folder", "tags"}, Limit: 10, - }, nil) + }, nil, nil) require.NoError(t, err) require.NotNil(t, resp) require.Equal(t, int64(3), resp.TotalHits) // Both doc1, doc2, and doc3 should have doc now @@ -319,7 +319,7 @@ func runTestResourceIndex(t *testing.T, backend resource.SearchBackend, nsPrefix Query: "", Fields: []string{"title"}, Limit: 10, - }, nil) + }, nil, nil) require.NoError(t, err) require.NotNil(t, resp) require.Equal(t, int64(1), resp.TotalHits) // Only dash1 should have lib-panel-1