diff --git a/pkg/storage/unified/apistore/restoptions.go b/pkg/storage/unified/apistore/restoptions.go index 53ce1d2bb2a..b48fbc6deaf 100644 --- a/pkg/storage/unified/apistore/restoptions.go +++ b/pkg/storage/unified/apistore/restoptions.go @@ -3,13 +3,11 @@ package apistore import ( - "context" "os" "path/filepath" "time" - "gocloud.dev/blob/fileblob" - "gocloud.dev/blob/memblob" + badger "github.com/dgraph-io/badger/v4" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apiserver/pkg/registry/generic" @@ -53,18 +51,30 @@ func NewRESTOptionsGetterForClient( } func NewRESTOptionsGetterMemory(originalStorageConfig storagebackend.Config, secrets secret.InlineSecureValueSupport) (*RESTOptionsGetter, error) { - backend, err := resource.NewCDKBackend(context.Background(), resource.CDKBackendOptions{ - Bucket: memblob.OpenBucket(&memblob.Options{}), + // 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, + WithExperimentalClusterScope: true, }) 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, @@ -83,25 +93,27 @@ func NewRESTOptionsGetterForFileXX(path string, path = filepath.Join(os.TempDir(), "grafana-apiserver") } - bucket, err := fileblob.OpenBucket(filepath.Join(path, "resource"), &fileblob.Options{ - CreateDir: true, - Metadata: fileblob.MetadataDontWrite, // skip - }) + db, err := badger.Open(badger.DefaultOptions(filepath.Join(path, "badger")). + WithLogger(nil)) if err != nil { return nil, err } - backend, err := resource.NewCDKBackend(context.Background(), resource.CDKBackendOptions{ - Bucket: bucket, + + kv := resource.NewBadgerKV(db) + backend, err := resource.NewKVStorageBackend(resource.KVBackendOptions{ + KvStore: kv, }) 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 afd2a77b08b..d23c8700c90 100644 --- a/pkg/storage/unified/apistore/watcher_test.go +++ b/pkg/storage/unified/apistore/watcher_test.go @@ -8,15 +8,13 @@ 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" @@ -105,24 +103,20 @@ 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: - backend, err := resource.NewCDKBackend(ctx, resource.CDKBackendOptions{ - Bucket: bucket, + // 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, }) require.NoError(t, err) diff --git a/pkg/storage/unified/client.go b/pkg/storage/unified/client.go index b0e7f6cf275..0a4d3e630ff 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,19 +111,22 @@ func newClient(opts options.StorageOptions, if opts.DataPath == "" { opts.DataPath = filepath.Join(cfg.DataPath, "grafana-apiserver") } - bucket, err := fileblob.OpenBucket(filepath.Join(opts.DataPath, "resource"), &fileblob.Options{ - CreateDir: true, - Metadata: fileblob.MetadataDontWrite, // skip - }) + + // Create BadgerDB instance + db, err := badger.Open(badger.DefaultOptions(filepath.Join(opts.DataPath, "badger")). + WithLogger(nil)) if err != nil { return nil, err } - backend, err := resource.NewCDKBackend(ctx, resource.CDKBackendOptions{ - Bucket: bucket, + + kv := resource.NewBadgerKV(db) + backend, err := resource.NewKVStorageBackend(resource.KVBackendOptions{ + KvStore: kv, }) 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 deleted file mode 100644 index 0cb9628758a..00000000000 --- a/pkg/storage/unified/resource/cdk_backend.go +++ /dev/null @@ -1,418 +0,0 @@ -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 59ae4e1dc83..76c3e3c2a8c 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 len(rsp.Items) >= int(req.Limit) || pageBytes >= maxPageBytes { + if (req.Limit > 0 && 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 66536b4d5d8..391a931696d 100644 --- a/pkg/storage/unified/resource/server_test.go +++ b/pkg/storage/unified/resource/server_test.go @@ -4,20 +4,17 @@ 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" @@ -41,20 +38,19 @@ func TestSimpleServer(t *testing.T) { } ctx := authlib.WithAuthInfo(context.Background(), testUserA) - bucket := memblob.OpenBucket(nil) - if false { - tmp, err := os.MkdirTemp("", "xxx-*") + // 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() require.NoError(t, err) + }() - 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, + kv := NewBadgerKV(db) + store, err := NewKVStorageBackend(KVBackendOptions{ + KvStore: kv, }) require.NoError(t, err) diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index 503123ed309..2ee08dd85d9 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -310,6 +310,39 @@ func (k *kvStorageBackend) ReadResource(ctx context.Context, req *resourcepb.Rea namespace := convertEmptyToClusterNamespace(req.Key.Namespace, k.withExperimentalClusterScope) + // 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, @@ -365,8 +398,15 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis resourceVersion = token.ResourceVersion } - // We set the listRV to the current time. + // We set the listRV to the last event resource version. + // If no events exist yet, we generate a new snowflake. 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 } @@ -390,7 +430,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 len(keys) >= int(req.Limit+1) { + if req.Limit > 0 && 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 027a9b4fd7e..772362aceef 100644 --- a/pkg/storage/unified/resource/storage_backend_test.go +++ b/pkg/storage/unified/resource/storage_backend_test.go @@ -336,6 +336,36 @@ 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 8ce39d55e33..440f2af8a26 100644 --- a/pkg/storage/unified/testing/storage_backend.go +++ b/pkg/storage/unified/testing/storage_backend.go @@ -387,6 +387,30 @@ 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, @@ -757,6 +781,30 @@ 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) {