Skip to content

Commit 3fc951c

Browse files
authored
Adding resource name in cloudevent (#139)
We need this so work agent can handle unmanaged manifestwork correctly. Signed-off-by: Jian Qiu <jqiu@redhat.com>
1 parent d4c9f78 commit 3fc951c

9 files changed

Lines changed: 135 additions & 34 deletions

File tree

pkg/cloudevents/clients/work/agent/codec/manifestbundle.go

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,7 @@ func (c *ManifestBundleCodec) Encode(source string, eventType types.CloudEventsT
6262
evt := types.NewEventBuilder(source, eventType).
6363
WithResourceID(string(work.UID)).
6464
WithStatusUpdateSequenceID(sequenceGenerator.Generate().String()).
65+
WithResourceName(work.Name).
6566
WithResourceVersion(resourceVersion).
6667
WithClusterName(work.Namespace).
6768
WithOriginalSource(originalSource).
@@ -104,6 +105,17 @@ func (c *ManifestBundleCodec) Decode(evt *cloudevents.Event) (*workv1.ManifestWo
104105
return nil, fmt.Errorf("failed to get resourceid extension: %v", err)
105106
}
106107

108+
var resourceName string
109+
if v, ok := evtExtensions[types.ExtensionResourceName]; ok {
110+
resourceName, err = cloudeventstypes.ToString(v)
111+
if err != nil {
112+
return nil, fmt.Errorf("failed to get resourcename extension: %v", err)
113+
}
114+
} else {
115+
// fall back to set resourceName to resourceID
116+
resourceName = resourceID
117+
}
118+
107119
resourceVersion, err := cloudeventstypes.ToInteger(evtExtensions[types.ExtensionResourceVersion])
108120
if err != nil {
109121
return nil, fmt.Errorf("failed to get resourceversion extension: %v", err)
@@ -123,7 +135,7 @@ func (c *ManifestBundleCodec) Decode(evt *cloudevents.Event) (*workv1.ManifestWo
123135
UID: kubetypes.UID(resourceID),
124136
Generation: int64(resourceVersion),
125137
ResourceVersion: fmt.Sprintf("%d", resourceVersion),
126-
Name: resourceID,
138+
Name: resourceName,
127139
Namespace: clusterName,
128140
Labels: map[string]string{
129141
common.CloudEventsOriginalSourceLabelKey: evt.Source(),

pkg/cloudevents/clients/work/agent/codec/manifestbundle_test.go

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -159,6 +159,7 @@ func TestManifestBundleDecode(t *testing.T) {
159159
evt.SetType("io.open-cluster-management.works.v1alpha1.manifestbundles.spec.test")
160160
evt.SetExtension("resourceid", "test")
161161
evt.SetExtension("resourceversion", "13")
162+
evt.SetExtension("resourcename", "work1")
162163
return &evt
163164
}(),
164165
expectedErr: true,
@@ -171,6 +172,7 @@ func TestManifestBundleDecode(t *testing.T) {
171172
evt.SetType("io.open-cluster-management.works.v1alpha1.manifestbundles.spec.test")
172173
evt.SetExtension("resourceid", "test")
173174
evt.SetExtension("resourceversion", "13")
175+
evt.SetExtension("resourcename", "work1")
174176
return &evt
175177
}(),
176178
expectedErr: true,
@@ -184,6 +186,7 @@ func TestManifestBundleDecode(t *testing.T) {
184186
evt.SetExtension("resourceid", "test")
185187
evt.SetExtension("resourceversion", "13")
186188
evt.SetExtension("clustername", "cluster1")
189+
evt.SetExtension("resourcename", "work1")
187190
evt.SetExtension("deletiontimestamp", "1985-04-12T23:20:50.52Z")
188191
return &evt
189192
}(),
@@ -197,6 +200,7 @@ func TestManifestBundleDecode(t *testing.T) {
197200
evt.SetExtension("resourceid", "test")
198201
evt.SetExtension("resourceversion", "13")
199202
evt.SetExtension("clustername", "cluster1")
203+
evt.SetExtension("resourcename", "work1")
200204
if err := evt.SetData(cloudevents.ApplicationJSON, &payload.ManifestBundle{}); err != nil {
201205
t.Fatal(err)
202206
}
@@ -206,6 +210,30 @@ func TestManifestBundleDecode(t *testing.T) {
206210
},
207211
{
208212
name: "decode a cloudevent",
213+
event: func() *cloudevents.Event {
214+
evt := cloudevents.NewEvent()
215+
evt.SetSource("source1")
216+
evt.SetType("io.open-cluster-management.works.v1alpha1.manifestbundles.spec.test")
217+
evt.SetExtension("resourceid", "test")
218+
evt.SetExtension("resourceversion", "13")
219+
evt.SetExtension("clustername", "cluster1")
220+
evt.SetExtension("resourcename", "work1")
221+
if err := evt.SetData(cloudevents.ApplicationJSON, &payload.ManifestBundle{
222+
Manifests: []workv1.Manifest{
223+
{
224+
RawExtension: runtime.RawExtension{
225+
Raw: toConfigMap(t),
226+
},
227+
},
228+
},
229+
}); err != nil {
230+
t.Fatal(err)
231+
}
232+
return &evt
233+
}(),
234+
},
235+
{
236+
name: "decode a cloudevent with empty resourceName",
209237
event: func() *cloudevents.Event {
210238
evt := cloudevents.NewEvent()
211239
evt.SetSource("source1")

pkg/cloudevents/clients/work/source/codec/manifestbundle.go

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ func (c *ManifestBundleCodec) Encode(source string, eventType types.CloudEventsT
4747
evt := types.NewEventBuilder(source, eventType).
4848
WithClusterName(work.Namespace).
4949
WithResourceID(string(work.UID)).
50+
WithResourceName(work.Name).
5051
WithResourceVersion(int64(resourceVersion)).
5152
NewEvent()
5253

@@ -92,6 +93,17 @@ func (c *ManifestBundleCodec) Decode(evt *cloudevents.Event) (*workv1.ManifestWo
9293
return nil, fmt.Errorf("failed to get resourceid extension: %v", err)
9394
}
9495

96+
var resourceName string
97+
if v, ok := evtExtensions[types.ExtensionResourceName]; ok {
98+
resourceName, err = cloudeventstypes.ToString(v)
99+
if err != nil {
100+
return nil, fmt.Errorf("failed to get resourcename extension: %v", err)
101+
}
102+
} else {
103+
// fall back to set resourceName to resourceID
104+
resourceName = resourceID
105+
}
106+
95107
resourceVersion, err := cloudeventstypes.ToInteger(evtExtensions[types.ExtensionResourceVersion])
96108
if err != nil {
97109
return nil, fmt.Errorf("failed to get resourceversion extension: %v", err)
@@ -119,6 +131,7 @@ func (c *ManifestBundleCodec) Decode(evt *cloudevents.Event) (*workv1.ManifestWo
119131
}
120132

121133
metaObj.UID = kubetypes.UID(resourceID)
134+
metaObj.Name = resourceName
122135
metaObj.ResourceVersion = fmt.Sprintf("%d", resourceVersion)
123136
if metaObj.Annotations == nil {
124137
metaObj.Annotations = map[string]string{}

pkg/cloudevents/clients/work/source/codec/manifestbundle_test.go

Lines changed: 44 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -123,6 +123,46 @@ func TestManifestBundleDecode(t *testing.T) {
123123
}(),
124124
expectedErr: true,
125125
},
126+
{
127+
name: "decode a manifestbundle status cloudevent with unset resourceName",
128+
event: func() *cloudevents.Event {
129+
evt := cloudevents.NewEvent()
130+
evt.SetSource("source1")
131+
evt.SetType("io.open-cluster-management.works.v1alpha1.manifestbundles.status.test")
132+
evt.SetExtension("resourceid", "test")
133+
evt.SetExtension("resourceversion", "13")
134+
evt.SetExtension("sequenceid", "1834773391719010304")
135+
if err := evt.SetData(cloudevents.ApplicationJSON, &payload.ManifestBundleStatus{
136+
Conditions: []metav1.Condition{
137+
{
138+
Type: "Test",
139+
Status: metav1.ConditionTrue,
140+
},
141+
},
142+
}); err != nil {
143+
t.Fatal(err)
144+
}
145+
return &evt
146+
}(),
147+
expectedWork: &workv1.ManifestWork{
148+
ObjectMeta: metav1.ObjectMeta{
149+
UID: kubetypes.UID("test"),
150+
ResourceVersion: "13",
151+
Annotations: map[string]string{
152+
"cloudevents.open-cluster-management.io/sequenceid": "1834773391719010304",
153+
},
154+
Name: "test",
155+
},
156+
Status: workv1.ManifestWorkStatus{
157+
Conditions: []metav1.Condition{
158+
{
159+
Type: "Test",
160+
Status: metav1.ConditionTrue,
161+
},
162+
},
163+
},
164+
},
165+
},
126166
{
127167
name: "decode a manifestbundle status cloudevent",
128168
event: func() *cloudevents.Event {
@@ -131,6 +171,7 @@ func TestManifestBundleDecode(t *testing.T) {
131171
evt.SetType("io.open-cluster-management.works.v1alpha1.manifestbundles.status.test")
132172
evt.SetExtension("resourceid", "test")
133173
evt.SetExtension("resourceversion", "13")
174+
evt.SetExtension("resourcename", "work1")
134175
evt.SetExtension("sequenceid", "1834773391719010304")
135176
if err := evt.SetData(cloudevents.ApplicationJSON, &payload.ManifestBundleStatus{
136177
Conditions: []metav1.Condition{
@@ -151,6 +192,7 @@ func TestManifestBundleDecode(t *testing.T) {
151192
Annotations: map[string]string{
152193
"cloudevents.open-cluster-management.io/sequenceid": "1834773391719010304",
153194
},
195+
Name: "work1",
154196
},
155197
Status: workv1.ManifestWorkStatus{
156198
Conditions: []metav1.Condition{
@@ -181,6 +223,7 @@ func TestManifestBundleDecode(t *testing.T) {
181223
evt.SetSource("source1")
182224
evt.SetType("io.open-cluster-management.works.v1alpha1.manifestbundles.status.test")
183225
evt.SetExtension("resourceid", "test")
226+
evt.SetExtension("resourcename", "work1")
184227
evt.SetExtension("resourceversion", "13")
185228
evt.SetExtension("metadata", string(metaJson))
186229
evt.SetExtension("sequenceid", "1834773391719010304")
@@ -224,7 +267,7 @@ func TestManifestBundleDecode(t *testing.T) {
224267
ObjectMeta: metav1.ObjectMeta{
225268
UID: kubetypes.UID("test"),
226269
ResourceVersion: "13",
227-
Name: "test",
270+
Name: "work1",
228271
Namespace: "cluster1",
229272
Labels: map[string]string{"test1": "test1"},
230273
Annotations: map[string]string{

pkg/cloudevents/generic/types/types.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,9 @@ const (
5757
// ExtensionResourceID is the cloud event extension key of the resource ID.
5858
ExtensionResourceID = "resourceid"
5959

60+
// ExtensionResourceName is the cloud event extension key of the resource name.
61+
ExtensionResourceName = "resourcename"
62+
6063
// ExtensionResourceVersion is the cloud event extension key of the resource version.
6164
ExtensionResourceVersion = "resourceversion"
6265

@@ -228,6 +231,7 @@ type EventBuilder struct {
228231
clusterName string
229232
originalSource string
230233
resourceID string
234+
resourceName string
231235
sequenceID string
232236
resourceVersion *int64
233237
eventType CloudEventsType
@@ -246,6 +250,11 @@ func (b *EventBuilder) WithResourceID(resourceID string) *EventBuilder {
246250
return b
247251
}
248252

253+
func (b *EventBuilder) WithResourceName(resourceName string) *EventBuilder {
254+
b.resourceName = resourceName
255+
return b
256+
}
257+
249258
func (b *EventBuilder) WithResourceVersion(resourceVersion int64) *EventBuilder {
250259
b.resourceVersion = &resourceVersion
251260
return b
@@ -285,6 +294,13 @@ func (b *EventBuilder) NewEvent() cloudevents.Event {
285294
evt.SetExtension(ExtensionResourceID, b.resourceID)
286295
}
287296

297+
if len(b.resourceName) != 0 {
298+
evt.SetExtension(ExtensionResourceName, b.resourceName)
299+
} else {
300+
// if resourceName is not set, uses resourceID as the resourceName
301+
evt.SetExtension(ExtensionResourceName, b.resourceID)
302+
}
303+
288304
if b.resourceVersion != nil {
289305
evt.SetExtension(ExtensionResourceVersion, *b.resourceVersion)
290306
}

test/integration/cloudevents/garbagecollector_test.go

Lines changed: 6 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,6 @@ import (
1818

1919
workv1informers "open-cluster-management.io/api/client/work/informers/externalversions/work/v1"
2020

21-
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/common"
22-
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
2321
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work"
2422
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/agent/codec"
2523
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/garbagecollector"
@@ -161,21 +159,19 @@ var _ = ginkgo.Describe("Garbage Collector Test", func() {
161159
ginkgo.By("agent update the work status", func() {
162160
gomega.Eventually(func() error {
163161
workClient := agentClientHolder.ManifestWorks(clusterName)
164-
workID1 := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName1)
165-
if err := util.AddWorkFinalizer(ctx, workClient, workID1); err != nil {
162+
if err := util.AddWorkFinalizer(ctx, workClient, workName1); err != nil {
166163
return err
167164
}
168165

169-
if err := util.UpdateWorkStatus(ctx, workClient, workID1, util.WorkCreatedCondition); err != nil {
166+
if err := util.UpdateWorkStatus(ctx, workClient, workName1, util.WorkCreatedCondition); err != nil {
170167
return err
171168
}
172169

173-
workID2 := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName2)
174-
if err := util.AddWorkFinalizer(ctx, workClient, workID2); err != nil {
170+
if err := util.AddWorkFinalizer(ctx, workClient, workName2); err != nil {
175171
return err
176172
}
177173

178-
return util.UpdateWorkStatus(ctx, workClient, workID2, util.WorkCreatedCondition)
174+
return util.UpdateWorkStatus(ctx, workClient, workName2, util.WorkCreatedCondition)
179175
}, 10*time.Second, 1*time.Second).Should(gomega.Succeed())
180176
})
181177

@@ -198,8 +194,7 @@ var _ = ginkgo.Describe("Garbage Collector Test", func() {
198194

199195
ginkgo.By("agent delete the first work with single owner", func() {
200196
gomega.Eventually(func() error {
201-
workID1 := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName1)
202-
return util.RemoveWorkFinalizer(ctx, agentClientHolder.ManifestWorks(clusterName), workID1)
197+
return util.RemoveWorkFinalizer(ctx, agentClientHolder.ManifestWorks(clusterName), workName1)
203198
}, 30*time.Second, 1*time.Second).Should(gomega.Succeed())
204199
})
205200

@@ -230,8 +225,7 @@ var _ = ginkgo.Describe("Garbage Collector Test", func() {
230225

231226
ginkgo.By("agent delete the work with two owners", func() {
232227
gomega.Eventually(func() error {
233-
workID2 := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName2)
234-
return util.RemoveWorkFinalizer(ctx, agentClientHolder.ManifestWorks(clusterName), workID2)
228+
return util.RemoveWorkFinalizer(ctx, agentClientHolder.ManifestWorks(clusterName), workName2)
235229
}, 30*time.Second, 1*time.Second).Should(gomega.Succeed())
236230
})
237231

test/integration/cloudevents/manifestworkclients_informer_test.go

Lines changed: 5 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,6 @@ import (
1212
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
1313
"k8s.io/apimachinery/pkg/util/rand"
1414

15-
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/common"
16-
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/utils"
1715
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work"
1816
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/agent/codec"
1917
"open-cluster-management.io/sdk-go/test/integration/cloudevents/agent"
@@ -104,13 +102,12 @@ func crudManifestWork(
104102
ginkgo.By("agent update the work status", func() {
105103
gomega.Eventually(func() error {
106104
workClient := agentClientHolder.ManifestWorks(clusterName)
107-
workID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName)
108105

109-
if err := util.AddWorkFinalizer(ctx, workClient, workID); err != nil {
106+
if err := util.AddWorkFinalizer(ctx, workClient, workName); err != nil {
110107
return err
111108
}
112109

113-
return util.UpdateWorkStatus(ctx, workClient, workID, util.WorkCreatedCondition)
110+
return util.UpdateWorkStatus(ctx, workClient, workName, util.WorkCreatedCondition)
114111
}, 10*time.Second, 1*time.Second).Should(gomega.Succeed())
115112
})
116113

@@ -128,12 +125,11 @@ func crudManifestWork(
128125
ginkgo.By("agent update the work status again", func() {
129126
gomega.Eventually(func() error {
130127
workClient := agentClientHolder.ManifestWorks(clusterName)
131-
workID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName)
132-
if err := util.AssertUpdatedWork(ctx, workClient, workID); err != nil {
128+
if err := util.AssertUpdatedWork(ctx, workClient, workName); err != nil {
133129
return err
134130
}
135131

136-
return util.UpdateWorkStatus(ctx, workClient, workID, util.WorkUpdatedCondition)
132+
return util.UpdateWorkStatus(ctx, workClient, workName, util.WorkUpdatedCondition)
137133
}, 10*time.Second, 1*time.Second).Should(gomega.Succeed())
138134
})
139135

@@ -150,9 +146,8 @@ func crudManifestWork(
150146

151147
ginkgo.By("agent delete the work", func() {
152148
gomega.Eventually(func() error {
153-
workID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, workName)
154149
workClient := agentClientHolder.ManifestWorks(clusterName)
155-
return util.RemoveWorkFinalizer(ctx, workClient, workID)
150+
return util.RemoveWorkFinalizer(ctx, workClient, workName)
156151
}, 10*time.Second, 1*time.Second).Should(gomega.Succeed())
157152
})
158153

test/integration/cloudevents/manifestworkclients_reync_test.go

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -46,16 +46,18 @@ var _ = ginkgo.Describe("ManifestWork Clients Test - Resync", func() {
4646

4747
// add two works in the agent cache
4848
store := informer.Informer().GetStore()
49-
work1UID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, fmt.Sprintf("%s-1", workNamePrefix))
50-
work1 := util.NewManifestWorkWithStatus(clusterName, work1UID)
49+
work1Name := fmt.Sprintf("%s-1", workNamePrefix)
50+
work1UID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, work1Name)
51+
work1 := util.NewManifestWorkWithStatus(clusterName, work1Name)
5152
work1.UID = apitypes.UID(work1UID)
5253
work1.ResourceVersion = "1"
5354
work1.Labels = map[string]string{common.CloudEventsOriginalSourceLabelKey: sourceID}
5455
work1.Annotations = map[string]string{common.CloudEventsDataTypeAnnotationKey: payload.ManifestBundleEventDataType.String()}
5556
gomega.Expect(store.Add(work1)).ToNot(gomega.HaveOccurred())
5657

57-
work2UID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, fmt.Sprintf("%s-2", workNamePrefix))
58-
work2 := util.NewManifestWorkWithStatus(clusterName, work2UID)
58+
work2Name := fmt.Sprintf("%s-2", workNamePrefix)
59+
work2UID := utils.UID(sourceID, common.ManifestWorkGR.String(), clusterName, work2Name)
60+
work2 := util.NewManifestWorkWithStatus(clusterName, work2Name)
5961
work2.UID = apitypes.UID(work2UID)
6062
work2.ResourceVersion = "1"
6163
work2.Labels = map[string]string{common.CloudEventsOriginalSourceLabelKey: sourceID}

0 commit comments

Comments
 (0)