unified-storage: add UnixTimestamp support to the sqlkv implementation (#115651)
* unified-storage: add `UnixTimestamp` support to sqlkv implementation * unified-storage: improve tests and enable all of them on sqlkv
This commit is contained in:
@@ -11,6 +11,7 @@ import (
|
||||
"iter"
|
||||
"strings"
|
||||
"text/template"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/grafana/grafana/pkg/storage/unified/sql/db"
|
||||
@@ -556,7 +557,7 @@ func (k *sqlKV) BatchDelete(ctx context.Context, section string, keys []string)
|
||||
}
|
||||
|
||||
func (k *sqlKV) UnixTimestamp(ctx context.Context) (int64, error) {
|
||||
panic("not implemented!")
|
||||
return time.Now().Unix(), nil
|
||||
}
|
||||
|
||||
func closeRows[T any](rows db.Rows, yield func(T, error) bool) {
|
||||
|
||||
+140
-124
@@ -148,13 +148,15 @@ func runTestKVGet(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
|
||||
func runTestKVSave(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second))
|
||||
nsPrefix += "-save"
|
||||
|
||||
t.Run("save new key", func(t *testing.T) {
|
||||
newKey := namespacedKey(nsPrefix, "new-key")
|
||||
testValue := "new test value"
|
||||
saveKVHelper(t, kv, ctx, testSection, "new-key", strings.NewReader(testValue))
|
||||
saveKVHelper(t, kv, ctx, testSection, newKey, strings.NewReader(testValue))
|
||||
|
||||
// Verify it was saved
|
||||
reader, err := kv.Get(ctx, testSection, "new-key")
|
||||
reader, err := kv.Get(ctx, testSection, newKey)
|
||||
require.NoError(t, err)
|
||||
|
||||
value, err := io.ReadAll(reader)
|
||||
@@ -165,15 +167,17 @@ func runTestKVSave(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
})
|
||||
|
||||
t.Run("save overwrite existing key", func(t *testing.T) {
|
||||
overwriteKey := namespacedKey(nsPrefix, "overwrite-key")
|
||||
|
||||
// First save
|
||||
saveKVHelper(t, kv, ctx, testSection, "overwrite-key", strings.NewReader("old value"))
|
||||
saveKVHelper(t, kv, ctx, testSection, overwriteKey, strings.NewReader("old value"))
|
||||
|
||||
// Overwrite
|
||||
newValue := "new value"
|
||||
saveKVHelper(t, kv, ctx, testSection, "overwrite-key", strings.NewReader(newValue))
|
||||
saveKVHelper(t, kv, ctx, testSection, overwriteKey, strings.NewReader(newValue))
|
||||
|
||||
// Verify it was updated
|
||||
reader, err := kv.Get(ctx, testSection, "overwrite-key")
|
||||
reader, err := kv.Get(ctx, testSection, overwriteKey)
|
||||
require.NoError(t, err)
|
||||
|
||||
value, err := io.ReadAll(reader)
|
||||
@@ -185,15 +189,17 @@ func runTestKVSave(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
|
||||
t.Run("save overwrite existing key (datastore)", func(t *testing.T) {
|
||||
section := "unified/data"
|
||||
overwriteKey := namespacedKey(nsPrefix, "overwrite-key")
|
||||
|
||||
// First save
|
||||
saveKVHelper(t, kv, ctx, section, "overwrite-key", strings.NewReader("old value"))
|
||||
saveKVHelper(t, kv, ctx, section, overwriteKey, strings.NewReader("old value"))
|
||||
|
||||
// Overwrite
|
||||
newValue := "new value"
|
||||
saveKVHelper(t, kv, ctx, section, "overwrite-key", strings.NewReader(newValue))
|
||||
saveKVHelper(t, kv, ctx, section, overwriteKey, strings.NewReader(newValue))
|
||||
|
||||
// Verify it was updated
|
||||
reader, err := kv.Get(ctx, section, "overwrite-key")
|
||||
reader, err := kv.Get(ctx, section, overwriteKey)
|
||||
require.NoError(t, err)
|
||||
|
||||
value, err := io.ReadAll(reader)
|
||||
@@ -210,11 +216,13 @@ func runTestKVSave(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
})
|
||||
|
||||
t.Run("save binary data", func(t *testing.T) {
|
||||
binaryKey := namespacedKey(nsPrefix, "binary-key")
|
||||
|
||||
binaryData := []byte{0x00, 0x01, 0x02, 0x03, 0xFF, 0xFE, 0xFD}
|
||||
saveKVHelper(t, kv, ctx, testSection, "binary-key", bytes.NewReader(binaryData))
|
||||
saveKVHelper(t, kv, ctx, testSection, binaryKey, bytes.NewReader(binaryData))
|
||||
|
||||
// Verify binary data
|
||||
reader, err := kv.Get(ctx, testSection, "binary-key")
|
||||
reader, err := kv.Get(ctx, testSection, binaryKey)
|
||||
require.NoError(t, err)
|
||||
|
||||
value, err := io.ReadAll(reader)
|
||||
@@ -225,11 +233,13 @@ func runTestKVSave(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
})
|
||||
|
||||
t.Run("save key with no data", func(t *testing.T) {
|
||||
emptyKey := namespacedKey(nsPrefix, "empty-key")
|
||||
|
||||
// Save a key with empty data
|
||||
saveKVHelper(t, kv, ctx, testSection, "empty-key", strings.NewReader(""))
|
||||
saveKVHelper(t, kv, ctx, testSection, emptyKey, strings.NewReader(""))
|
||||
|
||||
// Verify it was saved with empty data
|
||||
reader, err := kv.Get(ctx, testSection, "empty-key")
|
||||
reader, err := kv.Get(ctx, testSection, emptyKey)
|
||||
require.NoError(t, err)
|
||||
|
||||
value, err := io.ReadAll(reader)
|
||||
@@ -495,128 +505,134 @@ func runTestKVKeysWithSort(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
|
||||
func runTestKVConcurrent(t *testing.T, kv resource.KV, nsPrefix string) {
|
||||
ctx := testutil.NewTestContext(t, time.Now().Add(60*time.Second))
|
||||
section := nsPrefix + "-concurrent"
|
||||
nsPrefix += "-concurrent"
|
||||
|
||||
t.Run("concurrent save and get operations", func(t *testing.T) {
|
||||
const numGoroutines = 10
|
||||
const numOperations = 20
|
||||
// Test concurrent operations for both sections, as they have different behaviours
|
||||
// in the sqlkv implementation.
|
||||
for _, testSection := range []string{"unified/data", "unified/events"} {
|
||||
t.Run(testSection, func(t *testing.T) {
|
||||
t.Run("concurrent save and get operations", func(t *testing.T) {
|
||||
const numGoroutines = 10
|
||||
const numOperations = 20
|
||||
|
||||
done := make(chan error, numGoroutines)
|
||||
done := make(chan error, numGoroutines)
|
||||
|
||||
for i := 0; i < numGoroutines; i++ {
|
||||
go func(goroutineID int) {
|
||||
var err error
|
||||
defer func() { done <- err }()
|
||||
for goroutineID := range numGoroutines {
|
||||
go func() {
|
||||
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)
|
||||
for j := range numOperations {
|
||||
key := namespacedKey(nsPrefix, fmt.Sprintf("concurrent-key-%d-%d", goroutineID, j))
|
||||
value := fmt.Sprintf("concurrent-value-%d-%d", goroutineID, j)
|
||||
|
||||
// Save
|
||||
writer, err := kv.Save(ctx, section, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer func() {
|
||||
err := writer.Close()
|
||||
require.NoError(t, err)
|
||||
// Save
|
||||
writer, err := kv.Save(ctx, testSection, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer func() {
|
||||
err := writer.Close()
|
||||
require.NoError(t, err)
|
||||
}()
|
||||
_, err = io.Copy(writer, strings.NewReader(value))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = writer.Close()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Get immediately
|
||||
reader, err := kv.Get(ctx, testSection, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
readValue, err := io.ReadAll(reader)
|
||||
require.NoError(t, err)
|
||||
err = reader.Close()
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, value, string(readValue))
|
||||
}
|
||||
}()
|
||||
_, err = io.Copy(writer, strings.NewReader(value))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = writer.Close()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Get immediately
|
||||
reader, err := kv.Get(ctx, section, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
readValue, err := io.ReadAll(reader)
|
||||
// Wait for all goroutines to complete
|
||||
for range numGoroutines {
|
||||
err := <-done
|
||||
require.NoError(t, err)
|
||||
err = reader.Close()
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("concurrent save, delete, and list operations", func(t *testing.T) {
|
||||
const numGoroutines = 5
|
||||
done := make(chan error, numGoroutines)
|
||||
|
||||
for i := range numGoroutines {
|
||||
go func(goroutineID int) {
|
||||
var err error
|
||||
defer func() { done <- err }()
|
||||
|
||||
key := namespacedKey(nsPrefix, fmt.Sprintf("concurrent-ops-key-%d", goroutineID))
|
||||
value := fmt.Sprintf("concurrent-ops-value-%d", goroutineID)
|
||||
|
||||
// Save
|
||||
writer, err := kv.Save(ctx, testSection, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer func() {
|
||||
err := writer.Close()
|
||||
require.NoError(t, err)
|
||||
}()
|
||||
_, err = io.Copy(writer, strings.NewReader(value))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = writer.Close()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// List to verify it exists
|
||||
found := false
|
||||
for k, err := range kv.Keys(ctx, testSection, 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, testSection, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Verify it's deleted
|
||||
_, err = kv.Get(ctx, testSection, key)
|
||||
require.ErrorIs(t, resource.ErrNotFound, err)
|
||||
err = nil // Expected error, so clear it
|
||||
}(i)
|
||||
}
|
||||
|
||||
// Wait for all goroutines to complete
|
||||
for range numGoroutines {
|
||||
err := <-done
|
||||
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
|
||||
writer, err := kv.Save(ctx, section, key)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer func() {
|
||||
err := writer.Close()
|
||||
require.NoError(t, err)
|
||||
}()
|
||||
_, err = io.Copy(writer, strings.NewReader(value))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = writer.Close()
|
||||
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) {
|
||||
|
||||
@@ -44,11 +44,5 @@ func TestSQLKV(t *testing.T) {
|
||||
kv, err := resource.NewSQLKV(eDB)
|
||||
require.NoError(t, err)
|
||||
return kv
|
||||
}, &KVTestOptions{
|
||||
NSPrefix: "sql-kv-test",
|
||||
SkipTests: map[string]bool{
|
||||
TestKVConcurrent: true,
|
||||
TestKVUnixTimestamp: true,
|
||||
},
|
||||
})
|
||||
}, &KVTestOptions{NSPrefix: "sql-kv-test"})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user