From bc9540fadbfb18b16bef192d0e5fe5496a4aa7e6 Mon Sep 17 00:00:00 2001 From: Georges Chaudy Date: Mon, 27 Oct 2025 12:27:31 +0100 Subject: [PATCH] kvstore: use batch delete to cleanup old events (#112737) * use batchdelete for cleaning up old events * comment --- pkg/storage/unified/resource/datastore.go | 2 +- pkg/storage/unified/resource/eventstore.go | 42 ++++++++++++---- .../unified/resource/eventstore_test.go | 50 +++++++++++++++++++ 3 files changed, 82 insertions(+), 12 deletions(-) diff --git a/pkg/storage/unified/resource/datastore.go b/pkg/storage/unified/resource/datastore.go index 8ecb878b675..1b23a84e922 100644 --- a/pkg/storage/unified/resource/datastore.go +++ b/pkg/storage/unified/resource/datastore.go @@ -364,7 +364,7 @@ func (d *dataStore) Get(ctx context.Context, key DataKey) (io.ReadCloser, error) // BatchGet retrieves multiple data objects in batches. // It returns an iterator that yields DataObj results for the given keys. -// Keys are processed in batches (default 50) to balance between efficiency and memory usage. +// Keys are processed in batches (default 50). // Non-existent entries will not appear in the result. func (d *dataStore) BatchGet(ctx context.Context, keys []DataKey) iter.Seq2[DataObj, error] { return func(yield func(DataObj, error) bool) { diff --git a/pkg/storage/unified/resource/eventstore.go b/pkg/storage/unified/resource/eventstore.go index 04a35606f81..828e3cececf 100644 --- a/pkg/storage/unified/resource/eventstore.go +++ b/pkg/storage/unified/resource/eventstore.go @@ -15,7 +15,8 @@ import ( ) const ( - eventsSection = "unified/events" + eventsSection = "unified/events" + deleteEventBatchSize = 50 ) // eventStore is a store for events. @@ -231,24 +232,43 @@ func (n *eventStore) ListSince(ctx context.Context, sinceRV int64) iter.Seq2[Eve // CleanupOldEvents deletes events older than the specified retention period. func (n *eventStore) CleanupOldEvents(ctx context.Context, cutoff time.Time) (int, error) { - deletedCount := 0 - // Keys are stored in the format of "resource_version~namespace~group~resource~name" // With a start key of "1" and an end key of the cutoff time we can get all expired events. endKey := fmt.Sprintf("%d", snowflakeFromTime(cutoff)) + + // Collect keys to delete + keysToDelete := make([]string, 0, deleteEventBatchSize) for key, err := range n.kv.Keys(ctx, eventsSection, ListOptions{StartKey: "1", EndKey: endKey}) { if err != nil { - return deletedCount, fmt.Errorf("failed to list event keys: %w", err) + return 0, fmt.Errorf("failed to list event keys: %w", err) } - - // TODO should use batch deletes here when available - if err := n.kv.Delete(ctx, eventsSection, key); err != nil { - return deletedCount, fmt.Errorf("failed to delete event key %s: %w", key, err) - } - deletedCount++ + keysToDelete = append(keysToDelete, key) } - return deletedCount, nil + // Use batch delete + if err := n.batchDelete(ctx, keysToDelete); err != nil { + return 0, fmt.Errorf("failed to batch delete events: %w", err) + } + + return len(keysToDelete), nil +} + +// batchDelete deletes multiple events in batches. +// Keys are processed in batches (default 50). +func (n *eventStore) batchDelete(ctx context.Context, keys []string) error { + for len(keys) > 0 { + batch := keys + if len(batch) > deleteEventBatchSize { + batch = batch[:deleteEventBatchSize] + } + keys = keys[len(batch):] + + if err := n.kv.BatchDelete(ctx, eventsSection, batch); err != nil { + return err + } + } + + return nil } // snowflake id with last two sections set to 0 (machine id and sequence) diff --git a/pkg/storage/unified/resource/eventstore_test.go b/pkg/storage/unified/resource/eventstore_test.go index b846536e864..5744d81c985 100644 --- a/pkg/storage/unified/resource/eventstore_test.go +++ b/pkg/storage/unified/resource/eventstore_test.go @@ -610,3 +610,53 @@ func TestEventStore_CleanupOldEvents_EmptyStore(t *testing.T) { require.NoError(t, err) assert.Equal(t, 0, deletedCount, "Should not have deleted any events from empty store") } + +func TestEventStore_BatchDelete(t *testing.T) { + ctx := context.Background() + store := setupTestEventStore(t) + + // Create multiple events (more than batch size to test batching) + eventKeys := make([]string, 75) + for i := 0; i < 75; i++ { + event := Event{ + Namespace: "default", + Group: "apps", + Resource: "deployments", + Name: "test-deployment", + ResourceVersion: int64(1000 + i), + Action: DataActionCreated, + Folder: "test-folder", + PreviousRV: int64(999 + i), + } + err := store.Save(ctx, event) + require.NoError(t, err) + + eventKeys[i] = EventKey{ + Namespace: event.Namespace, + Group: event.Group, + Resource: event.Resource, + Name: event.Name, + ResourceVersion: event.ResourceVersion, + Action: event.Action, + Folder: event.Folder, + }.String() + } + + // Batch delete all events + err := store.batchDelete(ctx, eventKeys) + require.NoError(t, err) + + // Verify all events were deleted + for i := 0; i < 75; i++ { + _, err := store.Get(ctx, EventKey{ + Namespace: "default", + Group: "apps", + Resource: "deployments", + Name: "test-deployment", + ResourceVersion: int64(1000 + i), + Action: DataActionCreated, + Folder: "test-folder", + }) + require.Error(t, err, "Event should have been deleted") + } +}