Skip to content

Commit aa513d7

Browse files
mkolesnikclaude
andcommitted
Add NewWorkApplierWithRuntimeClient constructor
The README documents a NewWorkApplierWithRuntimeClient implementation using sigs.k8s.io/controller-runtime/pkg/client, but users cannot implement it externally because all WorkApplier fields are unexported. Promote it to an exported constructor. Closes #226 Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> Signed-off-by: Mike Kolesnik <mkolesni@redhat.com>
1 parent 3ea63fb commit aa513d7

3 files changed

Lines changed: 165 additions & 45 deletions

File tree

pkg/apis/work/v1/README.md

Lines changed: 3 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -89,52 +89,10 @@ The input `index` is the count of the existing manifestWorks.
8989

9090
1. Create a `WorkApplier` instance.
9191

92-
There is a default `WorkApplier` instance with ManifestWork typed client.
92+
There are two constructors:
9393

94-
One is `NewWorkApplierWithTypedClient(workClient workv1client.Interface, workLister worklister.ManifestWorkLister)`
95-
with manifestWork typed client.
96-
97-
You can also define the instance with runtime client, for example:
98-
```go
99-
func NewWorkApplierWithRuntimeClient(workClient client.Client) *WorkApplier {
100-
return &WorkApplier{
101-
cache: newWorkCache(),
102-
getWork: func(ctx context.Context, namespace, name string) (*workapiv1.ManifestWork, error) {
103-
work := &workapiv1.ManifestWork{}
104-
err := workClient.Get(ctx, types.NamespacedName{Namespace: namespace, Name: name}, work)
105-
return work, err
106-
},
107-
deleteWork: func(ctx context.Context, namespace, name string) error {
108-
work := &workapiv1.ManifestWork{
109-
ObjectMeta: metav1.ObjectMeta{
110-
Name: name,
111-
Namespace: namespace,
112-
},
113-
}
114-
return workClient.Delete(ctx, work)
115-
},
116-
patchWork: func(ctx context.Context, namespace, name string, pt types.PatchType, data []byte) (*workapiv1.ManifestWork, error) {
117-
work := &workapiv1.ManifestWork{
118-
ObjectMeta: metav1.ObjectMeta{
119-
Name: name,
120-
Namespace: namespace,
121-
},
122-
}
123-
if err := workClient.Patch(ctx, work, client.RawPatch(pt, data)); err != nil {
124-
return nil, err
125-
}
126-
if err := workClient.Get(ctx, types.NamespacedName{Namespace: namespace, Name: name}, work); err != nil {
127-
return nil, err
128-
}
129-
return work, nil
130-
},
131-
createWork: func(ctx context.Context, work *workapiv1.ManifestWork) (*workapiv1.ManifestWork, error) {
132-
err := workClient.Create(ctx, work)
133-
return work, err
134-
},
135-
}
136-
}
137-
```
94+
- `NewWorkApplierWithTypedClient(workClient workv1client.Interface, workLister worklister.ManifestWorkLister)`, for use with a typed ManifestWork client and a lister (e.g., from a SharedInformerFactory). Reads are served from the lister cache.
95+
- `NewWorkApplierWithRuntimeClient(workClient client.Client)`, for use with a controller-runtime `client.Client` (e.g., in a controller-runtime based controller). Reads go through the client's cache.
13896
13997
2. Apply a manifestWork.
14098

pkg/apis/work/v1/applier/workapplier.go

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@ import (
1414
"k8s.io/apimachinery/pkg/types"
1515
"k8s.io/klog/v2"
1616

17+
"sigs.k8s.io/controller-runtime/pkg/client"
18+
1719
workv1client "open-cluster-management.io/api/client/work/clientset/versioned"
1820
worklister "open-cluster-management.io/api/client/work/listers/work/v1"
1921
workapiv1 "open-cluster-management.io/api/work/v1"
@@ -46,6 +48,45 @@ func NewWorkApplierWithTypedClient(workClient workv1client.Interface,
4648
}
4749
}
4850

51+
func NewWorkApplierWithRuntimeClient(workClient client.Client) *WorkApplier {
52+
return &WorkApplier{
53+
cache: newWorkCache(),
54+
getWork: func(ctx context.Context, namespace, name string) (*workapiv1.ManifestWork, error) {
55+
work := &workapiv1.ManifestWork{}
56+
err := workClient.Get(ctx, types.NamespacedName{Namespace: namespace, Name: name}, work)
57+
return work, err
58+
},
59+
deleteWork: func(ctx context.Context, namespace, name string) error {
60+
work := &workapiv1.ManifestWork{
61+
ObjectMeta: metav1.ObjectMeta{
62+
Name: name,
63+
Namespace: namespace,
64+
},
65+
}
66+
return workClient.Delete(ctx, work)
67+
},
68+
patchWork: func(ctx context.Context, namespace, name string, pt types.PatchType, data []byte) (*workapiv1.ManifestWork, error) {
69+
work := &workapiv1.ManifestWork{
70+
ObjectMeta: metav1.ObjectMeta{
71+
Name: name,
72+
Namespace: namespace,
73+
},
74+
}
75+
if err := workClient.Patch(ctx, work, client.RawPatch(pt, data)); err != nil {
76+
return nil, err
77+
}
78+
if err := workClient.Get(ctx, types.NamespacedName{Namespace: namespace, Name: name}, work); err != nil {
79+
return nil, err
80+
}
81+
return work, nil
82+
},
83+
createWork: func(ctx context.Context, work *workapiv1.ManifestWork) (*workapiv1.ManifestWork, error) {
84+
err := workClient.Create(ctx, work)
85+
return work, err
86+
},
87+
}
88+
}
89+
4990
func (w *WorkApplier) Apply(ctx context.Context, work *workapiv1.ManifestWork) (*workapiv1.ManifestWork, error) {
5091
existingWork, err := w.getWork(ctx, work.Namespace, work.Name)
5192
existingWork = existingWork.DeepCopy()

pkg/apis/work/v1/applier/workapplier_test.go

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,28 @@
11
package applier
22

33
import (
4+
"bytes"
45
"context"
56
"testing"
7+
68
"time"
79

10+
"github.com/google/go-cmp/cmp"
11+
812
v1 "k8s.io/api/apps/v1"
913
apiequality "k8s.io/apimachinery/pkg/api/equality"
1014
apierrors "k8s.io/apimachinery/pkg/api/errors"
1115
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
1216
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
1317
"k8s.io/apimachinery/pkg/runtime"
1418
"k8s.io/apimachinery/pkg/runtime/serializer"
19+
"k8s.io/apimachinery/pkg/types"
1520
clienttesting "k8s.io/client-go/testing"
1621
fakework "open-cluster-management.io/api/client/work/clientset/versioned/fake"
1722
workinformers "open-cluster-management.io/api/client/work/informers/externalversions"
1823
workapiv1 "open-cluster-management.io/api/work/v1"
24+
"sigs.k8s.io/controller-runtime/pkg/client"
25+
"sigs.k8s.io/controller-runtime/pkg/client/fake"
1926
)
2027

2128
// assertActions asserts the actual actions have the expected action verb
@@ -50,6 +57,7 @@ func newUnstructured(apiVersion, kind, namespace, name string) *unstructured.Uns
5057

5158
func newFakeWork(name, namespace string, obj runtime.Object) *workapiv1.ManifestWork {
5259
rawObject, _ := runtime.Encode(unstructured.UnstructuredJSONScheme, obj)
60+
rawObject = bytes.TrimRight(rawObject, "\n")
5361

5462
return &workapiv1.ManifestWork{
5563
ObjectMeta: metav1.ObjectMeta{
@@ -183,6 +191,119 @@ func TestWorkApplierWithTypedClient(t *testing.T) {
183191
assertActions(t, fakeWorkClient.Actions(), "delete")
184192
}
185193

194+
func getWork(t *testing.T, c client.Client, namespace, name string) *workapiv1.ManifestWork {
195+
t.Helper()
196+
work := &workapiv1.ManifestWork{}
197+
if err := c.Get(context.TODO(), types.NamespacedName{Namespace: namespace, Name: name}, work); err != nil {
198+
t.Fatalf("failed to get work %s/%s: %v", namespace, name, err)
199+
}
200+
return work
201+
}
202+
203+
func assertWorkState(t *testing.T, c client.Client, namespace, name string, desired *workapiv1.ManifestWork) {
204+
t.Helper()
205+
actual := getWork(t, c, namespace, name)
206+
if diff := cmp.Diff(desired.Spec, actual.Spec); diff != "" {
207+
t.Fatalf("spec of %s/%s mismatch (-want +got):\n%s", namespace, name, diff)
208+
}
209+
if diff := cmp.Diff(desired.Annotations, actual.Annotations); diff != "" {
210+
t.Fatalf("annotations of %s/%s mismatch (-want +got):\n%s", namespace, name, diff)
211+
}
212+
}
213+
214+
func TestWorkApplierWithRuntimeClient(t *testing.T) {
215+
scheme := runtime.NewScheme()
216+
if err := workapiv1.Install(scheme); err != nil {
217+
t.Fatalf("failed to add work scheme: %v", err)
218+
}
219+
220+
fakeClient := fake.NewClientBuilder().WithScheme(scheme).Build()
221+
workApplier := NewWorkApplierWithRuntimeClient(fakeClient)
222+
ctx := context.TODO()
223+
224+
baseWork := newFakeWork("test", "test", newUnstructured("batch/v1", "Job", "default", "test"))
225+
desired := baseWork.DeepCopy()
226+
227+
// Create: apply a MW that doesn't exist, verify it's persisted
228+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
229+
t.Fatalf("failed to create work: %v", err)
230+
}
231+
assertWorkState(t, fakeClient, "test", "test", desired)
232+
233+
// Update: change the desired spec, verify the object is patched
234+
desired.Spec.DeleteOption = &workapiv1.DeleteOption{PropagationPolicy: workapiv1.DeletePropagationPolicyTypeForeground}
235+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
236+
t.Fatalf("failed to update work: %v", err)
237+
}
238+
assertWorkState(t, fakeClient, "test", "test", desired)
239+
240+
// Annotation add: verify annotations are patched
241+
desired.SetAnnotations(map[string]string{workapiv1.ManifestConfigSpecHashAnnotationKey: "hash"})
242+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
243+
t.Fatalf("failed to add annotation: %v", err)
244+
}
245+
assertWorkState(t, fakeClient, "test", "test", desired)
246+
247+
// Annotation remove: verify annotations are cleared
248+
desired.Annotations = nil
249+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
250+
t.Fatalf("failed to remove annotation: %v", err)
251+
}
252+
assertWorkState(t, fakeClient, "test", "test", desired)
253+
254+
// Cache hit: same desired, verify no write via unchanged resourceVersion
255+
rvBefore := getWork(t, fakeClient, "test", "test").ResourceVersion
256+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
257+
t.Fatalf("failed to re-apply unchanged work: %v", err)
258+
}
259+
rvAfter := getWork(t, fakeClient, "test", "test").ResourceVersion
260+
if rvBefore != rvAfter {
261+
t.Fatalf("expected no write, but resourceVersion changed from %s to %s", rvBefore, rvAfter)
262+
}
263+
264+
// Cache hit when generation unchanged: externally modify the spec without
265+
// bumping generation. The cache still sees matching generation + desired
266+
// hash, so it skips the apply.
267+
tampered := getWork(t, fakeClient, "test", "test").DeepCopy()
268+
tampered.Spec.DeleteOption = &workapiv1.DeleteOption{PropagationPolicy: workapiv1.DeletePropagationPolicyTypeOrphan}
269+
if err := fakeClient.Update(ctx, tampered); err != nil {
270+
t.Fatalf("failed to externally modify work: %v", err)
271+
}
272+
rvBefore = getWork(t, fakeClient, "test", "test").ResourceVersion
273+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
274+
t.Fatalf("failed to re-apply after external modification without generation bump: %v", err)
275+
}
276+
rvAfter = getWork(t, fakeClient, "test", "test").ResourceVersion
277+
if rvBefore != rvAfter {
278+
t.Fatalf("expected no write when generation unchanged, but resourceVersion changed from %s to %s", rvBefore, rvAfter)
279+
}
280+
281+
// External modification with generation bump: simulate a real API server
282+
// spec change (e.g., kubectl edit or another addon manager).
283+
// The applier should revert the work back to its desired state.
284+
tampered.Generation++
285+
if err := fakeClient.Update(ctx, tampered); err != nil {
286+
t.Fatalf("failed to externally modify work: %v", err)
287+
}
288+
if _, err := workApplier.Apply(ctx, desired.DeepCopy()); err != nil {
289+
t.Fatalf("failed to restore drifted work: %v", err)
290+
}
291+
assertWorkState(t, fakeClient, "test", "test", desired)
292+
293+
// Delete: verify object is removed
294+
if err := workApplier.Delete(ctx, "test", "test"); err != nil {
295+
t.Fatalf("failed to delete work: %v", err)
296+
}
297+
if err := fakeClient.Get(ctx, types.NamespacedName{Name: "test", Namespace: "test"}, &workapiv1.ManifestWork{}); !apierrors.IsNotFound(err) {
298+
t.Fatalf("expected NotFound after delete, got: %v", err)
299+
}
300+
301+
// Delete nonexistent: verify idempotency (no error)
302+
if err := workApplier.Delete(ctx, "test", "nonexistent"); err != nil {
303+
t.Fatalf("expected no error deleting nonexistent work, got: %v", err)
304+
}
305+
}
306+
186307
var deploymentJson = `{
187308
"apiVersion": "apps/v1",
188309
"kind": "Deployment",

0 commit comments

Comments
 (0)