Collect and log search stats. (#114107)

* Collect and log search stats.

* Fix compilation problems.
This commit is contained in:
Peter Štibraný
2025-11-19 16:52:47 +01:00
committed by GitHub
parent 1fcf15d05f
commit 42db6324c9
9 changed files with 225 additions and 70 deletions
+91 -27
View File
@@ -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)
}
+10 -10
View File
@@ -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)
}
+65
View File
@@ -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
}
}
+28 -3
View File
@@ -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
}
@@ -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)
@@ -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 {
+21 -21
View File
@@ -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)
}
+4 -3
View File
@@ -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))
@@ -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