Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
96 changes: 31 additions & 65 deletions pkg/cloudevents/clients/store/informer.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ package store

import (
"context"
"fmt"

"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
Expand All @@ -11,7 +10,6 @@ import (
"k8s.io/client-go/tools/cache"
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
)

// AgentInformerWatcherStore extends the BaseClientWatchStore.
Expand Down Expand Up @@ -49,80 +47,48 @@ func (s *AgentInformerWatcherStore[T]) Delete(resource runtime.Object) error {
return s.Store.Delete(resource)
}

func (s *AgentInformerWatcherStore[T]) HandleReceivedResource(ctx context.Context, action types.ResourceAction, resource T) error {
switch action {
case types.Added:
newObj, err := utils.ToRuntimeObject(resource)
if err != nil {
return err
}

return s.Add(newObj)
case types.Modified:
accessor, err := meta.Accessor(resource)
if err != nil {
return err
}

lastObj, exists, err := s.Get(accessor.GetNamespace(), accessor.GetName())
if err != nil {
return err
}
if !exists {
return fmt.Errorf("the resource %s/%s does not exist", accessor.GetNamespace(), accessor.GetName())
}

// if resource is deleting, keep the deletion timestamp
if !lastObj.GetDeletionTimestamp().IsZero() {
accessor.SetDeletionTimestamp(lastObj.GetDeletionTimestamp())
}

updated, err := utils.ToRuntimeObject(resource)
if err != nil {
return err
}

return s.Update(updated)
case types.Deleted:
newObj, err := meta.Accessor(resource)
if err != nil {
return err
}

if newObj.GetDeletionTimestamp().IsZero() {
return nil
}
func (s *AgentInformerWatcherStore[T]) HandleReceivedResource(ctx context.Context, resource T) error {
runtimeObj, err := utils.ToRuntimeObject(resource)
if err != nil {
return err
}

last, exists, err := s.Get(newObj.GetNamespace(), newObj.GetName())
if err != nil {
return err
}
if !exists {
return nil
}
metaObj, err := meta.Accessor(runtimeObj)
if err != nil {
return err
}

deletingObj, err := utils.ToRuntimeObject(last)
if err != nil {
return err
}
lastResource, exists, err := s.Get(metaObj.GetNamespace(), metaObj.GetName())
if err != nil {
return err
}
if !exists {
return s.Add(runtimeObj)
}

if !metaObj.GetDeletionTimestamp().IsZero() {
// trigger an update event if the object is deleting.
// Only need to update generation/finalizer/deletionTimeStamp of the object.
if len(newObj.GetFinalizers()) != 0 {
accessor, err := meta.Accessor(deletingObj)
if len(metaObj.GetFinalizers()) != 0 {
deletingObj, err := meta.Accessor(lastResource)
if err != nil {
return err
}
deletingObj.SetDeletionTimestamp(metaObj.GetDeletionTimestamp())
deletingObj.SetFinalizers(metaObj.GetFinalizers())
deletingObj.SetGeneration(metaObj.GetGeneration())
runtimeObj, err := utils.ToRuntimeObject(deletingObj)
if err != nil {
return err
}
accessor.SetDeletionTimestamp(newObj.GetDeletionTimestamp())
accessor.SetFinalizers(newObj.GetFinalizers())
accessor.SetGeneration(newObj.GetGeneration())
return s.Update(deletingObj)

return s.Update(runtimeObj)
}
Comment thread
skeeey marked this conversation as resolved.

return s.Delete(deletingObj)
default:
return fmt.Errorf("unsupported resource action %s", action)
return s.Delete(runtimeObj)
}

return s.Update(runtimeObj)
}

func (s *AgentInformerWatcherStore[T]) GetWatcher(namespace string, opts metav1.ListOptions) (watch.Interface, error) {
Expand Down
7 changes: 3 additions & 4 deletions pkg/cloudevents/clients/store/informer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/watch"
clusterv1 "open-cluster-management.io/api/cluster/v1"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
)

func TestGet(t *testing.T) {
Expand Down Expand Up @@ -174,10 +173,10 @@ func TestWatch(t *testing.T) {
}
}()

if err := watchStore.HandleReceivedResource(ctx, types.Added, &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "test0"}}); err != nil {
if err := watchStore.HandleReceivedResource(ctx, &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "test0"}}); err != nil {
t.Error(err)
}
if err := watchStore.HandleReceivedResource(ctx, types.Modified, &clusterv1.ManagedCluster{
if err := watchStore.HandleReceivedResource(ctx, &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{
Name: "test1",
},
Expand All @@ -186,7 +185,7 @@ func TestWatch(t *testing.T) {
}}); err != nil {
t.Error(err)
}
if err := watchStore.HandleReceivedResource(ctx, types.Deleted, &clusterv1.ManagedCluster{
if err := watchStore.HandleReceivedResource(ctx, &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{
Name: "test1",
DeletionTimestamp: &metav1.Time{Time: time.Now()},
Expand Down
3 changes: 1 addition & 2 deletions pkg/cloudevents/clients/store/interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@ import (
"k8s.io/klog/v2"

"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
)

const syncedPollPeriod = 100 * time.Millisecond
Expand All @@ -34,7 +33,7 @@ type ClientWatcherStore[T generic.ResourceObject] interface {
GetWatcher(namespace string, opts metav1.ListOptions) (watch.Interface, error)

// HandleReceivedResource handles the client received resource events.
HandleReceivedResource(ctx context.Context, action types.ResourceAction, resource T) error
HandleReceivedResource(ctx context.Context, resource T) error

// Add will be called by resource client when adding resources. The implementation is based on the specific
// watcher store, in some case, it does not need to update a store, but just send a watch event.
Expand Down
87 changes: 21 additions & 66 deletions pkg/cloudevents/clients/store/simplestore.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,18 +2,15 @@ package store

import (
"context"
"fmt"

"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/watch"
"k8s.io/client-go/tools/cache"
"k8s.io/klog/v2"

"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
)

// SimpleStore extends the BaseClientWatchStore.
Expand Down Expand Up @@ -51,73 +48,31 @@ func (s *SimpleStore[T]) HasInitiated() bool {
return true
}

func (s *SimpleStore[T]) HandleReceivedResource(ctx context.Context, action types.ResourceAction, resource T) error {
logger := klog.FromContext(ctx)

switch action {
case types.Added:
newObj, err := utils.ToRuntimeObject(resource)
if err != nil {
return err
}

return s.Add(newObj)
case types.Modified:
newObj, err := meta.Accessor(resource)
if err != nil {
return err
}

lastObj, exists, err := s.Get(newObj.GetNamespace(), newObj.GetName())
if err != nil {
return err
}
if !exists {
return fmt.Errorf("the resource %s/%s does not exist", newObj.GetNamespace(), newObj.GetName())
}

// prevent the resource from being updated if it is deleting
if !lastObj.GetDeletionTimestamp().IsZero() {
logger.Info("the resource is deleting, ignore the update",
"resourceNamespace", newObj.GetNamespace(), "resourceName", newObj.GetName())
return nil
}

updated, err := utils.ToRuntimeObject(resource)
if err != nil {
return err
}

return s.Update(updated)
case types.Deleted:
newObj, err := meta.Accessor(resource)
if err != nil {
return err
}
func (s *SimpleStore[T]) HandleReceivedResource(ctx context.Context, resource T) error {
runtimeObj, err := utils.ToRuntimeObject(resource)
if err != nil {
return err
}

if newObj.GetDeletionTimestamp().IsZero() {
return nil
}
metaObj, err := meta.Accessor(runtimeObj)
if err != nil {
return err
}

if len(newObj.GetFinalizers()) != 0 {
return nil
}
_, exists, err := s.Get(metaObj.GetNamespace(), metaObj.GetName())
if err != nil {
return err
}
if !exists {
return s.Add(runtimeObj)
}

last, exists, err := s.Get(newObj.GetNamespace(), newObj.GetName())
if err != nil {
return err
}
if !exists {
if !metaObj.GetDeletionTimestamp().IsZero() {
if len(metaObj.GetFinalizers()) != 0 {
return nil
}

deletingObj, err := utils.ToRuntimeObject(last)
if err != nil {
return err
}

return s.Delete(deletingObj)
default:
return fmt.Errorf("unsupported resource action %s", action)
return s.Delete(runtimeObj)
}

return s.Update(runtimeObj)
}
13 changes: 6 additions & 7 deletions pkg/cloudevents/clients/store/simplestore_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,7 @@ import (

coordv1 "k8s.io/api/coordination/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"

"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
"k8s.io/apimachinery/pkg/watch"
)

func TestHandleReceivedResource(t *testing.T) {
Expand All @@ -33,13 +32,13 @@ func TestHandleReceivedResource(t *testing.T) {

cases := []struct {
name string
action types.ResourceAction
action watch.EventType
received *coordv1.Lease
validate func(t *testing.T, namespace, name string)
}{
{
name: "add resource",
action: types.Added,
action: watch.Added,
received: &coordv1.Lease{
ObjectMeta: metav1.ObjectMeta{
Name: "new",
Expand All @@ -59,7 +58,7 @@ func TestHandleReceivedResource(t *testing.T) {
},
{
name: "update resource",
action: types.Modified,
action: watch.Modified,
received: &coordv1.Lease{
ObjectMeta: metav1.ObjectMeta{
Name: "update",
Expand All @@ -86,7 +85,7 @@ func TestHandleReceivedResource(t *testing.T) {
},
{
name: "delete resource",
action: types.Deleted,
action: watch.Deleted,
received: &coordv1.Lease{
ObjectMeta: metav1.ObjectMeta{
Name: "deletion",
Expand All @@ -109,7 +108,7 @@ func TestHandleReceivedResource(t *testing.T) {

for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
err := store.HandleReceivedResource(context.Background(), c.action, c.received)
err := store.HandleReceivedResource(context.Background(), c.received)
if err != nil {
t.Error(err)
}
Expand Down
10 changes: 2 additions & 8 deletions pkg/cloudevents/clients/work/store/base.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@ import (
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/common"
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/store"
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
)

type baseSourceStore struct {
Expand All @@ -27,13 +26,8 @@ type baseSourceStore struct {
receivedWorks workqueue.TypedRateLimitingInterface[*workv1.ManifestWork]
}

func (bs *baseSourceStore) HandleReceivedResource(_ context.Context, action types.ResourceAction, work *workv1.ManifestWork) error {
switch action {
case types.StatusModified:
bs.receivedWorks.Add(work)
default:
return fmt.Errorf("unsupported resource action %s", action)
}
func (bs *baseSourceStore) HandleReceivedResource(_ context.Context, work *workv1.ManifestWork) error {
bs.receivedWorks.Add(work)
return nil
}

Expand Down
Loading