diff --git a/pkg/storage/unified/resource/event.go b/pkg/storage/unified/resource/event.go index 4a1ceb06a21..4d1edb5d27d 100644 --- a/pkg/storage/unified/resource/event.go +++ b/pkg/storage/unified/resource/event.go @@ -2,6 +2,7 @@ package resource import ( "context" + "fmt" "github.com/grafana/grafana/pkg/apimachinery/utils" "github.com/grafana/grafana/pkg/storage/unified/resourcepb" @@ -26,6 +27,26 @@ type WriteEvent struct { ObjectOld utils.GrafanaMetaAccessor } +func (e *WriteEvent) Validate() error { + if e.Object == nil { + return fmt.Errorf("object is nil") + } + + if e.Key == nil { + return fmt.Errorf("key is nil") + } + + if e.Value == nil { + return fmt.Errorf("value is nil") + } + + if e.Type == resourcepb.WatchEvent_UNKNOWN { + return fmt.Errorf("watch event type is unknown") + } + + return nil +} + // WrittenEvent is a WriteEvent reported with a resource version. type WrittenEvent struct { Type resourcepb.WatchEvent_Type diff --git a/pkg/storage/unified/resource/kv_test.go b/pkg/storage/unified/resource/kv_test.go index a999495d5f7..8ecf6a54751 100644 --- a/pkg/storage/unified/resource/kv_test.go +++ b/pkg/storage/unified/resource/kv_test.go @@ -1,15 +1,13 @@ package resource import ( - "bytes" "context" - "fmt" "io" + "strings" "testing" badger "github.com/dgraph-io/badger/v4" - "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -27,139 +25,9 @@ func setupTestBadgerDB(t *testing.T) *badger.DB { func setupTestKV(t *testing.T) KV { db := setupTestBadgerDB(t) - t.Cleanup(func() { - err := db.Close() - require.NoError(t, err) - }) return NewBadgerKV(db) } -func TestBadgerKV_Get(t *testing.T) { - db := setupTestBadgerDB(t) - - kv := NewBadgerKV(db) - ctx := context.Background() - - // Setup test data - err := db.Update(func(txn *badger.Txn) error { - return txn.Set([]byte("section/key1"), []byte("value1")) - }) - require.NoError(t, err) - - t.Run("Get existing key", func(t *testing.T) { - obj, err := kv.Get(ctx, "section", "key1") - require.NoError(t, err) - assert.Equal(t, "key1", obj.Key) - - // Read the value from the Reader - value, err := io.ReadAll(obj.Value) - require.NoError(t, err) - assert.Equal(t, []byte("value1"), value) - }) - - t.Run("Get non-existent key", func(t *testing.T) { - _, err := kv.Get(ctx, "section", "nonexistent") - assert.Error(t, err) - assert.Equal(t, ErrNotFound, err) - }) -} - -func TestBadgerKV_Save(t *testing.T) { - db := setupTestBadgerDB(t) - - kv := NewBadgerKV(db) - ctx := context.Background() - - t.Run("Save new key", func(t *testing.T) { - err := kv.Save(ctx, "section", "key1", bytes.NewReader([]byte("value1"))) - require.NoError(t, err) - - // Verify the value was saved - obj, err := kv.Get(ctx, "section", "key1") - require.NoError(t, err) - assert.Equal(t, "key1", obj.Key) - - value, err := io.ReadAll(obj.Value) - require.NoError(t, err) - assert.Equal(t, []byte("value1"), value) - }) - - t.Run("Save overwrite existing key", func(t *testing.T) { - // First save - err := kv.Save(ctx, "section", "key1", bytes.NewReader([]byte("oldvalue"))) - require.NoError(t, err) - - // Overwrite - err = kv.Save(ctx, "section", "key1", bytes.NewReader([]byte("newvalue"))) - require.NoError(t, err) - - // Verify the value was updated - obj, err := kv.Get(ctx, "section", "key1") - require.NoError(t, err) - assert.Equal(t, "key1", obj.Key) - - value, err := io.ReadAll(obj.Value) - require.NoError(t, err) - assert.Equal(t, []byte("newvalue"), value) - }) -} - -func TestBadgerKV_Delete(t *testing.T) { - db := setupTestBadgerDB(t) - - kv := NewBadgerKV(db) - ctx := context.Background() - - t.Run("Delete existing key", func(t *testing.T) { - // First create a key - err := kv.Save(ctx, "section", "key1", bytes.NewReader([]byte("value1"))) - require.NoError(t, err) - - // Delete it - err = kv.Delete(ctx, "section", "key1") - require.NoError(t, err) - - // Verify it's gone - _, err = kv.Get(ctx, "section", "key1") - assert.Error(t, err) - assert.Equal(t, ErrNotFound, err) - }) - - t.Run("Delete non-existent key", func(t *testing.T) { - err := kv.Delete(ctx, "section", "nonexistent") - assert.Error(t, err) - assert.Equal(t, ErrNotFound, err) - }) -} - -// setupIteratorTestData creates a test environment with common test data -func setupIteratorTestData(t *testing.T) (*badgerKV, context.Context) { - db := setupTestBadgerDB(t) - t.Cleanup(func() { - err := db.Close() - require.NoError(t, err) - }) - - kv := NewBadgerKV(db) - ctx := context.Background() - - // Setup test data - keys := []string{"a1", "a2", "b1", "b2", "c1"} - for _, k := range keys { - err := kv.Save(ctx, "section", k, bytes.NewReader([]byte("value"+k))) - require.NoError(t, err) - } - - return kv, ctx -} - -// iteratorTestCase represents a test case for iteration methods -type iteratorTestCase struct { - name string - options ListOptions - expectedKeys []string -} - func TestPrefixRangeEnd(t *testing.T) { require.Equal(t, "b", PrefixRangeEnd("a")) require.Equal(t, "a/c", PrefixRangeEnd("a/b")) @@ -167,95 +35,190 @@ func TestPrefixRangeEnd(t *testing.T) { require.Equal(t, "", PrefixRangeEnd("")) } -func TestBadgerKV_Keys(t *testing.T) { - for _, tc := range []iteratorTestCase{ - { - name: "all items", - options: ListOptions{}, - expectedKeys: []string{"a1", "a2", "b1", "b2", "c1"}, - }, - { - name: "with limit", - options: ListOptions{Limit: 2}, - expectedKeys: []string{"a1", "a2"}, - }, - { - name: "with range", - options: ListOptions{StartKey: "a", EndKey: "b"}, - expectedKeys: []string{"a1", "a2"}, - }, - { - name: "with prefix", - options: ListOptions{StartKey: "a", EndKey: PrefixRangeEnd("a")}, - expectedKeys: []string{"a1", "a2"}, - }, - { - name: "in descending order", - options: ListOptions{Sort: SortOrderDesc}, - expectedKeys: []string{"c1", "b2", "b1", "a2", "a1"}, - }, - { - name: "in descending order with prefix", - options: ListOptions{StartKey: "a", EndKey: PrefixRangeEnd("a"), Sort: SortOrderDesc}, - expectedKeys: []string{"a2", "a1"}, - }, - } { - t.Run("Keys "+tc.name, func(t *testing.T) { - kv, ctx := setupIteratorTestData(t) +func TestBadgerKVSmoke(t *testing.T) { + // Simple smoke test to ensure the basic badger KV implementation works + kv := setupTestKV(t) + ctx := context.Background() - var keys []string - for k, err := range kv.Keys(ctx, "section", tc.options) { - require.NoError(t, err) - keys = append(keys, k) - } - assert.Equal(t, tc.expectedKeys, keys) - }) - } + // Test unix timestamp works + timestamp, err := kv.UnixTimestamp(ctx) + require.NoError(t, err) + require.Greater(t, timestamp, int64(0)) + + // Test get non-existent key returns proper error + _, err = kv.Get(ctx, "test-section", "non-existent") + require.Error(t, err) + require.Equal(t, ErrNotFound, err) } -func TestBadgerKV_Concurrent(t *testing.T) { +func TestBadgerKV_UnderlyingStorage(t *testing.T) { + // Test internal key storage format and structure db := setupTestBadgerDB(t) - kv := NewBadgerKV(db) ctx := context.Background() - t.Run("Concurrent operations", func(t *testing.T) { - const numGoroutines = 10 - done := make(chan struct{}) + t.Run("keys are stored with section prefix", func(t *testing.T) { + section := "test-section" + key := "test-key" + value := "test-value" + expectedInternalKey := section + "/" + key - for i := 0; i < numGoroutines; i++ { - go func(i int) { - defer func() { done <- struct{}{} }() + // Save through KV interface + err := kv.Save(ctx, section, key, strings.NewReader(value)) + require.NoError(t, err) - key := fmt.Sprintf("key%d", i) - value := []byte(fmt.Sprintf("value%d", i)) + // Verify the raw key exists in badger with correct format + err = db.View(func(txn *badger.Txn) error { + item, err := txn.Get([]byte(expectedInternalKey)) + require.NoError(t, err) - // Save - err := kv.Save(ctx, "section", key, bytes.NewReader(value)) - require.NoError(t, err) + // Verify the value is correct + valueBytes, err := item.ValueCopy(nil) + require.NoError(t, err) + require.Equal(t, value, string(valueBytes)) - // Get - obj, err := kv.Get(ctx, "section", key) - require.NoError(t, err) - assert.Equal(t, key, obj.Key) + return nil + }) + require.NoError(t, err) + }) - readValue, err := io.ReadAll(obj.Value) - require.NoError(t, err) - assert.Equal(t, value, readValue) + t.Run("sections are properly isolated", func(t *testing.T) { + section1 := "section1" + section2 := "section2" + key := "same-key" + value1 := "value-from-section1" + value2 := "value-from-section2" - // Delete - err = kv.Delete(ctx, "section", key) - require.NoError(t, err) + // Save same key in different sections + err := kv.Save(ctx, section1, key, strings.NewReader(value1)) + require.NoError(t, err) + err = kv.Save(ctx, section2, key, strings.NewReader(value2)) + require.NoError(t, err) - // Verify deleted - _, err = kv.Get(ctx, "section", key) - assert.Error(t, err) - }(i) + // Verify both keys exist in badger with different internal keys + err = db.View(func(txn *badger.Txn) error { + // Check section1 key + item1, err := txn.Get([]byte(section1 + "/" + key)) + require.NoError(t, err) + value1Bytes, err := item1.ValueCopy(nil) + require.NoError(t, err) + require.Equal(t, value1, string(value1Bytes)) + + // Check section2 key + item2, err := txn.Get([]byte(section2 + "/" + key)) + require.NoError(t, err) + value2Bytes, err := item2.ValueCopy(nil) + require.NoError(t, err) + require.Equal(t, value2, string(value2Bytes)) + + return nil + }) + require.NoError(t, err) + + // Verify KV interface returns correct values for each section + obj1, err := kv.Get(ctx, section1, key) + require.NoError(t, err) + val1, err := io.ReadAll(obj1.Value) + require.NoError(t, err) + require.Equal(t, value1, string(val1)) + err = obj1.Value.Close() + require.NoError(t, err) + + obj2, err := kv.Get(ctx, section2, key) + require.NoError(t, err) + val2, err := io.ReadAll(obj2.Value) + require.NoError(t, err) + require.Equal(t, value2, string(val2)) + err = obj2.Value.Close() + require.NoError(t, err) + }) + + t.Run("delete removes correct internal key", func(t *testing.T) { + section := "delete-section" + key := "delete-key" + value := "delete-value" + internalKey := section + "/" + key + + // Save and verify it exists + err := kv.Save(ctx, section, key, strings.NewReader(value)) + require.NoError(t, err) + + // Verify it exists in badger + err = db.View(func(txn *badger.Txn) error { + _, err := txn.Get([]byte(internalKey)) + return err + }) + require.NoError(t, err) + + // Delete through KV interface + err = kv.Delete(ctx, section, key) + require.NoError(t, err) + + // Verify it's gone from badger + err = db.View(func(txn *badger.Txn) error { + _, err := txn.Get([]byte(internalKey)) + return err + }) + require.Error(t, err) + require.Equal(t, badger.ErrKeyNotFound, err) + }) + + t.Run("keys iteration respects section boundaries", func(t *testing.T) { + section1 := "alpha" + section2 := "beta" + + // Add keys to both sections + keys1 := []string{"a1", "a2", "a3"} + keys2 := []string{"b1", "b2", "b3"} + + for _, k := range keys1 { + err := kv.Save(ctx, section1, k, strings.NewReader("value"+k)) + require.NoError(t, err) + } + for _, k := range keys2 { + err := kv.Save(ctx, section2, k, strings.NewReader("value"+k)) + require.NoError(t, err) } - // Wait for all goroutines to complete - for i := 0; i < numGoroutines; i++ { - <-done + // List keys from section1 only + var foundKeys1 []string + for k, err := range kv.Keys(ctx, section1, ListOptions{}) { + require.NoError(t, err) + foundKeys1 = append(foundKeys1, k) + } + require.Equal(t, keys1, foundKeys1) + + // List keys from section2 only + var foundKeys2 []string + for k, err := range kv.Keys(ctx, section2, ListOptions{}) { + require.NoError(t, err) + foundKeys2 = append(foundKeys2, k) + } + require.Equal(t, keys2, foundKeys2) + + // Verify raw badger contains all keys with proper prefixes + var allRawKeys []string + err := db.View(func(txn *badger.Txn) error { + opts := badger.DefaultIteratorOptions + opts.PrefetchValues = false + iter := txn.NewIterator(opts) + defer iter.Close() + + for iter.Rewind(); iter.Valid(); iter.Next() { + item := iter.Item() + allRawKeys = append(allRawKeys, string(item.Key())) + } + return nil + }) + require.NoError(t, err) + + // Check that all expected internal keys exist + expectedInternalKeys := []string{ + "alpha/a1", "alpha/a2", "alpha/a3", + "beta/b1", "beta/b2", "beta/b3", + } + for _, expectedKey := range expectedInternalKeys { + require.Contains(t, allRawKeys, expectedKey, "Expected internal key %s should exist", expectedKey) } }) } diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go new file mode 100644 index 00000000000..53a898e50b2 --- /dev/null +++ b/pkg/storage/unified/resource/storage_backend.go @@ -0,0 +1,798 @@ +package resource + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "math/rand/v2" + "net/http" + "sort" + "strings" + "time" + + "github.com/bwmarrin/snowflake" + "github.com/grafana/grafana-app-sdk/logging" + "github.com/grafana/grafana/pkg/apimachinery/utils" + "github.com/grafana/grafana/pkg/storage/unified/resourcepb" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +const ( + defaultListBufferSize = 100 +) + +// Unified storage backend based on KV storage. +type kvStorageBackend struct { + snowflake *snowflake.Node + kv KV + dataStore *dataStore + metaStore *metadataStore + eventStore *eventStore + notifier *notifier + builder DocumentBuilder + log logging.Logger +} + +var _ StorageBackend = &kvStorageBackend{} + +func NewKvStorageBackend(kv KV) *kvStorageBackend { + s, err := snowflake.NewNode(rand.Int64N(1024)) + if err != nil { + panic(err) + } + eventStore := newEventStore(kv) + return &kvStorageBackend{ + kv: kv, + dataStore: newDataStore(kv), + metaStore: newMetadataStore(kv), + eventStore: eventStore, + notifier: newNotifier(eventStore, notifierOptions{}), + snowflake: s, + builder: StandardDocumentBuilder(), // For now we use the standard document builder. + log: &logging.NoOpLogger{}, // Make this configurable + } +} + +// WriteEvent writes a resource event (create/update/delete) to the storage backend. +func (k *kvStorageBackend) WriteEvent(ctx context.Context, event WriteEvent) (int64, error) { + if err := event.Validate(); err != nil { + return 0, fmt.Errorf("invalid event: %w", err) + } + rv := k.snowflake.Generate().Int64() + + // Write data. + var action DataAction + switch event.Type { + case resourcepb.WatchEvent_ADDED: + action = DataActionCreated + // Check if resource already exists for create operations + _, err := k.metaStore.GetLatestResourceKey(ctx, MetaGetRequestKey{ + Namespace: event.Key.Namespace, + Group: event.Key.Group, + Resource: event.Key.Resource, + Name: event.Key.Name, + }) + if err == nil { + // Resource exists, return already exists error + return 0, ErrResourceAlreadyExists + } + if !errors.Is(err, ErrNotFound) { + // Some other error occurred + return 0, fmt.Errorf("failed to check if resource exists: %w", err) + } + case resourcepb.WatchEvent_MODIFIED: + action = DataActionUpdated + case resourcepb.WatchEvent_DELETED: + action = DataActionDeleted + default: + return 0, fmt.Errorf("invalid event type: %d", event.Type) + } + + // 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, + Group: event.Key.Group, + Resource: event.Key.Resource, + Name: event.Key.Name, + ResourceVersion: rv, + Action: action, + }, 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: event.Object.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, + Group: event.Key.Group, + Resource: event.Key.Resource, + Name: event.Key.Name, + ResourceVersion: rv, + Action: action, + Folder: event.Object.GetFolder(), + PreviousRV: event.PreviousRV, + }) + if err != nil { + return 0, fmt.Errorf("failed to save event: %w", err) + } + return rv, nil +} + +func (k *kvStorageBackend) ReadResource(ctx context.Context, req *resourcepb.ReadRequest) *BackendReadResponse { + 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, + Group: req.Key.Group, + Resource: req.Key.Resource, + Name: req.Key.Name, + }, req.ResourceVersion) + if errors.Is(err, ErrNotFound) { + return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusNotFound, Message: "not found"}} + } else if err != nil { + 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, + Name: req.Key.Name, + ResourceVersion: meta.ResourceVersion, + Action: meta.Action, + }) + if err != nil || data == nil { + return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusInternalServerError, Message: err.Error()}} + } + value, err := readAndClose(data) + if err != nil { + return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusInternalServerError, Message: err.Error()}} + } + return &BackendReadResponse{ + Key: req.Key, + ResourceVersion: meta.ResourceVersion, + Value: value, + Folder: meta.Folder, + } +} + +// ListIterator returns an iterator for listing resources. +func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.ListRequest, cb func(ListIterator) error) (int64, error) { + if req.Options == nil || req.Options.Key == nil { + return 0, fmt.Errorf("missing options or key in ListRequest") + } + // Parse continue token if provided + offset := int64(0) + resourceVersion := req.ResourceVersion + if req.NextPageToken != "" { + token, err := GetContinueToken(req.NextPageToken) + if err != nil { + return 0, fmt.Errorf("invalid continue token: %w", err) + } + offset = token.StartOffset + resourceVersion = token.ResourceVersion + } + + // We set the listRV to the current time. + listRV := k.snowflake.Generate().Int64() + if resourceVersion > 0 { + listRV = resourceVersion + } + + // Fetch the latest objects + keys := make([]MetaDataKey, 0, min(defaultListBufferSize, req.Limit+1)) + idx := 0 + for metaKey, err := range k.metaStore.ListResourceKeysAtRevision(ctx, MetaListRequestKey{ + Namespace: req.Options.Key.Namespace, + Group: req.Options.Key.Group, + Resource: req.Options.Key.Resource, + Name: req.Options.Key.Name, + }, resourceVersion) { + if err != nil { + return 0, err + } + // Skip the first offset items. This is not efficient, but it's a simple way to implement it for now. + if idx < int(offset) { + idx++ + continue + } + keys = append(keys, metaKey) + // Only fetch the first limit items + 1 to get the next token. + if len(keys) >= int(req.Limit+1) { + break + } + } + iter := kvListIterator{ + keys: keys, + currentIndex: -1, + ctx: ctx, + listRV: listRV, + offset: offset, + limit: req.Limit + 1, // TODO: for now we need at least one more item. Fix the caller + dataStore: k.dataStore, + } + err := cb(&iter) + if err != nil { + return 0, err + } + + return listRV, nil +} + +// kvListIterator implements ListIterator for KV storage +type kvListIterator struct { + ctx context.Context + keys []MetaDataKey + currentIndex int + dataStore *dataStore + listRV int64 + offset int64 + limit int64 + + // current + rv int64 + err error + value []byte +} + +func (i *kvListIterator) Next() bool { + i.currentIndex++ + + if i.currentIndex >= len(i.keys) { + return false + } + + if int64(i.currentIndex) >= i.limit { + return false + } + + 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, + }) + if err != nil { + i.err = err + return false + } + + i.value, i.err = readAndClose(data) + if i.err != nil { + return false + } + + // increment the offset + i.offset++ + + return true +} + +func (i *kvListIterator) Error() error { + return nil +} + +func (i *kvListIterator) ContinueToken() string { + return ContinueToken{ + StartOffset: i.offset, + ResourceVersion: i.listRV, + }.String() +} + +func (i *kvListIterator) ResourceVersion() int64 { + return i.rv +} + +func (i *kvListIterator) Namespace() string { + return i.keys[i.currentIndex].Namespace +} + +func (i *kvListIterator) Name() string { + return i.keys[i.currentIndex].Name +} + +func (i *kvListIterator) Folder() string { + return i.keys[i.currentIndex].Folder +} + +func (i *kvListIterator) Value() []byte { + return i.value +} + +func validateListHistoryRequest(req *resourcepb.ListRequest) error { + if req.Options == nil || req.Options.Key == nil { + return fmt.Errorf("missing options or key in ListRequest") + } + key := req.Options.Key + if key.Group == "" { + return fmt.Errorf("group is required") + } + if key.Resource == "" { + return fmt.Errorf("resource is required") + } + if key.Namespace == "" { + return fmt.Errorf("namespace is required") + } + if key.Name == "" { + return fmt.Errorf("name is required") + } + return nil +} + +// filterHistoryKeysByVersion filters history keys based on version match criteria +func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListRequest) ([]DataKey, error) { + switch req.GetVersionMatchV2() { + case resourcepb.ResourceVersionMatchV2_Exact: + if req.ResourceVersion <= 0 { + return nil, fmt.Errorf("expecting an explicit resource version query when using Exact matching") + } + var exactKeys []DataKey + for _, key := range historyKeys { + if key.ResourceVersion == req.ResourceVersion { + exactKeys = append(exactKeys, key) + } + } + return exactKeys, nil + case resourcepb.ResourceVersionMatchV2_NotOlderThan: + if req.ResourceVersion > 0 { + var filteredKeys []DataKey + for _, key := range historyKeys { + if key.ResourceVersion >= req.ResourceVersion { + filteredKeys = append(filteredKeys, key) + } + } + return filteredKeys, nil + } + default: + if req.ResourceVersion > 0 { + var filteredKeys []DataKey + for _, key := range historyKeys { + if key.ResourceVersion <= req.ResourceVersion { + filteredKeys = append(filteredKeys, key) + } + } + return filteredKeys, nil + } + } + return historyKeys, nil +} + +// applyLiveHistoryFilter applies "live" history logic by ignoring events before the last delete +func applyLiveHistoryFilter(filteredKeys []DataKey, req *resourcepb.ListRequest) []DataKey { + useLatestDeletionAsMinRV := req.ResourceVersion == 0 && req.Source != resourcepb.ListRequest_TRASH && req.GetVersionMatchV2() != resourcepb.ResourceVersionMatchV2_Exact + if !useLatestDeletionAsMinRV { + return filteredKeys + } + + latestDeleteRV := int64(0) + for _, key := range filteredKeys { + if key.Action == DataActionDeleted && key.ResourceVersion > latestDeleteRV { + latestDeleteRV = key.ResourceVersion + } + } + if latestDeleteRV > 0 { + var liveKeys []DataKey + for _, key := range filteredKeys { + if key.ResourceVersion > latestDeleteRV { + liveKeys = append(liveKeys, key) + } + } + return liveKeys + } + return filteredKeys +} + +// sortByResourceVersion sorts the history keys based on the sortAscending flag +func sortByResourceVersion(filteredKeys []DataKey, sortAscending bool) { + if sortAscending { + sort.Slice(filteredKeys, func(i, j int) bool { + return filteredKeys[i].ResourceVersion < filteredKeys[j].ResourceVersion + }) + } else { + sort.Slice(filteredKeys, func(i, j int) bool { + return filteredKeys[i].ResourceVersion > filteredKeys[j].ResourceVersion + }) + } +} + +// applyPagination filters keys based on pagination parameters +func applyPagination(keys []DataKey, lastSeenRV int64, sortAscending bool) []DataKey { + if lastSeenRV == 0 { + return keys + } + + var pagedKeys []DataKey + for _, key := range keys { + if sortAscending && key.ResourceVersion > lastSeenRV { + pagedKeys = append(pagedKeys, key) + } else if !sortAscending && key.ResourceVersion < lastSeenRV { + pagedKeys = append(pagedKeys, key) + } + } + return pagedKeys +} + +// ListHistory is like ListIterator, but it returns the history of a resource. +func (k *kvStorageBackend) ListHistory(ctx context.Context, req *resourcepb.ListRequest, fn func(ListIterator) error) (int64, error) { + if err := validateListHistoryRequest(req); err != nil { + return 0, err + } + key := req.Options.Key + // Parse continue token if provided + lastSeenRV := int64(0) + sortAscending := req.GetVersionMatchV2() == resourcepb.ResourceVersionMatchV2_NotOlderThan + if req.NextPageToken != "" { + token, err := GetContinueToken(req.NextPageToken) + if err != nil { + return 0, fmt.Errorf("invalid continue token: %w", err) + } + lastSeenRV = token.ResourceVersion + sortAscending = token.SortAscending + } + + // Generate a new resource version for the list + listRV := k.snowflake.Generate().Int64() + + // Get all history entries by iterating through datastore keys + historyKeys := make([]DataKey, 0, min(defaultListBufferSize, req.Limit+1)) + + // Use datastore.Keys to get all data keys for this specific resource + for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{ + Namespace: key.Namespace, + Group: key.Group, + Resource: key.Resource, + Name: key.Name, + }) { + if err != nil { + return 0, err + } + historyKeys = append(historyKeys, dataKey) + } + + // Check if context has been cancelled + if ctx.Err() != nil { + return 0, ctx.Err() + } + + // Handle trash differently from regular history + if req.Source == resourcepb.ListRequest_TRASH { + return k.processTrashEntries(ctx, req, fn, historyKeys, lastSeenRV, sortAscending, listRV) + } + + // Apply filtering based on version match + filteredKeys, filterErr := filterHistoryKeysByVersion(historyKeys, req) + if filterErr != nil { + return 0, filterErr + } + + // Apply "live" history logic: ignore events before the last delete + filteredKeys = applyLiveHistoryFilter(filteredKeys, req) + + // Sort the entries if not already sorted correctly + sortByResourceVersion(filteredKeys, sortAscending) + + // Pagination: filter out items up to and including lastSeenRV + pagedKeys := applyPagination(filteredKeys, lastSeenRV, sortAscending) + + iter := kvHistoryIterator{ + keys: pagedKeys, + currentIndex: -1, + ctx: ctx, + listRV: listRV, + sortAscending: sortAscending, + dataStore: k.dataStore, + } + + err := fn(&iter) + if err != nil { + return 0, err + } + + return listRV, nil +} + +// 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 + for _, key := range historyKeys { + if key.Action == DataActionDeleted { + deletedKeys = append(deletedKeys, key) + } + } + + // 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, + Group: req.Options.Key.Group, + Resource: req.Options.Key.Resource, + Name: req.Options.Key.Name, + }) + + var trashKeys []DataKey + if errors.Is(err, ErrNotFound) { + // Resource doesn't exist currently, so we can return the latest delete + // Find the latest delete event + var latestDelete *DataKey + for _, key := range deletedKeys { + if latestDelete == nil || key.ResourceVersion > latestDelete.ResourceVersion { + latestDelete = &key + } + } + if latestDelete != nil { + trashKeys = append(trashKeys, *latestDelete) + } + } + // If err != ErrNotFound, the resource exists, so no trash entries should be returned + + // Apply version filtering + filteredKeys, err := filterHistoryKeysByVersion(trashKeys, req) + if err != nil { + return 0, err + } + + // Sort the entries + sortByResourceVersion(filteredKeys, sortAscending) + + // Pagination: filter out items up to and including lastSeenRV + pagedKeys := applyPagination(filteredKeys, lastSeenRV, sortAscending) + + iter := kvHistoryIterator{ + keys: pagedKeys, + currentIndex: -1, + ctx: ctx, + listRV: listRV, + sortAscending: sortAscending, + dataStore: k.dataStore, + } + + err = fn(&iter) + if err != nil { + return 0, err + } + + return listRV, nil +} + +// kvHistoryIterator implements ListIterator for KV storage history +type kvHistoryIterator struct { + ctx context.Context + keys []DataKey + currentIndex int + listRV int64 + sortAscending bool + dataStore *dataStore + + // current + rv int64 + err error + value []byte + folder string +} + +func (i *kvHistoryIterator) Next() bool { + i.currentIndex++ + + if i.currentIndex >= len(i.keys) { + return false + } + + key := i.keys[i.currentIndex] + i.rv = key.ResourceVersion + + // Read the value from the ReadCloser + data, err := i.dataStore.Get(i.ctx, key) + if err != nil { + i.err = err + return false + } + if data == nil { + i.err = fmt.Errorf("data is nil") + return false + } + i.value, i.err = readAndClose(data) + if i.err != nil { + return false + } + + // Extract the folder from the meta data + partial := &metav1.PartialObjectMetadata{} + err = json.Unmarshal(i.value, partial) + if err != nil { + i.err = err + return false + } + + meta, err := utils.MetaAccessor(partial) + if err != nil { + i.err = err + return false + } + i.folder = meta.GetFolder() + i.err = nil + + return true +} + +func (i *kvHistoryIterator) Error() error { + return i.err +} + +func (i *kvHistoryIterator) ContinueToken() string { + if i.currentIndex < 0 || i.currentIndex >= len(i.keys) { + return "" + } + token := ContinueToken{ + StartOffset: i.rv, + ResourceVersion: i.keys[i.currentIndex].ResourceVersion, + SortAscending: i.sortAscending, + } + return token.String() +} + +func (i *kvHistoryIterator) ResourceVersion() int64 { + return i.rv +} + +func (i *kvHistoryIterator) Namespace() string { + if i.currentIndex >= 0 && i.currentIndex < len(i.keys) { + return i.keys[i.currentIndex].Namespace + } + return "" +} + +func (i *kvHistoryIterator) Name() string { + if i.currentIndex >= 0 && i.currentIndex < len(i.keys) { + return i.keys[i.currentIndex].Name + } + return "" +} + +func (i *kvHistoryIterator) Folder() string { + return i.folder +} + +func (i *kvHistoryIterator) Value() []byte { + return i.value +} + +// WatchWriteEvents returns a channel that receives write events. +func (k *kvStorageBackend) WatchWriteEvents(ctx context.Context) (<-chan *WrittenEvent, error) { + // Create a channel to receive events + events := make(chan *WrittenEvent, 10000) // TODO: make this configurable + + notifierEvents := k.notifier.Watch(ctx, defaultWatchOptions()) + go func() { + for event := range notifierEvents { + // fetch the data + dataReader, err := k.dataStore.Get(ctx, DataKey{ + Namespace: event.Namespace, + Group: event.Group, + Resource: event.Resource, + Name: event.Name, + ResourceVersion: event.ResourceVersion, + Action: event.Action, + }) + if err != nil || dataReader == nil { + k.log.Error("failed to get data for event", "error", err) + continue + } + data, err := readAndClose(dataReader) + if err != nil { + k.log.Error("failed to read and close data for event", "error", err) + continue + } + var t resourcepb.WatchEvent_Type + switch event.Action { + case DataActionCreated: + t = resourcepb.WatchEvent_ADDED + case DataActionUpdated: + t = resourcepb.WatchEvent_MODIFIED + case DataActionDeleted: + t = resourcepb.WatchEvent_DELETED + } + + events <- &WrittenEvent{ + Key: &resourcepb.ResourceKey{ + Namespace: event.Namespace, + Group: event.Group, + Resource: event.Resource, + Name: event.Name, + }, + Type: t, + Folder: event.Folder, + Value: data, + ResourceVersion: event.ResourceVersion, + PreviousRV: event.PreviousRV, + Timestamp: event.ResourceVersion / time.Second.Nanoseconds(), // convert to seconds + } + } + close(events) + }() + return events, nil +} + +// 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 +} + +// readAndClose reads all data from a ReadCloser and ensures it's closed, +// combining any errors from both operations. +func readAndClose(r io.ReadCloser) ([]byte, error) { + data, err := io.ReadAll(r) + return data, errors.Join(err, r.Close()) +} diff --git a/pkg/storage/unified/resource/storage_backend_test.go b/pkg/storage/unified/resource/storage_backend_test.go new file mode 100644 index 00000000000..899f7b2de25 --- /dev/null +++ b/pkg/storage/unified/resource/storage_backend_test.go @@ -0,0 +1,1007 @@ +package resource + +import ( + "context" + "fmt" + "io" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + + "github.com/grafana/grafana/pkg/apimachinery/utils" + "github.com/grafana/grafana/pkg/storage/unified/resourcepb" +) + +func setupTestStorageBackend(t *testing.T) *kvStorageBackend { + kv := setupTestKV(t) + return NewKvStorageBackend(kv) +} + +func TestNewKvStorageBackend(t *testing.T) { + backend := setupTestStorageBackend(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) +} + +func TestKvStorageBackend_WriteEvent_Success(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + tests := []struct { + name string + eventType resourcepb.WatchEvent_Type + }{ + { + name: "write ADDED event", + eventType: resourcepb.WatchEvent_ADDED, + }, + { + name: "write MODIFIED event", + eventType: resourcepb.WatchEvent_MODIFIED, + }, + { + name: "write DELETED event", + eventType: resourcepb.WatchEvent_DELETED, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + testObj, err := createTestObject() + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: tt.eventType, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 100, + } + + rv, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + assert.Greater(t, rv, int64(0), "resource version should be positive") + + // Verify data was written to dataStore + var expectedAction DataAction + switch tt.eventType { + case resourcepb.WatchEvent_ADDED: + expectedAction = DataActionCreated + case resourcepb.WatchEvent_MODIFIED: + expectedAction = DataActionUpdated + case resourcepb.WatchEvent_DELETED: + expectedAction = DataActionDeleted + default: + t.Fatalf("unexpected event type: %v", tt.eventType) + } + + dataKey := DataKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + ResourceVersion: rv, + Action: expectedAction, + } + + dataReader, err := backend.dataStore.Get(ctx, dataKey) + require.NoError(t, err) + dataValue, err := io.ReadAll(dataReader) + require.NoError(t, err) + 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", + Group: "apps", + Resource: "resources", + Name: "test-resource", + ResourceVersion: rv, + } + + _, err = backend.eventStore.Get(ctx, eventKey) + require.NoError(t, err) + }) + } +} + +func TestKvStorageBackend_WriteEvent_ResourceAlreadyExists(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create a test resource first + testObj, err := createTestObject() + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + // First create should succeed + rv1, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + require.Greater(t, rv1, int64(0)) + + // Try to create the same resource again - should fail with ErrResourceAlreadyExists + writeEvent.PreviousRV = 0 // Reset previous RV to simulate a fresh create attempt + rv2, err := backend.WriteEvent(ctx, writeEvent) + require.Error(t, err) + require.Equal(t, int64(0), rv2) + require.ErrorIs(t, err, ErrResourceAlreadyExists) +} + +func TestKvStorageBackend_ReadResource_Success(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // First, write a resource to read + testObj, rv := createAndWriteTestObject(t, backend) + + // Now test reading the resource + readReq := &resourcepb.ReadRequest{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + ResourceVersion: 0, // Read latest version + } + + response := backend.ReadResource(ctx, readReq) + require.Nil(t, response.Error, "ReadResource should succeed") + require.NotNil(t, response.Key, "Response should have a key") + require.Equal(t, "test-resource", response.Key.Name) + require.Equal(t, "default", response.Key.Namespace) + require.Equal(t, "apps", response.Key.Group) + require.Equal(t, "resources", response.Key.Resource) + require.Equal(t, rv, response.ResourceVersion) + require.Equal(t, objectToJSONBytes(t, testObj), response.Value) +} + +func TestKvStorageBackend_ReadResource_SpecificVersion(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create initial version + testObj, rv1 := createAndWriteTestObject(t, backend) + + // Update the resource + testObj.Object["spec"].(map[string]any)["value"] = "updated data" + rv2, err := writeObject(t, backend, testObj, resourcepb.WatchEvent_MODIFIED, rv1) + require.NoError(t, err) + + // Read the first version specifically + readReq := &resourcepb.ReadRequest{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + ResourceVersion: rv1, + } + + response := backend.ReadResource(ctx, readReq) + require.Nil(t, response.Error, "ReadResource should succeed for specific version") + require.Equal(t, rv1, response.ResourceVersion) + + // Verify we got the original data, not the updated data + originalObj, err := createTestObject() + require.NoError(t, err) + require.Equal(t, objectToJSONBytes(t, originalObj), response.Value) + + // Read the latest version + readReq.ResourceVersion = 0 + response = backend.ReadResource(ctx, readReq) + require.Nil(t, response.Error, "ReadResource should succeed for latest version") + require.Equal(t, rv2, response.ResourceVersion) + require.Equal(t, objectToJSONBytes(t, testObj), response.Value) +} + +func TestKvStorageBackend_ReadResource_NotFound(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + readReq := &resourcepb.ReadRequest{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "nonexistent-resource", + }, + ResourceVersion: 0, + } + + response := backend.ReadResource(ctx, readReq) + require.NotNil(t, response.Error, "ReadResource should return error for nonexistent resource") + require.Equal(t, int32(404), response.Error.Code) + require.Equal(t, "not found", response.Error.Message) + require.Nil(t, response.Key) + require.Equal(t, int64(0), response.ResourceVersion) + require.Nil(t, response.Value) +} + +func TestKvStorageBackend_ReadResource_MissingKey(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + readReq := &resourcepb.ReadRequest{ + Key: nil, // Missing key + ResourceVersion: 0, + } + + response := backend.ReadResource(ctx, readReq) + require.NotNil(t, response.Error, "ReadResource should return error for missing key") + require.Equal(t, int32(400), response.Error.Code) + require.Equal(t, "missing key", response.Error.Message) + require.Nil(t, response.Key) + require.Equal(t, int64(0), response.ResourceVersion) + require.Nil(t, response.Value) +} + +func TestKvStorageBackend_ReadResource_DeletedResource(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // First, create a resource + testObj, rv1 := createAndWriteTestObject(t, backend) + + // Delete the resource + _, err := writeObject(t, backend, testObj, resourcepb.WatchEvent_DELETED, rv1) + require.NoError(t, err) + + // Try to read the latest version (should be deleted and return not found) + readReq := &resourcepb.ReadRequest{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + ResourceVersion: 0, + } + + response := backend.ReadResource(ctx, readReq) + require.NotNil(t, response.Error, "ReadResource should return not found for deleted resource") + require.Equal(t, int32(404), response.Error.Code) + require.Equal(t, "not found", response.Error.Message) + + // Try to read the original version (should still work) + readReq.ResourceVersion = rv1 + response = backend.ReadResource(ctx, readReq) + require.Nil(t, response.Error, "ReadResource should succeed for specific version before deletion") + require.Equal(t, rv1, response.ResourceVersion) + require.Equal(t, objectToJSONBytes(t, testObj), response.Value) +} + +func TestKvStorageBackend_ListIterator_Success(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create multiple test resources + resources := []struct { + name string + group string + value string + }{ + {"resource-1", "apps", "data-1"}, + {"resource-2", "apps", "data-2"}, + {"resource-3", "core", "data-3"}, + } + + for _, res := range resources { + testObj, err := createTestObjectWithName(res.name, res.group, res.value) + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: res.group, + Resource: "resources", + Name: res.name, + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + _, err = backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + } + + // Test listing all resources in "apps" group + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + }, + }, + Limit: 10, + } + + var collectedItems []struct { + name string + namespace string + resourceVersion int64 + value []byte + } + + rv, err := backend.ListIterator(ctx, listReq, func(iter ListIterator) error { + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + collectedItems = append(collectedItems, struct { + name string + namespace string + resourceVersion int64 + value []byte + }{ + name: iter.Name(), + namespace: iter.Namespace(), + resourceVersion: iter.ResourceVersion(), + value: iter.Value(), + }) + } + return iter.Error() + }) + + require.NoError(t, err) + require.Greater(t, rv, int64(0)) + require.Len(t, collectedItems, 2) // Only resources in "apps" group + + // Verify the items contain expected data + names := make([]string, len(collectedItems)) + for i, item := range collectedItems { + names[i] = item.name + require.Equal(t, "default", item.namespace) + require.Greater(t, item.resourceVersion, int64(0)) + require.NotEmpty(t, item.value) + } + require.Equal(t, []string{"resource-1", "resource-2"}, names) +} + +func TestKvStorageBackend_ListIterator_WithPagination(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create multiple test resources + for i := 1; i <= 5; i++ { + testObj, err := createTestObjectWithName(fmt.Sprintf("resource-%d", i), "apps", fmt.Sprintf("data-%d", i)) + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: fmt.Sprintf("resource-%d", i), + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + _, err = backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + } + + // First page with limit 2 + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + }, + }, + Limit: 2, + } + + var firstPageItems []string + var continueToken string + + rv, err := backend.ListIterator(ctx, listReq, func(iter ListIterator) error { + count := 0 + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + firstPageItems = append(firstPageItems, iter.Name()) + count++ + // Simulate pagination by getting continue token after limit items + if count >= int(listReq.Limit) { + continueToken = iter.ContinueToken() + break + } + } + return iter.Error() + }) + + require.NoError(t, err) + require.Greater(t, rv, int64(0)) + require.Len(t, firstPageItems, 2) + require.Equal(t, []string{"resource-1", "resource-2"}, firstPageItems) + require.NotEmpty(t, continueToken) + + // Second page using continue token + listReq.NextPageToken = continueToken + var secondPageItems []string + var continueToken2 string + + _, err = backend.ListIterator(ctx, listReq, func(iter ListIterator) error { + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + secondPageItems = append(secondPageItems, iter.Name()) + } + // Capture continue token for potential third page + continueToken2 = iter.ContinueToken() + return iter.Error() + }) + // TODO: fix the ListIterator to respect the limit. This require a change to the resource server. + require.NoError(t, err) + require.Equal(t, 3, len(secondPageItems)) + require.Equal(t, []string{"resource-3", "resource-4", "resource-5"}, secondPageItems) + require.NotEmpty(t, continueToken2) +} +func TestKvStorageBackend_ListIterator_EmptyResult(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "nonexistent", + Group: "apps", + Resource: "resources", + }, + }, + Limit: 10, + } + + var collectedItems []string + rv, err := backend.ListIterator(ctx, listReq, func(iter ListIterator) error { + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + collectedItems = append(collectedItems, iter.Name()) + } + return iter.Error() + }) + + require.NoError(t, err) + require.Greater(t, rv, int64(0)) + require.Empty(t, collectedItems) +} + +func TestKvStorageBackend_ListIterator_MissingOptions(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + tests := []struct { + name string + request *resourcepb.ListRequest + }{ + { + name: "nil options", + request: &resourcepb.ListRequest{ + Options: nil, + Limit: 10, + }, + }, + { + name: "nil key", + request: &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: nil, + }, + Limit: 10, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, err := backend.ListIterator(ctx, tt.request, func(iter ListIterator) error { + return nil + }) + require.Error(t, err) + require.Contains(t, err.Error(), "missing options or key") + }) + } +} + +func TestKvStorageBackend_ListIterator_InvalidContinueToken(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + }, + }, + Limit: 10, + NextPageToken: "invalid-token", + } + + _, err := backend.ListIterator(ctx, listReq, func(iter ListIterator) error { + return nil + }) + require.Error(t, err) + require.Contains(t, err.Error(), "invalid continue token") +} + +func TestKvStorageBackend_ListIterator_SpecificResourceVersion(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create a resource + testObj, err := createTestObjectWithName("test-resource", "apps", "initial-data") + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + rv1, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // Update the resource + testObj.Object["spec"].(map[string]any)["value"] = "updated-data" + writeEvent.Type = resourcepb.WatchEvent_MODIFIED + writeEvent.Value = objectToJSONBytes(t, testObj) + writeEvent.PreviousRV = rv1 + + _, err = backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // List at specific resource version + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + }, + }, + ResourceVersion: rv1, + Limit: 10, + } + + var collectedItems [][]byte + rv, err := backend.ListIterator(ctx, listReq, func(iter ListIterator) error { + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + collectedItems = append(collectedItems, iter.Value()) + } + return iter.Error() + }) + + require.NoError(t, err) + require.Equal(t, rv1, rv) + require.Len(t, collectedItems, 1) + + // Verify we got the original data, not the updated data + originalObj, err := createTestObjectWithName("test-resource", "apps", "initial-data") + require.NoError(t, err) + require.Equal(t, objectToJSONBytes(t, originalObj), collectedItems[0]) +} + +func TestKvStorageBackend_ListHistory_Success(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create initial resource + testObj, err := createTestObjectWithName("test-resource", "apps", "initial-data") + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + rv1, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // Update the resource + testObj.Object["spec"].(map[string]any)["value"] = "updated-data" + writeEvent.Type = resourcepb.WatchEvent_MODIFIED + writeEvent.Value = objectToJSONBytes(t, testObj) + writeEvent.PreviousRV = rv1 + + rv2, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // Update again + testObj.Object["spec"].(map[string]any)["value"] = "final-data" + writeEvent.Value = objectToJSONBytes(t, testObj) + writeEvent.PreviousRV = rv2 + + rv3, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // List the history + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + }, + Source: resourcepb.ListRequest_HISTORY, + Limit: 10, + } + + var historyItems []struct { + resourceVersion int64 + value []byte + } + + rv, err := backend.ListHistory(ctx, listReq, func(iter ListIterator) error { + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + historyItems = append(historyItems, struct { + resourceVersion int64 + value []byte + }{ + resourceVersion: iter.ResourceVersion(), + value: iter.Value(), + }) + } + return iter.Error() + }) + + require.NoError(t, err) + require.Greater(t, rv, int64(0)) + require.Len(t, historyItems, 3) // Should have all 3 versions + + // Verify the history is sorted (newest first by default) + require.Equal(t, rv3, historyItems[0].resourceVersion) + require.Equal(t, rv2, historyItems[1].resourceVersion) + require.Equal(t, rv1, historyItems[2].resourceVersion) + + // Verify the content matches expectations for all versions + finalObj, err := createTestObjectWithName("test-resource", "apps", "final-data") + require.NoError(t, err) + require.Equal(t, objectToJSONBytes(t, finalObj), historyItems[0].value) + + updatedObj, err := createTestObjectWithName("test-resource", "apps", "updated-data") + require.NoError(t, err) + require.Equal(t, objectToJSONBytes(t, updatedObj), historyItems[1].value) + + initialObj, err := createTestObjectWithName("test-resource", "apps", "initial-data") + require.NoError(t, err) + require.Equal(t, objectToJSONBytes(t, initialObj), historyItems[2].value) +} + +func TestKvStorageBackend_ListTrash_Success(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create a resource + testObj, err := createTestObjectWithName("test-resource", "apps", "test-data") + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + rv1, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // Delete the resource + writeEvent.Type = resourcepb.WatchEvent_DELETED + writeEvent.PreviousRV = rv1 + + rv2, err := backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + + // List the trash (deleted items) + listReq := &resourcepb.ListRequest{ + Options: &resourcepb.ListOptions{ + Key: &resourcepb.ResourceKey{ + Namespace: "default", + Group: "apps", + Resource: "resources", + Name: "test-resource", + }, + }, + Source: resourcepb.ListRequest_TRASH, + Limit: 10, + } + + var trashItems []struct { + name string + resourceVersion int64 + value []byte + } + + rv, err := backend.ListHistory(ctx, listReq, func(iter ListIterator) error { + for iter.Next() { + if err := iter.Error(); err != nil { + return err + } + trashItems = append(trashItems, struct { + name string + resourceVersion int64 + value []byte + }{ + name: iter.Name(), + resourceVersion: iter.ResourceVersion(), + value: iter.Value(), + }) + } + return iter.Error() + }) + + require.NoError(t, err) + require.Greater(t, rv, int64(0)) + require.Len(t, trashItems, 1) // Should have the deleted item + + // Verify the trash item + require.Equal(t, "test-resource", trashItems[0].name) + require.Equal(t, rv2, trashItems[0].resourceVersion) + require.Equal(t, objectToJSONBytes(t, testObj), trashItems[0].value) +} + +func TestKvStorageBackend_GetResourceStats_Success(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + // Create resources in different groups and namespaces + resources := []struct { + namespace string + group string + resource string + name string + }{ + {"default", "apps", "resources", "app1"}, + {"default", "apps", "resources", "app2"}, + {"default", "core", "services", "svc1"}, + {"kube-system", "apps", "resources", "system-app"}, + {"kube-system", "core", "configmaps", "config1"}, + } + + for _, res := range resources { + testObj, err := createTestObjectWithName(res.name, res.group, "test-data") + require.NoError(t, err) + + metaAccessor, err := utils.MetaAccessor(testObj) + require.NoError(t, err) + + writeEvent := WriteEvent{ + Type: resourcepb.WatchEvent_ADDED, + Key: &resourcepb.ResourceKey{ + Namespace: res.namespace, + Group: res.group, + Resource: res.resource, + Name: res.name, + }, + Value: objectToJSONBytes(t, testObj), + Object: metaAccessor, + PreviousRV: 0, + } + + _, err = backend.WriteEvent(ctx, writeEvent) + require.NoError(t, err) + } + + // Get stats for default namespace + stats, err := backend.GetResourceStats(ctx, "default", 0) + require.NoError(t, err) + require.Len(t, stats, 2) // Should have stats for 2 resource types in default namespace + + // Verify the stats contain expected resource types + resourceTypes := make(map[string]int64) + for _, stat := range stats { + key := fmt.Sprintf("%s/%s/%s", stat.Namespace, stat.Group, stat.Resource) + resourceTypes[key] = stat.Count + require.Greater(t, stat.ResourceVersion, int64(0)) + } + + require.Equal(t, int64(2), resourceTypes["default/apps/resources"]) + require.Equal(t, int64(1), resourceTypes["default/core/services"]) + + // Get stats for all namespaces (empty string) + allStats, err := backend.GetResourceStats(ctx, "", 0) + require.NoError(t, err) + require.Len(t, allStats, 4) // Should have stats for all 4 resource types across namespaces + + // Get stats with minCount filter + filteredStats, err := backend.GetResourceStats(ctx, "", 1) + require.NoError(t, err) + require.Len(t, filteredStats, 1) // Only resources in default namespace has count > 1 + + require.Equal(t, "default", filteredStats[0].Namespace) + require.Equal(t, "apps", filteredStats[0].Group) + require.Equal(t, "resources", filteredStats[0].Resource) + require.Equal(t, int64(2), filteredStats[0].Count) +} + +// createTestObject creates a test unstructured object with standard values +func createTestObject() (*unstructured.Unstructured, error) { + return createTestObjectWithName("test-resource", "apps", "test data") +} + +// objectToJSONBytes converts an unstructured object to JSON bytes +func objectToJSONBytes(t *testing.T, obj *unstructured.Unstructured) []byte { + jsonBytes, err := obj.MarshalJSON() + require.NoError(t, err) + return jsonBytes +} + +// createTestObjectWithName creates a test unstructured object with specific name, group and value +func createTestObjectWithName(name, group, value string) (*unstructured.Unstructured, error) { + u := &unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": group + "/v1", + "kind": "resource", + "metadata": map[string]any{ + "name": name, + "namespace": "default", + }, + "spec": map[string]any{ + "value": value, + }, + }, + } + return u, nil +} + +// writeObject writes an unstructured object to the backend using the provided event type and previous resource version +func writeObject(t *testing.T, backend *kvStorageBackend, obj *unstructured.Unstructured, eventType resourcepb.WatchEvent_Type, previousRV int64) (int64, error) { + metaAccessor, err := utils.MetaAccessor(obj) + require.NoError(t, err) + + // Extract resource information from the object + namespace := metaAccessor.GetNamespace() + if namespace == "" { + namespace = "default" + } + + // Extract group from apiVersion (e.g., "apps/v1" -> "apps") + apiVersion := obj.GetAPIVersion() + group := "" + if parts := strings.Split(apiVersion, "/"); len(parts) > 1 { + group = parts[0] + } + + // Use standard resource type for tests + resource := "resources" + if group == "core" { + resource = "services" + } + + writeEvent := WriteEvent{ + Type: eventType, + Key: &resourcepb.ResourceKey{ + Namespace: namespace, + Group: group, + Resource: resource, + Name: metaAccessor.GetName(), + }, + Value: objectToJSONBytes(t, obj), + Object: metaAccessor, + PreviousRV: previousRV, + } + + return backend.WriteEvent(context.Background(), writeEvent) +} + +// createAndWriteTestObject creates a basic test object and writes it to the backend +func createAndWriteTestObject(t *testing.T, backend *kvStorageBackend) (*unstructured.Unstructured, int64) { + testObj, err := createTestObject() + require.NoError(t, err) + + rv, err := writeObject(t, backend, testObj, resourcepb.WatchEvent_ADDED, 0) + require.NoError(t, err) + + return testObj, rv +} diff --git a/pkg/storage/unified/testing/kv.go b/pkg/storage/unified/testing/kv.go new file mode 100644 index 00000000000..da91a8647d2 --- /dev/null +++ b/pkg/storage/unified/testing/kv.go @@ -0,0 +1,490 @@ +package test + +import ( + "bytes" + "context" + "fmt" + "io" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/grafana/grafana/pkg/storage/unified/resource" + "github.com/grafana/grafana/pkg/util/testutil" +) + +// Test names for the KV test suite +const ( + TestKVGet = "get operations" + TestKVSave = "save operations" + TestKVDelete = "delete operations" + TestKVKeys = "keys listing" + TestKVKeysWithLimits = "keys with limits and ranges" + TestKVKeysWithSort = "keys with sorting" + TestKVConcurrent = "concurrent operations" + TestKVUnixTimestamp = "unix timestamp" +) + +// NewKVFunc is a function that creates a new KV instance for testing +type NewKVFunc func(ctx context.Context) resource.KV + +// KVTestOptions configures which tests to run +type KVTestOptions struct { + NSPrefix string // namespace prefix for isolation +} + +// GenerateRandomKVPrefix creates a random namespace prefix for test isolation +func GenerateRandomKVPrefix() string { + return fmt.Sprintf("kvtest-%d", time.Now().UnixNano()) +} + +// RunKVTest runs the KV test suite +func RunKVTest(t *testing.T, newKV NewKVFunc, opts *KVTestOptions) { + if testing.Short() { + t.Skip("skipping integration test") + } + + if opts == nil { + opts = &KVTestOptions{} + } + + if opts.NSPrefix == "" { + opts.NSPrefix = GenerateRandomKVPrefix() + } + + t.Logf("Running KV tests with namespace prefix: %s", opts.NSPrefix) + + cases := []struct { + name string + fn func(*testing.T, resource.KV, string) + }{ + {TestKVGet, runTestKVGet}, + {TestKVSave, runTestKVSave}, + {TestKVDelete, runTestKVDelete}, + {TestKVKeys, runTestKVKeys}, + {TestKVKeysWithLimits, runTestKVKeysWithLimits}, + {TestKVKeysWithSort, runTestKVKeysWithSort}, + {TestKVConcurrent, runTestKVConcurrent}, + {TestKVUnixTimestamp, runTestKVUnixTimestamp}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + tc.fn(t, newKV(context.Background()), opts.NSPrefix) + }) + } +} + +func runTestKVGet(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-get" + + t.Run("get existing key", func(t *testing.T) { + // First save a key + testValue := "test value for get" + err := kv.Save(ctx, section, "existing-key", strings.NewReader(testValue)) + require.NoError(t, err) + + // Now get it + obj, err := kv.Get(ctx, section, "existing-key") + require.NoError(t, err) + assert.Equal(t, "existing-key", obj.Key) + + // Read the value + value, err := io.ReadAll(obj.Value) + require.NoError(t, err) + assert.Equal(t, testValue, string(value)) + + // Close the value reader + err = obj.Value.Close() + require.NoError(t, err) + }) + + t.Run("get non-existent key", func(t *testing.T) { + _, err := kv.Get(ctx, section, "non-existent-key") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + }) + + t.Run("get with empty section", func(t *testing.T) { + _, err := kv.Get(ctx, "", "some-key") + assert.Error(t, err) + assert.Contains(t, err.Error(), "section is required") + }) +} + +func runTestKVSave(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-save" + + t.Run("save new key", func(t *testing.T) { + testValue := "new test value" + err := kv.Save(ctx, section, "new-key", strings.NewReader(testValue)) + require.NoError(t, err) + + // Verify it was saved + obj, err := kv.Get(ctx, section, "new-key") + require.NoError(t, err) + assert.Equal(t, "new-key", obj.Key) + + value, err := io.ReadAll(obj.Value) + require.NoError(t, err) + assert.Equal(t, testValue, string(value)) + err = obj.Value.Close() + require.NoError(t, err) + }) + + t.Run("save overwrite existing key", func(t *testing.T) { + // First save + err := kv.Save(ctx, section, "overwrite-key", strings.NewReader("old value")) + require.NoError(t, err) + + // Overwrite + newValue := "new value" + err = kv.Save(ctx, section, "overwrite-key", strings.NewReader(newValue)) + require.NoError(t, err) + + // Verify it was updated + obj, err := kv.Get(ctx, section, "overwrite-key") + require.NoError(t, err) + + value, err := io.ReadAll(obj.Value) + require.NoError(t, err) + assert.Equal(t, newValue, string(value)) + err = obj.Value.Close() + require.NoError(t, err) + }) + + t.Run("save with empty section", func(t *testing.T) { + err := kv.Save(ctx, "", "some-key", strings.NewReader("some value")) + assert.Error(t, err) + assert.Contains(t, err.Error(), "section is required") + }) + + t.Run("save binary data", func(t *testing.T) { + binaryData := []byte{0x00, 0x01, 0x02, 0x03, 0xFF, 0xFE, 0xFD} + err := kv.Save(ctx, section, "binary-key", bytes.NewReader(binaryData)) + require.NoError(t, err) + + // Verify binary data + obj, err := kv.Get(ctx, section, "binary-key") + require.NoError(t, err) + + value, err := io.ReadAll(obj.Value) + require.NoError(t, err) + assert.Equal(t, binaryData, value) + err = obj.Value.Close() + require.NoError(t, err) + }) +} + +func runTestKVDelete(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-delete" + + t.Run("delete existing key", func(t *testing.T) { + // First create a key + err := kv.Save(ctx, section, "delete-key", strings.NewReader("delete me")) + require.NoError(t, err) + + // Verify it exists + _, err = kv.Get(ctx, section, "delete-key") + require.NoError(t, err) + + // Delete it + err = kv.Delete(ctx, section, "delete-key") + require.NoError(t, err) + + // Verify it's gone + _, err = kv.Get(ctx, section, "delete-key") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + }) + + t.Run("delete non-existent key", func(t *testing.T) { + err := kv.Delete(ctx, section, "non-existent-delete-key") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + }) + + t.Run("delete with empty section", func(t *testing.T) { + err := kv.Delete(ctx, "", "some-key") + assert.Error(t, err) + assert.Contains(t, err.Error(), "section is required") + }) +} + +func runTestKVKeys(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-keys" + + // Setup test data + testKeys := []string{"a1", "a2", "b1", "b2", "c1"} + for _, key := range testKeys { + err := kv.Save(ctx, section, key, strings.NewReader("value"+key)) + require.NoError(t, err) + } + + t.Run("list all keys", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{}) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, testKeys, keys) + }) + + t.Run("list keys with empty section", func(t *testing.T) { + var keys []string + var errors []error + for k, err := range kv.Keys(ctx, "", resource.ListOptions{}) { + if err != nil { + errors = append(errors, err) + break + } + keys = append(keys, k) + } + assert.Len(t, errors, 1) + assert.Contains(t, errors[0].Error(), "section is required") + assert.Empty(t, keys) + }) +} + +func runTestKVKeysWithLimits(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-keys-limits" + + // Setup test data + testKeys := []string{"a1", "a2", "b1", "b2", "c1", "c2", "d1", "d2"} + for _, key := range testKeys { + err := kv.Save(ctx, section, key, strings.NewReader("value"+key)) + require.NoError(t, err) + } + + t.Run("keys with limit", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{Limit: 3}) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"a1", "a2", "b1"}, keys) + }) + + t.Run("keys with range", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{StartKey: "b", EndKey: "d"}) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"b1", "b2", "c1", "c2"}, keys) + }) + + t.Run("keys with prefix", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{ + StartKey: "c", + EndKey: resource.PrefixRangeEnd("c"), + }) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"c1", "c2"}, keys) + }) + + t.Run("keys with limit and range", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{ + StartKey: "a", + EndKey: "c", + Limit: 2, + }) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"a1", "a2"}, keys) + }) +} + +func runTestKVKeysWithSort(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-keys-sort" + + // Setup test data + testKeys := []string{"a1", "a2", "b1", "b2", "c1"} + for _, key := range testKeys { + err := kv.Save(ctx, section, key, strings.NewReader("value"+key)) + require.NoError(t, err) + } + + t.Run("keys in ascending order (default)", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{Sort: resource.SortOrderAsc}) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"a1", "a2", "b1", "b2", "c1"}, keys) + }) + + t.Run("keys in descending order", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{Sort: resource.SortOrderDesc}) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"c1", "b2", "b1", "a2", "a1"}, keys) + }) + + t.Run("keys descending with prefix", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{ + StartKey: "a", + EndKey: resource.PrefixRangeEnd("a"), + Sort: resource.SortOrderDesc, + }) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"a2", "a1"}, keys) + }) + + t.Run("keys descending with limit", func(t *testing.T) { + var keys []string + for k, err := range kv.Keys(ctx, section, resource.ListOptions{ + Sort: resource.SortOrderDesc, + Limit: 3, + }) { + require.NoError(t, err) + keys = append(keys, k) + } + assert.Equal(t, []string{"c1", "b2", "b1"}, keys) + }) +} + +func runTestKVConcurrent(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(60*time.Second)) + section := nsPrefix + "-concurrent" + + t.Run("concurrent save and get operations", func(t *testing.T) { + const numGoroutines = 10 + const numOperations = 20 + + done := make(chan error, numGoroutines) + + for i := 0; i < numGoroutines; i++ { + go func(goroutineID int) { + var err error + defer func() { done <- err }() + + for j := 0; j < numOperations; j++ { + key := fmt.Sprintf("concurrent-key-%d-%d", goroutineID, j) + value := fmt.Sprintf("concurrent-value-%d-%d", goroutineID, j) + + // Save + err = kv.Save(ctx, section, key, strings.NewReader(value)) + if err != nil { + return + } + + // Get immediately + obj, err := kv.Get(ctx, section, key) + if err != nil { + return + } + + readValue, err := io.ReadAll(obj.Value) + require.NoError(t, err) + err = obj.Value.Close() + require.NoError(t, err) + assert.Equal(t, value, string(readValue)) + } + }(i) + } + + // Wait for all goroutines to complete + for i := 0; i < numGoroutines; i++ { + err := <-done + require.NoError(t, err) + } + }) + + t.Run("concurrent save, delete, and list operations", func(t *testing.T) { + const numGoroutines = 5 + done := make(chan error, numGoroutines) + + for i := 0; i < numGoroutines; i++ { + go func(goroutineID int) { + var err error + defer func() { done <- err }() + + key := fmt.Sprintf("concurrent-ops-key-%d", goroutineID) + value := fmt.Sprintf("concurrent-ops-value-%d", goroutineID) + + // Save + err = kv.Save(ctx, section, key, strings.NewReader(value)) + if err != nil { + return + } + + // List to verify it exists + found := false + for k, err := range kv.Keys(ctx, section, resource.ListOptions{}) { + if err != nil { + return + } + if k == key { + found = true + break + } + } + if !found { + err = fmt.Errorf("key %s not found in list", key) + return + } + + // Delete + err = kv.Delete(ctx, section, key) + if err != nil { + return + } + + // Verify it's deleted + _, err = kv.Get(ctx, section, key) + require.ErrorIs(t, resource.ErrNotFound, err) + err = nil // Expected error, so clear it + }(i) + } + + // Wait for all goroutines to complete + for i := 0; i < numGoroutines; i++ { + err := <-done + require.NoError(t, err) + } + }) +} + +func runTestKVUnixTimestamp(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + + t.Run("unix timestamp returns reasonable value", func(t *testing.T) { + timestamp, err := kv.UnixTimestamp(ctx) + require.NoError(t, err) + + now := time.Now().Unix() + // Allow for some time difference (up to 5 seconds) + assert.InDelta(t, now, timestamp, 5) + }) + + t.Run("unix timestamp is consistent", func(t *testing.T) { + timestamp1, err := kv.UnixTimestamp(ctx) + require.NoError(t, err) + + timestamp2, err := kv.UnixTimestamp(ctx) + require.NoError(t, err) + + // Should be very close (within 1 second) + require.InDelta(t, timestamp1, timestamp2, 1) + }) +} diff --git a/pkg/storage/unified/testing/kv_test.go b/pkg/storage/unified/testing/kv_test.go new file mode 100644 index 00000000000..1e9b1a16c45 --- /dev/null +++ b/pkg/storage/unified/testing/kv_test.go @@ -0,0 +1,28 @@ +package test + +import ( + "context" + "testing" + + badger "github.com/dgraph-io/badger/v4" + "github.com/stretchr/testify/require" + + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +func TestBadgerKV(t *testing.T) { + RunKVTest(t, func(ctx context.Context) resource.KV { + opts := badger.DefaultOptions("").WithInMemory(true).WithLogger(nil) + db, err := badger.Open(opts) + require.NoError(t, err) + + t.Cleanup(func() { + err := db.Close() + require.NoError(t, err) + }) + + return resource.NewBadgerKV(db) + }, &KVTestOptions{ + NSPrefix: "badger-kv-test", + }) +} diff --git a/pkg/storage/unified/testing/storage_backend.go b/pkg/storage/unified/testing/storage_backend.go index 45e756436fc..b99a9708665 100644 --- a/pkg/storage/unified/testing/storage_backend.go +++ b/pkg/storage/unified/testing/storage_backend.go @@ -55,10 +55,6 @@ func GenerateRandomNSPrefix() string { // RunStorageBackendTest runs the storage backend test suite func RunStorageBackendTest(t *testing.T, newBackend NewBackendFunc, opts *TestOptions) { - if testing.Short() { - t.Skip("skipping integration test") - } - if opts == nil { opts = &TestOptions{} } @@ -987,10 +983,10 @@ func runTestIntegrationBackendCreateNewResource(t *testing.T, backend resource.S Key: &resourcepb.ResourceKey{ Namespace: "default", Group: "test.grafana", - Resource: "Test", + Resource: "tests", Name: "test", }, - Value: []byte(`{"apiVersion":"test.grafana/v0alpha1","kind":"Test","metadata":{"name":"test","namespace":"default"}}`), + Value: []byte(`{"apiVersion":"test.grafana/v0alpha1","kind":"Test","metadata":{"name":"test","namespace":"default","uid":"test-uid-123"}}`), } response, err := server.Create(ctx, request) diff --git a/pkg/storage/unified/testing/storage_backend_test.go b/pkg/storage/unified/testing/storage_backend_test.go new file mode 100644 index 00000000000..a1da3b82e72 --- /dev/null +++ b/pkg/storage/unified/testing/storage_backend_test.go @@ -0,0 +1,29 @@ +package test + +import ( + "context" + "testing" + + badger "github.com/dgraph-io/badger/v4" + "github.com/stretchr/testify/require" + + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +func TestBadgerKVStorageBackend(t *testing.T) { + RunStorageBackendTest(t, func(ctx context.Context) resource.StorageBackend { + opts := badger.DefaultOptions("").WithInMemory(true).WithLogger(nil) + db, err := badger.Open(opts) + require.NoError(t, err) + t.Cleanup(func() { + _ = db.Close() + }) + return resource.NewKvStorageBackend(resource.NewBadgerKV(db)) + }, &TestOptions{ + NSPrefix: "kvstorage-test", + SkipTests: map[string]bool{ + // TODO: fix these tests and remove this skip + TestBlobSupport: true, + }, + }) +}