From 861af005e0a044699bfd8b420748d8df7c709897 Mon Sep 17 00:00:00 2001 From: Will Assis Date: Wed, 14 Jan 2026 19:24:15 -0300 Subject: [PATCH] unified-storage: fix listModifiedSince in storage_backend.go and enable its tests for both badger and sqlkv --- pkg/storage/unified/resource/eventstore.go | 8 ++-- .../unified/resource/eventstore_test.go | 10 ++--- pkg/storage/unified/resource/notifier.go | 2 +- .../unified/resource/storage_backend.go | 37 +++++++++++-------- .../unified/testing/storage_backend.go | 8 ++-- .../unified/testing/storage_backend_test.go | 2 - 6 files changed, 35 insertions(+), 32 deletions(-) diff --git a/pkg/storage/unified/resource/eventstore.go b/pkg/storage/unified/resource/eventstore.go index e0c01afb550..809dd61cafa 100644 --- a/pkg/storage/unified/resource/eventstore.go +++ b/pkg/storage/unified/resource/eventstore.go @@ -184,9 +184,9 @@ func (n *eventStore) Get(ctx context.Context, key EventKey) (Event, error) { } // ListSince returns a sequence of events since the given resource version. -func (n *eventStore) ListKeysSince(ctx context.Context, sinceRV int64) iter.Seq2[string, error] { +func (n *eventStore) ListKeysSince(ctx context.Context, sinceRV int64, sortOrder SortOrder) iter.Seq2[string, error] { opts := ListOptions{ - Sort: SortOrderAsc, + Sort: sortOrder, StartKey: fmt.Sprintf("%d", sinceRV), } return func(yield func(string, error) bool) { @@ -202,9 +202,9 @@ func (n *eventStore) ListKeysSince(ctx context.Context, sinceRV int64) iter.Seq2 } } -func (n *eventStore) ListSince(ctx context.Context, sinceRV int64) iter.Seq2[Event, error] { +func (n *eventStore) ListSince(ctx context.Context, sinceRV int64, sortOrder SortOrder) iter.Seq2[Event, error] { return func(yield func(Event, error) bool) { - for evtKey, err := range n.ListKeysSince(ctx, sinceRV) { + for evtKey, err := range n.ListKeysSince(ctx, sinceRV, sortOrder) { if err != nil { yield(Event{}, err) return diff --git a/pkg/storage/unified/resource/eventstore_test.go b/pkg/storage/unified/resource/eventstore_test.go index 270db1ddd3f..8a6e27d4327 100644 --- a/pkg/storage/unified/resource/eventstore_test.go +++ b/pkg/storage/unified/resource/eventstore_test.go @@ -369,7 +369,7 @@ func testEventStoreListKeysSince(t *testing.T, ctx context.Context, store *event // List events since RV 1500 (should get events with RV 2000 and 3000) retrievedEvents := make([]string, 0, 2) - for eventKey, err := range store.ListKeysSince(ctx, 1500) { + for eventKey, err := range store.ListKeysSince(ctx, 1500, SortOrderAsc) { require.NoError(t, err) retrievedEvents = append(retrievedEvents, eventKey) } @@ -429,7 +429,7 @@ func testEventStoreListSince(t *testing.T, ctx context.Context, store *eventStor // List events since RV 1500 (should get events with RV 2000 and 3000) retrievedEvents := make([]Event, 0, 2) - for event, err := range store.ListSince(ctx, 1500) { + for event, err := range store.ListSince(ctx, 1500, SortOrderAsc) { require.NoError(t, err) retrievedEvents = append(retrievedEvents, event) } @@ -453,7 +453,7 @@ func TestEventStore_ListSince_Empty(t *testing.T) { func testEventStoreListSinceEmpty(t *testing.T, ctx context.Context, store *eventStore) { // List events when store is empty retrievedEvents := make([]Event, 0) - for event, err := range store.ListSince(ctx, 0) { + for event, err := range store.ListSince(ctx, 0, SortOrderAsc) { require.NoError(t, err) retrievedEvents = append(retrievedEvents, event) } @@ -825,7 +825,7 @@ func testListKeysSinceWithSnowflakeTime(t *testing.T, ctx context.Context, store // List events since 90 minutes ago using subtractDurationFromSnowflake sinceRV := subtractDurationFromSnowflake(snowflakeFromTime(now), 90*time.Minute) retrievedEvents := make([]string, 0) - for eventKey, err := range store.ListKeysSince(ctx, sinceRV) { + for eventKey, err := range store.ListKeysSince(ctx, sinceRV, SortOrderAsc) { require.NoError(t, err) retrievedEvents = append(retrievedEvents, eventKey) } @@ -842,7 +842,7 @@ func testListKeysSinceWithSnowflakeTime(t *testing.T, ctx context.Context, store // List events since 30 minutes ago using subtractDurationFromSnowflake sinceRV = subtractDurationFromSnowflake(snowflakeFromTime(now), 30*time.Minute) retrievedEvents = make([]string, 0) - for eventKey, err := range store.ListKeysSince(ctx, sinceRV) { + for eventKey, err := range store.ListKeysSince(ctx, sinceRV, SortOrderAsc) { require.NoError(t, err) retrievedEvents = append(retrievedEvents, eventKey) } diff --git a/pkg/storage/unified/resource/notifier.go b/pkg/storage/unified/resource/notifier.go index 12cefa5add5..3ecf8b81044 100644 --- a/pkg/storage/unified/resource/notifier.go +++ b/pkg/storage/unified/resource/notifier.go @@ -119,7 +119,7 @@ func (n *pollingNotifier) Watch(ctx context.Context, opts watchOptions) <-chan E return case <-time.After(currentInterval): foundEvents := false - for evt, err := range n.eventStore.ListSince(ctx, subtractDurationFromSnowflake(lastRV, opts.LookbackPeriod)) { + for evt, err := range n.eventStore.ListSince(ctx, subtractDurationFromSnowflake(lastRV, opts.LookbackPeriod), SortOrderAsc) { if err != nil { n.log.Error("Failed to list events since", "error", err) continue diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index f285daf0cb8..ae5824acbe1 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -801,8 +801,20 @@ func (k *kvStorageBackend) ListModifiedSince(ctx context.Context, key Namespaced } } - // Generate a new resource version for the list - listRV := k.snowflake.Generate().Int64() + latestEvent, err := k.eventStore.LastEventKey(ctx) + if err != nil { + if errors.Is(err, ErrNotFound) { + return sinceRv, func(yield func(*ModifiedResource, error) bool) { /* nothing to return */ } + } + + return 0, func(yield func(*ModifiedResource, error) bool) { + yield(nil, fmt.Errorf("error trying to retrieve last event key: %s", err)) + } + } + + if latestEvent.ResourceVersion == sinceRv { + return sinceRv, func(yield func(*ModifiedResource, error) bool) { /* nothing to return */ } + } // Check if sinceRv is older than 1 hour sinceRvTimestamp := snowflake.ID(sinceRv).Time() @@ -811,11 +823,11 @@ func (k *kvStorageBackend) ListModifiedSince(ctx context.Context, key Namespaced if sinceRvAge > time.Hour { k.log.Debug("ListModifiedSince using data store", "sinceRv", sinceRv, "sinceRvAge", sinceRvAge) - return listRV, k.listModifiedSinceDataStore(ctx, key, sinceRv) + return latestEvent.ResourceVersion, k.listModifiedSinceDataStore(ctx, key, sinceRv) } k.log.Debug("ListModifiedSince using event store", "sinceRv", sinceRv, "sinceRvAge", sinceRvAge) - return listRV, k.listModifiedSinceEventStore(ctx, key, sinceRv) + return latestEvent.ResourceVersion, k.listModifiedSinceEventStore(ctx, key, sinceRv) } func convertEventType(action DataAction) resourcepb.WatchEvent_Type { @@ -916,9 +928,9 @@ func (k *kvStorageBackend) listModifiedSinceDataStore(ctx context.Context, key N func (k *kvStorageBackend) listModifiedSinceEventStore(ctx context.Context, key NamespacedResource, sinceRv int64) iter.Seq2[*ModifiedResource, error] { return func(yield func(*ModifiedResource, error) bool) { - // store all events ordered by RV for the given tenant here - eventKeys := make([]EventKey, 0) - for evtKeyStr, err := range k.eventStore.ListKeysSince(ctx, subtractDurationFromSnowflake(sinceRv, defaultLookbackPeriod)) { + // we only care about the latest revision of every resource in the list + seen := make(map[string]struct{}) + for evtKeyStr, err := range k.eventStore.ListKeysSince(ctx, subtractDurationFromSnowflake(sinceRv, defaultLookbackPeriod), SortOrderDesc) { if err != nil { yield(&ModifiedResource{}, err) return @@ -938,18 +950,11 @@ func (k *kvStorageBackend) listModifiedSinceEventStore(ctx context.Context, key continue } - eventKeys = append(eventKeys, evtKey) - } - - // we only care about the latest revision of every resource in the list - seen := make(map[string]struct{}) - for i := len(eventKeys) - 1; i >= 0; i -= 1 { - evtKey := eventKeys[i] if _, ok := seen[evtKey.Name]; ok { continue } - seen[evtKey.Name] = struct{}{} + seen[evtKey.Name] = struct{}{} value, err := k.getValueFromDataStore(ctx, DataKey(evtKey)) if err != nil { yield(&ModifiedResource{}, err) @@ -1307,7 +1312,7 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings if setting.RebuildCollection { for _, key := range setting.Collection { events := make([]string, 0) - for evtKeyStr, err := range b.eventStore.ListKeysSince(ctx, 1) { + for evtKeyStr, err := range b.eventStore.ListKeysSince(ctx, 1, SortOrderAsc) { if err != nil { b.log.Error("failed to list event: %s", err) return rsp diff --git a/pkg/storage/unified/testing/storage_backend.go b/pkg/storage/unified/testing/storage_backend.go index 29279a16e97..dcf1d94521a 100644 --- a/pkg/storage/unified/testing/storage_backend.go +++ b/pkg/storage/unified/testing/storage_backend.go @@ -555,7 +555,7 @@ func runTestIntegrationBackendListModifiedSince(t *testing.T, backend resource.S Resource: "resource", } latestRv, seq := backend.ListModifiedSince(ctx, key, rvCreated) - require.Greater(t, latestRv, rvCreated) + require.Equal(t, latestRv, rvDeleted) counter := 0 for res, err := range seq { @@ -629,11 +629,11 @@ func runTestIntegrationBackendListModifiedSince(t *testing.T, backend resource.S rvCreated3, _ := writeEvent(ctx, backend, "bItem", resourcepb.WatchEvent_ADDED, WithNamespace(ns)) latestRv, seq := backend.ListModifiedSince(ctx, key, rvCreated1-1) - require.Greater(t, latestRv, rvCreated3) + require.Equal(t, latestRv, rvCreated3) counter := 0 - names := []string{"aItem", "bItem", "cItem"} - rvs := []int64{rvCreated2, rvCreated3, rvCreated1} + names := []string{"bItem", "aItem", "cItem"} + rvs := []int64{rvCreated3, rvCreated2, rvCreated1} for res, err := range seq { require.NoError(t, err) require.Equal(t, key.Namespace, res.Key.Namespace) diff --git a/pkg/storage/unified/testing/storage_backend_test.go b/pkg/storage/unified/testing/storage_backend_test.go index 96b9d5716db..56f492c64fa 100644 --- a/pkg/storage/unified/testing/storage_backend_test.go +++ b/pkg/storage/unified/testing/storage_backend_test.go @@ -30,7 +30,6 @@ func TestBadgerKVStorageBackend(t *testing.T) { SkipTests: map[string]bool{ // TODO: fix these tests and remove this skip TestBlobSupport: true, - TestListModifiedSince: true, // Badger does not support bulk import yet. TestGetResourceLastImportTime: true, }, @@ -42,7 +41,6 @@ func TestIntegrationSQLKVStorageBackend(t *testing.T) { skipTests := map[string]bool{ TestBlobSupport: true, - TestListModifiedSince: true, TestGetResourceLastImportTime: true, }