unistore: replace CDK backend with KV store backend (again) (#113184)
* Reapply "unistore: replace CDK backend with KV store backend"" (#113132)
This reverts commit 7127b2538c.
* enable cluster scope
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user