unistore: add kv based storage backend (#107305)

* Add datastore

* too many slashes

* lint

* add metadata store

* simplify meta

* Add eventstore

* golint

* lint

* Add datastore

* too many slashes

* lint

* pr comments

* extract ParseKey

* readcloser

* remove get prefix

* use dedicated keys

* parsekey

* sameresource

* unrelated

* name

* renmae tests

* add key validation

* fix tests

* refactor a bit

* lint

* allow empty ns

* get keys instead of list

* rename the functions

* refactor yield candidate

* update test

* unistore: add LastResourceVersion to datastore

* lint

* use map string

* missing err check

* fix

* Add storage backend

* remove hasmore

* fix tests

* small refactor

* pre-alloc

* extract the folder

* lint

* refactor

* handle context canceled in ListHistory to pass the tests

* fix the resource test

* unistore: provide generic tests for the kv interface (#107443)

unistore: move the kv tests to the testing package

* Update pkg/storage/unified/resource/storage_backend_test.go

Co-authored-by: Peter Štibraný <pstibrany@gmail.com>

* address comments

* comments

* comments

* comments

* normalise the names and add helper method

* events comments

* rename function

---------

Co-authored-by: Peter Štibraný <pstibrany@gmail.com>
This commit is contained in:
Georges Chaudy
2025-07-02 10:57:37 +00:00
committed by GitHub
co-authored by Peter Štibraný
parent f0a5829eb7
commit 696657bdd1
8 changed files with 2544 additions and 212 deletions
+21
View File
@@ -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
+169 -206
View File
@@ -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)
}
})
}
@@ -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())
}
File diff suppressed because it is too large Load Diff
+490
View File
@@ -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)
})
}
+28
View File
@@ -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",
})
}
@@ -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)
@@ -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,
},
})
}