From 9a154ac15f8a925c09c23606d01037116ec98088 Mon Sep 17 00:00:00 2001 From: Georges Chaudy Date: Tue, 21 Oct 2025 21:32:03 +0200 Subject: [PATCH] kvstore: add batch delete (#112723) -add batch delete to the grpc --- pkg/storage/unified/resource/kv.go | 29 +++++++ pkg/storage/unified/testing/kv.go | 118 +++++++++++++++++++++++++++++ 2 files changed, 147 insertions(+) diff --git a/pkg/storage/unified/resource/kv.go b/pkg/storage/unified/resource/kv.go index 8700c9538aa..043ea45f695 100644 --- a/pkg/storage/unified/resource/kv.go +++ b/pkg/storage/unified/resource/kv.go @@ -53,6 +53,10 @@ type KV interface { // Delete a value Delete(ctx context.Context, section string, key string) error + // BatchDelete removes multiple keys from the store. + // Non-existent keys will be skipped silently without error. + BatchDelete(ctx context.Context, section string, keys []string) error + // UnixTimestamp returns the current time in seconds since Epoch. // This is used to ensure the server and client are not too far apart in time. UnixTimestamp(ctx context.Context) (int64, error) @@ -331,3 +335,28 @@ func IsValidKey(key string) bool { } return validKeyRegex.MatchString(key) } + +func (k *badgerKV) BatchDelete(ctx context.Context, section string, keys []string) error { + if k.db.IsClosed() { + return fmt.Errorf("database is closed") + } + + if section == "" { + return fmt.Errorf("section is required") + } + + txn := k.db.NewTransaction(true) + defer txn.Discard() + + for _, key := range keys { + keyWithSection := section + "/" + key + + // Delete the key (BadgerDB's Delete is idempotent - succeeds even if key doesn't exist) + err := txn.Delete([]byte(keyWithSection)) + if err != nil { + return err + } + } + + return txn.Commit() +} diff --git a/pkg/storage/unified/testing/kv.go b/pkg/storage/unified/testing/kv.go index 00f5a43e5e3..eab9aa9c845 100644 --- a/pkg/storage/unified/testing/kv.go +++ b/pkg/storage/unified/testing/kv.go @@ -27,6 +27,7 @@ const ( TestKVConcurrent = "concurrent operations" TestKVUnixTimestamp = "unix timestamp" TestKVBatchGet = "batch get operations" + TestKVBatchDelete = "batch delete operations" ) // NewKVFunc is a function that creates a new KV instance for testing @@ -67,6 +68,7 @@ func RunKVTest(t *testing.T, newKV NewKVFunc, opts *KVTestOptions) { {TestKVConcurrent, runTestKVConcurrent}, {TestKVUnixTimestamp, runTestKVUnixTimestamp}, {TestKVBatchGet, runTestKVBatchGet}, + {TestKVBatchDelete, runTestKVBatchDelete}, } for _, tc := range cases { @@ -673,6 +675,122 @@ func runTestKVBatchGet(t *testing.T, kv resource.KV, nsPrefix string) { }) } +func runTestKVBatchDelete(t *testing.T, kv resource.KV, nsPrefix string) { + ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second)) + section := nsPrefix + "-batchdelete" + + t.Run("batch delete existing keys", func(t *testing.T) { + // Setup test data + testData := map[string]string{ + "key1": "value1", + "key2": "value2", + "key3": "value3", + } + + // Save test data + for key, value := range testData { + saveKVHelper(t, kv, ctx, section, key, strings.NewReader(value)) + } + + // Verify keys exist before deletion + for key := range testData { + _, err := kv.Get(ctx, section, key) + require.NoError(t, err) + } + + // Batch delete all keys + keys := []string{"key1", "key2", "key3"} + err := kv.BatchDelete(ctx, section, keys) + require.NoError(t, err) + + // Verify all keys are deleted + for _, key := range keys { + _, err := kv.Get(ctx, section, key) + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + } + }) + + t.Run("batch delete with non-existent keys", func(t *testing.T) { + // Setup some test data + saveKVHelper(t, kv, ctx, section, "existing-key-1", strings.NewReader("value1")) + saveKVHelper(t, kv, ctx, section, "existing-key-2", strings.NewReader("value2")) + + // Batch delete with mix of existing and non-existent keys + keys := []string{"existing-key-1", "non-existent-1", "existing-key-2", "non-existent-2"} + err := kv.BatchDelete(ctx, section, keys) + require.NoError(t, err) + + // Verify existing keys are deleted + _, err = kv.Get(ctx, section, "existing-key-1") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + + _, err = kv.Get(ctx, section, "existing-key-2") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + }) + + t.Run("batch delete with all non-existent keys", func(t *testing.T) { + // Batch delete keys that don't exist + keys := []string{"non-existent-1", "non-existent-2", "non-existent-3"} + err := kv.BatchDelete(ctx, section, keys) + require.NoError(t, err) + }) + + t.Run("batch delete with empty keys list", func(t *testing.T) { + keys := []string{} + err := kv.BatchDelete(ctx, section, keys) + require.NoError(t, err) + }) + + t.Run("batch delete with empty section", func(t *testing.T) { + keys := []string{"some-key"} + err := kv.BatchDelete(ctx, "", keys) + assert.Error(t, err) + assert.Contains(t, err.Error(), "section is required") + }) + + t.Run("batch delete preserves other keys", func(t *testing.T) { + // Setup test data + saveKVHelper(t, kv, ctx, section, "keep-key-1", strings.NewReader("keep-value-1")) + saveKVHelper(t, kv, ctx, section, "delete-key-1", strings.NewReader("delete-value-1")) + saveKVHelper(t, kv, ctx, section, "keep-key-2", strings.NewReader("keep-value-2")) + saveKVHelper(t, kv, ctx, section, "delete-key-2", strings.NewReader("delete-value-2")) + + // Batch delete specific keys + keys := []string{"delete-key-1", "delete-key-2"} + err := kv.BatchDelete(ctx, section, keys) + require.NoError(t, err) + + // Verify deleted keys are gone + _, err = kv.Get(ctx, section, "delete-key-1") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + + _, err = kv.Get(ctx, section, "delete-key-2") + assert.Error(t, err) + assert.Equal(t, resource.ErrNotFound, err) + + // Verify kept keys still exist + reader, err := kv.Get(ctx, section, "keep-key-1") + require.NoError(t, err) + value, err := io.ReadAll(reader) + require.NoError(t, err) + assert.Equal(t, "keep-value-1", string(value)) + err = reader.Close() + require.NoError(t, err) + + reader, err = kv.Get(ctx, section, "keep-key-2") + require.NoError(t, err) + value, err = io.ReadAll(reader) + require.NoError(t, err) + assert.Equal(t, "keep-value-2", string(value)) + err = reader.Close() + require.NoError(t, err) + }) +} + // saveKVHelper is a helper function to save data to KV store using the new WriteCloser interface func saveKVHelper(t *testing.T, kv resource.KV, ctx context.Context, section, key string, value io.Reader) { t.Helper()