Skip to content

Commit 1be7def

Browse files
committed
handle resource in informer
Signed-off-by: Wei Liu <liuweixa@redhat.com>
1 parent 7e7c96b commit 1be7def

16 files changed

Lines changed: 154 additions & 401 deletions

File tree

pkg/cloudevents/clients/store/informer.go

Lines changed: 31 additions & 65 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ package store
22

33
import (
44
"context"
5-
"fmt"
65

76
"k8s.io/apimachinery/pkg/api/meta"
87
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -11,7 +10,6 @@ import (
1110
"k8s.io/client-go/tools/cache"
1211
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
1312
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
14-
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
1513
)
1614

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

52-
func (s *AgentInformerWatcherStore[T]) HandleReceivedResource(ctx context.Context, action types.ResourceAction, resource T) error {
53-
switch action {
54-
case types.Added:
55-
newObj, err := utils.ToRuntimeObject(resource)
56-
if err != nil {
57-
return err
58-
}
59-
60-
return s.Add(newObj)
61-
case types.Modified:
62-
accessor, err := meta.Accessor(resource)
63-
if err != nil {
64-
return err
65-
}
66-
67-
lastObj, exists, err := s.Get(accessor.GetNamespace(), accessor.GetName())
68-
if err != nil {
69-
return err
70-
}
71-
if !exists {
72-
return fmt.Errorf("the resource %s/%s does not exist", accessor.GetNamespace(), accessor.GetName())
73-
}
74-
75-
// if resource is deleting, keep the deletion timestamp
76-
if !lastObj.GetDeletionTimestamp().IsZero() {
77-
accessor.SetDeletionTimestamp(lastObj.GetDeletionTimestamp())
78-
}
79-
80-
updated, err := utils.ToRuntimeObject(resource)
81-
if err != nil {
82-
return err
83-
}
84-
85-
return s.Update(updated)
86-
case types.Deleted:
87-
newObj, err := meta.Accessor(resource)
88-
if err != nil {
89-
return err
90-
}
91-
92-
if newObj.GetDeletionTimestamp().IsZero() {
93-
return nil
94-
}
50+
func (s *AgentInformerWatcherStore[T]) HandleReceivedResource(ctx context.Context, resource T) error {
51+
runtimeObj, err := utils.ToRuntimeObject(resource)
52+
if err != nil {
53+
return err
54+
}
9555

96-
last, exists, err := s.Get(newObj.GetNamespace(), newObj.GetName())
97-
if err != nil {
98-
return err
99-
}
100-
if !exists {
101-
return nil
102-
}
56+
metaObj, err := meta.Accessor(runtimeObj)
57+
if err != nil {
58+
return err
59+
}
10360

104-
deletingObj, err := utils.ToRuntimeObject(last)
105-
if err != nil {
106-
return err
107-
}
61+
lastResource, exists, err := s.Get(metaObj.GetNamespace(), metaObj.GetName())
62+
if err != nil {
63+
return err
64+
}
65+
if !exists {
66+
return s.Add(runtimeObj)
67+
}
10868

69+
if !metaObj.GetDeletionTimestamp().IsZero() {
10970
// trigger an update event if the object is deleting.
11071
// Only need to update generation/finalizer/deletionTimeStamp of the object.
111-
if len(newObj.GetFinalizers()) != 0 {
112-
accessor, err := meta.Accessor(deletingObj)
72+
if len(metaObj.GetFinalizers()) != 0 {
73+
deletingObj, err := meta.Accessor(lastResource)
74+
if err != nil {
75+
return err
76+
}
77+
deletingObj.SetDeletionTimestamp(metaObj.GetDeletionTimestamp())
78+
deletingObj.SetFinalizers(metaObj.GetFinalizers())
79+
deletingObj.SetGeneration(metaObj.GetGeneration())
80+
runtimeObj, err := utils.ToRuntimeObject(deletingObj)
11381
if err != nil {
11482
return err
11583
}
116-
accessor.SetDeletionTimestamp(newObj.GetDeletionTimestamp())
117-
accessor.SetFinalizers(newObj.GetFinalizers())
118-
accessor.SetGeneration(newObj.GetGeneration())
119-
return s.Update(deletingObj)
84+
85+
return s.Update(runtimeObj)
12086
}
12187

122-
return s.Delete(deletingObj)
123-
default:
124-
return fmt.Errorf("unsupported resource action %s", action)
88+
return s.Delete(runtimeObj)
12589
}
90+
91+
return s.Update(runtimeObj)
12692
}
12793

12894
func (s *AgentInformerWatcherStore[T]) GetWatcher(namespace string, opts metav1.ListOptions) (watch.Interface, error) {

pkg/cloudevents/clients/store/informer_test.go

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ import (
1010
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
1111
"k8s.io/apimachinery/pkg/watch"
1212
clusterv1 "open-cluster-management.io/api/cluster/v1"
13-
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
1413
)
1514

1615
func TestGet(t *testing.T) {
@@ -174,10 +173,10 @@ func TestWatch(t *testing.T) {
174173
}
175174
}()
176175

177-
if err := watchStore.HandleReceivedResource(ctx, types.Added, &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "test0"}}); err != nil {
176+
if err := watchStore.HandleReceivedResource(ctx, &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "test0"}}); err != nil {
178177
t.Error(err)
179178
}
180-
if err := watchStore.HandleReceivedResource(ctx, types.Modified, &clusterv1.ManagedCluster{
179+
if err := watchStore.HandleReceivedResource(ctx, &clusterv1.ManagedCluster{
181180
ObjectMeta: metav1.ObjectMeta{
182181
Name: "test1",
183182
},
@@ -186,7 +185,7 @@ func TestWatch(t *testing.T) {
186185
}}); err != nil {
187186
t.Error(err)
188187
}
189-
if err := watchStore.HandleReceivedResource(ctx, types.Deleted, &clusterv1.ManagedCluster{
188+
if err := watchStore.HandleReceivedResource(ctx, &clusterv1.ManagedCluster{
190189
ObjectMeta: metav1.ObjectMeta{
191190
Name: "test1",
192191
DeletionTimestamp: &metav1.Time{Time: time.Now()},

pkg/cloudevents/clients/store/interface.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,6 @@ import (
1111
"k8s.io/klog/v2"
1212

1313
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
14-
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
1514
)
1615

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

3635
// HandleReceivedResource handles the client received resource events.
37-
HandleReceivedResource(ctx context.Context, action types.ResourceAction, resource T) error
36+
HandleReceivedResource(ctx context.Context, resource T) error
3837

3938
// Add will be called by resource client when adding resources. The implementation is based on the specific
4039
// watcher store, in some case, it does not need to update a store, but just send a watch event.

pkg/cloudevents/clients/store/simplestore.go

Lines changed: 21 additions & 66 deletions
Original file line numberDiff line numberDiff line change
@@ -2,18 +2,15 @@ package store
22

33
import (
44
"context"
5-
"fmt"
65

76
"k8s.io/apimachinery/pkg/api/meta"
87
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
98
"k8s.io/apimachinery/pkg/runtime"
109
"k8s.io/apimachinery/pkg/watch"
1110
"k8s.io/client-go/tools/cache"
12-
"k8s.io/klog/v2"
1311

1412
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
1513
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
16-
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
1714
)
1815

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

54-
func (s *SimpleStore[T]) HandleReceivedResource(ctx context.Context, action types.ResourceAction, resource T) error {
55-
logger := klog.FromContext(ctx)
56-
57-
switch action {
58-
case types.Added:
59-
newObj, err := utils.ToRuntimeObject(resource)
60-
if err != nil {
61-
return err
62-
}
63-
64-
return s.Add(newObj)
65-
case types.Modified:
66-
newObj, err := meta.Accessor(resource)
67-
if err != nil {
68-
return err
69-
}
70-
71-
lastObj, exists, err := s.Get(newObj.GetNamespace(), newObj.GetName())
72-
if err != nil {
73-
return err
74-
}
75-
if !exists {
76-
return fmt.Errorf("the resource %s/%s does not exist", newObj.GetNamespace(), newObj.GetName())
77-
}
78-
79-
// prevent the resource from being updated if it is deleting
80-
if !lastObj.GetDeletionTimestamp().IsZero() {
81-
logger.Info("the resource is deleting, ignore the update",
82-
"resourceNamespace", newObj.GetNamespace(), "resourceName", newObj.GetName())
83-
return nil
84-
}
85-
86-
updated, err := utils.ToRuntimeObject(resource)
87-
if err != nil {
88-
return err
89-
}
90-
91-
return s.Update(updated)
92-
case types.Deleted:
93-
newObj, err := meta.Accessor(resource)
94-
if err != nil {
95-
return err
96-
}
51+
func (s *SimpleStore[T]) HandleReceivedResource(ctx context.Context, resource T) error {
52+
runtimeObj, err := utils.ToRuntimeObject(resource)
53+
if err != nil {
54+
return err
55+
}
9756

98-
if newObj.GetDeletionTimestamp().IsZero() {
99-
return nil
100-
}
57+
metaObj, err := meta.Accessor(runtimeObj)
58+
if err != nil {
59+
return err
60+
}
10161

102-
if len(newObj.GetFinalizers()) != 0 {
103-
return nil
104-
}
62+
_, exists, err := s.Get(metaObj.GetNamespace(), metaObj.GetName())
63+
if err != nil {
64+
return err
65+
}
66+
if !exists {
67+
return s.Add(runtimeObj)
68+
}
10569

106-
last, exists, err := s.Get(newObj.GetNamespace(), newObj.GetName())
107-
if err != nil {
108-
return err
109-
}
110-
if !exists {
70+
if !metaObj.GetDeletionTimestamp().IsZero() {
71+
if len(metaObj.GetFinalizers()) != 0 {
11172
return nil
11273
}
113-
114-
deletingObj, err := utils.ToRuntimeObject(last)
115-
if err != nil {
116-
return err
117-
}
118-
119-
return s.Delete(deletingObj)
120-
default:
121-
return fmt.Errorf("unsupported resource action %s", action)
74+
return s.Delete(runtimeObj)
12275
}
76+
77+
return s.Update(runtimeObj)
12378
}

pkg/cloudevents/clients/store/simplestore_test.go

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,7 @@ import (
77

88
coordv1 "k8s.io/api/coordination/v1"
99
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
10-
11-
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
10+
"k8s.io/apimachinery/pkg/watch"
1211
)
1312

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

3433
cases := []struct {
3534
name string
36-
action types.ResourceAction
35+
action watch.EventType
3736
received *coordv1.Lease
3837
validate func(t *testing.T, namespace, name string)
3938
}{
4039
{
4140
name: "add resource",
42-
action: types.Added,
41+
action: watch.Added,
4342
received: &coordv1.Lease{
4443
ObjectMeta: metav1.ObjectMeta{
4544
Name: "new",
@@ -59,7 +58,7 @@ func TestHandleReceivedResource(t *testing.T) {
5958
},
6059
{
6160
name: "update resource",
62-
action: types.Modified,
61+
action: watch.Modified,
6362
received: &coordv1.Lease{
6463
ObjectMeta: metav1.ObjectMeta{
6564
Name: "update",
@@ -86,7 +85,7 @@ func TestHandleReceivedResource(t *testing.T) {
8685
},
8786
{
8887
name: "delete resource",
89-
action: types.Deleted,
88+
action: watch.Deleted,
9089
received: &coordv1.Lease{
9190
ObjectMeta: metav1.ObjectMeta{
9291
Name: "deletion",
@@ -109,7 +108,7 @@ func TestHandleReceivedResource(t *testing.T) {
109108

110109
for _, c := range cases {
111110
t.Run(c.name, func(t *testing.T) {
112-
err := store.HandleReceivedResource(context.Background(), c.action, c.received)
111+
err := store.HandleReceivedResource(context.Background(), c.received)
113112
if err != nil {
114113
t.Error(err)
115114
}

pkg/cloudevents/clients/work/store/base.go

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@ import (
1717
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/common"
1818
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/store"
1919
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
20-
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
2120
)
2221

2322
type baseSourceStore struct {
@@ -27,13 +26,8 @@ type baseSourceStore struct {
2726
receivedWorks workqueue.TypedRateLimitingInterface[*workv1.ManifestWork]
2827
}
2928

30-
func (bs *baseSourceStore) HandleReceivedResource(_ context.Context, action types.ResourceAction, work *workv1.ManifestWork) error {
31-
switch action {
32-
case types.StatusModified:
33-
bs.receivedWorks.Add(work)
34-
default:
35-
return fmt.Errorf("unsupported resource action %s", action)
36-
}
29+
func (bs *baseSourceStore) HandleReceivedResource(_ context.Context, work *workv1.ManifestWork) error {
30+
bs.receivedWorks.Add(work)
3731
return nil
3832
}
3933

0 commit comments

Comments
 (0)