diff --git a/pkg/storage/unified/apistore/restoptions.go b/pkg/storage/unified/apistore/restoptions.go index 58f702e78c1..53ce1d2bb2a 100644 --- a/pkg/storage/unified/apistore/restoptions.go +++ b/pkg/storage/unified/apistore/restoptions.go @@ -3,11 +3,13 @@ package apistore import ( + "context" "os" "path/filepath" "time" - badger "github.com/dgraph-io/badger/v4" + "gocloud.dev/blob/fileblob" + "gocloud.dev/blob/memblob" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apiserver/pkg/registry/generic" @@ -51,29 +53,18 @@ func NewRESTOptionsGetterForClient( } func NewRESTOptionsGetterMemory(originalStorageConfig storagebackend.Config, secrets secret.InlineSecureValueSupport) (*RESTOptionsGetter, error) { - // Create BadgerDB with in-memory mode - db, err := badger.Open(badger.DefaultOptions(""). - WithInMemory(true). - WithLogger(nil)) - if err != nil { - return nil, err - } - - kv := resource.NewBadgerKV(db) - backend, err := resource.NewKVStorageBackend(resource.KVBackendOptions{ - KvStore: kv, + backend, err := resource.NewCDKBackend(context.Background(), resource.CDKBackendOptions{ + Bucket: memblob.OpenBucket(&memblob.Options{}), }) if err != nil { return nil, err } - server, err := resource.NewResourceServer(resource.ResourceServerOptions{ Backend: backend, }) if err != nil { return nil, err } - return NewRESTOptionsGetterForClient( resource.NewLocalResourceClient(server), secrets, @@ -92,27 +83,25 @@ func NewRESTOptionsGetterForFileXX(path string, path = filepath.Join(os.TempDir(), "grafana-apiserver") } - db, err := badger.Open(badger.DefaultOptions(filepath.Join(path, "badger")). - WithLogger(nil)) + bucket, err := fileblob.OpenBucket(filepath.Join(path, "resource"), &fileblob.Options{ + CreateDir: true, + Metadata: fileblob.MetadataDontWrite, // skip + }) if err != nil { return nil, err } - - kv := resource.NewBadgerKV(db) - backend, err := resource.NewKVStorageBackend(resource.KVBackendOptions{ - KvStore: kv, + backend, err := resource.NewCDKBackend(context.Background(), resource.CDKBackendOptions{ + Bucket: bucket, }) if err != nil { return nil, err } - server, err := resource.NewResourceServer(resource.ResourceServerOptions{ Backend: backend, }) if err != nil { return nil, err } - return NewRESTOptionsGetterForClient( resource.NewLocalResourceClient(server), nil, // secrets diff --git a/pkg/storage/unified/apistore/watcher_test.go b/pkg/storage/unified/apistore/watcher_test.go index d23c8700c90..afd2a77b08b 100644 --- a/pkg/storage/unified/apistore/watcher_test.go +++ b/pkg/storage/unified/apistore/watcher_test.go @@ -8,13 +8,15 @@ package apistore_test import ( "context" "fmt" + "os" "strings" "testing" "time" - badger "github.com/dgraph-io/badger/v4" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "gocloud.dev/blob/fileblob" + "gocloud.dev/blob/memblob" "k8s.io/apimachinery/pkg/api/apitesting" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" @@ -103,20 +105,24 @@ func testSetup(t testing.TB, opts ...setupOption) (context.Context, storage.Inte Resource: "pods", } + bucket := memblob.OpenBucket(nil) + if true { + tmp, err := os.MkdirTemp("", "xxx-*") + require.NoError(t, err) + + bucket, err = fileblob.OpenBucket(tmp, &fileblob.Options{ + CreateDir: true, + Metadata: fileblob.MetadataDontWrite, // skip + }) + require.NoError(t, err) + } ctx := storagetesting.NewContext() var server resource.ResourceServer switch setupOpts.storageType { case StorageTypeFile: - // Create in-memory BadgerDB for testing - db, err := badger.Open(badger.DefaultOptions(""). - WithInMemory(true). - WithLogger(nil)) - require.NoError(t, err) - - kv := resource.NewBadgerKV(db) - backend, err := resource.NewKVStorageBackend(resource.KVBackendOptions{ - KvStore: kv, + backend, err := resource.NewCDKBackend(ctx, resource.CDKBackendOptions{ + Bucket: bucket, }) require.NoError(t, err) diff --git a/pkg/storage/unified/client.go b/pkg/storage/unified/client.go index 0a4d3e630ff..b0e7f6cf275 100644 --- a/pkg/storage/unified/client.go +++ b/pkg/storage/unified/client.go @@ -6,12 +6,12 @@ import ( "path/filepath" "time" - badger "github.com/dgraph-io/badger/v4" otgrpc "github.com/opentracing-contrib/go-grpc" "github.com/opentracing/opentracing-go" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promauto" "go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc" + "gocloud.dev/blob/fileblob" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/keepalive" @@ -111,22 +111,19 @@ func newClient(opts options.StorageOptions, if opts.DataPath == "" { opts.DataPath = filepath.Join(cfg.DataPath, "grafana-apiserver") } - - // Create BadgerDB instance - db, err := badger.Open(badger.DefaultOptions(filepath.Join(opts.DataPath, "badger")). - WithLogger(nil)) + bucket, err := fileblob.OpenBucket(filepath.Join(opts.DataPath, "resource"), &fileblob.Options{ + CreateDir: true, + Metadata: fileblob.MetadataDontWrite, // skip + }) if err != nil { return nil, err } - - kv := resource.NewBadgerKV(db) - backend, err := resource.NewKVStorageBackend(resource.KVBackendOptions{ - KvStore: kv, + backend, err := resource.NewCDKBackend(ctx, resource.CDKBackendOptions{ + Bucket: bucket, }) if err != nil { return nil, err } - server, err := resource.NewResourceServer(resource.ResourceServerOptions{ Backend: backend, Blob: resource.BlobConfig{ diff --git a/pkg/storage/unified/resource/cdk_backend.go b/pkg/storage/unified/resource/cdk_backend.go new file mode 100644 index 00000000000..0cb9628758a --- /dev/null +++ b/pkg/storage/unified/resource/cdk_backend.go @@ -0,0 +1,418 @@ +package resource + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "iter" + "net/http" + "sort" + "strconv" + "strings" + "sync" + "sync/atomic" + "time" + + "go.opentelemetry.io/otel/trace" + "go.opentelemetry.io/otel/trace/noop" + "gocloud.dev/blob" + _ "gocloud.dev/blob/fileblob" + _ "gocloud.dev/blob/memblob" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + + "github.com/grafana/grafana/pkg/apimachinery/utils" + "github.com/grafana/grafana/pkg/storage/unified/resourcepb" +) + +type CDKBackendOptions struct { + Tracer trace.Tracer + Bucket CDKBucket + RootFolder string +} + +func NewCDKBackend(ctx context.Context, opts CDKBackendOptions) (StorageBackend, error) { + if opts.Tracer == nil { + opts.Tracer = noop.NewTracerProvider().Tracer("cdk-appending-store") + } + + if opts.Bucket == nil { + return nil, fmt.Errorf("missing bucket") + } + + found, _, err := opts.Bucket.ListPage(ctx, blob.FirstPageToken, 1, &blob.ListOptions{ + Prefix: opts.RootFolder, + Delimiter: "/", + }) + if err != nil { + return nil, err + } + if found == nil { + return nil, fmt.Errorf("the root folder does not exist") + } + + backend := &cdkBackend{ + tracer: opts.Tracer, + bucket: opts.Bucket, + root: opts.RootFolder, + } + backend.rv.Swap(time.Now().UnixMilli()) + return backend, nil +} + +type cdkBackend struct { + tracer trace.Tracer + bucket CDKBucket + root string + + mutex sync.Mutex + rv atomic.Int64 + + // Simple watch stream -- NOTE, this only works for single tenant! + broadcaster Broadcaster[*WrittenEvent] + stream chan<- *WrittenEvent +} + +func (s *cdkBackend) GetResourceLastImportTimes(ctx context.Context) iter.Seq2[ResourceLastImportTime, error] { + return func(yield func(ResourceLastImportTime, error) bool) { + yield(ResourceLastImportTime{}, errors.New("not implemented")) + } +} + +func (s *cdkBackend) ListModifiedSince(ctx context.Context, key NamespacedResource, sinceRv int64) (int64, iter.Seq2[*ModifiedResource, error]) { + return 0, func(yield func(*ModifiedResource, error) bool) { + yield(nil, errors.New("not implemented")) + } +} + +func (s *cdkBackend) getPath(key *resourcepb.ResourceKey, rv int64) string { + var buffer bytes.Buffer + buffer.WriteString(s.root) + + if key.Group == "" { + return buffer.String() + } + buffer.WriteString(key.Group) + + if key.Resource == "" { + return buffer.String() + } + buffer.WriteString("/") + buffer.WriteString(key.Resource) + + if key.Namespace == "" { + if key.Name == "" { + return buffer.String() + } + buffer.WriteString("/__cluster__") + } else { + buffer.WriteString("/") + buffer.WriteString(key.Namespace) + } + + if key.Name == "" { + return buffer.String() + } + buffer.WriteString("/") + buffer.WriteString(key.Name) + + if rv > 0 { + buffer.WriteString(fmt.Sprintf("/%d.json", rv)) + } + return buffer.String() +} + +// GetResourceStats implements Backend. +func (s *cdkBackend) GetResourceStats(ctx context.Context, namespace string, minCount int) ([]ResourceStats, error) { + return nil, fmt.Errorf("not implemented") +} + +func (s *cdkBackend) WriteEvent(ctx context.Context, event WriteEvent) (rv int64, err error) { + if event.Type == resourcepb.WatchEvent_ADDED { + // ReadResource deals with deleted values (i.e. a file exists but has generation -999). + resp := s.ReadResource(ctx, &resourcepb.ReadRequest{Key: event.Key}) + if resp.Error != nil && resp.Error.Code != http.StatusNotFound { + return 0, GetError(resp.Error) + } + if resp.Value != nil { + return 0, ErrResourceAlreadyExists + } + } + + // Scope the lock + { + s.mutex.Lock() + defer s.mutex.Unlock() + + rv = s.rv.Add(1) + err = s.bucket.WriteAll(ctx, s.getPath(event.Key, rv), event.Value, &blob.WriterOptions{ + ContentType: "application/json", + }) + } + + // notify all subscribers + if s.stream != nil { + write := &WrittenEvent{ + Type: event.Type, + Key: event.Key, + PreviousRV: event.PreviousRV, + Value: event.Value, + Timestamp: time.Now().UnixMilli(), + ResourceVersion: rv, + } + s.stream <- write + } + return rv, err +} + +func (s *cdkBackend) ReadResource(ctx context.Context, req *resourcepb.ReadRequest) *BackendReadResponse { + rv := req.ResourceVersion + + path := s.getPath(req.Key, rv) + if rv < 1 { + iter := s.bucket.List(&blob.ListOptions{Prefix: path + "/", Delimiter: "/"}) + for { + obj, err := iter.Next(ctx) + if errors.Is(err, io.EOF) { + break + } + if strings.HasSuffix(obj.Key, ".json") { + idx := strings.LastIndex(obj.Key, "/") + 1 + edx := strings.LastIndex(obj.Key, ".") + if idx > 0 { + v, err := strconv.ParseInt(obj.Key[idx:edx], 10, 64) + if err == nil && v > rv { + rv = v + path = obj.Key // find the path with biggest resource version + } + } + } + } + } + + raw, err := s.bucket.ReadAll(ctx, path) + if raw == nil && req.ResourceVersion > 0 { + if req.ResourceVersion > s.rv.Load() { + return &BackendReadResponse{ + Error: &resourcepb.ErrorResult{ + Code: http.StatusGatewayTimeout, + Reason: string(metav1.StatusReasonTimeout), // match etcd behavior + Message: "ResourceVersion is larger than max", + Details: &resourcepb.ErrorDetails{ + Causes: []*resourcepb.ErrorCause{ + { + Reason: string(metav1.CauseTypeResourceVersionTooLarge), + Message: fmt.Sprintf("requested: %d, current %d", req.ResourceVersion, s.rv.Load()), + }, + }, + }, + }, + } + } + + // If the there was an explicit request, get the latest + rsp := s.ReadResource(ctx, &resourcepb.ReadRequest{Key: req.Key}) + if rsp != nil && len(rsp.Value) > 0 { + raw = rsp.Value + rv = rsp.ResourceVersion + err = nil + } + } + if err == nil && isDeletedValue(raw) { + raw = nil + } + if raw == nil { + return &BackendReadResponse{Error: NewNotFoundError(req.Key)} + } + return &BackendReadResponse{ + Key: req.Key, + Folder: "", // TODO: implement this + ResourceVersion: rv, + Value: raw, + } +} + +func isDeletedValue(raw []byte) bool { + if bytes.Contains(raw, []byte(`"generation":-999`)) { + tmp := &unstructured.Unstructured{} + err := tmp.UnmarshalJSON(raw) + if err == nil && tmp.GetGeneration() == utils.DeletedGeneration { + return true + } + } + return false +} + +func (s *cdkBackend) ListIterator(ctx context.Context, req *resourcepb.ListRequest, cb func(ListIterator) error) (int64, error) { + resources, err := buildTree(ctx, s, req.Options.Key) + if err != nil { + return 0, err + } + err = cb(resources) + return resources.listRV, err +} + +func (s *cdkBackend) ListHistory(ctx context.Context, req *resourcepb.ListRequest, cb func(ListIterator) error) (int64, error) { + return 0, fmt.Errorf("listing from history not supported in CDK backend") +} + +func (s *cdkBackend) WatchWriteEvents(ctx context.Context) (<-chan *WrittenEvent, error) { + s.mutex.Lock() + defer s.mutex.Unlock() + + if s.broadcaster == nil { + var err error + s.broadcaster, err = NewBroadcaster(context.Background(), func(c chan<- *WrittenEvent) error { + s.stream = c + return nil + }) + if err != nil { + return nil, err + } + } + return s.broadcaster.Subscribe(ctx) +} + +// group > resource > namespace > name > versions +type cdkResource struct { + prefix string + versions []cdkVersion +} +type cdkVersion struct { + rv int64 + key string +} + +type cdkListIterator struct { + bucket CDKBucket + ctx context.Context + err error + + listRV int64 + resources []cdkResource + index int + + currentRV int64 + currentKey string + currentVal []byte +} + +// Next implements ListIterator. +func (c *cdkListIterator) Next() bool { + if c.err != nil { + return false + } + for { + c.currentVal = nil + c.index += 1 + if c.index >= len(c.resources) { + return false + } + + item := c.resources[c.index] + latest := item.versions[0] + raw, err := c.bucket.ReadAll(c.ctx, latest.key) + if err != nil { + c.err = err + return false + } + if !isDeletedValue(raw) { + c.currentRV = latest.rv + c.currentKey = latest.key + c.currentVal = raw + return true + } + } +} + +// Error implements ListIterator. +func (c *cdkListIterator) Error() error { + return c.err +} + +// ResourceVersion implements ListIterator. +func (c *cdkListIterator) ResourceVersion() int64 { + return c.currentRV +} + +// Value implements ListIterator. +func (c *cdkListIterator) Value() []byte { + return c.currentVal +} + +// ContinueToken implements ListIterator. +func (c *cdkListIterator) ContinueToken() string { + return fmt.Sprintf("index:%d/key:%s", c.index, c.currentKey) +} + +// Name implements ListIterator. +func (c *cdkListIterator) Name() string { + return c.currentKey // TODO (parse name from key) +} + +// Namespace implements ListIterator. +func (c *cdkListIterator) Namespace() string { + return c.currentKey // TODO (parse namespace from key) +} + +func (c *cdkListIterator) Folder() string { + return "" // TODO: implement this +} + +var _ ListIterator = (*cdkListIterator)(nil) + +func buildTree(ctx context.Context, s *cdkBackend, key *resourcepb.ResourceKey) (*cdkListIterator, error) { + byPrefix := make(map[string]*cdkResource) + path := s.getPath(key, 0) + iter := s.bucket.List(&blob.ListOptions{Prefix: path, Delimiter: ""}) // "" is recursive + for { + obj, err := iter.Next(ctx) + if errors.Is(err, io.EOF) { + break + } + if strings.HasSuffix(obj.Key, ".json") { + idx := strings.LastIndex(obj.Key, "/") + 1 + edx := strings.LastIndex(obj.Key, ".") + if idx > 0 { + rv, err := strconv.ParseInt(obj.Key[idx:edx], 10, 64) + if err == nil { + prefix := obj.Key[:idx] + res, ok := byPrefix[prefix] + if !ok { + res = &cdkResource{prefix: prefix} + byPrefix[prefix] = res + } + + res.versions = append(res.versions, cdkVersion{ + rv: rv, + key: obj.Key, + }) + } + } + } + } + + // Now sort all versions + resources := make([]cdkResource, 0, len(byPrefix)) + for _, res := range byPrefix { + sort.Slice(res.versions, func(i, j int) bool { + return res.versions[i].rv > res.versions[j].rv + }) + resources = append(resources, *res) + } + sort.Slice(resources, func(i, j int) bool { + a := resources[i].prefix + b := resources[j].prefix + return a < b + }) + + return &cdkListIterator{ + ctx: ctx, + bucket: s.bucket, + resources: resources, + listRV: s.rv.Load(), + index: -1, // must call next first + }, nil +} diff --git a/pkg/storage/unified/resource/server.go b/pkg/storage/unified/resource/server.go index 76c3e3c2a8c..59ae4e1dc83 100644 --- a/pkg/storage/unified/resource/server.go +++ b/pkg/storage/unified/resource/server.go @@ -1072,7 +1072,7 @@ func (s *server) List(ctx context.Context, req *resourcepb.ListRequest) (*resour pageBytes += len(item.Value) rsp.Items = append(rsp.Items, item) - if (req.Limit > 0 && len(rsp.Items) >= int(req.Limit)) || pageBytes >= maxPageBytes { + if len(rsp.Items) >= int(req.Limit) || pageBytes >= maxPageBytes { t := iter.ContinueToken() if iter.Next() { rsp.NextPageToken = t diff --git a/pkg/storage/unified/resource/server_test.go b/pkg/storage/unified/resource/server_test.go index 391a931696d..66536b4d5d8 100644 --- a/pkg/storage/unified/resource/server_test.go +++ b/pkg/storage/unified/resource/server_test.go @@ -4,17 +4,20 @@ import ( "context" "encoding/json" "errors" + "fmt" "log/slog" "net/http" + "os" "strings" "sync" "testing" "time" - badger "github.com/dgraph-io/badger/v4" "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "gocloud.dev/blob/fileblob" + "gocloud.dev/blob/memblob" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" authlib "github.com/grafana/authlib/types" @@ -38,19 +41,20 @@ func TestSimpleServer(t *testing.T) { } ctx := authlib.WithAuthInfo(context.Background(), testUserA) - // Create in-memory BadgerDB for testing - db, err := badger.Open(badger.DefaultOptions(""). - WithInMemory(true). - WithLogger(nil)) - require.NoError(t, err) - defer func() { - err := db.Close() + bucket := memblob.OpenBucket(nil) + if false { + tmp, err := os.MkdirTemp("", "xxx-*") require.NoError(t, err) - }() - kv := NewBadgerKV(db) - store, err := NewKVStorageBackend(KVBackendOptions{ - KvStore: kv, + bucket, err = fileblob.OpenBucket(tmp, &fileblob.Options{ + CreateDir: true, + Metadata: fileblob.MetadataDontWrite, // skip + }) + require.NoError(t, err) + fmt.Printf("ROOT: %s\n\n", tmp) + } + store, err := NewCDKBackend(ctx, CDKBackendOptions{ + Bucket: bucket, }) require.NoError(t, err) diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index 4de28e758d5..24bb98c1f8b 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -283,40 +283,6 @@ func (k *kvStorageBackend) ReadResource(ctx context.Context, req *resourcepb.Rea if req.Key == nil { return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusBadRequest, Message: "missing key"}} } - - // If a specific resource version is requested, validate that it's not too high - if req.ResourceVersion > 0 { - // Fetch the latest RV - latestRV := k.snowflake.Generate().Int64() - if lastEventKey, err := k.eventStore.LastEventKey(ctx); err == nil { - latestRV = lastEventKey.ResourceVersion - } else if !errors.Is(err, ErrNotFound) { - return &BackendReadResponse{Error: &resourcepb.ErrorResult{ - Code: http.StatusInternalServerError, - Message: fmt.Sprintf("failed to fetch latest resource version: %v", err), - }} - } - - // Check if the requested RV is higher than the latest available RV - if req.ResourceVersion > latestRV { - return &BackendReadResponse{ - Error: &resourcepb.ErrorResult{ - Code: http.StatusGatewayTimeout, - Reason: string(metav1.StatusReasonTimeout), // match etcd behavior - Message: "ResourceVersion is larger than max", - Details: &resourcepb.ErrorDetails{ - Causes: []*resourcepb.ErrorCause{ - { - Reason: string(metav1.CauseTypeResourceVersionTooLarge), - Message: fmt.Sprintf("requested: %d, current %d", req.ResourceVersion, latestRV), - }, - }, - }, - }, - } - } - } - meta, err := k.dataStore.GetResourceKeyAtRevision(ctx, GetRequestKey{ Group: req.Key.Group, Resource: req.Key.Resource, @@ -369,15 +335,8 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis resourceVersion = token.ResourceVersion } - // We set the listRV to the last event resource version. - // If no events exist yet, we generate a new snowflake. + // We set the listRV to the current time. listRV := k.snowflake.Generate().Int64() - if lastEventKey, err := k.eventStore.LastEventKey(ctx); err == nil { - listRV = lastEventKey.ResourceVersion - } else if !errors.Is(err, ErrNotFound) { - return 0, fmt.Errorf("failed to fetch last event: %w", err) - } - if resourceVersion > 0 { listRV = resourceVersion } @@ -401,7 +360,7 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis } keys = append(keys, dataKey) // Only fetch the first limit items + 1 to get the next token. - if req.Limit > 0 && len(keys) >= int(req.Limit+1) { + if len(keys) >= int(req.Limit+1) { break } } diff --git a/pkg/storage/unified/resource/storage_backend_test.go b/pkg/storage/unified/resource/storage_backend_test.go index 05920a603c4..0a1c263d216 100644 --- a/pkg/storage/unified/resource/storage_backend_test.go +++ b/pkg/storage/unified/resource/storage_backend_test.go @@ -323,36 +323,6 @@ func TestKvStorageBackend_ReadResource_DeletedResource(t *testing.T) { require.Equal(t, objectToJSONBytes(t, testObj), response.Value) } -func TestKvStorageBackend_ReadResource_TooHighResourceVersion(t *testing.T) { - backend := setupTestStorageBackend(t) - ctx := context.Background() - - // First, create a resource - _, rv := createAndWriteTestObject(t, backend) - - // Try to read with a resource version that's way too high - readReq := &resourcepb.ReadRequest{ - Key: &resourcepb.ResourceKey{ - Namespace: "default", - Group: "apps", - Resource: "resources", - Name: "test-resource", - }, - ResourceVersion: rv + 1000000000000, // Way in the future - } - - response := backend.ReadResource(ctx, readReq) - require.NotNil(t, response.Error, "ReadResource should return error for too high resource version") - require.Equal(t, int32(504), response.Error.Code) // http.StatusGatewayTimeout - require.Equal(t, "Timeout", response.Error.Reason) - require.Equal(t, "ResourceVersion is larger than max", response.Error.Message) - require.NotNil(t, response.Error.Details) - require.Len(t, response.Error.Details.Causes, 1) - require.Equal(t, "ResourceVersionTooLarge", response.Error.Details.Causes[0].Reason) - require.Contains(t, response.Error.Details.Causes[0].Message, "requested:") - require.Contains(t, response.Error.Details.Causes[0].Message, "current") -} - func TestKvStorageBackend_ListIterator_Success(t *testing.T) { backend := setupTestStorageBackend(t) ctx := context.Background() diff --git a/pkg/storage/unified/testing/storage_backend.go b/pkg/storage/unified/testing/storage_backend.go index 440f2af8a26..8ce39d55e33 100644 --- a/pkg/storage/unified/testing/storage_backend.go +++ b/pkg/storage/unified/testing/storage_backend.go @@ -387,30 +387,6 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend require.Empty(t, res.NextPageToken) }) - t.Run("fetch all with limit 0", func(t *testing.T) { - res, err := server.List(ctx, &resourcepb.ListRequest{ - Limit: 0, - Options: &resourcepb.ListOptions{ - Key: &resourcepb.ResourceKey{ - Namespace: ns, - Group: "group", - Resource: "resource", - }, - }, - }) - require.NoError(t, err) - require.Nil(t, res.Error) - require.Len(t, res.Items, 5) - // should be sorted by key ASC - 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) - }) - t.Run("list latest first page ", func(t *testing.T) { res, err := server.List(ctx, &resourcepb.ListRequest{ Limit: 3, @@ -781,30 +757,6 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage require.Contains(t, string(secondPageRes.Items[i].Value), "item1 MODIFIED") } }) - - // Test with limit=0 (should return all items) - t.Run("fetch all history with limit 0", func(t *testing.T) { - res, err := server.List(ctx, &resourcepb.ListRequest{ - Limit: 0, - Source: resourcepb.ListRequest_HISTORY, - Options: &resourcepb.ListOptions{ - Key: baseKey, - }, - }) - require.NoError(t, err) - require.Nil(t, res.Error) - require.Len(t, res.Items, 6) // Should return all 6 history items (1 ADDED + 5 MODIFIED) - - // Should be in descending order (default for history) - require.Equal(t, rvHistory5, res.Items[0].ResourceVersion) - require.Equal(t, rvHistory4, res.Items[1].ResourceVersion) - require.Equal(t, rvHistory3, res.Items[2].ResourceVersion) - require.Equal(t, rvHistory2, res.Items[3].ResourceVersion) - require.Equal(t, rvHistory1, res.Items[4].ResourceVersion) - require.Equal(t, rv1, res.Items[5].ResourceVersion) - - require.Empty(t, res.NextPageToken) - }) }) t.Run("fetch second page of history at revision", func(t *testing.T) {