From e162c69c34ca26d6236ae937d2163c3f8602ab75 Mon Sep 17 00:00:00 2001 From: Georges Chaudy Date: Mon, 12 May 2025 07:56:25 +0200 Subject: [PATCH] search: fix document missing at startup (#105198) * fix document missing at startup * go-lint * fix tests * fix tests * fix integration tests now that we are storing real values --- pkg/storage/unified/resource/search.go | 21 +- pkg/storage/unified/search/bleve.go | 3 + .../unified/sql/test/integration_test.go | 60 ++++- pkg/storage/unified/testing/benchmark.go | 79 ++++-- .../unified/testing/search_and_storage.go | 249 ++++++++++++++++++ pkg/storage/unified/testing/search_backend.go | 64 ++++- .../unified/testing/storage_backend.go | 100 ++++--- 7 files changed, 496 insertions(+), 80 deletions(-) create mode 100644 pkg/storage/unified/testing/search_and_storage.go diff --git a/pkg/storage/unified/resource/search.go b/pkg/storage/unified/resource/search.go index 0a05ee97567..ad1b982d93a 100644 --- a/pkg/storage/unified/resource/search.go +++ b/pkg/storage/unified/resource/search.go @@ -546,17 +546,15 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size s.log.Debug("Building index", "resource", nsr.Resource, "size", size, "rv", rv) - key := &ResourceKey{ - Group: nsr.Group, - Resource: nsr.Resource, - Namespace: nsr.Namespace, - } - index, err := s.search.BuildIndex(ctx, nsr, size, rv, fields, func(index ResourceIndex) (int64, error) { rv, err = s.storage.ListIterator(ctx, &ListRequest{ Limit: 1000000000000, // big number Options: &ListOptions{ - Key: key, + Key: &ResourceKey{ + Group: nsr.Group, + Resource: nsr.Resource, + Namespace: nsr.Namespace, + }, }, }, func(iter ListIterator) error { // Collect all documents in a single bulk request @@ -568,7 +566,12 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size } // Update the key name - key.Name = iter.Name() + key := &ResourceKey{ + Group: nsr.Group, + Resource: nsr.Resource, + Namespace: nsr.Namespace, + Name: iter.Name(), + } // Convert it to an indexable document doc, err := builder.BuildDocument(ctx, key, iter.ResourceVersion(), iter.Value()) @@ -607,7 +610,7 @@ func (s *searchSupport) build(ctx context.Context, nsr NamespacedResource, size s.log.Warn("error getting doc count", "error", err) } if s.indexMetrics != nil { - s.indexMetrics.IndexedKinds.WithLabelValues(key.Resource).Add(float64(docCount)) + s.indexMetrics.IndexedKinds.WithLabelValues(nsr.Resource).Add(float64(docCount)) } // rv is the last RV we read. when watching, we must add all events since that time diff --git a/pkg/storage/unified/search/bleve.go b/pkg/storage/unified/search/bleve.go index 83fa1d9fbb3..0bb24127aca 100644 --- a/pkg/storage/unified/search/bleve.go +++ b/pkg/storage/unified/search/bleve.go @@ -324,6 +324,9 @@ func (b *bleveIndex) BulkIndex(req *resource.BulkIndexRequest) error { for _, item := range req.Items { switch item.Action { case resource.ActionIndex: + if item.Doc == nil { + return fmt.Errorf("missing document") + } doc := item.Doc.UpdateCopyFields() doc.References = nil // remove references (for now!) diff --git a/pkg/storage/unified/sql/test/integration_test.go b/pkg/storage/unified/sql/test/integration_test.go index a9c2f3a82d4..8a3b0435dda 100644 --- a/pkg/storage/unified/sql/test/integration_test.go +++ b/pkg/storage/unified/sql/test/integration_test.go @@ -2,6 +2,7 @@ package test import ( "context" + "os" "testing" "time" @@ -12,12 +13,13 @@ import ( "github.com/grafana/authlib/authn" "github.com/grafana/authlib/types" "github.com/grafana/dskit/services" - infraDB "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/infra/db" "github.com/grafana/grafana/pkg/infra/tracing" "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified" "github.com/grafana/grafana/pkg/storage/unified/resource" + "github.com/grafana/grafana/pkg/storage/unified/search" "github.com/grafana/grafana/pkg/storage/unified/sql" "github.com/grafana/grafana/pkg/storage/unified/sql/db/dbimpl" unitest "github.com/grafana/grafana/pkg/storage/unified/testing" @@ -30,11 +32,11 @@ func TestMain(m *testing.M) { } func TestIntegrationStorageServer(t *testing.T) { - if infraDB.IsTestDBSpanner() { + if db.IsTestDBSpanner() { t.Skip("skipping integration test") } unitest.RunStorageServerTest(t, func(ctx context.Context) resource.StorageBackend { - dbstore := infraDB.InitTestDB(t) + dbstore := db.InitTestDB(t) eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil) require.NoError(t, err) require.NotNil(t, eDB) @@ -53,13 +55,13 @@ func TestIntegrationStorageServer(t *testing.T) { // TestStorageBackend is a test for the StorageBackend interface. func TestIntegrationSQLStorageBackend(t *testing.T) { - if infraDB.IsTestDBSpanner() { + if db.IsTestDBSpanner() { t.Skip("skipping integration test") } t.Run("IsHA (polling notifier)", func(t *testing.T) { unitest.RunStorageBackendTest(t, func(ctx context.Context) resource.StorageBackend { - dbstore := infraDB.InitTestDB(t) + dbstore := db.InitTestDB(t) eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil) require.NoError(t, err) require.NotNil(t, eDB) @@ -78,7 +80,7 @@ func TestIntegrationSQLStorageBackend(t *testing.T) { t.Run("NotHA (in process notifier)", func(t *testing.T) { unitest.RunStorageBackendTest(t, func(ctx context.Context) resource.StorageBackend { - dbstore := infraDB.InitTestDB(t) + dbstore := db.InitTestDB(t) eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil) require.NoError(t, err) require.NotNil(t, eDB) @@ -96,16 +98,56 @@ func TestIntegrationSQLStorageBackend(t *testing.T) { }) } +func TestIntegrationSearchAndStorage(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test in short mode") + } + if db.IsTestDBSpanner() { + t.Skip("Skipping benchmark on Spanner") + } + + ctx := context.Background() + + tempDir := t.TempDir() + t.Cleanup(func() { + _ = os.RemoveAll(tempDir) + }) + // Create a new bleve backend + search, err := search.NewBleveBackend(search.BleveOptions{ + FileThreshold: 0, + Root: tempDir, + }, tracing.NewNoopTracerService(), featuremgmt.WithFeatures(featuremgmt.FlagUnifiedStorageSearchPermissionFiltering), nil) + require.NoError(t, err) + require.NotNil(t, search) + + // Create a new resource backend + dbstore := db.InitTestDB(t) + eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil) + require.NoError(t, err) + require.NotNil(t, eDB) + + storage, err := sql.NewBackend(sql.BackendOptions{ + DBProvider: eDB, + IsHA: false, + }) + require.NoError(t, err) + require.NotNil(t, storage) + + err = storage.Init(ctx) + require.NoError(t, err) + unitest.RunTestSearchAndStorage(t, ctx, storage, search) +} + func TestClientServer(t *testing.T) { - if infraDB.IsTestDbSQLite() { + if db.IsTestDbSQLite() { t.Skip("TODO: test blocking, skipping to unblock Enterprise until we fix this") } - if infraDB.IsTestDBSpanner() { + if db.IsTestDBSpanner() { t.Skip("skipping integration test") } ctx := testutil.NewTestContext(t, time.Now().Add(5*time.Second)) - dbstore := infraDB.InitTestDB(t) + dbstore := db.InitTestDB(t) cfg := setting.NewCfg() cfg.GRPCServer.Address = "localhost:0" // get a free address diff --git a/pkg/storage/unified/testing/benchmark.go b/pkg/storage/unified/testing/benchmark.go index 11c577625ce..e3dab5acbeb 100644 --- a/pkg/storage/unified/testing/benchmark.go +++ b/pkg/storage/unified/testing/benchmark.go @@ -11,6 +11,7 @@ import ( "github.com/grafana/grafana/pkg/storage/unified/resource" "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" ) @@ -56,7 +57,7 @@ func initializeBackend(ctx context.Context, backend resource.StorageBackend, opt WithNamespace(namespace), WithGroup(group), WithResource(resourceType), - WithValue([]byte("init"))) + WithValue("init")) if err != nil { return fmt.Errorf("failed to initialize backend: %w", err) } @@ -107,7 +108,7 @@ func runStorageBackendBenchmark(ctx context.Context, backend resource.StorageBac WithNamespace(namespace), WithGroup(group), WithResource(resourceType), - WithValue([]byte(strings.Repeat("abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789", 20)))) // ~1.21 KiB + WithValue(strings.Repeat("abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789", 20))) // ~1.21 KiB if err != nil { errors <- err @@ -336,12 +337,18 @@ func BenchmarkSearchBackend(tb testing.TB, backend resource.SearchBackend, opts func BenchmarkIndexServer(tb testing.TB, ctx context.Context, backend resource.StorageBackend, searchBackend resource.SearchBackend, opts *BenchmarkOptions) { events := make(chan *resource.IndexEvent, opts.NumResources) + groupsResources := make(map[string]string) + for g := 0; g < opts.NumGroups; g++ { + for r := 0; r < opts.NumResourceTypes; r++ { + groupsResources[fmt.Sprintf("group-%d", g)] = fmt.Sprintf("resource-%d", r) + } + } server, err := resource.NewResourceServer(resource.ResourceServerOptions{ Backend: backend, Search: resource.SearchOptions{ Backend: searchBackend, IndexEventsChan: events, - Resources: &testDocumentBuilderSupplier{opts: opts}, + Resources: &testDocumentBuilderSupplier{groupsResources: groupsResources}, }, }) require.NoError(tb, err) @@ -414,37 +421,67 @@ func BenchmarkIndexServer(tb testing.TB, ctx context.Context, backend resource.S type testDocumentBuilder struct{} func (b *testDocumentBuilder) BuildDocument(ctx context.Context, key *resource.ResourceKey, rv int64, value []byte) (*resource.IndexableDocument, error) { + // convert value to unstructured.Unstructured + var u unstructured.Unstructured + if err := u.UnmarshalJSON(value); err != nil { + return nil, fmt.Errorf("failed to unmarshal value: %w", err) + } + + title := "" + tags := []string{} + val := "" + + spec, ok, _ := unstructured.NestedMap(u.Object, "spec") + if ok { + if v, ok := spec["title"]; ok { + title = v.(string) + } + if v, ok := spec["tags"]; ok { + if tagSlice, ok := v.([]interface{}); ok { + tags = make([]string, len(tagSlice)) + for i, tag := range tagSlice { + if strTag, ok := tag.(string); ok { + tags[i] = strTag + } + } + } + } + if v, ok := spec["value"]; ok { + val = v.(string) + } + } return &resource.IndexableDocument{ - Key: key, - Title: fmt.Sprintf("Document %s", key.Name), - Tags: []string{"test", "benchmark"}, + Key: &resource.ResourceKey{ + Namespace: key.Namespace, + Group: key.Group, + Resource: key.Resource, + Name: u.GetName(), + }, + Title: title, + Tags: tags, Fields: map[string]interface{}{ - "value": string(value), + "value": val, }, }, nil } // testDocumentBuilderSupplier implements DocumentBuilderSupplier for testing type testDocumentBuilderSupplier struct { - opts *BenchmarkOptions + groupsResources map[string]string } func (s *testDocumentBuilderSupplier) GetDocumentBuilders() ([]resource.DocumentBuilderInfo, error) { - builders := make([]resource.DocumentBuilderInfo, 0, s.opts.NumGroups*s.opts.NumResourceTypes) + builders := make([]resource.DocumentBuilderInfo, 0, len(s.groupsResources)) // Add builders for all possible group/resource combinations - for g := 0; g < s.opts.NumGroups; g++ { - group := fmt.Sprintf("group-%d", g) - for r := 0; r < s.opts.NumResourceTypes; r++ { - resourceType := fmt.Sprintf("resource-%d", r) - builders = append(builders, resource.DocumentBuilderInfo{ - GroupResource: schema.GroupResource{ - Group: group, - Resource: resourceType, - }, - Builder: &testDocumentBuilder{}, - }) - } + for group, resourceType := range s.groupsResources { + builders = append(builders, resource.DocumentBuilderInfo{ + GroupResource: schema.GroupResource{ + Group: group, + Resource: resourceType, + }, + Builder: &testDocumentBuilder{}, + }) } return builders, nil diff --git a/pkg/storage/unified/testing/search_and_storage.go b/pkg/storage/unified/testing/search_and_storage.go new file mode 100644 index 00000000000..ba0b83dbf19 --- /dev/null +++ b/pkg/storage/unified/testing/search_and_storage.go @@ -0,0 +1,249 @@ +package test + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + + claims "github.com/grafana/authlib/types" + "github.com/grafana/grafana/pkg/apimachinery/identity" + "github.com/grafana/grafana/pkg/apimachinery/utils" + "github.com/grafana/grafana/pkg/storage/unified/resource" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +func RunTestSearchAndStorage(t *testing.T, ctx context.Context, backend resource.StorageBackend, searchBackend resource.SearchBackend) { + // Create a test user with admin permissions + testUser := &identity.StaticRequester{ + Type: claims.TypeUser, + Login: "testuser", + UserID: 123, + UserUID: "u123", + OrgRole: identity.RoleAdmin, + IsGrafanaAdmin: true, + } + ctx = claims.WithAuthInfo(ctx, testUser) + + nsPrefix := "test-ns" + + var server resource.ResourceServer + + t.Run("Create initial resources in storage", func(t *testing.T) { + initialResources := []struct { + name string + title string + tags []string + }{ + { + name: "initial1", + title: "First Initial Document", + tags: []string{"tag0", "initial"}, + }, + { + name: "initial2", + title: "Second Initial Document", + tags: []string{"tag0", "initial"}, + }, + } + + for _, doc := range initialResources { + key := &resource.ResourceKey{ + Group: "test.grafana.app", + Resource: "testresources", + Namespace: nsPrefix, + Name: doc.name, + } + + // Create document using unstructured + obj := &unstructured.Unstructured{ + Object: map[string]interface{}{ + "apiVersion": "test.grafana.app/v1", + "kind": "testresources", + "metadata": map[string]interface{}{ + "name": doc.name, + "namespace": nsPrefix, + }, + "spec": map[string]interface{}{ + "title": doc.title, + "tags": doc.tags, + }, + }, + } + + meta, err := utils.MetaAccessor(obj) + require.NoError(t, err) + + // Convert unstructured to bytes + value, err := obj.MarshalJSON() + require.NoError(t, err) + + // Create document + rv, err := backend.WriteEvent(ctx, resource.WriteEvent{ + Type: resource.WatchEvent_ADDED, + Key: key, + Value: value, + Object: meta, + }) + require.NoError(t, err) + require.Greater(t, rv, int64(0)) + } + }) + ch := make(chan *resource.IndexEvent) + + t.Run("Create a resource server with both backends", func(t *testing.T) { + // Create a resource server with both backends + var err error + server, err = resource.NewResourceServer(resource.ResourceServerOptions{ + Backend: backend, + Search: resource.SearchOptions{ + Backend: searchBackend, + Resources: &testDocumentBuilderSupplier{ + groupsResources: map[string]string{ + "test.grafana.app": "testresources", + }, + }, + IndexEventsChan: ch, + }, + }) + require.NoError(t, err) + }) + + t.Run("Search for initial resources", func(t *testing.T) { + // Test 1: Search for initial resources + searchResp, err := server.Search(ctx, &resource.ResourceSearchRequest{ + Options: &resource.ListOptions{ + Key: &resource.ResourceKey{ + Group: "test.grafana.app", + Resource: "testresources", + Namespace: nsPrefix, + }, + }, + Query: "initial", + Limit: 10, + }) + require.NoError(t, err) + require.NotNil(t, searchResp) + require.Nil(t, searchResp.Error) + require.Equal(t, int64(2), searchResp.TotalHits) + }) + + t.Run("Add search documents", func(t *testing.T) { + testDocs := []struct { + name string + title string + tags []string + }{ + { + name: "doc1", + title: "First Test Document", + tags: []string{"hello", "tag1"}, + }, + { + name: "doc2", + title: "Second Test Document", + tags: []string{"hello", "tag2"}, + }, + { + name: "doc3", + title: "Third Test Document", + tags: []string{"hello", "tag3"}, + }, + } + + for _, doc := range testDocs { + key := &resource.ResourceKey{ + Group: "test.grafana.app", + Resource: "testresources", + Namespace: nsPrefix, + Name: doc.name, + } + + // Create document using unstructured + obj := &unstructured.Unstructured{ + Object: map[string]interface{}{ + "apiVersion": "test.grafana.app/v1", + "kind": "testresources", + "metadata": map[string]interface{}{ + "name": doc.name, + "namespace": nsPrefix, + }, + "spec": map[string]interface{}{ + "title": doc.title, + "tags": doc.tags, + }, + }, + } + + // Convert unstructured to bytes + value, err := obj.MarshalJSON() + require.NoError(t, err) + + // Create document + createResp, err := server.Create(ctx, &resource.CreateRequest{ + Key: key, + Value: value, + }) + require.NoError(t, err) + require.NotNil(t, createResp) + require.Nil(t, createResp.Error) + + ev := <-ch + require.NotNil(t, ev) + } + }) + + t.Run("Search for documents", func(t *testing.T) { + searchResp, err := server.Search(ctx, &resource.ResourceSearchRequest{ + Options: &resource.ListOptions{ + Key: &resource.ResourceKey{ + Group: "test.grafana.app", + Resource: "testresources", + Namespace: nsPrefix, + }, + }, + Query: "Document", + Limit: 10, + }) + require.NoError(t, err) + require.NotNil(t, searchResp) + require.Nil(t, searchResp.Error) + require.Equal(t, int64(5), searchResp.TotalHits) + }) + + t.Run("Search with tags", func(t *testing.T) { + searchResp, err := server.Search(ctx, &resource.ResourceSearchRequest{ + Options: &resource.ListOptions{ + Key: &resource.ResourceKey{ + Group: "test.grafana.app", + Resource: "testresources", + Namespace: nsPrefix, + }, + }, + Query: "hello", + Limit: 10, + }) + require.NoError(t, err) + require.NotNil(t, searchResp) + require.Nil(t, searchResp.Error) + require.Equal(t, int64(3), searchResp.TotalHits) + }) + + t.Run("Search with specific tag", func(t *testing.T) { + searchResp, err := server.Search(ctx, &resource.ResourceSearchRequest{ + Options: &resource.ListOptions{ + Key: &resource.ResourceKey{ + Group: "test.grafana.app", + Resource: "testresources", + Namespace: nsPrefix, + }, + }, + Query: "tag1", + Limit: 10, + }) + require.NoError(t, err) + require.NotNil(t, searchResp) + require.Nil(t, searchResp.Error) + require.Equal(t, int64(1), searchResp.TotalHits) + }) +} diff --git a/pkg/storage/unified/testing/search_backend.go b/pkg/storage/unified/testing/search_backend.go index edb03073273..55b7722b756 100644 --- a/pkg/storage/unified/testing/search_backend.go +++ b/pkg/storage/unified/testing/search_backend.go @@ -160,7 +160,7 @@ func runTestResourceIndex(t *testing.T, backend resource.SearchBackend, nsPrefix require.NotNil(t, index) t.Run("Search", func(t *testing.T) { - req := &resource.ResourceSearchRequest{ + resp, err := index.Search(ctx, nil, &resource.ResourceSearchRequest{ Options: &resource.ListOptions{ Key: &resource.ResourceKey{ Namespace: ns.Namespace, @@ -171,10 +171,68 @@ func runTestResourceIndex(t *testing.T, backend resource.SearchBackend, nsPrefix Fields: []string{"title", "folder", "tags"}, Query: "tag3", Limit: 10, - } - resp, err := index.Search(ctx, nil, req, nil) + }, nil) require.NoError(t, err) require.NotNil(t, resp) require.Equal(t, int64(1), resp.TotalHits) // Only doc3 should have tag3 now + + // Search for Document + resp, err = index.Search(ctx, nil, &resource.ResourceSearchRequest{ + Options: &resource.ListOptions{ + Key: &resource.ResourceKey{ + Namespace: ns.Namespace, + Group: ns.Group, + Resource: ns.Resource, + }, + }, + Query: "Document", + Fields: []string{"title", "folder", "tags"}, + Limit: 10, + }, nil) + require.NoError(t, err) + require.NotNil(t, resp) + require.Equal(t, int64(2), resp.TotalHits) // Both doc1 and doc2 should have doc now + }) + + t.Run("Add a new document", func(t *testing.T) { + // Add a new document + err := index.BulkIndex(&resource.BulkIndexRequest{ + Items: []*resource.BulkIndexItem{ + { + Action: resource.ActionIndex, + Doc: &resource.IndexableDocument{ + Key: &resource.ResourceKey{ + Namespace: ns.Namespace, + Group: ns.Group, + Resource: ns.Resource, + Name: "doc3", + }, + Title: "Document 3", + Tags: []string{"tag3", "tag4"}, + Fields: map[string]interface{}{ + "field1": 3, + "field2": "value3", + }, + }, + }, + }, + }) + require.NoError(t, err) + // Search for Document + resp, err := index.Search(ctx, nil, &resource.ResourceSearchRequest{ + Options: &resource.ListOptions{ + Key: &resource.ResourceKey{ + Namespace: ns.Namespace, + Group: ns.Group, + Resource: ns.Resource, + }, + }, + Query: "Document", + Fields: []string{"title", "folder", "tags"}, + Limit: 10, + }, 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 }) } diff --git a/pkg/storage/unified/testing/storage_backend.go b/pkg/storage/unified/testing/storage_backend.go index df2b2f5833f..71849648943 100644 --- a/pkg/storage/unified/testing/storage_backend.go +++ b/pkg/storage/unified/testing/storage_backend.go @@ -144,7 +144,7 @@ func runTestIntegrationBackendHappyPath(t *testing.T, backend resource.StorageBa }) require.Nil(t, resp.Error) require.Equal(t, rv4, resp.ResourceVersion) - require.Equal(t, "item2 MODIFIED", string(resp.Value)) + require.Contains(t, string(resp.Value), "item2 MODIFIED") require.Equal(t, "folderuid", resp.Folder) }) @@ -160,7 +160,7 @@ func runTestIntegrationBackendHappyPath(t *testing.T, backend resource.StorageBa }) require.Nil(t, resp.Error) require.Equal(t, rv2, resp.ResourceVersion) - require.Equal(t, "item2 ADDED", string(resp.Value)) + require.Contains(t, string(resp.Value), "item2 ADDED") }) t.Run("List latest", func(t *testing.T) { @@ -176,8 +176,8 @@ func runTestIntegrationBackendHappyPath(t *testing.T, backend resource.StorageBa require.NoError(t, err) require.Nil(t, resp.Error) require.Len(t, resp.Items, 2) - require.Equal(t, "item2 MODIFIED", string(resp.Items[0].Value)) - require.Equal(t, "item3 ADDED", string(resp.Items[1].Value)) + require.Contains(t, string(resp.Items[0].Value), "item2 MODIFIED") + require.Contains(t, string(resp.Items[1].Value), "item3 ADDED") require.GreaterOrEqual(t, resp.ResourceVersion, rv5) // rv5 is the latest resource version }) @@ -373,11 +373,11 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend require.Nil(t, res.Error) require.Len(t, res.Items, 5) // should be sorted by key ASC - require.Equal(t, "item1 ADDED", string(res.Items[0].Value)) - require.Equal(t, "item2 MODIFIED", string(res.Items[1].Value)) - require.Equal(t, "item4 ADDED", string(res.Items[2].Value)) - require.Equal(t, "item5 ADDED", string(res.Items[3].Value)) - require.Equal(t, "item6 ADDED", string(res.Items[4].Value)) + require.Contains(t, string(res.Items[0].Value), "item1 ADDED") + require.Contains(t, string(res.Items[1].Value), "item2 MODIFIED") + require.Contains(t, string(res.Items[2].Value), "item4 ADDED") + require.Contains(t, string(res.Items[3].Value), "item5 ADDED") + require.Contains(t, string(res.Items[4].Value), "item6 ADDED") require.Empty(t, res.NextPageToken) }) @@ -397,9 +397,9 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend require.Len(t, res.Items, 3) continueToken, err := resource.GetContinueToken(res.NextPageToken) require.NoError(t, err) - require.Equal(t, "item1 ADDED", string(res.Items[0].Value)) - require.Equal(t, "item2 MODIFIED", string(res.Items[1].Value)) - require.Equal(t, "item4 ADDED", string(res.Items[2].Value)) + require.Contains(t, string(res.Items[0].Value), "item1 ADDED") + require.Contains(t, string(res.Items[1].Value), "item2 MODIFIED") + require.Contains(t, string(res.Items[2].Value), "item4 ADDED") require.GreaterOrEqual(t, continueToken.ResourceVersion, rv8) }) @@ -416,10 +416,10 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend require.NoError(t, err) require.Nil(t, res.Error) require.Len(t, res.Items, 4) - require.Equal(t, "item1 ADDED", string(res.Items[0].Value)) - require.Equal(t, "item2 ADDED", string(res.Items[1].Value)) - require.Equal(t, "item3 ADDED", string(res.Items[2].Value)) - require.Equal(t, "item4 ADDED", string(res.Items[3].Value)) + require.Contains(t, string(res.Items[0].Value), "item1 ADDED") + require.Contains(t, string(res.Items[1].Value), "item2 ADDED") + require.Contains(t, string(res.Items[2].Value), "item3 ADDED") + require.Contains(t, string(res.Items[3].Value), "item4 ADDED") require.Empty(t, res.NextPageToken) }) @@ -439,9 +439,9 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend require.Nil(t, res.Error) require.Len(t, res.Items, 3) t.Log(res.Items) - require.Equal(t, "item1 ADDED", string(res.Items[0].Value)) - require.Equal(t, "item2 MODIFIED", string(res.Items[1].Value)) - require.Equal(t, "item4 ADDED", string(res.Items[2].Value)) + require.Contains(t, string(res.Items[0].Value), "item1 ADDED") + require.Contains(t, string(res.Items[1].Value), "item2 MODIFIED") + require.Contains(t, string(res.Items[2].Value), "item4 ADDED") continueToken, err := resource.GetContinueToken(res.NextPageToken) require.NoError(t, err) @@ -468,8 +468,8 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend require.Nil(t, res.Error) require.Len(t, res.Items, 2) t.Log(res.Items) - require.Equal(t, "item4 ADDED", string(res.Items[0].Value)) - require.Equal(t, "item5 ADDED", string(res.Items[1].Value)) + require.Contains(t, string(res.Items[0].Value), "item4 ADDED") + require.Contains(t, string(res.Items[1].Value), "item5 ADDED") continueToken, err = resource.GetContinueToken(res.NextPageToken) require.NoError(t, err) @@ -580,7 +580,7 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage // Check versions and values match expectations for i := 0; i < 3; i++ { require.Equal(t, tc.expectedVersions[i], res.Items[i].ResourceVersion) - require.Equal(t, tc.expectedValues[i], string(res.Items[i].Value)) + require.Contains(t, string(res.Items[i].Value), tc.expectedValues[i]) } // Check resource version in response @@ -628,7 +628,7 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage expectedRVs := []int64{rvHistory3, rvHistory4, rvHistory5} for i, expectedRV := range expectedRVs { require.Equal(t, expectedRV, secondPageRes.Items[i].ResourceVersion) - require.Equal(t, "item1 MODIFIED", string(secondPageRes.Items[i].Value)) + require.Contains(t, string(secondPageRes.Items[i].Value), "item1 MODIFIED") } }) }) @@ -656,9 +656,9 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage require.Nil(t, res.Error) require.Len(t, res.Items, 2) t.Log(res.Items) - require.Equal(t, "item1 MODIFIED", string(res.Items[0].Value)) + require.Contains(t, string(res.Items[0].Value), "item1 MODIFIED") require.Equal(t, rvHistory2, res.Items[0].ResourceVersion) - require.Equal(t, "item1 MODIFIED", string(res.Items[1].Value)) + require.Contains(t, string(res.Items[1].Value), "item1 MODIFIED") require.Equal(t, rvHistory1, res.Items[1].ResourceVersion) }) @@ -770,11 +770,11 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage // Verify the first item is the initial ADDED event require.Equal(t, initialRV, allItems[0].ResourceVersion, "First item should be the initial ADDED event") - require.Equal(t, "paged-item ADDED", string(allItems[0].Value)) + require.Contains(t, string(allItems[0].Value), "paged-item ADDED") // Verify all other items are MODIFIED events and correspond to our recorded resource versions for i := 1; i < len(allItems); i++ { - require.Equal(t, "paged-item MODIFIED", string(allItems[i].Value)) + require.Contains(t, string(allItems[i].Value), "paged-item MODIFIED") require.Equal(t, resourceVersions[i], allItems[i].ResourceVersion) } }) @@ -835,9 +835,9 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage require.NoError(t, err) require.Nil(t, res.Error) require.Len(t, res.Items, 2) - require.Equal(t, "deleted-item MODIFIED", string(res.Items[0].Value)) + require.Contains(t, string(res.Items[0].Value), "deleted-item MODIFIED") require.Equal(t, rv2, res.Items[0].ResourceVersion) - require.Equal(t, "deleted-item ADDED", string(res.Items[1].Value)) + require.Contains(t, string(res.Items[1].Value), "deleted-item ADDED") require.Equal(t, rv1, res.Items[1].ResourceVersion) }) } @@ -933,11 +933,11 @@ func runTestIntegrationBlobSupport(t *testing.T, backend resource.StorageBackend // Check that we can still access both values found, err := store.GetResourceBlob(ctx, key, &utils.BlobInfo{UID: b1.Uid}, true) require.NoError(t, err) - require.Equal(t, []byte("hello 11111"), found.Value) + require.Contains(t, string(found.Value), "hello 11111") found, err = store.GetResourceBlob(ctx, key, &utils.BlobInfo{UID: b2.Uid}, true) require.NoError(t, err) - require.Equal(t, []byte("hello 22222"), found.Value) + require.Contains(t, string(found.Value), "hello 22222") // Save a resource with annotation obj := &unstructured.Unstructured{} @@ -959,13 +959,13 @@ func runTestIntegrationBlobSupport(t *testing.T, backend resource.StorageBackend res, err := server.GetBlob(ctx, &resource.GetBlobRequest{Resource: key}) require.NoError(t, err) require.Nil(t, out.Error) - require.Equal(t, "hello 22222", string(res.Value)) + require.Contains(t, string(res.Value), "hello 22222") // But we can still get an older version with an explicit UID res, err = server.GetBlob(ctx, &resource.GetBlobRequest{Resource: key, Uid: b1.Uid}) require.NoError(t, err) require.Nil(t, out.Error) - require.Equal(t, "hello 11111", string(res.Value)) + require.Contains(t, string(res.Value), "hello 11111") }) } @@ -1037,9 +1037,22 @@ func WithFolder(folder string) WriteEventOption { } // WithValue sets the value for the write event -func WithValue(value []byte) WriteEventOption { +func WithValue(value string) WriteEventOption { return func(o *WriteEventOptions) { - o.Value = value + u := unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": o.Group + "/v1", + "kind": o.Resource, + "metadata": map[string]any{ + "name": "name", + "namespace": "ns", + }, + "spec": map[string]any{ + "value": value, + }, + }, + } + o.Value, _ = u.MarshalJSON() } } @@ -1064,10 +1077,21 @@ func writeEvent(ctx context.Context, store resource.StorageBackend, name string, for _, opt := range opts { opt(&options) } - - // Set default value if not provided if options.Value == nil { - options.Value = []byte(name + " " + resource.WatchEvent_Type_name[int32(action)]) + u := unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": options.Group + "/v1", + "kind": options.Resource, + "metadata": map[string]any{ + "name": name, + "namespace": options.Namespace, + }, + "spec": map[string]any{ + "value": name + " " + resource.WatchEvent_Type_name[int32(action)], + }, + }, + } + options.Value, _ = u.MarshalJSON() } res := &unstructured.Unstructured{