From c251ebf4d5ecbb92afbde63516c62a0a707ce95a Mon Sep 17 00:00:00 2001 From: Georges Chaudy Date: Tue, 16 Sep 2025 12:16:29 +0200 Subject: [PATCH] kvstore: merge the metadata store into the datastore (#110334) * migrate eventstore to datastore * Add folder to event key * lint * lint * lint * lint * remove foundkye * refactor the Keys methods to move the Sort outside the ListKey method * remove bad import * fix missing params * lint * fix test * perf improvement --- pkg/storage/unified/resource/datastore.go | 485 ++++- .../unified/resource/datastore_test.go | 1655 +++++++++++++++-- pkg/storage/unified/resource/eventstore.go | 14 +- .../unified/resource/eventstore_test.go | 22 +- pkg/storage/unified/resource/metadata.go | 391 ---- pkg/storage/unified/resource/metadata_test.go | 1355 -------------- .../unified/resource/storage_backend.go | 149 +- .../unified/resource/storage_backend_test.go | 30 +- 8 files changed, 1943 insertions(+), 2158 deletions(-) delete mode 100644 pkg/storage/unified/resource/metadata.go delete mode 100644 pkg/storage/unified/resource/metadata_test.go diff --git a/pkg/storage/unified/resource/datastore.go b/pkg/storage/unified/resource/datastore.go index f939f18f1df..e45e1b74101 100644 --- a/pkg/storage/unified/resource/datastore.go +++ b/pkg/storage/unified/resource/datastore.go @@ -5,23 +5,31 @@ import ( "fmt" "io" "iter" + "math" "regexp" "strconv" "strings" + "time" + + gocache "github.com/patrickmn/go-cache" ) const ( dataSection = "unified/data" + // cache + groupResourcesCacheKey = "group-resources" ) // dataStore is a data store that uses a KV store to store data. type dataStore struct { - kv KV + kv KV + cache *gocache.Cache } func newDataStore(kv KV) *dataStore { return &dataStore{ - kv: kv, + kv: kv, + cache: gocache.New(time.Hour, 10*time.Minute), // 1 hour expiration, 10 minute cleanup } } @@ -37,6 +45,13 @@ type DataKey struct { Name string ResourceVersion int64 Action DataAction + Folder string +} + +// GroupResource represents a unique group/resource combination +type GroupResource struct { + Group string + Resource string } var ( @@ -46,40 +61,34 @@ var ( ) func (k DataKey) String() string { - return fmt.Sprintf("%s/%s/%s/%s/%d~%s", k.Namespace, k.Group, k.Resource, k.Name, k.ResourceVersion, k.Action) + return fmt.Sprintf("%s/%s/%s/%s/%d~%s~%s", k.Group, k.Resource, k.Namespace, k.Name, k.ResourceVersion, k.Action, k.Folder) } func (k DataKey) Equals(other DataKey) bool { - return k.Namespace == other.Namespace && k.Group == other.Group && k.Resource == other.Resource && k.Name == other.Name && k.ResourceVersion == other.ResourceVersion && k.Action == other.Action + return k.Group == other.Group && k.Resource == other.Resource && k.Namespace == other.Namespace && k.Name == other.Name && k.ResourceVersion == other.ResourceVersion && k.Action == other.Action && k.Folder == other.Folder } func (k DataKey) Validate() error { - if k.Namespace == "" { - if k.Group != "" || k.Resource != "" || k.Name != "" { - return fmt.Errorf("namespace is required when group, resource, or name are provided") - } - return fmt.Errorf("namespace cannot be empty") - } if k.Group == "" { - if k.Resource != "" || k.Name != "" { - return fmt.Errorf("group is required when resource or name are provided") - } - return fmt.Errorf("group cannot be empty") + return fmt.Errorf("group is required") } if k.Resource == "" { - if k.Name != "" { - return fmt.Errorf("resource is required when name is provided") - } - return fmt.Errorf("resource cannot be empty") + return fmt.Errorf("resource is required") + } + if k.Namespace == "" { + return fmt.Errorf("namespace is required") } if k.Name == "" { - return fmt.Errorf("name cannot be empty") + return fmt.Errorf("name is required") + } + if k.ResourceVersion <= 0 { + return fmt.Errorf("resource version must be positive") } if k.Action == "" { - return fmt.Errorf("action cannot be empty") + return fmt.Errorf("action is required") } - // Validate each field against the naming rules + // Validate naming conventions for all required fields if !validNameRegex.MatchString(k.Namespace) { return fmt.Errorf("namespace '%s' is invalid", k.Namespace) } @@ -93,6 +102,12 @@ func (k DataKey) Validate() error { return fmt.Errorf("name '%s' is invalid", k.Name) } + // Validate folder field if provided (optional field) + if k.Folder != "" && !validNameRegex.MatchString(k.Folder) { + return fmt.Errorf("folder '%s' is invalid", k.Folder) + } + + // Validate action is one of the valid values switch k.Action { case DataActionCreated, DataActionUpdated, DataActionDeleted: return nil @@ -102,47 +117,23 @@ func (k DataKey) Validate() error { } type ListRequestKey struct { - Namespace string Group string Resource string - Name string - Sort SortOrder + Namespace string + Name string // optional for listing multiple resources } func (k ListRequestKey) Validate() error { - // Check hierarchical validation - if a field is empty, more specific fields should also be empty - if k.Namespace == "" { - if k.Group != "" || k.Resource != "" || k.Name != "" { - return fmt.Errorf("namespace is required when group, resource, or name are provided") - } - return nil // Empty namespace is allowed for ListRequestKey - } if k.Group == "" { - if k.Resource != "" || k.Name != "" { - return fmt.Errorf("group is required when resource or name are provided") - } - // Only validate namespace if it's provided - if !validNameRegex.MatchString(k.Namespace) { - return fmt.Errorf("namespace '%s' is invalid", k.Namespace) - } - return nil + return fmt.Errorf("group is required") } if k.Resource == "" { - if k.Name != "" { - return fmt.Errorf("resource is required when name is provided") - } - // Validate namespace and group if they're provided - if !validNameRegex.MatchString(k.Namespace) { - return fmt.Errorf("namespace '%s' is invalid", k.Namespace) - } - if !validNameRegex.MatchString(k.Group) { - return fmt.Errorf("group '%s' is invalid", k.Group) - } - return nil + return fmt.Errorf("resource is required") } - - // All fields are provided, validate each one - if !validNameRegex.MatchString(k.Namespace) { + if k.Namespace == "" && k.Name != "" { + return fmt.Errorf("name must be empty when namespace is empty") + } + if k.Namespace != "" && !validNameRegex.MatchString(k.Namespace) { return fmt.Errorf("namespace '%s' is invalid", k.Namespace) } if !validNameRegex.MatchString(k.Group) { @@ -154,24 +145,62 @@ func (k ListRequestKey) Validate() error { if k.Name != "" && !validNameRegex.MatchString(k.Name) { return fmt.Errorf("name '%s' is invalid", k.Name) } - return nil } func (k ListRequestKey) Prefix() string { if k.Namespace == "" { - return "" - } - if k.Group == "" { - return fmt.Sprintf("%s/", k.Namespace) - } - if k.Resource == "" { - return fmt.Sprintf("%s/%s/", k.Namespace, k.Group) + return fmt.Sprintf("%s/%s/", k.Group, k.Resource) } if k.Name == "" { - return fmt.Sprintf("%s/%s/%s/", k.Namespace, k.Group, k.Resource) + return fmt.Sprintf("%s/%s/%s/", k.Group, k.Resource, k.Namespace) } - return fmt.Sprintf("%s/%s/%s/%s/", k.Namespace, k.Group, k.Resource, k.Name) + return fmt.Sprintf("%s/%s/%s/%s/", k.Group, k.Resource, k.Namespace, k.Name) +} + +// GetRequestKey is used for getting a specific data object by latest version +type GetRequestKey struct { + Group string + Resource string + Namespace string + Name string +} + +// Validate validates the get request key +func (k GetRequestKey) Validate() error { + if k.Group == "" { + return fmt.Errorf("group is required") + } + if k.Resource == "" { + return fmt.Errorf("resource is required") + } + if k.Namespace == "" { + return fmt.Errorf("namespace is required") + } + if k.Name == "" { + return fmt.Errorf("name is required") + } + + // Validate naming conventions + if !validNameRegex.MatchString(k.Namespace) { + return fmt.Errorf("namespace '%s' is invalid", k.Namespace) + } + if !validNameRegex.MatchString(k.Group) { + return fmt.Errorf("group '%s' is invalid", k.Group) + } + if !validNameRegex.MatchString(k.Resource) { + return fmt.Errorf("resource '%s' is invalid", k.Resource) + } + if !validNameRegex.MatchString(k.Name) { + return fmt.Errorf("name '%s' is invalid", k.Name) + } + + return nil +} + +// Prefix returns the prefix for getting a specific data object +func (k GetRequestKey) Prefix() string { + return fmt.Sprintf("%s/%s/%s/%s/", k.Group, k.Resource, k.Namespace, k.Name) } type DataAction string @@ -183,19 +212,18 @@ const ( ) // Keys returns all keys for a given key by iterating through the KV store -func (d *dataStore) Keys(ctx context.Context, key ListRequestKey) iter.Seq2[DataKey, error] { +func (d *dataStore) Keys(ctx context.Context, key ListRequestKey, sort SortOrder) iter.Seq2[DataKey, error] { if err := key.Validate(); err != nil { return func(yield func(DataKey, error) bool) { yield(DataKey{}, err) } } - prefix := key.Prefix() return func(yield func(DataKey, error) bool) { for k, err := range d.kv.Keys(ctx, dataSection, ListOptions{ StartKey: prefix, EndKey: PrefixRangeEnd(prefix), - Sort: key.Sort, + Sort: sort, }) { if err != nil { yield(DataKey{}, err) @@ -236,6 +264,123 @@ func (d *dataStore) LastResourceVersion(ctx context.Context, key ListRequestKey) return DataKey{}, ErrNotFound } +// GetLatestResourceKey retrieves the data key for the latest version of a resource. +// Returns the key with the highest resource version that is not deleted. +func (d *dataStore) GetLatestResourceKey(ctx context.Context, key GetRequestKey) (DataKey, error) { + return d.GetResourceKeyAtRevision(ctx, key, 0) +} + +// GetResourceKeyAtRevision retrieves the data key for a resource at a specific revision. +// If rv is 0, it returns the latest version. Returns the highest version <= rv that is not deleted. +func (d *dataStore) GetResourceKeyAtRevision(ctx context.Context, key GetRequestKey, rv int64) (DataKey, error) { + if err := key.Validate(); err != nil { + return DataKey{}, fmt.Errorf("invalid get request key: %w", err) + } + + if rv == 0 { + rv = math.MaxInt64 + } + + listKey := ListRequestKey(key) + + iter := d.ListResourceKeysAtRevision(ctx, listKey, rv) + for dataKey, err := range iter { + if err != nil { + return DataKey{}, err + } + return dataKey, nil + } + return DataKey{}, ErrNotFound +} + +// 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) +} + +// 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 { + return func(yield func(DataKey, error) bool) { + yield(DataKey{}, fmt.Errorf("invalid list request key: %w", err)) + } + } + + 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, + }) + + return func(yield func(DataKey, error) bool) { + var candidateKey *DataKey // The current candidate key we are iterating over + + // yieldCandidate is a helper function to yield results. + // Won't yield if the resource was last deleted. + yieldCandidate := func() bool { + if candidateKey.Action == DataActionDeleted { + // Skip because the resource was last deleted. + return true + } + return yield(*candidateKey, nil) + } + + for key, err := range iter { + if err != nil { + yield(DataKey{}, err) + return + } + + dataKey, err := ParseKey(key) + if err != nil { + yield(DataKey{}, err) + return + } + + if candidateKey == nil { + // Skip until we have our first candidate + if dataKey.ResourceVersion <= rv { + // New candidate found. + candidateKey = &dataKey + } + continue + } + // Should yield if either: + // - We reached the next resource. + // - We reached a resource version greater than the target resource version. + if !dataKey.SameResource(*candidateKey) || dataKey.ResourceVersion > rv { + if !yieldCandidate() { + return + } + // If we moved to a different resource and the resource version matches, make it the new candidate + if !dataKey.SameResource(*candidateKey) && dataKey.ResourceVersion <= rv { + candidateKey = &dataKey + } else { + // If we moved to a different resource and the resource version does not match, reset the candidate + candidateKey = nil + } + } else { + // Update candidate to the current key (same resource, valid version) + candidateKey = &dataKey + } + } + if candidateKey != nil { + // Yield the last selected object + if !yieldCandidate() { + return + } + } + } +} + func (d *dataStore) Get(ctx context.Context, key DataKey) (io.ReadCloser, error) { if err := key.Validate(); err != nil { return nil, fmt.Errorf("invalid data key: %w", err) @@ -276,20 +421,212 @@ func ParseKey(key string) (DataKey, error) { if len(parts) != 5 { return DataKey{}, fmt.Errorf("invalid key: %s", key) } - uidActionParts := strings.Split(parts[4], "~") - if len(uidActionParts) != 2 { + rvActionFolderParts := strings.Split(parts[4], "~") + if len(rvActionFolderParts) != 3 { return DataKey{}, fmt.Errorf("invalid key: %s", key) } - rv, err := strconv.ParseInt(uidActionParts[0], 10, 64) + rv, err := strconv.ParseInt(rvActionFolderParts[0], 10, 64) if err != nil { - return DataKey{}, fmt.Errorf("invalid resource version: %s", uidActionParts[0]) + return DataKey{}, fmt.Errorf("invalid resource version '%s' in key %s: %w", rvActionFolderParts[0], key, err) } return DataKey{ - Namespace: parts[0], - Group: parts[1], - Resource: parts[2], + Group: parts[0], + Resource: parts[1], + Namespace: parts[2], Name: parts[3], ResourceVersion: rv, - Action: DataAction(uidActionParts[1]), + Action: DataAction(rvActionFolderParts[1]), + Folder: rvActionFolderParts[2], }, nil } + +// SameResource checks if this key represents the same resource as another key. +// It compares the identifying fields: Group, Resource, Namespace, and Name. +// ResourceVersion, Action, and Folder are ignored as they don't identify the resource itself. +func (k DataKey) SameResource(other DataKey) bool { + return k.Group == other.Group && + k.Resource == other.Resource && + k.Namespace == other.Namespace && + k.Name == other.Name +} + +// GetResourceStats returns resource stats within the data store by first discovering +// all group/resource combinations, then issuing targeted list operations for each one. +// If namespace is provided, only keys matching that namespace are considered. +func (d *dataStore) GetResourceStats(ctx context.Context, namespace string, minCount int) ([]ResourceStats, error) { + // First, get all unique group/resource combinations in the store + groupResources, err := d.getGroupResources(ctx) + if err != nil { + return nil, fmt.Errorf("failed to get group resources: %w", err) + } + + var stats []ResourceStats + + // Process each group/resource combination + for _, groupResource := range groupResources { + groupStats, err := d.processGroupResourceStats(ctx, groupResource, namespace, minCount) + if err != nil { + return nil, fmt.Errorf("failed to process stats for %s/%s: %w", groupResource.Group, groupResource.Resource, err) + } + stats = append(stats, groupStats...) + } + + return stats, nil +} + +// processGroupResourceStats processes stats for a specific group/resource combination +func (d *dataStore) processGroupResourceStats(ctx context.Context, groupResource GroupResource, namespace string, minCount int) ([]ResourceStats, error) { + // Use ListRequestKey to construct the appropriate prefix + listKey := ListRequestKey{ + Group: groupResource.Group, + Resource: groupResource.Resource, + Namespace: namespace, // Empty string if not specified, which will list all namespaces + } + + // Maps to track counts per namespace for this group/resource + namespaceCounts := make(map[string]int64) // namespace -> count of existing resources + namespaceVersions := make(map[string]int64) // namespace -> latest resource version + + // Track current resource being processed + var currentResourceKey string + var lastDataKey *DataKey + + // Helper function to process the last seen resource + processLastResource := func() { + if lastDataKey != nil { + // Initialize namespace version if not exists + if _, exists := namespaceVersions[lastDataKey.Namespace]; !exists { + namespaceVersions[lastDataKey.Namespace] = 0 + } + + // If resource exists (not deleted), increment the count for this namespace + if lastDataKey.Action != DataActionDeleted { + namespaceCounts[lastDataKey.Namespace]++ + } + + // Update to latest resource version seen + if lastDataKey.ResourceVersion > namespaceVersions[lastDataKey.Namespace] { + namespaceVersions[lastDataKey.Namespace] = lastDataKey.ResourceVersion + } + } + } + + // List all keys using the existing Keys method + for dataKey, err := range d.Keys(ctx, listKey, SortOrderAsc) { + if err != nil { + return nil, err + } + + // Create unique resource identifier (namespace/group/resource/name) + resourceKey := fmt.Sprintf("%s/%s/%s/%s", dataKey.Namespace, dataKey.Group, dataKey.Resource, dataKey.Name) + + // If we've moved to a different resource, process the previous one + if currentResourceKey != "" && resourceKey != currentResourceKey { + processLastResource() + } + + // Update tracking variables for the current resource + currentResourceKey = resourceKey + lastDataKey = &dataKey + } + + // Process the final resource + processLastResource() + + // Convert namespace counts to ResourceStats + stats := make([]ResourceStats, 0, len(namespaceCounts)) + for ns, count := range namespaceCounts { + // Skip if count is below or equal to minimum + if count <= int64(minCount) { + continue + } + + stats = append(stats, ResourceStats{ + NamespacedResource: NamespacedResource{ + Namespace: ns, + Group: groupResource.Group, + Resource: groupResource.Resource, + }, + Count: count, + ResourceVersion: namespaceVersions[ns], + }) + } + + return stats, nil +} + +// getGroupResources returns all unique group/resource combinations in the data store. +// It efficiently discovers these by using the key ordering and PrefixRangeEnd to jump +// between different group/resource prefixes without iterating through all keys. +// Results are cached to improve performance. +func (d *dataStore) getGroupResources(ctx context.Context) ([]GroupResource, error) { + // Check cache first + if cached, found := d.cache.Get(groupResourcesCacheKey); found { + if cachedResults, ok := cached.([]GroupResource); ok { + return cachedResults, nil + } + } + + // Cache miss or invalid data, compute the results + results := make([]GroupResource, 0) + seenGroupResources := make(map[string]bool) // "group/resource" -> seen + + startKey := "" + + for { + // List with limit 1 to get the next key + var foundKey string + + for key, err := range d.kv.Keys(ctx, dataSection, ListOptions{ + StartKey: startKey, + Limit: 1, + Sort: SortOrderAsc, + }) { + if err != nil { + return nil, err + } + foundKey = key + break // Only process the first (and only) key + } + + // If no key found, we're done + if foundKey == "" { + break + } + + // Parse the key to extract group and resource + dataKey, err := ParseKey(foundKey) + if err != nil { + return nil, fmt.Errorf("failed to parse key %s: %w", foundKey, err) + } + + // Create the group/resource identifier + groupResourceKey := fmt.Sprintf("%s/%s", dataKey.Group, dataKey.Resource) + + // Add to results if we haven't seen this group/resource combination before + if !seenGroupResources[groupResourceKey] { + seenGroupResources[groupResourceKey] = true + //nolint:staticcheck // SA4010: wrongly assumes that this result of append is never used + results = append(results, GroupResource{ + Group: dataKey.Group, + Resource: dataKey.Resource, + }) + } + + // Compute the next starting point by finding the end of this group/resource prefix + groupResourcePrefix := fmt.Sprintf("%s/%s/", dataKey.Group, dataKey.Resource) + nextStartKey := PrefixRangeEnd(groupResourcePrefix) + + // If we've reached the end of the key space, we're done + if nextStartKey == "" { + break + } + + startKey = nextStartKey + } + + // Cache the results using the default expiration (1 hour) + d.cache.Set(groupResourcesCacheKey, results, gocache.DefaultExpiration) + + return results, nil +} diff --git a/pkg/storage/unified/resource/datastore_test.go b/pkg/storage/unified/resource/datastore_test.go index d82396f1ee7..1bdb97fb949 100644 --- a/pkg/storage/unified/resource/datastore_test.go +++ b/pkg/storage/unified/resource/datastore_test.go @@ -33,37 +33,40 @@ func TestDataKey_String(t *testing.T) { { name: "created key", key: DataKey{ - Namespace: "test-namespace", Group: "test-group", Resource: "test-resource", + Namespace: "test-namespace", Name: "test-name", ResourceVersion: rv, Action: DataActionCreated, + Folder: "test-folder", }, - expected: "test-namespace/test-group/test-resource/test-name/1934555792099250176~created", + expected: "test-group/test-resource/test-namespace/test-name/1934555792099250176~created~test-folder", }, { name: "updated key", key: DataKey{ - Namespace: "test-namespace", Group: "test-group", Resource: "test-resource", + Namespace: "test-namespace", Name: "test-name", ResourceVersion: rv, Action: DataActionUpdated, + Folder: "test-folder", }, - expected: "test-namespace/test-group/test-resource/test-name/1934555792099250176~updated", + expected: "test-group/test-resource/test-namespace/test-name/1934555792099250176~updated~test-folder", }, { name: "deleted key", key: DataKey{ - Namespace: "test-namespace", Group: "test-group", Resource: "test-resource", + Namespace: "test-namespace", Name: "test-name", ResourceVersion: rv, Action: DataActionDeleted, + Folder: "test-folder", }, - expected: "test-namespace/test-group/test-resource/test-name/1934555792099250176~deleted", + expected: "test-group/test-resource/test-namespace/test-name/1934555792099250176~deleted~test-folder", }, } @@ -78,15 +81,6 @@ func TestDataKey_String(t *testing.T) { func TestDataKey_Validate(t *testing.T) { rv := int64(1234567890) - validKey := DataKey{ - Namespace: "test-namespace", - Group: "test-group", - Resource: "test-resource", - Name: "test-name", - ResourceVersion: rv, - Action: DataActionCreated, - } - tests := []struct { name string key DataKey @@ -94,8 +88,15 @@ func TestDataKey_Validate(t *testing.T) { errorMsg string }{ { - name: "valid key with created action", - key: validKey, + name: "valid key with created action", + key: DataKey{ + Namespace: "test-namespace", + Group: "test-group", + Resource: "test-resource", + Name: "test-name", + ResourceVersion: rv, + Action: DataActionCreated, + }, expectError: false, }, { @@ -170,7 +171,7 @@ func TestDataKey_Validate(t *testing.T) { Action: DataActionCreated, }, expectError: true, - errorMsg: "namespace is required when group, resource, or name are provided", + errorMsg: "namespace is required", }, { name: "invalid - empty group", @@ -183,7 +184,7 @@ func TestDataKey_Validate(t *testing.T) { Action: DataActionCreated, }, expectError: true, - errorMsg: "group is required when resource or name are provided", + errorMsg: "group is required", }, { name: "invalid - empty resource", @@ -196,7 +197,7 @@ func TestDataKey_Validate(t *testing.T) { Action: DataActionCreated, }, expectError: true, - errorMsg: "resource is required when name is provided", + errorMsg: "resource is required", }, { name: "invalid - empty name", @@ -209,7 +210,7 @@ func TestDataKey_Validate(t *testing.T) { Action: DataActionCreated, }, expectError: true, - errorMsg: "name cannot be empty", + errorMsg: "name is required", }, { name: "invalid - empty action", @@ -222,7 +223,7 @@ func TestDataKey_Validate(t *testing.T) { Action: "", }, expectError: true, - errorMsg: "action cannot be empty", + errorMsg: "action is required", }, { name: "invalid - all fields empty", @@ -235,7 +236,7 @@ func TestDataKey_Validate(t *testing.T) { Action: "", }, expectError: true, - errorMsg: "namespace cannot be empty", + errorMsg: "group is required", }, // Invalid cases - uppercase characters { @@ -438,7 +439,7 @@ func TestParseKey(t *testing.T) { }{ { name: "valid normal key", - key: "test-namespace/test-group/test-resource/test-name/" + rv.String() + "~created", + key: "test-group/test-resource/test-namespace/test-name/" + rv.String() + "~created~team-folder", expected: DataKey{ Namespace: "test-namespace", Group: "test-group", @@ -446,11 +447,12 @@ func TestParseKey(t *testing.T) { Name: "test-name", ResourceVersion: rv.Int64(), Action: DataActionCreated, + Folder: "team-folder", }, }, { name: "valid deleted key", - key: "test-namespace/test-group/test-resource/test-name/" + rv.String() + "~deleted", + key: "test-group/test-resource/test-namespace/test-name/" + rv.String() + "~deleted~team-folder", expected: DataKey{ Namespace: "test-namespace", Group: "test-group", @@ -458,6 +460,7 @@ func TestParseKey(t *testing.T) { Name: "test-name", ResourceVersion: rv.Int64(), Action: DataActionDeleted, + Folder: "team-folder", }, }, { @@ -467,17 +470,12 @@ func TestParseKey(t *testing.T) { }, { name: "invalid key - too many slashes", - key: "test-namespace/test-group/test-resource/test-name/1934555792099250176~created/extra-slash", + key: "test-group/test-resource/test-namespace/test-name/1934555792099250176~created~team-folder/extra-slash", expectError: true, }, { - name: "invalid key - invalid uuid", - key: "test-namespace/test-group/test-resource/test-name/invalid-uuid", - expectError: true, - }, - { - name: "invalid key - too many dashes in uuid part", - key: "test-namespace/test-group/test-resource/test-name/uuid-part-extra-dash", + name: "invalid key - invalid rv", + key: "test-group/test-resource/test-namespace/test-name/invalid-rv~team-folder", expectError: true, }, } @@ -503,9 +501,9 @@ func TestDataStore_Save_And_Get(t *testing.T) { rv := node.Generate() testKey := DataKey{ - Namespace: "test-namespace", Group: "test-group", Resource: "test-resource", + Namespace: "test-namespace", Name: "test-name", ResourceVersion: rv.Int64(), Action: DataActionCreated, @@ -546,9 +544,9 @@ func TestDataStore_Save_And_Get(t *testing.T) { rv := node.Generate() nonExistentKey := DataKey{ - Namespace: "non-existent", Group: "test-group", Resource: "test-resource", + Namespace: "non-existent", Name: "test-name", ResourceVersion: rv.Int64(), Action: DataActionCreated, @@ -567,9 +565,9 @@ func TestDataStore_Delete(t *testing.T) { rv := node.Generate() testKey := DataKey{ - Namespace: "test-namespace", Group: "test-group", Resource: "test-resource", + Namespace: "test-namespace", Name: "test-name", ResourceVersion: rv.Int64(), Action: DataActionCreated, @@ -655,8 +653,8 @@ func TestDataStore_List(t *testing.T) { require.NoError(t, err) // List the data - var results []DataKey - for key, err := range ds.Keys(ctx, resourceKey) { + results := make([]DataKey, 0, 2) + for key, err := range ds.Keys(ctx, resourceKey, SortOrderAsc) { require.NoError(t, err) results = append(results, key) } @@ -689,8 +687,8 @@ func TestDataStore_List(t *testing.T) { Name: "empty-name", } - var results []DataKey - for key, err := range ds.Keys(ctx, emptyResourceKey) { + results := make([]DataKey, 0, 1) + for key, err := range ds.Keys(ctx, emptyResourceKey, SortOrderAsc) { require.NoError(t, err) results = append(results, key) } @@ -723,8 +721,8 @@ func TestDataStore_List(t *testing.T) { require.NoError(t, err) // List should include deleted keys - var results []DataKey - for key, err := range ds.Keys(ctx, deletedResourceKey) { + results := make([]DataKey, 0, 2) + for key, err := range ds.Keys(ctx, deletedResourceKey, SortOrderAsc) { require.NoError(t, err) results = append(results, key) } @@ -777,8 +775,8 @@ func TestDataStore_Integration(t *testing.T) { } // List all versions - var results []DataKey - for key, err := range ds.Keys(ctx, resourceKey) { + results := make([]DataKey, 0, 3) + for key, err := range ds.Keys(ctx, resourceKey, SortOrderAsc) { require.NoError(t, err) results = append(results, key) } @@ -804,7 +802,7 @@ func TestDataStore_Integration(t *testing.T) { // List should now have 2 items results = nil - for key, err := range ds.Keys(ctx, resourceKey) { + for key, err := range ds.Keys(ctx, resourceKey, SortOrderAsc) { require.NoError(t, err) results = append(results, key) } @@ -881,8 +879,8 @@ func TestDataStore_Keys(t *testing.T) { require.NoError(t, err) // Get keys - var keys []DataKey - for key, err := range ds.Keys(ctx, resourceKey) { + keys := make([]DataKey, 0, 2) + for key, err := range ds.Keys(ctx, resourceKey, SortOrderAsc) { require.NoError(t, err) keys = append(keys, key) } @@ -910,8 +908,8 @@ func TestDataStore_Keys(t *testing.T) { Name: "empty-name", } - var keys []DataKey - for key, err := range ds.Keys(ctx, emptyResourceKey) { + keys := make([]DataKey, 0, 1) + for key, err := range ds.Keys(ctx, emptyResourceKey, SortOrderAsc) { require.NoError(t, err) keys = append(keys, key) } @@ -941,8 +939,8 @@ func TestDataStore_Keys(t *testing.T) { err := ds.Save(ctx, dataKey4, io.NopCloser(bytes.NewReader([]byte("different-value")))) require.NoError(t, err) - var keys []DataKey - for key, err := range ds.Keys(ctx, partialKey) { + keys := make([]DataKey, 0, 4) + for key, err := range ds.Keys(ctx, partialKey, SortOrderAsc) { require.NoError(t, err) keys = append(keys, key) } @@ -954,62 +952,19 @@ func TestDataStore_Keys(t *testing.T) { require.Contains(t, keys, dataKey4) }) - t.Run("keys with namespace only prefix", func(t *testing.T) { - // Create keys with different groups but same namespace - namespaceOnlyKey := ListRequestKey{ - Namespace: resourceKey.Namespace, - // Group, Resource, Name are empty + t.Run("keys with group and resource only prefix", func(t *testing.T) { + groupAndResourceKey := ListRequestKey{ + Group: "test-group", + Resource: "test-resource", } - rv5 := node.Generate() - dataKey5 := DataKey{ - Namespace: resourceKey.Namespace, - Group: "different-group", - Resource: "different-resource", - Name: "different-name", - ResourceVersion: rv5.Int64(), - Action: DataActionCreated, - } - - err := ds.Save(ctx, dataKey5, io.NopCloser(bytes.NewReader([]byte("namespace-only-value")))) - require.NoError(t, err) - - var keys []DataKey - for key, err := range ds.Keys(ctx, namespaceOnlyKey) { + keys := make([]DataKey, 0, 4) + for key, err := range ds.Keys(ctx, groupAndResourceKey, SortOrderAsc) { require.NoError(t, err) keys = append(keys, key) } - // Should include all keys with matching namespace - require.Len(t, keys, 5) // 4 from previous tests + 1 new one - - // Verify the new key is included - require.Contains(t, keys, dataKey5) - }) - - t.Run("keys with empty namespace", func(t *testing.T) { - // Group, Resource, Name are provided but will be ignored due to validation - emptyNamespaceKey := ListRequestKey{ - Namespace: "", - Group: "test-group", - Resource: "test-resource", - Name: "test-name", - } - - var keys []DataKey - var hasError bool - for key, err := range ds.Keys(ctx, emptyNamespaceKey) { - if err != nil { - hasError = true - require.Error(t, err) - require.Contains(t, err.Error(), "namespace is required") - break - } - keys = append(keys, key) - } - - // Should get an error due to validation - require.True(t, hasError, "Expected an error due to empty namespace with other fields provided") + require.Len(t, keys, 4) }) } @@ -1100,22 +1055,15 @@ func TestListRequestKey_Validate(t *testing.T) { expectError: false, }, { - name: "valid - only namespace", + name: "valid - only group and resource", key: ListRequestKey{ - Namespace: "test-namespace", + Group: "test-group", + Resource: "test-resource", }, expectError: false, }, { - name: "valid - namespace and group", - key: ListRequestKey{ - Namespace: "test-namespace", - Group: "test-group", - }, - expectError: false, - }, - { - name: "valid - namespace, group, and resource", + name: "valid - namespace and group and resource", key: ListRequestKey{ Namespace: "test-namespace", Group: "test-group", @@ -1124,87 +1072,58 @@ func TestListRequestKey_Validate(t *testing.T) { expectError: false, }, { - name: "valid - all empty", + name: "invalid - all empty", key: ListRequestKey{}, - expectError: false, - }, - { - name: "valid - with dots and dashes", - key: ListRequestKey{ - Namespace: "test.namespace-123", - Group: "test-group.v1", - Resource: "test-resource", - Name: "test.name-456", - }, - expectError: false, + expectError: true, + errorMsg: "group is required", }, // Invalid hierarchical cases { - name: "invalid - group without namespace", + name: "invalid - group without resource", key: ListRequestKey{ Group: "test-group", }, expectError: true, - errorMsg: "namespace is required when group, resource, or name are provided", - }, - { - name: "invalid - resource without namespace", - key: ListRequestKey{ - Resource: "test-resource", - }, - expectError: true, - errorMsg: "namespace is required when group, resource, or name are provided", + errorMsg: "resource is required", }, { name: "invalid - name without namespace", key: ListRequestKey{ - Name: "test-name", + Name: "test-name", + Resource: "test-resource", + Group: "test-group", }, expectError: true, - errorMsg: "namespace is required when group, resource, or name are provided", + errorMsg: "name must be empty when namespace is empty", }, { - name: "invalid - resource without group", - key: ListRequestKey{ - Namespace: "test-namespace", - Resource: "test-resource", - }, - expectError: true, - errorMsg: "group is required when resource or name are provided", - }, - { - name: "invalid - name without group", + name: "invalid - name without group and resource", key: ListRequestKey{ Namespace: "test-namespace", Name: "test-name", }, expectError: true, - errorMsg: "group is required when resource or name are provided", - }, - { - name: "invalid - name without resource", - key: ListRequestKey{ - Namespace: "test-namespace", - Group: "test-group", - Name: "test-name", - }, - expectError: true, - errorMsg: "resource is required when name is provided", + errorMsg: "group is required", }, // Invalid naming cases { name: "invalid - uppercase in namespace", key: ListRequestKey{ Namespace: "Test-Namespace", + Group: "test-group", + Resource: "test-resource", + Name: "test-name", }, expectError: true, errorMsg: "namespace 'Test-Namespace' is invalid", }, { - name: "invalid - uppercase in group", + name: "invalid - uppercase in group and resource", key: ListRequestKey{ Namespace: "test-namespace", Group: "Test-Group", + Resource: "test-resource", + Name: "test-name", }, expectError: true, errorMsg: "group 'Test-Group' is invalid", @@ -1234,6 +1153,9 @@ func TestListRequestKey_Validate(t *testing.T) { name: "invalid - underscore in namespace", key: ListRequestKey{ Namespace: "test_namespace", + Group: "test-group", + Resource: "test-resource", + Name: "test-name", }, expectError: true, errorMsg: "namespace 'test_namespace' is invalid", @@ -1242,6 +1164,9 @@ func TestListRequestKey_Validate(t *testing.T) { name: "invalid - starts with dash", key: ListRequestKey{ Namespace: "-test-namespace", + Group: "test-group", + Resource: "test-resource", + Name: "test-name", }, expectError: true, errorMsg: "namespace '-test-namespace' is invalid", @@ -1251,6 +1176,8 @@ func TestListRequestKey_Validate(t *testing.T) { key: ListRequestKey{ Namespace: "test-namespace", Group: "test-group.", + Resource: "test-resource", + Name: "test-name", }, expectError: true, errorMsg: "group 'test-group.' is invalid", @@ -1286,7 +1213,7 @@ func TestListRequestKey_Prefix(t *testing.T) { Resource: "test-resource", Name: "test-name", }, - expected: "test-namespace/test-group/test-resource/test-name/", + expected: "test-group/test-resource/test-namespace/test-name/", }, { name: "name is empty", @@ -1296,37 +1223,17 @@ func TestListRequestKey_Prefix(t *testing.T) { Resource: "test-resource", Name: "", }, - expected: "test-namespace/test-group/test-resource/", + expected: "test-group/test-resource/test-namespace/", }, { - name: "resource is empty", + name: "namespace is empty", key: ListRequestKey{ - Namespace: "test-namespace", Group: "test-group", - Resource: "", - Name: "", - }, - expected: "test-namespace/test-group/", - }, - { - name: "only namespace provided", - key: ListRequestKey{ - Namespace: "test-namespace", - Group: "", - Resource: "", - Name: "", - }, - expected: "test-namespace/", - }, - { - name: "all fields empty", - key: ListRequestKey{ Namespace: "", - Group: "", - Resource: "", + Resource: "test-resource", Name: "", }, - expected: "", + expected: "test-group/test-resource/", }, { name: "fields with special characters", @@ -1336,14 +1243,7 @@ func TestListRequestKey_Prefix(t *testing.T) { Resource: "test-resource", Name: "test-name-with-multiple.special-chars", }, - expected: "test-namespace-with-dashes/test.group.with.dots/test-resource/test-name-with-multiple.special-chars/", - }, - { - name: "invalid key still produces prefix", - key: ListRequestKey{ - Namespace: "Test-Namespace", // invalid but we assume validity in Prefix - }, - expected: "Test-Namespace/", + expected: "test.group.with.dots/test-resource/test-namespace-with-dashes/test-name-with-multiple.special-chars/", }, } @@ -1456,3 +1356,1376 @@ func TestDataStore_LastResourceVersion(t *testing.T) { } }) } + +func TestDataStore_GetLatestResourceKey(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + key := GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + } + + // Create multiple versions with different timestamps + rv1 := node.Generate().Int64() + rv2 := node.Generate().Int64() + rv3 := node.Generate().Int64() + + // Save multiple versions (rv3 should be latest) + dataKey1 := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: rv1, + Action: DataActionCreated, + Folder: "test-folder", + } + dataKey2 := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: rv2, + Action: DataActionUpdated, + Folder: "test-folder", + } + dataKey3 := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: rv3, + Action: DataActionUpdated, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey1, bytes.NewReader([]byte("version1"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey2, bytes.NewReader([]byte("version2"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey3, bytes.NewReader([]byte("version3"))) + require.NoError(t, err) + + // GetLatestResourceKey should return rv3 + latestKey, err := ds.GetLatestResourceKey(ctx, key) + require.NoError(t, err) + + require.Equal(t, dataKey3, latestKey) + require.Equal(t, rv3, latestKey.ResourceVersion) + require.Equal(t, DataActionUpdated, latestKey.Action) +} + +func TestDataStore_GetLatestResourceKey_Deleted(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + key := GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + } + + dataKey := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: node.Generate().Int64(), + Action: DataActionDeleted, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey, bytes.NewReader([]byte("deleted"))) + require.NoError(t, err) + + _, err = ds.GetLatestResourceKey(ctx, key) + require.Equal(t, ErrNotFound, err) +} + +func TestDataStore_GetLatestResourceKey_NotFound(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + key := GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "non-existent", + } + + _, err := ds.GetLatestResourceKey(ctx, key) + require.Equal(t, ErrNotFound, err) +} + +func TestDataStore_GetResourceKeyAtRevision(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + key := GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + } + + // Create multiple versions + rv1 := node.Generate().Int64() + rv2 := node.Generate().Int64() + rv3 := node.Generate().Int64() + + dataKey1 := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: rv1, + Action: DataActionCreated, + Folder: "test-folder", + } + dataKey2 := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: rv2, + Action: DataActionUpdated, + Folder: "test-folder", + } + dataKey3 := DataKey{ + Group: key.Group, + Resource: key.Resource, + Namespace: key.Namespace, + Name: key.Name, + ResourceVersion: rv3, + Action: DataActionUpdated, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey1, bytes.NewReader([]byte("version1"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey2, bytes.NewReader([]byte("version2"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey3, bytes.NewReader([]byte("version3"))) + require.NoError(t, err) + + // Get key at rv2 should return rv2 + dataKey, err := ds.GetResourceKeyAtRevision(ctx, key, rv2) + require.NoError(t, err) + + require.Equal(t, rv2, dataKey.ResourceVersion) + require.Equal(t, DataActionUpdated, dataKey.Action) + + // Get key at rv1 should return rv1 + dataKey, err = ds.GetResourceKeyAtRevision(ctx, key, rv1) + require.NoError(t, err) + + require.Equal(t, rv1, dataKey.ResourceVersion) + require.Equal(t, DataActionCreated, dataKey.Action) + + // Get key at revision 0 should return latest (rv3) + dataKey, err = ds.GetResourceKeyAtRevision(ctx, key, 0) + require.NoError(t, err) + + require.Equal(t, rv3, dataKey.ResourceVersion) + require.Equal(t, DataActionUpdated, dataKey.Action) +} + +func TestDataStore_ListLatestResourceKeys(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + listKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + } + + // Save multiple versions - ListLatestResourceKeys should return only the latest + rv1 := node.Generate().Int64() + rv2 := node.Generate().Int64() + + dataKey1 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: listKey.Name, + ResourceVersion: rv1, + Action: DataActionCreated, + Folder: "test-folder", + } + dataKey2 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: listKey.Name, + ResourceVersion: rv2, + Action: DataActionUpdated, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey1, bytes.NewReader([]byte("version1"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey2, bytes.NewReader([]byte("version2"))) + require.NoError(t, err) + + // List latest resource keys + resultKeys := make([]DataKey, 0, 1) + for dataKey, err := range ds.ListLatestResourceKeys(ctx, listKey) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 1) + require.Equal(t, dataKey2, resultKeys[0]) + require.Equal(t, rv2, resultKeys[0].ResourceVersion) + require.Equal(t, DataActionUpdated, resultKeys[0].Action) +} + +func TestDataStore_ListLatestResourceKeys_Deleted(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + listKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + } + + // Save a resource and then delete it + rv1 := node.Generate().Int64() + rv2 := node.Generate().Int64() + + dataKey1 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: listKey.Name, + ResourceVersion: rv1, + Action: DataActionCreated, + Folder: "test-folder", + } + dataKey2 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: listKey.Name, + ResourceVersion: rv2, + Action: DataActionDeleted, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey1, bytes.NewReader([]byte("version1"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey2, bytes.NewReader([]byte("deleted"))) + require.NoError(t, err) + + // ListLatestResourceKeys should exclude deleted resources + resultKeys := make([]DataKey, 0, 1) + for dataKey, err := range ds.ListLatestResourceKeys(ctx, listKey) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 0) // Should be empty because resource was deleted +} + +func TestDataStore_ListLatestResourceKeys_Multiple(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + listKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + } + + // Save multiple resources with different names + rv1 := node.Generate().Int64() + rv2 := node.Generate().Int64() + rv3 := node.Generate().Int64() + + dataKey1 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: "resource-1", + ResourceVersion: rv1, + Action: DataActionCreated, + Folder: "test-folder", + } + dataKey2 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: "resource-2", + ResourceVersion: rv2, + Action: DataActionCreated, + Folder: "test-folder", + } + dataKey3 := DataKey{ + Group: listKey.Group, + Resource: listKey.Resource, + Namespace: listKey.Namespace, + Name: "resource-1", + ResourceVersion: rv3, + Action: DataActionUpdated, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey1, bytes.NewReader([]byte("resource-1-v1"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey2, bytes.NewReader([]byte("resource-2-v1"))) + require.NoError(t, err) + + err = ds.Save(ctx, dataKey3, bytes.NewReader([]byte("resource-1-v2"))) + require.NoError(t, err) + + // List latest resource keys for all resources + resultKeys := make([]DataKey, 0, 3) + for dataKey, err := range ds.ListLatestResourceKeys(ctx, listKey) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 2) // resource-1 (latest version) and resource-2 + + // Check we got the correct keys + names := make(map[string]int64) + for _, key := range resultKeys { + names[key.Name] = key.ResourceVersion + } + + require.Equal(t, rv3, names["resource-1"]) // Should be the updated version + require.Equal(t, rv2, names["resource-2"]) +} + +func TestDataStore_ListResourceKeysAtRevision(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + // Create multiple resources with different versions + rv1 := node.Generate().Int64() + rv2 := node.Generate().Int64() + rv3 := node.Generate().Int64() + rv4 := node.Generate().Int64() + rv5 := node.Generate().Int64() + + // Resource 1: Created at rv1, updated at rv3 + key1 := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "resource1", + ResourceVersion: rv1, + Action: DataActionCreated, + Folder: "test-folder", + } + err := ds.Save(ctx, key1, bytes.NewReader([]byte("resource1-v1"))) + require.NoError(t, err) + + key1Updated := key1 + key1Updated.ResourceVersion = rv3 + key1Updated.Action = DataActionUpdated + err = ds.Save(ctx, key1Updated, bytes.NewReader([]byte("resource1-v2"))) + require.NoError(t, err) + + // Resource 2: Created at rv2 + key2 := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "resource2", + ResourceVersion: rv2, + Action: DataActionCreated, + Folder: "test-folder", + } + err = ds.Save(ctx, key2, bytes.NewReader([]byte("resource2-v1"))) + require.NoError(t, err) + + // Resource 3: Created at rv4 + key3 := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "resource3", + ResourceVersion: rv4, + Action: DataActionCreated, + Folder: "test-folder", + } + err = ds.Save(ctx, key3, bytes.NewReader([]byte("resource3-v1"))) + require.NoError(t, err) + + // Resource 4: Created at rv2, deleted at rv5 + key4 := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "resource4", + ResourceVersion: rv2, + Action: DataActionCreated, + Folder: "test-folder", + } + err = ds.Save(ctx, key4, bytes.NewReader([]byte("resource4-v1"))) + require.NoError(t, err) + + key4Deleted := key4 + key4Deleted.ResourceVersion = rv5 + key4Deleted.Action = DataActionDeleted + err = ds.Save(ctx, key4Deleted, bytes.NewReader([]byte("resource4-deleted"))) + require.NoError(t, err) + + listKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + } + + 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) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 1) + require.Equal(t, "resource1", resultKeys[0].Name) + require.Equal(t, rv1, resultKeys[0].ResourceVersion) + require.Equal(t, DataActionCreated, resultKeys[0].Action) + }) + + 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) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 3) // resource1, resource2, resource4 + names := make(map[string]int64) + for _, result := range resultKeys { + names[result.Name] = result.ResourceVersion + } + + require.Equal(t, rv1, names["resource1"]) // Should be the original version + require.Equal(t, rv2, names["resource2"]) + require.Equal(t, rv2, names["resource4"]) + }) + + 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) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 3) // resource1 (updated), resource2, resource4 + names := make(map[string]int64) + actions := make(map[string]DataAction) + for _, result := range resultKeys { + names[result.Name] = result.ResourceVersion + actions[result.Name] = result.Action + } + + require.Equal(t, rv3, names["resource1"]) // Should be the updated version + require.Equal(t, DataActionUpdated, actions["resource1"]) + require.Equal(t, rv2, names["resource2"]) + require.Equal(t, rv2, names["resource4"]) + }) + + 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) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 4) // resource1 (updated), resource2, resource3, resource4 + names := make(map[string]int64) + for _, result := range resultKeys { + names[result.Name] = result.ResourceVersion + } + + require.Equal(t, rv3, names["resource1"]) + require.Equal(t, rv2, names["resource2"]) + require.Equal(t, rv4, names["resource3"]) + require.Equal(t, rv2, names["resource4"]) + }) + + 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) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 3) // resource1 (updated), resource2, resource3 (resource4 excluded because deleted) + names := make(map[string]bool) + for _, result := range resultKeys { + names[result.Name] = true + } + + require.True(t, names["resource1"]) + require.True(t, names["resource2"]) + require.True(t, names["resource3"]) + require.False(t, names["resource4"]) // Should be excluded because it's deleted + }) + + t.Run("list with specific resource name", func(t *testing.T) { + specificListKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "resource1", + } + + resultKeys := make([]DataKey, 0, 2) + for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, specificListKey, rv3) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 1) + require.Equal(t, "resource1", resultKeys[0].Name) + require.Equal(t, rv3, resultKeys[0].ResourceVersion) + require.Equal(t, DataActionUpdated, resultKeys[0].Action) + }) + + 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) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + // Should return all non-deleted resources at their latest versions + require.Len(t, resultKeys, 3) // resource1 (updated), resource2, resource3 + names := make(map[string]bool) + for _, result := range resultKeys { + names[result.Name] = true + } + + require.True(t, names["resource1"]) + require.True(t, names["resource2"]) + require.True(t, names["resource3"]) + require.False(t, names["resource4"]) // Excluded because deleted + }) +} + +func TestDataStore_ListResourceKeysAtRevision_ValidationErrors(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + tests := []struct { + name string + key ListRequestKey + }{ + { + name: "missing group", + key: ListRequestKey{ + Namespace: "default", + Resource: "resources", + }, + }, + { + name: "missing resource", + key: ListRequestKey{ + Namespace: "default", + Group: "apps", + }, + }, + { + name: "name without namespace", + key: ListRequestKey{ + Group: "apps", + Resource: "resources", + Name: "test-name", + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + for _, err := range ds.ListResourceKeysAtRevision(ctx, tt.key, 0) { + require.Error(t, err) + return + } + }) + } +} + +func TestDataStore_ListResourceKeysAtRevision_EmptyResults(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + listKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "empty", + } + + resultKeys := make([]DataKey, 0, 1) + for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, 0) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + require.Len(t, resultKeys, 0) +} + +func TestDataStore_ListResourceKeysAtRevision_ResourcesNewerThanRevision(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + // Create a resource with a high resource version + rv := node.Generate().Int64() + key := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "future-resource", + ResourceVersion: rv, + Action: DataActionCreated, + Folder: "test-folder", + } + err := ds.Save(ctx, key, bytes.NewReader([]byte("future-resource"))) + require.NoError(t, err) + + listKey := ListRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + } + + // List at a revision before the resource was created + resultKeys := make([]DataKey, 0, 1) + for dataKey, err := range ds.ListResourceKeysAtRevision(ctx, listKey, rv-1000) { + require.NoError(t, err) + resultKeys = append(resultKeys, dataKey) + } + + // Should return no results since the resource is newer than the target revision + require.Len(t, resultKeys, 0) +} + +func TestDataKey_Equals(t *testing.T) { + baseKey := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + } + + tests := []struct { + name string + key1 DataKey + key2 DataKey + expected bool + }{ + { + name: "identical keys", + key1: baseKey, + key2: baseKey, + expected: true, + }, + { + name: "different resource version", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 456, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different action", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionUpdated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different folder", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "other-folder", + }, + expected: false, + }, + { + name: "different namespace", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "other-namespace", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different group", + key1: baseKey, + key2: DataKey{ + Group: "extensions", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different resource", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "services", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different name", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "other-deployment", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "empty keys", + key1: DataKey{}, + key2: DataKey{}, + expected: true, + }, + { + name: "one empty key", + key1: baseKey, + key2: DataKey{}, + expected: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + result := tt.key1.Equals(tt.key2) + require.Equal(t, tt.expected, result) + + // Test symmetry: Equals should be commutative + reverseResult := tt.key2.Equals(tt.key1) + require.Equal(t, result, reverseResult, "Equals method should be commutative") + }) + } +} + +func TestDataKey_SameResource(t *testing.T) { + baseKey := DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + } + + tests := []struct { + name string + key1 DataKey + key2 DataKey + expected bool + }{ + { + name: "identical keys", + key1: baseKey, + key2: baseKey, + expected: true, + }, + { + name: "same identifying fields, different resource version", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 456, // Different resource version + Action: DataActionUpdated, + Folder: "other-folder", + }, + expected: true, // Should still be equal as ResourceVersion, Action, and Folder don't matter + }, + { + name: "different namespace", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "other-namespace", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different group", + key1: baseKey, + key2: DataKey{ + Group: "extensions", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different resource", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "services", + Namespace: "default", + Name: "test-resource", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "different name", + key1: baseKey, + key2: DataKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "other-deployment", + ResourceVersion: 123, + Action: DataActionCreated, + Folder: "test-folder", + }, + expected: false, + }, + { + name: "empty keys", + key1: DataKey{}, + key2: DataKey{}, + expected: true, + }, + { + name: "one empty key", + key1: baseKey, + key2: DataKey{}, + expected: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + result := tt.key1.SameResource(tt.key2) + require.Equal(t, tt.expected, result) + + // Test symmetry: SameResource should be commutative + reverseResult := tt.key2.SameResource(tt.key1) + require.Equal(t, result, reverseResult, "SameResource method should be commutative") + }) + } +} + +func TestGetRequestKey_Validate(t *testing.T) { + tests := []struct { + name string + key GetRequestKey + expectErr bool + wantError string + }{ + { + name: "valid key", + key: GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + }, + expectErr: false, + }, + { + name: "valid key with dots and dashes", + key: GetRequestKey{ + Group: "apps.v1", + Resource: "deployment-configs", + Namespace: "default-ns", + Name: "test-resource.v1", + }, + expectErr: false, + }, + { + name: "missing group", + key: GetRequestKey{ + Resource: "resources", + Namespace: "default", + Name: "test-resource", + }, + expectErr: true, + wantError: "group is required", + }, + { + name: "missing resource", + key: GetRequestKey{ + Group: "apps", + Namespace: "default", + Name: "test-resource", + }, + expectErr: true, + wantError: "resource is required", + }, + { + name: "missing namespace", + key: GetRequestKey{ + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + expectErr: true, + wantError: "namespace is required", + }, + { + name: "missing name", + key: GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + }, + expectErr: true, + wantError: "name is required", + }, + { + name: "invalid namespace - uppercase", + key: GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "Default", + Name: "test-resource", + }, + expectErr: true, + wantError: "namespace 'Default' is invalid", + }, + { + name: "invalid group - underscore", + key: GetRequestKey{ + Group: "apps_v1", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + }, + expectErr: true, + wantError: "group 'apps_v1' is invalid", + }, + { + name: "invalid resource - starts with dash", + key: GetRequestKey{ + Group: "apps", + Resource: "-resources", + Namespace: "default", + Name: "test-resource", + }, + expectErr: true, + wantError: "resource '-resources' is invalid", + }, + { + name: "invalid name - ends with dot", + key: GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource.", + }, + expectErr: true, + wantError: "name 'test-resource.' is invalid", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + err := tt.key.Validate() + if tt.expectErr { + require.Error(t, err) + if tt.wantError != "" { + require.Contains(t, err.Error(), tt.wantError) + } + } else { + require.NoError(t, err) + } + }) + } +} + +func TestGetRequestKey_Prefix(t *testing.T) { + tests := []struct { + name string + key GetRequestKey + expectedPrefix string + }{ + { + name: "standard key", + key: GetRequestKey{ + Group: "apps", + Resource: "resources", + Namespace: "default", + Name: "test-resource", + }, + expectedPrefix: "apps/resources/default/test-resource/", + }, + { + name: "key with special characters", + key: GetRequestKey{ + Group: "apps.v1", + Resource: "deployment-configs", + Namespace: "system-namespace", + Name: "my-app.v2", + }, + expectedPrefix: "apps.v1/deployment-configs/system-namespace/my-app.v2/", + }, + { + name: "key with single character fields", + key: GetRequestKey{ + Group: "a", + Resource: "b", + Namespace: "c", + Name: "d", + }, + expectedPrefix: "a/b/c/d/", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + prefix := tt.key.Prefix() + require.Equal(t, tt.expectedPrefix, prefix) + }) + } +} + +func TestDataStore_GetResourceStats_Comprehensive(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + // Test setup: 3 namespaces × 3 groups × 3 resources × 3 names × 3 versions = 243 total entries + // But each name will have only 1 latest version that counts, so 3 × 3 × 3 × 3 = 81 non-deleted resources + namespaces := []string{"ns1", "ns2", "ns3"} + groups := []string{"apps", "extensions", "networking"} + resources := []string{"deployments", "services", "ingresses"} + names := []string{"item1", "item2", "item3"} + + // Create all the test data + totalEntries := 0 + for _, ns := range namespaces { + for _, group := range groups { + for _, resource := range resources { + for _, name := range names { + // Create 3 versions for each resource name + for version := 1; version <= 3; version++ { + rv := node.Generate().Int64() + + var action DataAction + switch version { + case 1: + action = DataActionCreated + case 2, 3: + action = DataActionUpdated + } + + dataKey := DataKey{ + Namespace: ns, + Group: group, + Resource: resource, + Name: name, + ResourceVersion: rv, + Action: action, + Folder: "test-folder", + } + + content := fmt.Sprintf("%s/%s/%s/%s-v%d", ns, group, resource, name, version) + err := ds.Save(ctx, dataKey, bytes.NewReader([]byte(content))) + require.NoError(t, err) + totalEntries++ + } + } + } + } + } + + // Verify we created the expected number of entries + require.Equal(t, 243, totalEntries) // 3×3×3×3×3 = 243 total entries + + t.Run("get stats for all namespaces", func(t *testing.T) { + stats, err := ds.GetResourceStats(ctx, "", 0) + require.NoError(t, err) + + // Should have 27 resource types (3 namespaces × 3 groups × 3 resources) + require.Len(t, stats, 27) + + // Each resource type should have exactly 3 items (3 names per resource type) + for _, stat := range stats { + require.Equal(t, int64(3), stat.Count, "Resource %s/%s/%s should have 3 items", stat.Namespace, stat.Group, stat.Resource) + require.Greater(t, stat.ResourceVersion, int64(0), "ResourceVersion should be positive") + } + + // Verify all expected combinations are present + expectedCombinations := make(map[string]bool) + for _, ns := range namespaces { + for _, group := range groups { + for _, resource := range resources { + key := fmt.Sprintf("%s/%s/%s", ns, group, resource) + expectedCombinations[key] = false + } + } + } + + for _, stat := range stats { + key := fmt.Sprintf("%s/%s/%s", stat.Namespace, stat.Group, stat.Resource) + expectedCombinations[key] = true + } + + for key, found := range expectedCombinations { + require.True(t, found, "Expected combination not found: %s", key) + } + }) + + t.Run("get stats for specific namespace ns1", func(t *testing.T) { + stats, err := ds.GetResourceStats(ctx, "ns1", 0) + require.NoError(t, err) + + // Should have 9 resource types (3 groups × 3 resources for ns1) + require.Len(t, stats, 9) + + // All stats should be for ns1 + for _, stat := range stats { + require.Equal(t, "ns1", stat.Namespace) + require.Equal(t, int64(3), stat.Count) // 3 names per resource type + } + + // Verify we have all expected groups and resources for ns1 + foundCombinations := make(map[string]bool) + for _, stat := range stats { + key := fmt.Sprintf("%s/%s", stat.Group, stat.Resource) + foundCombinations[key] = true + } + + expectedCount := len(groups) * len(resources) // 3×3=9 + require.Equal(t, expectedCount, len(foundCombinations)) + }) + + t.Run("get stats for specific namespace ns2", func(t *testing.T) { + stats, err := ds.GetResourceStats(ctx, "ns2", 0) + require.NoError(t, err) + + // Should have 9 resource types (3 groups × 3 resources for ns2) + require.Len(t, stats, 9) + + // All stats should be for ns2 + for _, stat := range stats { + require.Equal(t, "ns2", stat.Namespace) + require.Equal(t, int64(3), stat.Count) + } + }) + + t.Run("get stats with minCount filter", func(t *testing.T) { + // With minCount=0, all resources should be included (each has 3 items > 0) + stats, err := ds.GetResourceStats(ctx, "", 0) + require.NoError(t, err) + require.Len(t, stats, 27) // All 27 resource types should be included + + // With minCount=2, all resources should still be included (each has 3 items > 2) + stats, err = ds.GetResourceStats(ctx, "", 2) + require.NoError(t, err) + require.Len(t, stats, 27) // All 27 resource types should still be included + + // With minCount=3, no resources should be included (each has exactly 3 items, not > 3) + stats, err = ds.GetResourceStats(ctx, "", 3) + require.NoError(t, err) + require.Len(t, stats, 0) + + // With minCount=4, no resources should be included (each has only 3 items < 4) + stats, err = ds.GetResourceStats(ctx, "", 4) + require.NoError(t, err) + require.Len(t, stats, 0) + }) + + t.Run("get stats for non-existent namespace", func(t *testing.T) { + stats, err := ds.GetResourceStats(ctx, "non-existent", 0) + require.NoError(t, err) + require.Len(t, stats, 0) + }) + + t.Run("add deleted resources and verify counts", func(t *testing.T) { + // Delete one resource from ns1/apps/deployments/item1 + rv := node.Generate().Int64() + deletedKey := DataKey{ + Namespace: "ns1", + Group: "apps", + Resource: "deployments", + Name: "item1", + ResourceVersion: rv, + Action: DataActionDeleted, + Folder: "test-folder", + } + + err := ds.Save(ctx, deletedKey, bytes.NewReader([]byte("deleted"))) + require.NoError(t, err) + + // Get stats for ns1 - apps/deployments should now have 2 items instead of 3 + stats, err := ds.GetResourceStats(ctx, "ns1", 0) + require.NoError(t, err) + + // Find the apps/deployments stat + var appsDeploymentsCount int64 = -1 + for _, stat := range stats { + if stat.Group == "apps" && stat.Resource == "deployments" { + appsDeploymentsCount = stat.Count + break + } + } + + require.Equal(t, int64(2), appsDeploymentsCount, "apps/deployments should have 2 items after deletion") + + // Other resource types in ns1 should still have 3 items + otherResourceCount := 0 + for _, stat := range stats { + if stat.Group != "apps" || stat.Resource != "deployments" { + require.Equal(t, int64(3), stat.Count, "Other resources should still have 3 items") + otherResourceCount++ + } + } + require.Equal(t, 8, otherResourceCount) // 9 total - 1 apps/deployments = 8 + }) + + t.Run("verify resource versions are meaningful", func(t *testing.T) { + stats, err := ds.GetResourceStats(ctx, "ns1", 0) + require.NoError(t, err) + + // All ResourceVersions should be positive and reasonable + for _, stat := range stats { + require.Greater(t, stat.ResourceVersion, int64(0)) + // ResourceVersion should be a snowflake ID, so it should be quite large + require.Greater(t, stat.ResourceVersion, int64(1000000)) + } + }) +} + +func TestDataStore_getGroupResources(t *testing.T) { + ds := setupTestDataStore(t) + ctx := context.Background() + + // Create test data with multiple group/resource combinations + testData := []struct { + group string + resource string + namespace string + name string + }{ + {"apps", "deployments", "default", "web-app"}, + {"apps", "deployments", "test", "api-server"}, + {"apps", "services", "default", "web-svc"}, + {"networking", "ingresses", "default", "web-ingress"}, + {"batch", "jobs", "default", "cleanup-job"}, + {"batch", "jobs", "test", "migration-job"}, + } + + // Save all test data + for i, data := range testData { + rv := node.Generate().Int64() + dataKey := DataKey{ + Namespace: data.namespace, + Group: data.group, + Resource: data.resource, + Name: data.name, + ResourceVersion: rv, + Action: DataActionCreated, + Folder: "test-folder", + } + + err := ds.Save(ctx, dataKey, bytes.NewReader([]byte(fmt.Sprintf("content-%d", i)))) + require.NoError(t, err) + } + + // Test GetGroupResources + results, err := ds.getGroupResources(ctx) + require.NoError(t, err) + + // Should find exactly 4 unique group/resource combinations + expectedCombinations := []string{ + "apps/deployments", + "apps/services", + "networking/ingresses", + "batch/jobs", + } + + require.Len(t, results, len(expectedCombinations)) + + // Verify all expected combinations are present and no duplicates + foundCombinations := make(map[string]bool) + for _, result := range results { + key := fmt.Sprintf("%s/%s", result.Group, result.Resource) + require.False(t, foundCombinations[key], "Duplicate group/resource found: %s", key) + foundCombinations[key] = true + } + + for _, expected := range expectedCombinations { + require.True(t, foundCombinations[expected], "Expected combination not found: %s", expected) + } +} diff --git a/pkg/storage/unified/resource/eventstore.go b/pkg/storage/unified/resource/eventstore.go index 651fcb52092..6558a81839c 100644 --- a/pkg/storage/unified/resource/eventstore.go +++ b/pkg/storage/unified/resource/eventstore.go @@ -28,10 +28,11 @@ type EventKey struct { Name string ResourceVersion int64 Action DataAction + Folder string } func (k EventKey) String() string { - return fmt.Sprintf("%d~%s~%s~%s~%s~%s", k.ResourceVersion, k.Namespace, k.Group, k.Resource, k.Name, k.Action) + return fmt.Sprintf("%d~%s~%s~%s~%s~%s~%s", k.ResourceVersion, k.Namespace, k.Group, k.Resource, k.Name, k.Action, k.Folder) } func (k EventKey) Validate() error { @@ -53,7 +54,9 @@ func (k EventKey) Validate() error { if k.Action == "" { return fmt.Errorf("action cannot be empty") } - + if k.Folder != "" && !validNameRegex.MatchString(k.Folder) { + return fmt.Errorf("folder '%s' is invalid", k.Folder) + } // Validate each field against the naming rules (reusing the regex from datastore.go) if !validNameRegex.MatchString(k.Namespace) { return fmt.Errorf("namespace '%s' is invalid", k.Namespace) @@ -67,7 +70,9 @@ func (k EventKey) Validate() error { if !validNameRegex.MatchString(k.Name) { return fmt.Errorf("name '%s' is invalid", k.Name) } - + if k.Folder != "" && !validNameRegex.MatchString(k.Folder) { + return fmt.Errorf("folder '%s' is invalid", k.Folder) + } switch k.Action { case DataActionCreated, DataActionUpdated, DataActionDeleted: default: @@ -97,7 +102,7 @@ func newEventStore(kv KV) *eventStore { // ParseEventKey parses a key string back into an EventKey struct func ParseEventKey(key string) (EventKey, error) { parts := strings.Split(key, "~") - if len(parts) != 6 { + if len(parts) != 7 { return EventKey{}, fmt.Errorf("invalid key format: expected 6 parts, got %d", len(parts)) } @@ -113,6 +118,7 @@ func ParseEventKey(key string) (EventKey, error) { Resource: parts[3], Name: parts[4], Action: DataAction(parts[5]), + Folder: parts[6], }, nil } diff --git a/pkg/storage/unified/resource/eventstore_test.go b/pkg/storage/unified/resource/eventstore_test.go index 63ddc999567..c2b549f5e0c 100644 --- a/pkg/storage/unified/resource/eventstore_test.go +++ b/pkg/storage/unified/resource/eventstore_test.go @@ -40,8 +40,9 @@ func TestEventKey_String(t *testing.T) { Name: "test-resource", ResourceVersion: 1000, Action: "created", + Folder: "test-folder", }, - expected: "1000~default~apps~resource~test-resource~created", + expected: "1000~default~apps~resource~test-resource~created~test-folder", }, { name: "empty namespace", @@ -52,8 +53,9 @@ func TestEventKey_String(t *testing.T) { Name: "test-resource", ResourceVersion: 2000, Action: "updated", + Folder: "test-folder", }, - expected: "2000~~apps~resource~test-resource~updated", + expected: "2000~~apps~resource~test-resource~updated~test-folder", }, { name: "special characters in name", @@ -64,8 +66,9 @@ func TestEventKey_String(t *testing.T) { Name: "test-resource-with-dashes", ResourceVersion: 3000, Action: "deleted", + Folder: "test-folder", }, - expected: "3000~test-ns~apps~resource~test-resource-with-dashes~deleted", + expected: "3000~test-ns~apps~resource~test-resource-with-dashes~deleted~test-folder", }, } @@ -86,7 +89,7 @@ func TestEventKey_Validate(t *testing.T) { }{ { name: "valid key", - key: "1000~default~apps~resource~test-resource~created", + key: "1000~default~apps~resource~test-resource~created~test-folder", expected: EventKey{ ResourceVersion: 1000, Namespace: "default", @@ -94,11 +97,12 @@ func TestEventKey_Validate(t *testing.T) { Resource: "resource", Name: "test-resource", Action: "created", + Folder: "test-folder", }, }, { name: "empty namespace", - key: "2000~~apps~resource~test-resource~updated", + key: "2000~~apps~resource~test-resource~updated~", expected: EventKey{ ResourceVersion: 2000, Namespace: "", @@ -110,7 +114,7 @@ func TestEventKey_Validate(t *testing.T) { }, { name: "special characters in name", - key: "3000~test-ns~apps~resource~test-resource-with-dashes~updated", + key: "3000~test-ns~apps~resource~test-resource-with-dashes~updated~", expected: EventKey{ ResourceVersion: 3000, Namespace: "test-ns", @@ -122,17 +126,17 @@ func TestEventKey_Validate(t *testing.T) { }, { name: "invalid key - too few parts", - key: "1000~default~apps~resource", + key: "1000~default~apps~resource~", expectError: true, }, { name: "invalid key - too many parts", - key: "1000~default~apps~resource~test~extra~parts", + key: "1000~default~apps~resource~test~extra~parts~", expectError: true, }, { name: "invalid resource version", - key: "invalid~default~apps~resource~test~cerated", + key: "invalid~default~apps~resource~test~cerated~", expectError: true, }, { diff --git a/pkg/storage/unified/resource/metadata.go b/pkg/storage/unified/resource/metadata.go deleted file mode 100644 index 83240f97727..00000000000 --- a/pkg/storage/unified/resource/metadata.go +++ /dev/null @@ -1,391 +0,0 @@ -package resource - -import ( - "context" - "encoding/json" - "fmt" - "iter" - "math" - "strconv" - "strings" -) - -const ( - metaSection = "unified/meta" -) - -// Metadata store stores search documents for resources in unified storage. -// The store keeps track of the latest versions of each resource. -type MetaData struct { - IndexableDocument -} - -type MetaDataKey struct { - Namespace string - Group string - Resource string - Name string - ResourceVersion int64 - Folder string - Action DataAction -} - -// String returns the string representation of the MetaDataKey used as the storage key -func (k MetaDataKey) String() string { - return fmt.Sprintf("%s/%s/%s/%s/%d~%s~%s", k.Group, k.Resource, k.Namespace, k.Name, k.ResourceVersion, k.Action, k.Folder) -} - -// Validate validates that all required fields are present and valid -func (k MetaDataKey) Validate() error { - if k.Group == "" { - return fmt.Errorf("group is required") - } - if k.Resource == "" { - return fmt.Errorf("resource is required") - } - if k.Namespace == "" { - return fmt.Errorf("namespace is required") - } - if k.Name == "" { - return fmt.Errorf("name is required") - } - if k.ResourceVersion <= 0 { - return fmt.Errorf("resource version must be positive") - } - if k.Action == "" { - return fmt.Errorf("action is required") - } - - // Validate naming conventions for all required fields - if !validNameRegex.MatchString(k.Namespace) { - return fmt.Errorf("namespace '%s' is invalid", k.Namespace) - } - if !validNameRegex.MatchString(k.Group) { - return fmt.Errorf("group '%s' is invalid", k.Group) - } - if !validNameRegex.MatchString(k.Resource) { - return fmt.Errorf("resource '%s' is invalid", k.Resource) - } - if !validNameRegex.MatchString(k.Name) { - return fmt.Errorf("name '%s' is invalid", k.Name) - } - - // Validate folder field if provided (optional field) - if k.Folder != "" && !validNameRegex.MatchString(k.Folder) { - return fmt.Errorf("folder '%s' is invalid", k.Folder) - } - - // Validate action is one of the valid values - switch k.Action { - case DataActionCreated, DataActionUpdated, DataActionDeleted: - return nil - default: - return fmt.Errorf("action '%s' is invalid: must be one of 'created', 'updated', or 'deleted'", k.Action) - } -} - -// MetaListRequestKey is used for listing metadata objects -type MetaListRequestKey struct { - Namespace string - Group string - Resource string - Name string // optional for listing multiple resources -} - -// Validate validates the list request key -func (k MetaListRequestKey) Validate() error { - if k.Group == "" { - return fmt.Errorf("group is required") - } - if k.Resource == "" { - return fmt.Errorf("resource is required") - } - - // If namespace is empty, name must also be empty - if k.Namespace == "" && k.Name != "" { - return fmt.Errorf("name must be empty when namespace is empty") - } - - // Validate naming conventions - if k.Namespace != "" && !validNameRegex.MatchString(k.Namespace) { - return fmt.Errorf("namespace '%s' is invalid", k.Namespace) - } - if !validNameRegex.MatchString(k.Group) { - return fmt.Errorf("group '%s' is invalid", k.Group) - } - if !validNameRegex.MatchString(k.Resource) { - return fmt.Errorf("resource '%s' is invalid", k.Resource) - } - if k.Name != "" && !validNameRegex.MatchString(k.Name) { - return fmt.Errorf("name '%s' is invalid", k.Name) - } - - return nil -} - -// Prefix returns the prefix for listing metadata objects -func (k MetaListRequestKey) Prefix() string { - if k.Name == "" { - if k.Namespace == "" { - return fmt.Sprintf("%s/%s/", k.Group, k.Resource) - } - return fmt.Sprintf("%s/%s/%s/", k.Group, k.Resource, k.Namespace) - } - if k.Namespace == "" { - return fmt.Sprintf("%s/%s/%s/", k.Group, k.Resource, k.Name) - } - return fmt.Sprintf("%s/%s/%s/%s/", k.Group, k.Resource, k.Namespace, k.Name) -} - -// MetaGetRequestKey is used for getting a specific metadata object by latest version -type MetaGetRequestKey struct { - Namespace string - Group string - Resource string - Name string -} - -// Validate validates the get request key -func (k MetaGetRequestKey) Validate() error { - if k.Group == "" { - return fmt.Errorf("group is required") - } - if k.Resource == "" { - return fmt.Errorf("resource is required") - } - if k.Namespace == "" { - return fmt.Errorf("namespace is required") - } - if k.Name == "" { - return fmt.Errorf("name is required") - } - - // Validate naming conventions - if !validNameRegex.MatchString(k.Namespace) { - return fmt.Errorf("namespace '%s' is invalid", k.Namespace) - } - if !validNameRegex.MatchString(k.Group) { - return fmt.Errorf("group '%s' is invalid", k.Group) - } - if !validNameRegex.MatchString(k.Resource) { - return fmt.Errorf("resource '%s' is invalid", k.Resource) - } - if !validNameRegex.MatchString(k.Name) { - return fmt.Errorf("name '%s' is invalid", k.Name) - } - - return nil -} - -// Prefix returns the prefix for getting a specific metadata object -func (k MetaGetRequestKey) Prefix() string { - return fmt.Sprintf("%s/%s/%s/%s/", k.Group, k.Resource, k.Namespace, k.Name) -} - -type MetaDataObj struct { - Key MetaDataKey - Value MetaData -} - -type metadataStore struct { - kv KV -} - -// newMetadataStore creates a new metadata store instance with the given key-value store backend. -func newMetadataStore(kv KV) *metadataStore { - return &metadataStore{ - kv: kv, - } -} - -// Get retrieves the metadata for a specific metadata key. -// It validates the key and returns the raw metadata content. -func (d *metadataStore) Get(ctx context.Context, key MetaDataKey) (MetaData, error) { - if err := key.Validate(); err != nil { - return MetaData{}, fmt.Errorf("invalid metadata key: %w", err) - } - - reader, err := d.kv.Get(ctx, metaSection, key.String()) - if err != nil { - return MetaData{}, err - } - defer func() { - _ = reader.Close() - }() - var meta MetaData - err = json.NewDecoder(reader).Decode(&meta) - return meta, err -} - -// GetLatestResourceKey retrieves the metadata key for the latest version of a resource. -// Returns the key with the highest resource version that is not deleted. -func (d *metadataStore) GetLatestResourceKey(ctx context.Context, key MetaGetRequestKey) (MetaDataKey, error) { - return d.GetResourceKeyAtRevision(ctx, key, 0) -} - -// GetResourceKeyAtRevision retrieves the metadata key for a resource at a specific revision. -// If rv is 0, it returns the latest version. Returns the highest version <= rv that is not deleted. -func (d *metadataStore) GetResourceKeyAtRevision(ctx context.Context, key MetaGetRequestKey, rv int64) (MetaDataKey, error) { - if err := key.Validate(); err != nil { - return MetaDataKey{}, fmt.Errorf("invalid get request key: %w", err) - } - - if rv == 0 { - rv = math.MaxInt64 - } - - listKey := MetaListRequestKey(key) - - iter := d.ListResourceKeysAtRevision(ctx, listKey, rv) - for metaKey, err := range iter { - if err != nil { - return MetaDataKey{}, err - } - return metaKey, nil - } - return MetaDataKey{}, ErrNotFound -} - -// ListLatestResourceKeys returns an iterator over the metadata keys for the latest versions of resources. -// Only returns keys for resources that are not deleted. -func (d *metadataStore) ListLatestResourceKeys(ctx context.Context, key MetaListRequestKey) iter.Seq2[MetaDataKey, error] { - return d.ListResourceKeysAtRevision(ctx, key, 0) -} - -// ListResourceKeysAtRevision returns an iterator over metadata 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 *metadataStore) ListResourceKeysAtRevision(ctx context.Context, key MetaListRequestKey, rv int64) iter.Seq2[MetaDataKey, error] { - if err := key.Validate(); err != nil { - return func(yield func(MetaDataKey, error) bool) { - yield(MetaDataKey{}, fmt.Errorf("invalid list request key: %w", err)) - } - } - - if rv == 0 { - rv = math.MaxInt64 - } - - prefix := key.Prefix() - // List all keys in the prefix. - iter := d.kv.Keys(ctx, metaSection, ListOptions{ - StartKey: prefix, - EndKey: PrefixRangeEnd(prefix), - Sort: SortOrderAsc, - }) - - return func(yield func(MetaDataKey, error) bool) { - var candidateKey *MetaDataKey // The current candidate key we are iterating over - - // yieldCandidate is a helper function to yield results. - // Won't yield if the resource was last deleted. - yieldCandidate := func() bool { - if candidateKey.Action == DataActionDeleted { - // Skip because the resource was last deleted. - return true - } - return yield(*candidateKey, nil) - } - - for key, err := range iter { - if err != nil { - yield(MetaDataKey{}, err) - return - } - - metaKey, err := parseMetaDataKey(key) - if err != nil { - yield(MetaDataKey{}, err) - return - } - - if candidateKey == nil { - // Skip until we have our first candidate - if metaKey.ResourceVersion <= rv { - // New candidate found. - candidateKey = &metaKey - } - continue - } - // Should yield if either: - // - We reached the next resource. - // - We reached a resource version greater than the target resource version. - if !metaKey.SameResource(*candidateKey) || metaKey.ResourceVersion > rv { - if !yieldCandidate() { - return - } - // If we moved to a different resource and the resource version matches, make it the new candidate - if !metaKey.SameResource(*candidateKey) && metaKey.ResourceVersion <= rv { - candidateKey = &metaKey - } else { - // If we moved to a different resource and the resource version does not match, reset the candidate - candidateKey = nil - } - } else { - // Update candidate to the current key (same resource, valid version) - candidateKey = &metaKey - } - } - if candidateKey != nil { - // Yield the last selected object - if !yieldCandidate() { - return - } - } - } -} - -// Save stores a metadata object in the store. -func (d *metadataStore) Save(ctx context.Context, obj MetaDataObj) error { - if err := obj.Key.Validate(); err != nil { - return fmt.Errorf("invalid metadata key: %w", err) - } - - writer, err := d.kv.Save(ctx, metaSection, obj.Key.String()) - if err != nil { - return err - } - encoder := json.NewEncoder(writer) - if err := encoder.Encode(obj.Value); err != nil { - _ = writer.Close() - return err - } - - return writer.Close() -} - -// parseMetaDataKey parses a string key into a MetaDataKey struct -func parseMetaDataKey(key string) (MetaDataKey, error) { - parts := strings.Split(key, "/") - if len(parts) != 5 { - return MetaDataKey{}, fmt.Errorf("invalid key: %s", key) - } - - rvActionFolderParts := strings.Split(parts[4], "~") - if len(rvActionFolderParts) != 3 { - return MetaDataKey{}, fmt.Errorf("invalid key: %s", key) - } - - rv, err := strconv.ParseInt(rvActionFolderParts[0], 10, 64) - if err != nil { - return MetaDataKey{}, fmt.Errorf("invalid resource version '%s' in key %s: %w", rvActionFolderParts[0], key, err) - } - return MetaDataKey{ - Namespace: parts[2], - Group: parts[0], - Resource: parts[1], - Name: parts[3], - ResourceVersion: rv, - Action: DataAction(rvActionFolderParts[1]), - Folder: rvActionFolderParts[2], - }, nil -} - -// SameResource checks if this key represents the same resource as another key. -// It compares the identifying fields: Namespace, Group, Resource, and Name. -// ResourceVersion, Action, and Folder are ignored as they don't identify the resource itself. -func (k MetaDataKey) SameResource(other MetaDataKey) bool { - return k.Namespace == other.Namespace && - k.Group == other.Group && - k.Resource == other.Resource && - k.Name == other.Name -} diff --git a/pkg/storage/unified/resource/metadata_test.go b/pkg/storage/unified/resource/metadata_test.go deleted file mode 100644 index 2e4816a3c32..00000000000 --- a/pkg/storage/unified/resource/metadata_test.go +++ /dev/null @@ -1,1355 +0,0 @@ -package resource - -import ( - "context" - "encoding/json" - "io" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -func setupTestMetadataStore(t *testing.T) *metadataStore { - kv := setupTestKV(t) - return newMetadataStore(kv) -} - -func TestNewMetadataStore(t *testing.T) { - store := setupTestMetadataStore(t) - - assert.NotNil(t, store) -} - -func TestMetaDataKey_String(t *testing.T) { - rv := int64(56500267212345678) - key := MetaDataKey{ - Folder: "test-folder", - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: rv, - Action: DataActionCreated, - } - - expectedKey := "apps/resource/default/test-resource/56500267212345678~created~test-folder" - actualKey := key.String() - - assert.Equal(t, expectedKey, actualKey) -} - -func TestParseMetaDataKey(t *testing.T) { - rv := node.Generate() - key := "apps/resource/default/test-resource/" + rv.String() + "~" + string(DataActionCreated) + "~test-folder" - - resourceKey, err := parseMetaDataKey(key) - - require.NoError(t, err) - assert.Equal(t, "apps", resourceKey.Group) - assert.Equal(t, "resource", resourceKey.Resource) - assert.Equal(t, "default", resourceKey.Namespace) - assert.Equal(t, "test-resource", resourceKey.Name) - assert.Equal(t, rv.Int64(), resourceKey.ResourceVersion) - assert.Equal(t, DataActionCreated, resourceKey.Action) - assert.Equal(t, "test-folder", resourceKey.Folder) -} - -func TestParseMetaDataKey_InvalidKey(t *testing.T) { - tests := []struct { - name string - key string - }{ - { - name: "too few parts", - key: "apps", - }, - { - name: "invalid uuid", - key: "apps/resource/default/test-resource/invalid-uuid", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - _, err := parseMetaDataKey(tt.key) - assert.Error(t, err) - }) - } -} - -func TestMetadataStore_Save(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: node.Generate().Int64(), - Action: DataActionCreated, - } - - metadata := MetaData{ - IndexableDocument: IndexableDocument{ - Title: "This is a test resource", - Description: "This is a test resource description", - Tags: []string{"tag1", "tag2"}, - Labels: map[string]string{ - "label1": "label1", - "label2": "label2", - }, - Folder: "test-folder", - }, - } - - err := store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata, - }) - require.NoError(t, err) - // Verify in the kv store that the metadata is saved - reader, err := store.kv.Get(ctx, metaSection, key.String()) - require.NoError(t, err) - var retrivedMeta MetaData - actualData, err := io.ReadAll(reader) - require.NoError(t, err) - err = reader.Close() - require.NoError(t, err) - err = json.Unmarshal(actualData, &retrivedMeta) - require.NoError(t, err) - assert.Equal(t, metadata, retrivedMeta) -} - -func TestMetadataStore_Get(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: node.Generate().Int64(), - Action: DataActionCreated, - } - - metadata := MetaData{ - IndexableDocument: IndexableDocument{ - Title: "This is a test resource", - Description: "This is a test resource description", - Tags: []string{"tag1", "tag2"}, - Labels: map[string]string{ - "label1": "label1", - "label2": "label2", - }, - Folder: "test-folder", - }, - } - - // Save first - err := store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata, - }) - require.NoError(t, err) - - // Get it back - retrievedMetadata, err := store.Get(ctx, key) - require.NoError(t, err) - - assert.Equal(t, metadata, retrievedMetadata) -} - -func TestMetadataStore_Get_NotFound(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: node.Generate().Int64(), - Action: DataActionCreated, - } - - _, err := store.Get(ctx, key) - assert.Equal(t, ErrNotFound, err) -} - -func TestMetadataStore_GetLatestResourceKey(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - } - - // Create multiple versions with different timestamps - rv1 := node.Generate().Int64() - rv2 := node.Generate().Int64() - rv3 := node.Generate().Int64() - - // Save multiple versions (rv3 should be latest) - metadata1 := MetaData{ - IndexableDocument: IndexableDocument{ - Title: "Initial version", - }, - } - metadata2 := MetaData{ - IndexableDocument: IndexableDocument{ - Title: "Updated version", - }, - } - metadata3 := MetaData{ - IndexableDocument: IndexableDocument{ - Title: "Latest version", - }, - } - - key.ResourceVersion = rv1 - key.Action = DataActionCreated - err := store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata1, - }) - require.NoError(t, err) - - key.ResourceVersion = rv2 - key.Action = DataActionUpdated - err = store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata2, - }) - require.NoError(t, err) - - key.ResourceVersion = rv3 - key.Action = DataActionCreated - err = store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata3, - }) - require.NoError(t, err) - - // GetLatestKey should return rv3 - latestKey, err := store.GetLatestResourceKey(ctx, MetaGetRequestKey{ - Group: key.Group, - Resource: key.Resource, - Namespace: key.Namespace, - Name: key.Name, - }) - require.NoError(t, err) - - assert.Equal(t, key, latestKey) - assert.Equal(t, rv3, latestKey.ResourceVersion) -} - -func TestMetadataStore_GetLatestKey_Deleted(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: node.Generate().Int64(), - Action: DataActionDeleted, - } - - metadata := MetaData{} - - err := store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata, - }) - require.NoError(t, err) - - _, err = store.GetLatestResourceKey(ctx, MetaGetRequestKey{ - Group: key.Group, - Resource: key.Resource, - Namespace: key.Namespace, - Name: key.Name, - }) - assert.Equal(t, ErrNotFound, err) -} - -func TestMetadataStore_GetResourceKeyAtRevision(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - } - - // Create multiple versions - rv1 := node.Generate().Int64() - rv2 := node.Generate().Int64() - rv3 := node.Generate().Int64() - - metadata1 := MetaData{} - metadata2 := MetaData{} - metadata3 := MetaData{} - - key.ResourceVersion = rv1 - key.Action = DataActionCreated - err := store.Save(ctx, MetaDataObj{Key: key, Value: metadata1}) - require.NoError(t, err) - - key.ResourceVersion = rv2 - key.Action = DataActionUpdated - err = store.Save(ctx, MetaDataObj{Key: key, Value: metadata2}) - require.NoError(t, err) - - key.ResourceVersion = rv3 - key.Action = DataActionUpdated - err = store.Save(ctx, MetaDataObj{Key: key, Value: metadata3}) - require.NoError(t, err) - - // Get key at rv2 should return rv2 - metaKey, err := store.GetResourceKeyAtRevision(ctx, MetaGetRequestKey{ - Group: key.Group, - Resource: key.Resource, - Namespace: key.Namespace, - Name: key.Name, - }, rv2) - require.NoError(t, err) - - assert.Equal(t, rv2, metaKey.ResourceVersion) - assert.Equal(t, DataActionUpdated, metaKey.Action) - - // Get key at rv1 should return rv1 - metaKey, err = store.GetResourceKeyAtRevision(ctx, MetaGetRequestKey{ - Group: key.Group, - Resource: key.Resource, - Namespace: key.Namespace, - Name: key.Name, - }, rv1) - require.NoError(t, err) - - assert.Equal(t, rv1, metaKey.ResourceVersion) - assert.Equal(t, DataActionCreated, metaKey.Action) -} - -func TestMetadataStore_ListLatestResourceKeys(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - } - - // Save multiple metadata objects - rv1 := node.Generate().Int64() - rv2 := node.Generate().Int64() - - metadata1 := MetaData{} - metadata2 := MetaData{} - - key.ResourceVersion = rv1 - key.Action = DataActionCreated - err := store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata1, - }) - require.NoError(t, err) - - key.ResourceVersion = rv2 - key.Action = DataActionCreated - err = store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata2, - }) - require.NoError(t, err) - - // List latest metadata keys - resultKeys := make([]MetaDataKey, 0, 1) - for metaKey, err := range store.ListLatestResourceKeys(ctx, MetaListRequestKey{ - Group: key.Group, - Resource: key.Resource, - Namespace: key.Namespace, - Name: key.Name, - }) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - assert.Len(t, resultKeys, 1) - assert.Equal(t, key, resultKeys[0]) - assert.Equal(t, rv2, resultKeys[0].ResourceVersion) - - // Get the metadata to verify - metadata, err := store.Get(ctx, resultKeys[0]) - require.NoError(t, err) - assert.Equal(t, metadata2, metadata) -} - -func TestMetadataStore_ListResourceKeysAtRevision(t *testing.T) { - store := newMetadataStore(setupTestKV(t)) - ctx := context.Background() - - // Create multiple resources with different versions - rv1 := node.Generate().Int64() - rv2 := node.Generate().Int64() - rv3 := node.Generate().Int64() - rv4 := node.Generate().Int64() - rv5 := node.Generate().Int64() - - // Resource 1: Created at rv1, updated at rv3 - key1 := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "resource1", - ResourceVersion: rv1, - Action: DataActionCreated, - } - metadata1 := MetaData{} - err := store.Save(ctx, MetaDataObj{Key: key1, Value: metadata1}) - require.NoError(t, err) - - key1Updated := key1 - key1Updated.ResourceVersion = rv3 - key1Updated.Action = DataActionUpdated - metadata1Updated := MetaData{} - err = store.Save(ctx, MetaDataObj{Key: key1Updated, Value: metadata1Updated}) - require.NoError(t, err) - - // Resource 2: Created at rv2 - key2 := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "resource2", - ResourceVersion: rv2, - Action: DataActionCreated, - } - metadata2 := MetaData{} - err = store.Save(ctx, MetaDataObj{Key: key2, Value: metadata2}) - require.NoError(t, err) - - // Resource 3: Created at rv4 - key3 := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "resource3", - ResourceVersion: rv4, - Action: DataActionCreated, - } - metadata3 := MetaData{} - err = store.Save(ctx, MetaDataObj{Key: key3, Value: metadata3}) - require.NoError(t, err) - - // Resource 4: Created at rv2, deleted at rv5 - key4 := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "resource4", - ResourceVersion: rv2, - Action: DataActionCreated, - } - metadata4 := MetaData{} - err = store.Save(ctx, MetaDataObj{Key: key4, Value: metadata4}) - require.NoError(t, err) - - key4Deleted := key4 - key4Deleted.ResourceVersion = rv5 - key4Deleted.Action = DataActionDeleted - err = store.Save(ctx, MetaDataObj{Key: key4Deleted, Value: metadata4}) - require.NoError(t, err) - - t.Run("list at revision rv1 - should return only resource1 initial version", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 1) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, rv1) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - require.Len(t, resultKeys, 1) - assert.Equal(t, "resource1", resultKeys[0].Name) - assert.Equal(t, rv1, resultKeys[0].ResourceVersion) - assert.Equal(t, DataActionCreated, resultKeys[0].Action) - }) - - t.Run("list at revision rv2 - should return resource1, resource2 and resource4", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 3) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, rv2) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - require.Len(t, resultKeys, 3) // resource1, resource2, resource4 - names := make(map[string]int64) - for _, result := range resultKeys { - names[result.Name] = result.ResourceVersion - } - - assert.Equal(t, rv1, names["resource1"]) // Should be the original version - assert.Equal(t, rv2, names["resource2"]) - assert.Equal(t, rv2, names["resource4"]) - }) - - t.Run("list at revision rv3 - should return resource1, resource2 and resource4", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 3) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, rv3) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - require.Len(t, resultKeys, 3) // resource1 (updated), resource2, resource4 - names := make(map[string]int64) - actions := make(map[string]DataAction) - for _, result := range resultKeys { - names[result.Name] = result.ResourceVersion - actions[result.Name] = result.Action - } - - assert.Equal(t, rv3, names["resource1"]) // Should be the updated version - assert.Equal(t, DataActionUpdated, actions["resource1"]) - assert.Equal(t, rv2, names["resource2"]) - assert.Equal(t, rv2, names["resource4"]) - }) - - t.Run("list at revision rv4 - should return", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 4) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, rv4) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - require.Len(t, resultKeys, 4) // resource1 (updated), resource2, resource3, resource4 - names := make(map[string]int64) - for _, result := range resultKeys { - names[result.Name] = result.ResourceVersion - } - - assert.Equal(t, rv3, names["resource1"]) - assert.Equal(t, rv2, names["resource2"]) - assert.Equal(t, rv4, names["resource3"]) - assert.Equal(t, rv2, names["resource4"]) - }) - - t.Run("list at revision rv5 - should exclude deleted resource4", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 3) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, rv5) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - require.Len(t, resultKeys, 3) // resource1 (updated), resource2, resource3 (resource4 excluded because deleted) - names := make(map[string]bool) - for _, result := range resultKeys { - names[result.Name] = true - } - - assert.True(t, names["resource1"]) - assert.True(t, names["resource2"]) - assert.True(t, names["resource3"]) - assert.False(t, names["resource4"]) // Should be excluded because it's deleted - }) - - t.Run("list with specific resource name", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 1) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "resource1", - }, rv3) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - require.Len(t, resultKeys, 1) - assert.Equal(t, "resource1", resultKeys[0].Name) - assert.Equal(t, rv3, resultKeys[0].ResourceVersion) - assert.Equal(t, DataActionUpdated, resultKeys[0].Action) - }) - - t.Run("list at revision 0 should use MaxInt64", func(t *testing.T) { - resultKeys := make([]MetaDataKey, 0, 3) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, 0) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - // Should return all non-deleted resources at their latest versions - require.Len(t, resultKeys, 3) // resource1 (updated), resource2, resource3 - names := make(map[string]bool) - for _, result := range resultKeys { - names[result.Name] = true - } - - assert.True(t, names["resource1"]) - assert.True(t, names["resource2"]) - assert.True(t, names["resource3"]) - assert.False(t, names["resource4"]) // Excluded because deleted - }) -} - -func TestMetadataStore_ListResourceKeysAtRevision_ValidationErrors(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - tests := []struct { - name string - key MetaListRequestKey - }{ - { - name: "missing namespace", - key: MetaListRequestKey{ - Group: "apps", - Resource: "resource", - }, - }, - { - name: "missing group", - key: MetaListRequestKey{ - Namespace: "default", - Resource: "resource", - }, - }, - { - name: "missing resource", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - }, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - var resultKeys []MetaDataKey - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, tt.key, 0) { - if err != nil { - assert.Error(t, err) - return - } - resultKeys = append(resultKeys, metaKey) - } - }) - } -} - -func TestMetadataStore_ListResourceKeysAtRevision_EmptyResults(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - resultKeys := make([]MetaDataKey, 0, 1) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, 0) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - assert.Len(t, resultKeys, 0) -} - -func TestMetadataStore_ListResourceKeysAtRevision_ResourcesNewerThanRevision(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - // Create a resource with a high resource version - rv := node.Generate().Int64() - key := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "future-resource", - ResourceVersion: rv, - Action: DataActionCreated, - } - metadata := MetaData{} - err := store.Save(ctx, MetaDataObj{Key: key, Value: metadata}) - require.NoError(t, err) - - // List at a revision before the resource was created - resultKeys := make([]MetaDataKey, 0, 1) - for metaKey, err := range store.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - }, rv-1000) { - require.NoError(t, err) - resultKeys = append(resultKeys, metaKey) - } - - // Should return no results since the resource is newer than the target revision - assert.Len(t, resultKeys, 0) -} - -func TestMetaDataKey_Validate_Valid(t *testing.T) { - validKey := MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "test-folder", - } - - err := validKey.Validate() - assert.NoError(t, err) -} - -func TestMetaDataKey_Validate_ValidEdgeCases(t *testing.T) { - tests := []struct { - name string - key MetaDataKey - }{ - { - name: "valid with empty folder", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "", // empty folder should be allowed - }, - }, - { - name: "valid with single character names", - key: MetaDataKey{ - Group: "a", - Resource: "r", - Namespace: "n", - Name: "t", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "f", - }, - }, - { - name: "valid with hyphens and dots", - key: MetaDataKey{ - Group: "my-group.v1", - Resource: "my-resource.v2", - Namespace: "my-namespace.test", - Name: "my-name.test", - ResourceVersion: 123, - Action: DataActionUpdated, - Folder: "my-folder.test", - }, - }, - { - name: "valid with all action types", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionDeleted, - }, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - err := tt.key.Validate() - assert.NoError(t, err) - }) - } -} - -func TestMetaDataKey_Validate_Invalid(t *testing.T) { - tests := []struct { - name string - key MetaDataKey - wantError string - }{ - { - name: "empty group", - key: MetaDataKey{ - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "group is required", - }, - { - name: "empty resource", - key: MetaDataKey{ - Group: "apps", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "resource is required", - }, - { - name: "empty namespace", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "namespace is required", - }, - { - name: "empty name", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "name is required", - }, - { - name: "zero resource version", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 0, - Action: DataActionCreated, - }, - wantError: "resource version must be positive", - }, - { - name: "negative resource version", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: -1, - Action: DataActionCreated, - }, - wantError: "resource version must be positive", - }, - { - name: "empty action", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - }, - wantError: "action is required", - }, - { - name: "invalid name with uppercase", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "Test-Resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "name 'Test-Resource' is invalid", - }, - { - name: "invalid namespace with special chars", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default_ns", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "namespace 'default_ns' is invalid", - }, - { - name: "invalid group with uppercase", - key: MetaDataKey{ - Group: "Apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "group 'Apps' is invalid", - }, - { - name: "invalid resource with special chars", - key: MetaDataKey{ - Group: "apps", - Resource: "resource_type", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - }, - wantError: "resource 'resource_type' is invalid", - }, - { - name: "invalid folder with special chars", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "invalid_folder", - }, - wantError: "folder 'invalid_folder' is invalid", - }, - { - name: "invalid action", - key: MetaDataKey{ - Group: "apps", - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: 123, - Action: "invalid", - }, - wantError: "action 'invalid' is invalid", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - err := tt.key.Validate() - assert.Error(t, err) - assert.Contains(t, err.Error(), tt.wantError) - }) - } -} - -func TestMetaDataKey_SameResource(t *testing.T) { - baseKey := MetaDataKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "test-resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "test-folder", - } - - tests := []struct { - name string - key1 MetaDataKey - key2 MetaDataKey - expected bool - }{ - { - name: "identical keys", - key1: baseKey, - key2: baseKey, - expected: true, - }, - { - name: "same identifying fields, different resource version", - key1: baseKey, - key2: MetaDataKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "test-resource", - ResourceVersion: 456, // Different resource version - Action: DataActionUpdated, - Folder: "other-folder", - }, - expected: true, // Should still be equal as ResourceVersion, Action, and Folder don't matter - }, - { - name: "different namespace", - key1: baseKey, - key2: MetaDataKey{ - Namespace: "other-namespace", - Group: "apps", - Resource: "resource", - Name: "test-Resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "test-folder", - }, - expected: false, - }, - { - name: "different group", - key1: baseKey, - key2: MetaDataKey{ - Namespace: "default", - Group: "extensions", - Resource: "resource", - Name: "test-Resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "test-folder", - }, - expected: false, - }, - { - name: "different resource", - key1: baseKey, - key2: MetaDataKey{ - Namespace: "default", - Group: "apps", - Resource: "daemonsets", - Name: "test-Resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "test-folder", - }, - expected: false, - }, - { - name: "different name", - key1: baseKey, - key2: MetaDataKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "other-Resource", - ResourceVersion: 123, - Action: DataActionCreated, - Folder: "test-folder", - }, - expected: false, - }, - { - name: "empty keys", - key1: MetaDataKey{}, - key2: MetaDataKey{}, - expected: true, - }, - { - name: "one empty key", - key1: baseKey, - key2: MetaDataKey{}, - expected: false, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - result := tt.key1.SameResource(tt.key2) - assert.Equal(t, tt.expected, result) - - // Test symmetry: SameResource should be commutative - reverseResult := tt.key2.SameResource(tt.key1) - assert.Equal(t, result, reverseResult, "SameResource method should be commutative") - }) - } -} - -func TestMetadataStore_Save_InvalidKey(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "", // invalid: empty group - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: node.Generate().Int64(), - Action: DataActionCreated, - } - - metadata := MetaData{} - - err := store.Save(ctx, MetaDataObj{ - Key: key, - Value: metadata, - }) - assert.Error(t, err) - assert.Contains(t, err.Error(), "invalid metadata key") - assert.Contains(t, err.Error(), "group is required") -} - -func TestMetadataStore_Get_InvalidKey(t *testing.T) { - store := setupTestMetadataStore(t) - ctx := context.Background() - - key := MetaDataKey{ - Group: "", // invalid: empty group - Resource: "resource", - Namespace: "default", - Name: "test-resource", - ResourceVersion: node.Generate().Int64(), - Action: DataActionCreated, - } - - _, err := store.Get(ctx, key) - assert.Error(t, err) - assert.Contains(t, err.Error(), "invalid metadata key") - assert.Contains(t, err.Error(), "group is required") -} - -func TestMetaListRequestKey_Validate(t *testing.T) { - tests := []struct { - name string - key MetaListRequestKey - expectErr bool - }{ - { - name: "valid key with all fields", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "test-resource", - }, - expectErr: false, - }, - { - name: "valid key without name (for listing multiple resources)", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - }, - expectErr: false, - }, - { - name: "invalid key: name provided when namespace is empty", - key: MetaListRequestKey{ - Group: "apps", - Resource: "resource", - Name: "test-resource", - }, - expectErr: true, - }, - { - name: "valid key with empty namespace and no name", - key: MetaListRequestKey{ - Namespace: "", - Group: "apps", - Resource: "resource", - }, - expectErr: false, - }, - { - name: "missing group", - key: MetaListRequestKey{ - Namespace: "default", - Resource: "resource", - Name: "test-Resource", - }, - expectErr: true, - }, - { - name: "missing resource", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - Name: "test-Resource", - }, - expectErr: true, - }, - { - name: "invalid namespace", - key: MetaListRequestKey{ - Namespace: "Default", - Group: "apps", - Resource: "resource", - }, - expectErr: true, - }, - { - name: "invalid name", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "Test-Resource", - }, - expectErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - err := tt.key.Validate() - if tt.expectErr { - assert.Error(t, err) - } else { - assert.NoError(t, err) - } - }) - } -} - -func TestMetaGetRequestKey_Validate(t *testing.T) { - tests := []struct { - name string - key MetaGetRequestKey - expectErr bool - }{ - { - name: "valid key", - key: MetaGetRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "test-resource", - }, - expectErr: false, - }, - { - name: "missing namespace", - key: MetaGetRequestKey{ - Group: "apps", - Resource: "resource", - Name: "test- ", - }, - expectErr: true, - }, - { - name: "missing group", - key: MetaGetRequestKey{ - Namespace: "default", - Resource: "resource", - Name: "test- ", - }, - expectErr: true, - }, - { - name: "missing resource", - key: MetaGetRequestKey{ - Namespace: "default", - Group: "apps", - Name: "test- ", - }, - expectErr: true, - }, - { - name: "missing name", - key: MetaGetRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - }, - expectErr: true, - }, - { - name: "invalid namespace", - key: MetaGetRequestKey{ - Namespace: "Default", - Group: "apps", - Resource: "resource", - Name: "test- ", - }, - expectErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - err := tt.key.Validate() - if tt.expectErr { - assert.Error(t, err) - } else { - assert.NoError(t, err) - } - }) - } -} - -func TestMetaListRequestKey_Prefix(t *testing.T) { - tests := []struct { - name string - key MetaListRequestKey - expectedPrefix string - }{ - { - name: "full key with name", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "test-resource", - }, - expectedPrefix: "apps/resource/default/test-resource/", - }, - { - name: "key without name", - key: MetaListRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - }, - expectedPrefix: "apps/resource/default/", - }, - { - name: "key without namespace and without name", - key: MetaListRequestKey{ - Namespace: "", - Group: "apps", - Resource: "resource", - }, - expectedPrefix: "apps/resource/", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - prefix := tt.key.Prefix() - assert.Equal(t, tt.expectedPrefix, prefix) - }) - } -} - -func TestMetaGetRequestKey_Prefix(t *testing.T) { - key := MetaGetRequestKey{ - Namespace: "default", - Group: "apps", - Resource: "resource", - Name: "test- ", - } - - expectedPrefix := "apps/resource/default/test- /" - prefix := key.Prefix() - assert.Equal(t, expectedPrefix, prefix) -} diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index 54bd33ba67f..61c2d59995d 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -11,7 +11,6 @@ import ( "math/rand/v2" "net/http" "sort" - "strings" "time" "github.com/bwmarrin/snowflake" @@ -36,7 +35,6 @@ type kvStorageBackend struct { snowflake *snowflake.Node kv KV dataStore *dataStore - metaStore *metadataStore eventStore *eventStore notifier *notifier builder DocumentBuilder @@ -83,7 +81,6 @@ func NewKvStorageBackend(opts KvBackendOptions) (StorageBackend, error) { backend := &kvStorageBackend{ kv: kv, dataStore: newDataStore(kv), - metaStore: newMetadataStore(kv), eventStore: eventStore, notifier: newNotifier(eventStore, notifierOptions{}), snowflake: s, @@ -139,16 +136,14 @@ func (k *kvStorageBackend) pruneEvents(ctx context.Context, key PruningKey) erro return fmt.Errorf("invalid pruning key, all fields must be set: %+v", key) } - listKey := ListRequestKey{ + counter := 0 + // iterate over all keys for the resource and delete versions beyond the latest 20 + for datakey, err := range k.dataStore.Keys(ctx, ListRequestKey{ Namespace: key.Namespace, Group: key.Group, Resource: key.Resource, Name: key.Name, - Sort: SortOrderDesc, - } - counter := 0 - // iterate over all keys for the resource and delete versions beyond the latest 20 - for datakey, err := range k.dataStore.Keys(ctx, listKey) { + }, SortOrderDesc) { if err != nil { return err } @@ -217,10 +212,10 @@ func (k *kvStorageBackend) WriteEvent(ctx context.Context, event WriteEvent) (in case resourcepb.WatchEvent_ADDED: action = DataActionCreated // Check if resource already exists for create operations - _, err := k.metaStore.GetLatestResourceKey(ctx, MetaGetRequestKey{ - Namespace: event.Key.Namespace, + _, err := k.dataStore.GetLatestResourceKey(ctx, GetRequestKey{ Group: event.Key.Group, Resource: event.Key.Resource, + Namespace: event.Key.Namespace, Name: event.Key.Name, }) if err == nil { @@ -244,44 +239,20 @@ func (k *kvStorageBackend) WriteEvent(ctx context.Context, event WriteEvent) (in return 0, fmt.Errorf("object is nil") } - // Build the search document - doc, err := k.builder.BuildDocument(ctx, event.Key, rv, event.Value) - if err != nil { - return 0, fmt.Errorf("failed to build document: %w", err) - } - // Write the data - err = k.dataStore.Save(ctx, DataKey{ - Namespace: event.Key.Namespace, + err := k.dataStore.Save(ctx, DataKey{ Group: event.Key.Group, Resource: event.Key.Resource, + Namespace: event.Key.Namespace, Name: event.Key.Name, ResourceVersion: rv, Action: action, + Folder: obj.GetFolder(), }, bytes.NewReader(event.Value)) if err != nil { return 0, fmt.Errorf("failed to write data: %w", err) } - // Write metadata - err = k.metaStore.Save(ctx, MetaDataObj{ - Key: MetaDataKey{ - Namespace: event.Key.Namespace, - Group: event.Key.Group, - Resource: event.Key.Resource, - Name: event.Key.Name, - ResourceVersion: rv, - Action: action, - Folder: obj.GetFolder(), - }, - Value: MetaData{ - IndexableDocument: *doc, - }, - }) - if err != nil { - return 0, fmt.Errorf("failed to write metadata: %w", err) - } - // Write event err = k.eventStore.Save(ctx, Event{ Namespace: event.Key.Namespace, @@ -311,10 +282,10 @@ 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"}} } - meta, err := k.metaStore.GetResourceKeyAtRevision(ctx, MetaGetRequestKey{ - Namespace: req.Key.Namespace, + meta, err := k.dataStore.GetResourceKeyAtRevision(ctx, GetRequestKey{ Group: req.Key.Group, Resource: req.Key.Resource, + Namespace: req.Key.Namespace, Name: req.Key.Name, }, req.ResourceVersion) if errors.Is(err, ErrNotFound) { @@ -323,12 +294,13 @@ func (k *kvStorageBackend) ReadResource(ctx context.Context, req *resourcepb.Rea return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusInternalServerError, Message: err.Error()}} } data, err := k.dataStore.Get(ctx, DataKey{ - Namespace: req.Key.Namespace, Group: req.Key.Group, Resource: req.Key.Resource, + Namespace: req.Key.Namespace, Name: req.Key.Name, ResourceVersion: meta.ResourceVersion, Action: meta.Action, + Folder: meta.Folder, }) if err != nil || data == nil { return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusInternalServerError, Message: err.Error()}} @@ -369,12 +341,12 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis } // Fetch the latest objects - keys := make([]MetaDataKey, 0, min(defaultListBufferSize, req.Limit+1)) + keys := make([]DataKey, 0, min(defaultListBufferSize, req.Limit+1)) idx := 0 - for metaKey, err := range k.metaStore.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ - Namespace: req.Options.Key.Namespace, + for dataKey, err := range k.dataStore.ListResourceKeysAtRevision(ctx, ListRequestKey{ Group: req.Options.Key.Group, Resource: req.Options.Key.Resource, + Namespace: req.Options.Key.Namespace, Name: req.Options.Key.Name, }, resourceVersion) { if err != nil { @@ -385,7 +357,7 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis idx++ continue } - keys = append(keys, metaKey) + keys = append(keys, dataKey) // Only fetch the first limit items + 1 to get the next token. if len(keys) >= int(req.Limit+1) { break @@ -411,7 +383,7 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis // kvListIterator implements ListIterator for KV storage type kvListIterator struct { ctx context.Context - keys []MetaDataKey + keys []DataKey currentIndex int dataStore *dataStore listRV int64 @@ -437,14 +409,7 @@ func (i *kvListIterator) Next() bool { i.rv, i.err = i.keys[i.currentIndex].ResourceVersion, nil - data, err := i.dataStore.Get(i.ctx, DataKey{ - Namespace: i.keys[i.currentIndex].Namespace, - Group: i.keys[i.currentIndex].Group, - Resource: i.keys[i.currentIndex].Resource, - Name: i.keys[i.currentIndex].Name, - ResourceVersion: i.keys[i.currentIndex].ResourceVersion, - Action: i.keys[i.currentIndex].Action, - }) + data, err := i.dataStore.Get(i.ctx, i.keys[i.currentIndex]) if err != nil { i.err = err return false @@ -519,7 +484,7 @@ func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListReque if req.ResourceVersion <= 0 { return nil, fmt.Errorf("expecting an explicit resource version query when using Exact matching") } - var exactKeys []DataKey + exactKeys := make([]DataKey, 0, len(historyKeys)) for _, key := range historyKeys { if key.ResourceVersion == req.ResourceVersion { exactKeys = append(exactKeys, key) @@ -528,7 +493,7 @@ func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListReque return exactKeys, nil case resourcepb.ResourceVersionMatchV2_NotOlderThan: if req.ResourceVersion > 0 { - var filteredKeys []DataKey + filteredKeys := make([]DataKey, 0, len(historyKeys)) for _, key := range historyKeys { if key.ResourceVersion >= req.ResourceVersion { filteredKeys = append(filteredKeys, key) @@ -538,7 +503,7 @@ func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListReque } default: if req.ResourceVersion > 0 { - var filteredKeys []DataKey + filteredKeys := make([]DataKey, 0, len(historyKeys)) for _, key := range historyKeys { if key.ResourceVersion <= req.ResourceVersion { filteredKeys = append(filteredKeys, key) @@ -564,7 +529,7 @@ func applyLiveHistoryFilter(filteredKeys []DataKey, req *resourcepb.ListRequest) } } if latestDeleteRV > 0 { - var liveKeys []DataKey + liveKeys := make([]DataKey, 0, len(filteredKeys)) for _, key := range filteredKeys { if key.ResourceVersion > latestDeleteRV { liveKeys = append(liveKeys, key) @@ -594,7 +559,7 @@ func applyPagination(keys []DataKey, lastSeenRV int64, sortAscending bool) []Dat return keys } - var pagedKeys []DataKey + pagedKeys := make([]DataKey, 0, len(keys)) for _, key := range keys { if sortAscending && key.ResourceVersion > lastSeenRV { pagedKeys = append(pagedKeys, key) @@ -666,7 +631,7 @@ func (k *kvStorageBackend) listModifiedSinceDataStore(ctx context.Context, key N return func(yield func(*ModifiedResource, error) bool) { var lastSeenResource *ModifiedResource var lastSeenDataKey DataKey - for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{Namespace: key.Namespace, Group: key.Group, Resource: key.Resource}) { + for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{Namespace: key.Namespace, Group: key.Group, Resource: key.Resource}, SortOrderAsc) { if err != nil { yield(&ModifiedResource{}, err) return @@ -767,7 +732,14 @@ func (k *kvStorageBackend) listModifiedSinceEventStore(ctx context.Context, key } seen[evtKey.Name] = struct{}{} - value, err := k.getValueFromDataStore(ctx, DataKey(evtKey)) + value, err := k.getValueFromDataStore(ctx, DataKey{ + Group: evtKey.Group, + Resource: evtKey.Resource, + Namespace: evtKey.Namespace, + Name: evtKey.Name, + ResourceVersion: evtKey.ResourceVersion, + Action: evtKey.Action, + }) if err != nil { yield(&ModifiedResource{}, err) return @@ -820,7 +792,7 @@ func (k *kvStorageBackend) ListHistory(ctx context.Context, req *resourcepb.List Group: key.Group, Resource: key.Resource, Name: key.Name, - }) { + }, SortOrderAsc) { if err != nil { return 0, err } @@ -872,7 +844,7 @@ func (k *kvStorageBackend) ListHistory(ctx context.Context, req *resourcepb.List // processTrashEntries handles the special case of listing deleted items (trash) func (k *kvStorageBackend) processTrashEntries(ctx context.Context, req *resourcepb.ListRequest, fn func(ListIterator) error, historyKeys []DataKey, lastSeenRV int64, sortAscending bool, listRV int64) (int64, error) { // Filter to only deleted entries - var deletedKeys []DataKey + deletedKeys := make([]DataKey, 0, len(historyKeys)) for _, key := range historyKeys { if key.Action == DataActionDeleted { deletedKeys = append(deletedKeys, key) @@ -881,14 +853,14 @@ func (k *kvStorageBackend) processTrashEntries(ctx context.Context, req *resourc // Check if the resource currently exists (is live) // If it exists, don't return any trash entries - _, err := k.metaStore.GetLatestResourceKey(ctx, MetaGetRequestKey{ - Namespace: req.Options.Key.Namespace, + _, err := k.dataStore.GetLatestResourceKey(ctx, GetRequestKey{ Group: req.Options.Key.Group, Resource: req.Options.Key.Resource, + Namespace: req.Options.Key.Namespace, Name: req.Options.Key.Name, }) - var trashKeys []DataKey + trashKeys := make([]DataKey, 0, 1) if errors.Is(err, ErrNotFound) { // Resource doesn't exist currently, so we can return the latest delete // Find the latest delete event @@ -1045,12 +1017,13 @@ func (k *kvStorageBackend) WatchWriteEvents(ctx context.Context) (<-chan *Writte for event := range notifierEvents { // fetch the data dataReader, err := k.dataStore.Get(ctx, DataKey{ - Namespace: event.Namespace, Group: event.Group, Resource: event.Resource, + Namespace: event.Namespace, Name: event.Name, ResourceVersion: event.ResourceVersion, Action: event.Action, + Folder: event.Folder, }) if err != nil || dataReader == nil { k.log.Error("failed to get data for event", "error", err) @@ -1092,48 +1065,8 @@ func (k *kvStorageBackend) WatchWriteEvents(ctx context.Context) (<-chan *Writte } // GetResourceStats returns resource stats within the storage backend. -// TODO: this isn't very efficient, we should use a more efficient algorithm. func (k *kvStorageBackend) GetResourceStats(ctx context.Context, namespace string, minCount int) ([]ResourceStats, error) { - stats := make([]ResourceStats, 0) - res := make(map[string]map[string]bool) - rvs := make(map[string]int64) - - // Use datastore.Keys to get all data keys for the namespace - for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{Namespace: namespace}) { - if err != nil { - return nil, err - } - key := fmt.Sprintf("%s/%s/%s", dataKey.Namespace, dataKey.Group, dataKey.Resource) - if _, ok := res[key]; !ok { - res[key] = make(map[string]bool) - rvs[key] = 1 - } - res[key][dataKey.Name] = dataKey.Action != DataActionDeleted - rvs[key] = dataKey.ResourceVersion - } - - for key, names := range res { - parts := strings.Split(key, "/") - count := int64(0) - for _, exists := range names { - if exists { - count++ - } - } - if count <= int64(minCount) { - continue - } - stats = append(stats, ResourceStats{ - NamespacedResource: NamespacedResource{ - Namespace: parts[0], - Group: parts[1], - Resource: parts[2], - }, - Count: count, - ResourceVersion: rvs[key], - }) - } - return stats, nil + return k.dataStore.GetResourceStats(ctx, namespace, minCount) } // readAndClose reads all data from a ReadCloser and ensures it's closed, diff --git a/pkg/storage/unified/resource/storage_backend_test.go b/pkg/storage/unified/resource/storage_backend_test.go index 65cd65354d3..6f0fe9e3cba 100644 --- a/pkg/storage/unified/resource/storage_backend_test.go +++ b/pkg/storage/unified/resource/storage_backend_test.go @@ -43,7 +43,7 @@ func TestNewKvStorageBackend(t *testing.T) { assert.NotNil(t, backend) assert.NotNil(t, backend.kv) assert.NotNil(t, backend.dataStore) - assert.NotNil(t, backend.metaStore) + assert.NotNil(t, backend.eventStore) assert.NotNil(t, backend.notifier) assert.NotNil(t, backend.snowflake) @@ -126,25 +126,6 @@ func TestKvStorageBackend_WriteEvent_Success(t *testing.T) { require.NoError(t, dataReader.Close()) assert.Equal(t, objectToJSONBytes(t, testObj), dataValue) - // Verify metadata was written to metaStore - metaKey := MetaDataKey{ - Namespace: "default", - Group: "apps", - Resource: "resources", - Name: "test-resource", - ResourceVersion: rv, - Action: expectedAction, - Folder: "", - } - - m, err := backend.metaStore.Get(ctx, metaKey) - require.NoError(t, err) - require.NotNil(t, m) - require.Equal(t, "test-resource", m.Key.Name) - require.Equal(t, "default", m.Key.Namespace) - require.Equal(t, "apps", m.Key.Group) - require.Equal(t, "resources", m.Key.Resource) - // Verify event was written to eventStore eventKey := EventKey{ Namespace: "default", @@ -1258,8 +1239,7 @@ func TestKvStorageBackend_PruneEvents(t *testing.T) { Group: "apps", Resource: "resources", Name: "test-resource", - Sort: SortOrderDesc, - }) { + }, SortOrderDesc) { require.NoError(t, err) require.NotEqual(t, rv1, datakey.ResourceVersion) counter++ @@ -1325,8 +1305,7 @@ func TestKvStorageBackend_PruneEvents(t *testing.T) { Group: "apps", Resource: "resources", Name: "test-resource", - Sort: SortOrderDesc, - }) { + }, SortOrderDesc) { require.NoError(t, err) counter++ } @@ -1393,8 +1372,7 @@ func TestKvStorageBackend_PruneEvents(t *testing.T) { Group: "apps", Resource: "resources", Name: "test-resource", - Sort: SortOrderDesc, - }) { + }, SortOrderDesc) { require.NoError(t, err) counter++ }