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
This commit is contained in:
Georges Chaudy
2025-05-12 07:56:25 +02:00
committed by GitHub
parent d91e4b0582
commit e162c69c34
7 changed files with 496 additions and 80 deletions
+12 -9
View File
@@ -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
+3
View File
@@ -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!)
@@ -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
+58 -21
View File
@@ -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
@@ -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)
})
}
+61 -3
View File
@@ -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
})
}
+62 -38
View File
@@ -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{