implement last import time
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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++
|
||||
}
|
||||
}
|
||||
|
||||
@@ -31,7 +31,7 @@ func TestBadgerKVStorageBackend(t *testing.T) {
|
||||
TestBlobSupport: true,
|
||||
TestListModifiedSince: true,
|
||||
// Badger does not support bulk import yet.
|
||||
TestGetResourceLastImportTime: true,
|
||||
// TestGetResourceLastImportTime: true,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user