Skip to content

Commit 2c22006

Browse files
committed
address feedback
Signed-off-by: Marcin Owsiany <porridge@redhat.com>
1 parent e8366c3 commit 2c22006

4 files changed

Lines changed: 63 additions & 66 deletions

File tree

internal/harness/harness.go

Lines changed: 33 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,11 @@ import (
1717

1818
docker "github.com/moby/moby/client"
1919
"gopkg.in/yaml.v2"
20+
corev1 "k8s.io/api/core/v1"
21+
apiErrors "k8s.io/apimachinery/pkg/api/errors"
2022
"k8s.io/apimachinery/pkg/labels"
2123
"k8s.io/apimachinery/pkg/runtime"
24+
"k8s.io/apimachinery/pkg/util/wait"
2225
"k8s.io/client-go/discovery"
2326
"k8s.io/client-go/rest"
2427
"k8s.io/client-go/tools/clientcmd"
@@ -289,13 +292,13 @@ func (h *Harness) Config() (*rest.Config, error) {
289292
}
290293

291294
func (h *Harness) waitForFunctionalCluster() error {
292-
err := kubernetes.WaitForSA(h.config, "default", "default")
295+
err := h.waitForSA("default", "default")
293296
if err == nil {
294297
return nil
295298
}
296299
// if there is a namespace provided but no "default"/"default" SA found, also check a SA in the provided NS
297300
if h.TestSuite.Namespace != "" {
298-
tempErr := kubernetes.WaitForSA(h.config, "default", h.TestSuite.Namespace)
301+
tempErr := h.waitForSA("default", h.TestSuite.Namespace)
299302
if tempErr == nil {
300303
return nil
301304
}
@@ -304,6 +307,34 @@ func (h *Harness) waitForFunctionalCluster() error {
304307
return err
305308
}
306309

310+
// waitForSA waits for a service account to be present.
311+
func (h *Harness) waitForSA(name, namespace string) error {
312+
c, err := kubernetes.NewRetryClient(h.config, client.Options{
313+
Scheme: kubernetes.Scheme(),
314+
})
315+
if err != nil {
316+
return err
317+
}
318+
319+
obj := &corev1.ServiceAccount{}
320+
321+
key := client.ObjectKey{
322+
Namespace: namespace,
323+
Name: name,
324+
}
325+
return wait.PollUntilContextTimeout(context.TODO(), 500*time.Millisecond, 60*time.Second, true, func(ctx context.Context) (done bool, err error) {
326+
err = c.Get(ctx, key, obj)
327+
if apiErrors.IsNotFound(err) {
328+
return false, nil
329+
}
330+
if err != nil {
331+
h.T.Logf("Error waiting for service account %s/%s (will retry): %v", namespace, name, err)
332+
return false, nil
333+
}
334+
return true, nil
335+
})
336+
}
337+
307338
// Client returns the current Kubernetes client for the test harness.
308339
func (h *Harness) Client(forceNew bool) (client.Client, error) {
309340
h.clientLock.Lock()

internal/kubernetes/wait.go

Lines changed: 15 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -2,54 +2,38 @@ package kubernetes
22

33
import (
44
"context"
5+
"fmt"
56
"time"
67

7-
corev1 "k8s.io/api/core/v1"
88
"k8s.io/apimachinery/pkg/api/errors"
99
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
10-
"k8s.io/apimachinery/pkg/runtime"
1110
"k8s.io/apimachinery/pkg/util/wait"
12-
"k8s.io/client-go/rest"
1311
"sigs.k8s.io/controller-runtime/pkg/client"
1412
)
1513

16-
// WaitForDelete waits for the provide runtime objects to be deleted from cluster
17-
func WaitForDelete(c *RetryClient, objs []runtime.Object) error {
18-
// Wait for resources to be deleted.
19-
return wait.PollUntilContextTimeout(context.TODO(), 100*time.Millisecond, 10*time.Second, true, func(ctx context.Context) (done bool, err error) {
20-
for _, obj := range objs {
14+
// WaitForDelete waits for the provided runtime objects to be deleted from cluster, up to duration.
15+
// Retries on transient errors.
16+
func WaitForDelete(cl client.Client, toDelete []client.Object, duration time.Duration) error {
17+
lastCheckMsg := ""
18+
err := wait.PollUntilContextTimeout(context.TODO(), 100*time.Millisecond, duration, true, func(ctx context.Context) (done bool, err error) {
19+
for _, obj := range toDelete {
2120
actual := &unstructured.Unstructured{}
2221
actual.SetGroupVersionKind(obj.GetObjectKind().GroupVersionKind())
23-
err = c.Get(ctx, ObjectKey(obj), actual)
24-
// Retry on transient API errors (return nil to keep polling) rather
25-
// than aborting the entire wait, which would surface as a misleading
26-
// "timed out" error well before the real deadline.
22+
err = cl.Get(ctx, ObjectKey(obj), actual)
23+
if err == nil {
24+
lastCheckMsg = fmt.Sprintf("%v %s still exists", obj.GetObjectKind().GroupVersionKind(), obj.GetName())
25+
return false, nil
26+
}
2727
if !errors.IsNotFound(err) {
28+
lastCheckMsg = fmt.Sprintf("checking existence of %v %s failed: %v", obj.GetObjectKind().GroupVersionKind(), obj.GetName(), err)
2829
return false, nil
2930
}
3031
}
3132

3233
return true, nil
3334
})
34-
}
35-
36-
// WaitForSA waits for a service account to be present
37-
func WaitForSA(config *rest.Config, name, namespace string) error {
38-
c, err := NewRetryClient(config, client.Options{
39-
Scheme: Scheme(),
40-
})
4135
if err != nil {
42-
return err
36+
return fmt.Errorf("timed out waiting for resource deletion (result of last check was: %q): %w", lastCheckMsg, err)
4337
}
44-
45-
obj := &corev1.ServiceAccount{}
46-
47-
key := client.ObjectKey{
48-
Namespace: namespace,
49-
Name: name,
50-
}
51-
return wait.PollUntilContextTimeout(context.TODO(), 500*time.Millisecond, 60*time.Second, true, func(ctx context.Context) (done bool, err error) {
52-
// Retry on all errors (not-found, transient) rather than aborting the wait.
53-
return c.Get(ctx, key, obj) == nil, nil
54-
})
38+
return nil
5539
}

internal/kubernetes/wait_test.go

Lines changed: 14 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -4,19 +4,19 @@ import (
44
"context"
55
"fmt"
66
"testing"
7+
"time"
78

89
"github.com/stretchr/testify/assert"
910
"github.com/stretchr/testify/require"
1011
k8serrors "k8s.io/apimachinery/pkg/api/errors"
1112
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
12-
"k8s.io/apimachinery/pkg/runtime"
1313
"k8s.io/apimachinery/pkg/runtime/schema"
1414
"sigs.k8s.io/controller-runtime/pkg/client"
1515
"sigs.k8s.io/controller-runtime/pkg/client/fake"
1616
"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
1717
)
1818

19-
func testObj() runtime.Object {
19+
func testObj() *unstructured.Unstructured {
2020
obj := &unstructured.Unstructured{}
2121
obj.SetGroupVersionKind(schema.GroupVersionKind{Group: "", Version: "v1", Kind: "ConfigMap"})
2222
obj.SetName("cm")
@@ -26,9 +26,8 @@ func testObj() runtime.Object {
2626

2727
func TestWaitForDelete_AlreadyGone(t *testing.T) {
2828
cl := fake.NewClientBuilder().Build()
29-
rc := &RetryClient{Client: cl}
3029

31-
err := WaitForDelete(rc, []runtime.Object{testObj()})
30+
err := WaitForDelete(cl, []client.Object{testObj()}, time.Second*2)
3231
require.NoError(t, err)
3332
}
3433

@@ -43,9 +42,8 @@ func TestWaitForDelete_TransientErrorThenGone(t *testing.T) {
4342
return k8serrors.NewNotFound(schema.GroupResource{Resource: "configmaps"}, "cm")
4443
},
4544
}).Build()
46-
rc := &RetryClient{Client: cl}
4745

48-
err := WaitForDelete(rc, []runtime.Object{testObj()})
46+
err := WaitForDelete(cl, []client.Object{testObj()}, time.Second*2)
4947
require.NoError(t, err)
5048
assert.Greater(t, callCount, 3)
5149
}
@@ -61,22 +59,28 @@ func TestWaitForDelete_StillExistsThenGone(t *testing.T) {
6159
return k8serrors.NewNotFound(schema.GroupResource{Resource: "configmaps"}, "cm")
6260
},
6361
}).Build()
64-
rc := &RetryClient{Client: cl}
6562

66-
err := WaitForDelete(rc, []runtime.Object{testObj()})
63+
err := WaitForDelete(cl, []client.Object{testObj()}, time.Second*2)
6764
require.NoError(t, err)
6865
assert.Greater(t, callCount, 2)
6966
}
7067

7168
func TestWaitForDelete_PersistentErrorTimesOut(t *testing.T) {
69+
callCount := 0
7270
cl := fake.NewClientBuilder().WithInterceptorFuncs(interceptor.Funcs{
7371
Get: func(context.Context, client.WithWatch, client.ObjectKey, client.Object, ...client.GetOption) error {
72+
callCount++
73+
if callCount <= 2 {
74+
return fmt.Errorf("initial transient error")
75+
}
7476
return fmt.Errorf("persistent API error")
7577
},
7678
}).Build()
77-
rc := &RetryClient{Client: cl}
7879

79-
err := WaitForDelete(rc, []runtime.Object{testObj()})
80+
err := WaitForDelete(cl, []client.Object{testObj()}, time.Second*2)
8081
assert.Error(t, err)
8182
assert.ErrorIs(t, err, context.DeadlineExceeded)
83+
assert.ErrorContains(t, err, "result of last check was:")
84+
assert.ErrorContains(t, err, "failed: persistent API error")
85+
assert.NotContains(t, err.Error(), "initial transient error")
8286
}

internal/step/step.go

Lines changed: 1 addition & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -166,29 +166,7 @@ func (s *Step) DeleteExisting(namespace string) error {
166166
s.Logger.Log(kubernetes.ResourceID(del), action)
167167
}
168168

169-
// Wait for resources to be deleted.
170-
lastCheckMsg := ""
171-
err = wait.PollUntilContextTimeout(context.TODO(), 100*time.Millisecond, time.Duration(s.GetTimeout())*time.Second, true, func(ctx context.Context) (done bool, err error) {
172-
for _, obj := range toDelete {
173-
actual := &unstructured.Unstructured{}
174-
actual.SetGroupVersionKind(obj.GetObjectKind().GroupVersionKind())
175-
err = cl.Get(ctx, kubernetes.ObjectKey(obj), actual)
176-
if err == nil {
177-
lastCheckMsg = fmt.Sprintf("%v %s still exists", obj.GetObjectKind().GroupVersionKind(), obj.GetName())
178-
return false, nil
179-
}
180-
if !k8serrors.IsNotFound(err) {
181-
lastCheckMsg = fmt.Sprintf("checking existence of %v %s failed: %v", obj.GetObjectKind().GroupVersionKind(), obj.GetName(), err)
182-
return false, nil
183-
}
184-
}
185-
186-
return true, nil
187-
})
188-
if err != nil {
189-
return fmt.Errorf("timed out waiting for resource deletion (result of last check was: %q): %w", lastCheckMsg, err)
190-
}
191-
return nil
169+
return kubernetes.WaitForDelete(cl, toDelete, time.Duration(s.GetTimeout())*time.Second)
192170
}
193171

194172
// Create applies all resources defined in the Apply list.

0 commit comments

Comments
 (0)