unified-storage: use name instead of offset in kvstore continue token (#113560)
* unified-storage: use name instead of offset in kvstore continue token --------- Co-authored-by: Georges Chaudy <chaudyg@gmail.com>
This commit is contained in:
co-authored by
Georges Chaudy
parent
bdf529c545
commit
047e6d45fa
@@ -6,10 +6,18 @@ import (
|
||||
"fmt"
|
||||
)
|
||||
|
||||
// ContinueToken represents a pagination token for list operations.
|
||||
type ContinueToken struct {
|
||||
StartOffset int64 `json:"o"`
|
||||
// Namespace is the namespace to continue from. Only set for cross-namespace list queries.
|
||||
Namespace string `json:"ns,omitempty"`
|
||||
// Name is the name to continue from. Required for list resources, empty for list history.
|
||||
Name string `json:"n,omitempty"`
|
||||
// ResourceVersion is the resource version for pagination.
|
||||
// For list resources: the RV the list was performed at.
|
||||
// For list history: the last seen RV for pagination.
|
||||
ResourceVersion int64 `json:"v"`
|
||||
SortAscending bool `json:"s"`
|
||||
// SortAscending indicates the sort order (used by list history).
|
||||
SortAscending bool `json:"s,omitempty"`
|
||||
}
|
||||
|
||||
func (c ContinueToken) String() string {
|
||||
|
||||
@@ -4,13 +4,54 @@ import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestContinueToken(t *testing.T) {
|
||||
token := &ContinueToken{
|
||||
ResourceVersion: 100,
|
||||
StartOffset: 50,
|
||||
SortAscending: false,
|
||||
}
|
||||
assert.Equal(t, "eyJvIjo1MCwidiI6MTAwLCJzIjpmYWxzZX0=", token.String())
|
||||
t.Run("round-trip with namespace (cross-namespace query)", func(t *testing.T) {
|
||||
original := ContinueToken{
|
||||
Namespace: "my-namespace",
|
||||
Name: "my-resource",
|
||||
ResourceVersion: 200,
|
||||
}
|
||||
decoded, err := GetContinueToken(original.String())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, original.Namespace, decoded.Namespace)
|
||||
assert.Equal(t, original.Name, decoded.Name)
|
||||
assert.Equal(t, original.ResourceVersion, decoded.ResourceVersion)
|
||||
})
|
||||
|
||||
t.Run("round-trip without namespace (single-namespace query)", func(t *testing.T) {
|
||||
original := ContinueToken{
|
||||
Name: "test-resource",
|
||||
ResourceVersion: 100,
|
||||
}
|
||||
decoded, err := GetContinueToken(original.String())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "", decoded.Namespace)
|
||||
assert.Equal(t, original.Name, decoded.Name)
|
||||
assert.Equal(t, original.ResourceVersion, decoded.ResourceVersion)
|
||||
})
|
||||
|
||||
t.Run("history token (no name, uses ResourceVersion for pagination)", func(t *testing.T) {
|
||||
original := ContinueToken{
|
||||
ResourceVersion: 500,
|
||||
SortAscending: true,
|
||||
}
|
||||
decoded, err := GetContinueToken(original.String())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "", decoded.Name)
|
||||
assert.Equal(t, int64(500), decoded.ResourceVersion)
|
||||
assert.True(t, decoded.SortAscending)
|
||||
})
|
||||
|
||||
t.Run("rejects invalid base64", func(t *testing.T) {
|
||||
_, err := GetContinueToken("not-valid-base64!")
|
||||
assert.Error(t, err)
|
||||
})
|
||||
|
||||
t.Run("rejects invalid json", func(t *testing.T) {
|
||||
_, err := GetContinueToken("bm90LWpzb24=") // "not-json" in base64
|
||||
assert.Error(t, err)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -303,7 +303,7 @@ func (d *dataStore) GetResourceKeyAtRevision(ctx context.Context, key GetRequest
|
||||
|
||||
listKey := ListRequestKey(key)
|
||||
|
||||
iter := d.ListResourceKeysAtRevision(ctx, listKey, rv)
|
||||
iter := d.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv})
|
||||
for dataKey, err := range iter {
|
||||
if err != nil {
|
||||
return DataKey{}, err
|
||||
@@ -313,32 +313,77 @@ func (d *dataStore) GetResourceKeyAtRevision(ctx context.Context, key GetRequest
|
||||
return DataKey{}, ErrNotFound
|
||||
}
|
||||
|
||||
type ListRequestOptions struct {
|
||||
// Key defines the range to query (Group/Resource/Namespace/Name prefix).
|
||||
Key ListRequestKey
|
||||
// ContinueNamespace is the namespace to continue from.
|
||||
// Only used when Key.Namespace is empty (cross-namespace query).
|
||||
ContinueNamespace string
|
||||
// ContinueName is the name to continue from.
|
||||
ContinueName string
|
||||
ResourceVersion int64
|
||||
}
|
||||
|
||||
// Validate checks that the ListRequestOptions are valid.
|
||||
func (o ListRequestOptions) Validate() error {
|
||||
if err := o.Key.Validate(); err != nil {
|
||||
return fmt.Errorf("invalid list request key: %w", err)
|
||||
}
|
||||
// ContinueNamespace is only valid for cross-namespace queries
|
||||
if o.ContinueNamespace != "" && o.Key.Namespace != "" {
|
||||
return fmt.Errorf("continue namespace %q not allowed when request namespace is set to %q", o.ContinueNamespace, o.Key.Namespace)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListLatestResourceKeys returns an iterator over the data keys for the latest versions of resources.
|
||||
// Only returns keys for resources that are not deleted.
|
||||
func (d *dataStore) ListLatestResourceKeys(ctx context.Context, key ListRequestKey) iter.Seq2[DataKey, error] {
|
||||
return d.ListResourceKeysAtRevision(ctx, key, 0)
|
||||
return d.ListResourceKeysAtRevision(ctx, ListRequestOptions{
|
||||
Key: key,
|
||||
})
|
||||
}
|
||||
|
||||
// ListResourceKeysAtRevision returns an iterator over data keys for resources at a specific revision.
|
||||
// If rv is 0, it returns the latest versions. Only returns keys for resources that are not deleted at the given revision.
|
||||
func (d *dataStore) ListResourceKeysAtRevision(ctx context.Context, key ListRequestKey, rv int64) iter.Seq2[DataKey, error] {
|
||||
if err := key.Validate(); err != nil {
|
||||
func (d *dataStore) ListResourceKeysAtRevision(ctx context.Context, options ListRequestOptions) iter.Seq2[DataKey, error] {
|
||||
if err := options.Validate(); err != nil {
|
||||
return func(yield func(DataKey, error) bool) {
|
||||
yield(DataKey{}, fmt.Errorf("invalid list request key: %w", err))
|
||||
yield(DataKey{}, err)
|
||||
}
|
||||
}
|
||||
|
||||
rv := options.ResourceVersion
|
||||
prefix := options.Key.Prefix()
|
||||
|
||||
startKey := prefix
|
||||
if options.ContinueName != "" {
|
||||
// Build the start key from the continue position
|
||||
continueKey := ListRequestKey{
|
||||
Group: options.Key.Group,
|
||||
Resource: options.Key.Resource,
|
||||
Namespace: options.Key.Namespace,
|
||||
Name: options.ContinueName,
|
||||
}
|
||||
// For cross-namespace queries, use the continue namespace
|
||||
if options.Key.Namespace == "" && options.ContinueNamespace != "" {
|
||||
continueKey.Namespace = options.ContinueNamespace
|
||||
}
|
||||
startKey = continueKey.Prefix()
|
||||
}
|
||||
|
||||
listOptions := ListOptions{
|
||||
StartKey: startKey,
|
||||
EndKey: PrefixRangeEnd(prefix),
|
||||
Sort: SortOrderAsc,
|
||||
}
|
||||
|
||||
if rv == 0 {
|
||||
rv = math.MaxInt64
|
||||
}
|
||||
|
||||
prefix := key.Prefix()
|
||||
// List all keys in the prefix.
|
||||
iter := d.kv.Keys(ctx, dataSection, ListOptions{
|
||||
StartKey: prefix,
|
||||
EndKey: PrefixRangeEnd(prefix),
|
||||
Sort: SortOrderAsc,
|
||||
})
|
||||
iter := d.kv.Keys(ctx, dataSection, listOptions)
|
||||
|
||||
return func(yield func(DataKey, error) bool) {
|
||||
var candidateKey *DataKey // The current candidate key we are iterating over
|
||||
|
||||
@@ -2022,7 +2022,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
|
||||
t.Run("list at revision rv1 - should return only resource1 initial version", func(t *testing.T) {
|
||||
resultKeys := make([]DataKey, 0, 2)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv1) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv1}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2035,7 +2035,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
|
||||
t.Run("list at revision rv2 - should return resource1, resource2 and resource4", func(t *testing.T) {
|
||||
resultKeys := make([]DataKey, 0, 3)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv2) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv2}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2053,7 +2053,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
|
||||
t.Run("list at revision rv3 - should return resource1, resource2 and resource4", func(t *testing.T) {
|
||||
resultKeys := make([]DataKey, 0, 3)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv3) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv3}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2074,7 +2074,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
|
||||
t.Run("list at revision rv4 - should return all resources", func(t *testing.T) {
|
||||
resultKeys := make([]DataKey, 0, 4)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv4) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv4}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2093,7 +2093,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
|
||||
t.Run("list at revision rv5 - should exclude deleted resource4", func(t *testing.T) {
|
||||
resultKeys := make([]DataKey, 0, 3)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv5) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv5}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2119,7 +2119,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
}
|
||||
|
||||
resultKeys := make([]DataKey, 0, 2)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, specificListKey, rv3) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: specificListKey, ResourceVersion: rv3}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2132,7 +2132,7 @@ func TestDataStore_ListResourceKeysAtRevision(t *testing.T) {
|
||||
|
||||
t.Run("list at revision 0 should use MaxInt64", func(t *testing.T) {
|
||||
resultKeys := make([]DataKey, 0, 4)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, 0) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: 0}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2185,7 +2185,7 @@ func TestDataStore_ListResourceKeysAtRevision_ValidationErrors(t *testing.T) {
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
for _, err := range ds.ListResourceKeysAtRevision(ctx, tt.key, 0) {
|
||||
for _, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: tt.key, ResourceVersion: 0}) {
|
||||
require.Error(t, err)
|
||||
return
|
||||
}
|
||||
@@ -2204,7 +2204,7 @@ func TestDataStore_ListResourceKeysAtRevision_EmptyResults(t *testing.T) {
|
||||
}
|
||||
|
||||
resultKeys := make([]DataKey, 0, 1)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, 0) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: 0}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
@@ -2238,7 +2238,7 @@ func TestDataStore_ListResourceKeysAtRevision_ResourcesNewerThanRevision(t *test
|
||||
|
||||
// List at a revision before the resource was created
|
||||
resultKeys := make([]DataKey, 0, 1)
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv-1000) {
|
||||
for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, ListRequestOptions{Key: listKey, ResourceVersion: rv - 1000}) {
|
||||
require.NoError(t, err)
|
||||
resultKeys = append(resultKeys, dataKey)
|
||||
}
|
||||
|
||||
@@ -472,15 +472,30 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
namespace := convertEmptyToClusterNamespace(req.Options.Key.Namespace, k.withExperimentalClusterScope)
|
||||
|
||||
// Parse continue token if provided
|
||||
offset := int64(0)
|
||||
resourceVersion := req.ResourceVersion
|
||||
listOptions := ListRequestOptions{
|
||||
Key: ListRequestKey{
|
||||
Group: req.Options.Key.Group,
|
||||
Resource: req.Options.Key.Resource,
|
||||
Namespace: namespace,
|
||||
Name: req.Options.Key.Name,
|
||||
},
|
||||
ResourceVersion: req.ResourceVersion,
|
||||
}
|
||||
|
||||
if req.NextPageToken != "" {
|
||||
token, err := GetContinueToken(req.NextPageToken)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("invalid continue token: %w", err)
|
||||
}
|
||||
offset = token.StartOffset
|
||||
resourceVersion = token.ResourceVersion
|
||||
if token.Name == "" {
|
||||
return 0, fmt.Errorf("invalid continue token: name is required for list resources")
|
||||
}
|
||||
// Only use token namespace for cross-namespace queries (when request namespace is empty)
|
||||
if req.Options.Key.Namespace == "" {
|
||||
listOptions.ContinueNamespace = token.Namespace
|
||||
}
|
||||
listOptions.ContinueName = token.Name
|
||||
listOptions.ResourceVersion = token.ResourceVersion
|
||||
}
|
||||
|
||||
// We set the listRV to the last event resource version.
|
||||
@@ -492,27 +507,17 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
return 0, fmt.Errorf("failed to fetch last event: %w", err)
|
||||
}
|
||||
|
||||
if resourceVersion > 0 {
|
||||
listRV = resourceVersion
|
||||
if listOptions.ResourceVersion > 0 {
|
||||
listRV = listOptions.ResourceVersion
|
||||
}
|
||||
|
||||
// Fetch the latest objects
|
||||
keys := make([]DataKey, 0, min(defaultListBufferSize, req.Limit+1))
|
||||
idx := 0
|
||||
for dataKey, err := range k.dataStore.ListResourceKeysAtRevision(ctx, ListRequestKey{
|
||||
Group: req.Options.Key.Group,
|
||||
Resource: req.Options.Key.Resource,
|
||||
Namespace: namespace,
|
||||
Name: req.Options.Key.Name,
|
||||
}, resourceVersion) {
|
||||
for dataKey, err := range k.dataStore.ListResourceKeysAtRevision(ctx, listOptions) {
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
// Skip the first offset items. This is not efficient, but it's a simple way to implement it for now.
|
||||
if idx < int(offset) {
|
||||
idx++
|
||||
continue
|
||||
}
|
||||
|
||||
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) {
|
||||
@@ -524,9 +529,9 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
defer stop()
|
||||
|
||||
iter := kvListIterator{
|
||||
listRV: listRV,
|
||||
offset: offset,
|
||||
next: next,
|
||||
listRV: listRV,
|
||||
isCrossNamespace: req.Options.Key.Namespace == "",
|
||||
next: next,
|
||||
}
|
||||
err := cb(&iter)
|
||||
if err != nil {
|
||||
@@ -538,38 +543,44 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
|
||||
// kvListIterator implements ListIterator for KV storage
|
||||
type kvListIterator struct {
|
||||
listRV int64
|
||||
offset int64
|
||||
listRV int64
|
||||
isCrossNamespace bool
|
||||
|
||||
// pull-style iterator
|
||||
next func() (DataObj, error, bool)
|
||||
|
||||
// current item state
|
||||
currentDataObj *DataObj
|
||||
started bool
|
||||
currentDataObj DataObj
|
||||
value []byte
|
||||
err error
|
||||
nextDataObj DataObj
|
||||
nextErr error
|
||||
hasMore bool
|
||||
}
|
||||
|
||||
func (i *kvListIterator) Next() bool {
|
||||
// Pull next item from the iterator
|
||||
dataObj, err, ok := i.next()
|
||||
if !ok {
|
||||
return false
|
||||
if !i.started {
|
||||
i.started = true
|
||||
|
||||
i.nextDataObj, i.nextErr, i.hasMore = i.next()
|
||||
}
|
||||
if err != nil {
|
||||
i.err = err
|
||||
|
||||
if !i.hasMore {
|
||||
return false
|
||||
}
|
||||
|
||||
i.currentDataObj = &dataObj
|
||||
|
||||
i.value, err = readAndClose(dataObj.Value)
|
||||
if err != nil {
|
||||
i.err = err
|
||||
i.currentDataObj, i.err = i.nextDataObj, i.nextErr
|
||||
if i.err != nil {
|
||||
return false
|
||||
}
|
||||
|
||||
i.offset++
|
||||
i.value, i.err = readAndClose(i.currentDataObj.Value)
|
||||
if i.err != nil {
|
||||
return false
|
||||
}
|
||||
|
||||
i.nextDataObj, i.nextErr, i.hasMore = i.next()
|
||||
|
||||
return true
|
||||
}
|
||||
@@ -579,38 +590,31 @@ func (i *kvListIterator) Error() error {
|
||||
}
|
||||
|
||||
func (i *kvListIterator) ContinueToken() string {
|
||||
return ContinueToken{
|
||||
StartOffset: i.offset,
|
||||
token := ContinueToken{
|
||||
Name: i.nextDataObj.Key.Name,
|
||||
ResourceVersion: i.listRV,
|
||||
}.String()
|
||||
}
|
||||
// Only store namespace in token for cross-namespace queries
|
||||
if i.isCrossNamespace {
|
||||
token.Namespace = i.nextDataObj.Key.Namespace
|
||||
}
|
||||
return token.String()
|
||||
}
|
||||
|
||||
func (i *kvListIterator) ResourceVersion() int64 {
|
||||
if i.currentDataObj != nil {
|
||||
return i.currentDataObj.Key.ResourceVersion
|
||||
}
|
||||
return 0
|
||||
return i.currentDataObj.Key.ResourceVersion
|
||||
}
|
||||
|
||||
func (i *kvListIterator) Namespace() string {
|
||||
if i.currentDataObj != nil {
|
||||
return convertClusterNamespaceToEmpty(i.currentDataObj.Key.Namespace)
|
||||
}
|
||||
return ""
|
||||
return convertClusterNamespaceToEmpty(i.currentDataObj.Key.Namespace)
|
||||
}
|
||||
|
||||
func (i *kvListIterator) Name() string {
|
||||
if i.currentDataObj != nil {
|
||||
return i.currentDataObj.Key.Name
|
||||
}
|
||||
return ""
|
||||
return i.currentDataObj.Key.Name
|
||||
}
|
||||
|
||||
func (i *kvListIterator) Folder() string {
|
||||
if i.currentDataObj != nil {
|
||||
return i.currentDataObj.Key.Folder
|
||||
}
|
||||
return ""
|
||||
return i.currentDataObj.Key.Folder
|
||||
}
|
||||
|
||||
func (i *kvListIterator) Value() []byte {
|
||||
@@ -1129,7 +1133,6 @@ func (i *kvHistoryIterator) ContinueToken() string {
|
||||
}
|
||||
rv := i.currentDataObj.Key.ResourceVersion
|
||||
token := ContinueToken{
|
||||
StartOffset: rv,
|
||||
ResourceVersion: rv,
|
||||
SortAscending: i.sortAscending,
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user