diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index be9e516f90e..d5c95ac63b6 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -58,6 +58,7 @@ type kvStorageBackend struct { bulkLock *BulkLock dataStore *dataStore eventStore *eventStore + metadataStore *metadataStore notifier *notifier builder DocumentBuilder log logging.Logger @@ -107,6 +108,7 @@ func NewKVStorageBackend(opts KVBackendOptions) (StorageBackend, error) { bulkLock: NewBulkLock(), dataStore: newDataStore(kv), eventStore: eventStore, + metadataStore: newMetadataStore(kv), notifier: newNotifier(eventStore, notifierOptions{}), snowflake: s, builder: StandardDocumentBuilder(), // For now we use the standard document builder. @@ -1235,23 +1237,65 @@ func (k *kvStorageBackend) GetResourceStats(ctx context.Context, nsr NamespacedR func (k *kvStorageBackend) GetResourceLastImportTimes(ctx context.Context) iter.Seq2[ResourceLastImportTime, error] { return func(yield func(ResourceLastImportTime, error) bool) { - yield(ResourceLastImportTime{}, fmt.Errorf("not implemented")) + for metadata, err := range k.metadataStore.GetAll(ctx) { + if err != nil { + yield(ResourceLastImportTime{}, err) + return + } + + if metadata.LastImportTime.IsZero() { + continue + } + + // TODO clear LastImportTime from metadata if > lastImportTimeMaxAge? + + if !yield(ResourceLastImportTime{ + NamespacedResource: NamespacedResource{ + Namespace: metadata.Namespace, + Group: metadata.Group, + Resource: metadata.Resource, + }, + LastImportTime: metadata.LastImportTime, + }, nil) { + return + } + } } } -func (b *kvStorageBackend) updateLastImportTime(ctx, now time.Time) error { - return nil +func (k *kvStorageBackend) updateLastImportTime(ctx context.Context, key *resourcepb.ResourceKey, now time.Time) error { + metadata, err := k.metadataStore.Get(ctx, MetadataKey{ + Namespace: key.Namespace, + Group: key.Group, + Resource: key.Resource, + }) + + if err != nil && err != ErrNotFound { + k.log.Error("Error retrieving metadata for namespace %s: %s", key.Namespace, err) + return err + } + + if err == ErrNotFound { + metadata = Metadata{ + Namespace: key.Namespace, + Group: key.Group, + Resource: key.Resource, + } + } + + metadata.LastImportTime = now.UTC() + return k.metadataStore.Save(ctx, metadata) } -func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings, iter BulkRequestIterator) *resourcepb.BulkResponse { +func (k *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings, iter BulkRequestIterator) *resourcepb.BulkResponse { // TODO cross-node lock - err := b.bulkLock.Start(setting.Collection) + err := k.bulkLock.Start(setting.Collection) if err != nil { return &resourcepb.BulkResponse{ Error: AsErrorResult(err), } } - defer b.bulkLock.Finish(setting.Collection) + defer k.bulkLock.Finish(setting.Collection) bulkRvGenerator := newBulkRV() summaries := make(map[string]*resourcepb.BulkResponse_Summary, len(setting.Collection)) @@ -1260,15 +1304,15 @@ 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 k.eventStore.ListKeysSince(ctx, 1) { if err != nil { - b.log.Error("failed to list event: %s", err) + k.log.Error("failed to list event: %s", err) return rsp } evtKey, err := ParseEventKey(evtKeyStr) if err != nil { - b.log.Error("error parsing event key: %s", err) + k.log.Error("error parsing event key: %s", err) return rsp } @@ -1279,20 +1323,20 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings events = append(events, evtKeyStr) } - if err := b.eventStore.batchDelete(ctx, events); err != nil { - b.log.Error("failed to delete events: %s", err) + if err := k.eventStore.batchDelete(ctx, events); err != nil { + k.log.Error("failed to delete events: %s", err) return rsp } historyKeys := make([]DataKey, 0) - for dataKey, err := range b.dataStore.Keys(ctx, ListRequestKey{ + for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{ Namespace: key.Namespace, Group: key.Group, Resource: key.Resource, }, SortOrderAsc) { if err != nil { - b.log.Error("failed to list collection before delete: %s", err) + k.log.Error("failed to list collection before delete: %s", err) return rsp } @@ -1300,8 +1344,8 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings } previousCount := int64(len(historyKeys)) - if err := b.dataStore.batchDelete(ctx, historyKeys); err != nil { - b.log.Error("failed to delete collection: %s", err) + if err := k.dataStore.batchDelete(ctx, historyKeys); err != nil { + k.log.Error("failed to delete collection: %s", err) return rsp } summaries[NSGR(key)] = &resourcepb.BulkResponse_Summary{ @@ -1326,9 +1370,9 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings saved := make([]DataKey, 0) rollback := func() { // we don't have transactions in the kv store, so we simply delete everything we created - err = b.dataStore.batchDelete(ctx, saved) + err = k.dataStore.batchDelete(ctx, saved) if err != nil { - b.log.Error("failed to delete during rollback: %s", err) + k.log.Error("failed to delete during rollback: %s", err) } } @@ -1352,7 +1396,7 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings case resourcepb.WatchEvent_ADDED: action = DataActionCreated // Check if resource already exists for create operations - _, err := b.dataStore.GetLatestResourceKey(ctx, GetRequestKey{ + _, err := k.dataStore.GetLatestResourceKey(ctx, GetRequestKey{ Group: req.Key.Group, Resource: req.Key.Resource, Namespace: req.Key.Namespace, @@ -1406,7 +1450,7 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings Action: action, Folder: req.Folder, } - err = b.dataStore.Save(ctx, dataKey, bytes.NewReader(req.Value)) + err = k.dataStore.Save(ctx, dataKey, bytes.NewReader(req.Value)) if err != nil { rsp.Rejected = append(rsp.Rejected, &resourcepb.BulkResponse_Rejected{ Key: req.Key, @@ -1419,7 +1463,13 @@ func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings saved = append(saved, dataKey) } - // TODO update last import time + for _, key := range setting.Collection { + if err := k.updateLastImportTime(ctx, key, time.Now()); err != nil { + rollback() + rsp.Error = AsErrorResult(err) + return rsp + } + } return rsp } diff --git a/pkg/storage/unified/resource/storage_backend_test.go b/pkg/storage/unified/resource/storage_backend_test.go index f8ea2717a1e..a1bbdfe8cd5 100644 --- a/pkg/storage/unified/resource/storage_backend_test.go +++ b/pkg/storage/unified/resource/storage_backend_test.go @@ -1997,3 +1997,68 @@ func TestKvStorageBackend_ClusterScopedResources(t *testing.T) { } } } + +func TestKvStorageBackend_ResourceLastImportTime(t *testing.T) { + backend := setupTestStorageBackend(t) + ctx := context.Background() + + lastImportTimes := []ResourceLastImportTime{ + { + NamespacedResource: NamespacedResource{ + Namespace: "stacks-1", + Group: "apps", + Resource: "resource", + }, + LastImportTime: time.Now().Add(-2 * time.Minute).UTC().Truncate(time.Microsecond), + }, + { + NamespacedResource: NamespacedResource{ + Namespace: "stacks-2", + Group: "apps", + Resource: "resource", + }, + LastImportTime: time.Now().Add(-7 * time.Minute).UTC().Truncate(time.Microsecond), + }, + { + NamespacedResource: NamespacedResource{ + Namespace: "stacks-3", + Group: "apps", + Resource: "resource", + }, + LastImportTime: time.Now().Add(-3 * time.Minute).UTC().Truncate(time.Microsecond), + }, + { + NamespacedResource: NamespacedResource{ + Namespace: "stacks-4", + Group: "apps", + Resource: "resource", + }, + LastImportTime: time.Now().Add(-24 * time.Minute).UTC().Truncate(time.Microsecond), + }, + { + NamespacedResource: NamespacedResource{ + Namespace: "stacks-5", + Group: "apps", + Resource: "resource", + }, + LastImportTime: time.Now().Add(-44 * time.Minute).UTC().Truncate(time.Microsecond), + }, + } + + for _, lastImportTime := range lastImportTimes { + err := backend.updateLastImportTime(ctx, &resourcepb.ResourceKey{ + Namespace: lastImportTime.NamespacedResource.Namespace, + Group: lastImportTime.NamespacedResource.Group, + Resource: lastImportTime.NamespacedResource.Resource, + }, lastImportTime.LastImportTime) + + require.NoError(t, err) + } + + var i int + for lastImportTime, err := range backend.GetResourceLastImportTimes(ctx) { + require.NoError(t, err) + require.Equal(t, lastImportTimes[i], lastImportTime) + i++ + } +} diff --git a/pkg/storage/unified/testing/storage_backend_test.go b/pkg/storage/unified/testing/storage_backend_test.go index 04f34e9102f..835fd4b31ec 100644 --- a/pkg/storage/unified/testing/storage_backend_test.go +++ b/pkg/storage/unified/testing/storage_backend_test.go @@ -31,7 +31,7 @@ func TestBadgerKVStorageBackend(t *testing.T) { TestBlobSupport: true, TestListModifiedSince: true, // Badger does not support bulk import yet. - TestGetResourceLastImportTime: true, + // TestGetResourceLastImportTime: true, }, }) }