From 41276676eb62912a58cc05684e29986f5731563c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Roberto=20Jim=C3=A9nez=20S=C3=A1nchez?= Date: Thu, 20 Nov 2025 15:12:07 +0100 Subject: [PATCH] Provisioning: add retry logic for transient errors in Kubernetes client (#114215) * feat: add retry logic for transient errors in Kubernetes client Add retry wrapper for dynamic.ResourceInterface that automatically retries transient errors using Kubernetes' wait.ExponentialBackoff utility. - Implements retry logic with exponential backoff for all Kubernetes API operations - Detects transient errors: ServiceUnavailable, ServerTimeout, TooManyRequests, InternalError, Timeout, and network errors - Uses wait.ExponentialBackoff from k8s.io/apimachinery/pkg/util/wait - Respects context cancellation - Includes comprehensive unit tests Part of https://github.com/grafana/git-ui-sync-project/issues/634 * docs: add detailed documentation for defaultRetryBackoff Document when retry attempts will happen, what errors trigger retries, and the retry behavior (attempts, delays, exponential backoff, jitter). * feat: add logging and increase retry attempts for Kubernetes client - Add context logger to track retry attempts (Info for retries, Warn for exhaustion) - Increase retry attempts from 5 to 8 steps (~10 seconds total retry window) - Document when all retry attempts will fail: * API server completely unavailable/unreachable * Network connectivity issues persist beyond retry window * Consistent transient errors for entire retry duration * Context cancellation before retries complete * chore: update retry client documentation * fix: resolve linting issues in retry client - Replace type assertions with errors.As for wrapped errors - Remove deprecated Temporary() check (deprecated since Go 1.18) - Update tests to remove temporary error test case * fix: resolve staticcheck S1008 linting issue in retry_client.go Simplify return statement to use errors.As directly instead of if-return pattern --- .../apis/provisioning/resources/client.go | 6 +- .../provisioning/resources/retry_client.go | 296 ++++++++++ .../resources/retry_client_test.go | 507 ++++++++++++++++++ 3 files changed, 807 insertions(+), 2 deletions(-) create mode 100644 pkg/registry/apis/provisioning/resources/retry_client.go create mode 100644 pkg/registry/apis/provisioning/resources/retry_client_test.go diff --git a/pkg/registry/apis/provisioning/resources/client.go b/pkg/registry/apis/provisioning/resources/client.go index f12bcf7b053..388e5311f72 100644 --- a/pkg/registry/apis/provisioning/resources/client.go +++ b/pkg/registry/apis/provisioning/resources/client.go @@ -221,10 +221,11 @@ func (c *resourceClients) ForKind(ctx context.Context, gvk schema.GroupVersionKi return nil, schema.GroupVersionResource{}, err } } + baseClient := dynamic.Resource(gvr).Namespace(c.namespace) info = &clientInfo{ gvk: gvk, gvr: gvr, - client: dynamic.Resource(gvr).Namespace(c.namespace), + client: newRetryResourceInterface(baseClient, defaultRetryBackoff()), } c.byKind[gvk] = info c.byResource[gvr] = info @@ -274,10 +275,11 @@ func (c *resourceClients) ForResource(ctx context.Context, gvr schema.GroupVersi return nil, schema.GroupVersionKind{}, fmt.Errorf("getting kind for resource for %s: %w", gvr.String(), err) } } + baseClient := dynamic.Resource(gvr).Namespace(c.namespace) info = &clientInfo{ gvk: gvk, gvr: gvr, - client: dynamic.Resource(gvr).Namespace(c.namespace), + client: newRetryResourceInterface(baseClient, defaultRetryBackoff()), } c.byKind[gvk] = info c.byResource[gvr] = info diff --git a/pkg/registry/apis/provisioning/resources/retry_client.go b/pkg/registry/apis/provisioning/resources/retry_client.go new file mode 100644 index 00000000000..cc09443e018 --- /dev/null +++ b/pkg/registry/apis/provisioning/resources/retry_client.go @@ -0,0 +1,296 @@ +package resources + +import ( + "context" + "errors" + "net" + "time" + + "github.com/grafana/grafana-app-sdk/logging" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/dynamic" +) + +// defaultRetryBackoff returns a default backoff configuration for retries. +// +// Retry attempts will happen when: +// - The Kubernetes API returns transient errors: ServiceUnavailable (503), ServerTimeout (504), +// TooManyRequests (429), InternalError (500), or Timeout errors +// - Network errors occur: connection timeouts, temporary network failures, or connection errors +// +// The retry behavior: +// - Total attempts: 8 (1 initial attempt + 7 retries) +// - Initial delay: 100ms before the first retry +// - Exponential backoff: delay doubles after each failed attempt (100ms → 200ms → 400ms → 800ms → 1.6s → 3.2s → 5s) +// - Maximum delay: capped at 5 seconds +// - Jitter: 10% randomization to prevent thundering herd problems +// - Total retry window: approximately 10 seconds from first attempt to last retry +// +// All attempts will fail when: +// - The Kubernetes API server is completely unavailable or unreachable +// - Network connectivity issues persist beyond the retry window (~10 seconds) +// - The API server returns transient errors consistently for the entire retry duration +// - Context cancellation occurs before retries complete +// +// Non-transient errors (e.g., NotFound, BadRequest, Forbidden) are not retried and returned immediately. +func defaultRetryBackoff() wait.Backoff { + return wait.Backoff{ + Duration: 100 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 8, // 1 initial attempt + 7 retries = 8 total attempts (~10s total retry window) + Cap: 5 * time.Second, + } +} + +// retryResourceInterface wraps a dynamic.ResourceInterface with retry logic for transient errors +type retryResourceInterface struct { + client dynamic.ResourceInterface + backoff wait.Backoff +} + +// newRetryResourceInterface creates a new ResourceInterface wrapper with retry logic +func newRetryResourceInterface(client dynamic.ResourceInterface, backoff wait.Backoff) dynamic.ResourceInterface { + return &retryResourceInterface{ + client: client, + backoff: backoff, + } +} + +// isTransientError determines if an error is transient and should be retried +func isTransientError(err error) bool { + if err == nil { + return false + } + + // Check for Kubernetes API transient errors + if apierrors.IsServiceUnavailable(err) { + return true + } + if apierrors.IsServerTimeout(err) { + return true + } + if apierrors.IsTooManyRequests(err) { + return true + } + if apierrors.IsInternalError(err) { + return true + } + if apierrors.IsTimeout(err) { + return true + } + + // Check for network errors + var netErr net.Error + if errors.As(err, &netErr) { + if netErr.Timeout() { + return true + } + } + + // Check for connection errors + var opErr *net.OpError + return errors.As(err, &opErr) +} + +// retryWithBackoff executes a function with exponential backoff retry logic using wait.ExponentialBackoff +func (r *retryResourceInterface) retryWithBackoff(ctx context.Context, fn func() error) error { + var lastErr error + attempt := 0 + logger := logging.FromContext(ctx) + + err := wait.ExponentialBackoff(r.backoff, func() (bool, error) { + attempt++ + + // Check if context is cancelled + if ctx.Err() != nil { + logger.Debug("Retry cancelled due to context cancellation", "attempt", attempt) + return false, ctx.Err() + } + + err := fn() + if err == nil { + if attempt > 1 { + logger.Debug("Operation succeeded after retry", "attempt", attempt) + } + return true, nil // success, stop retrying + } + + // If not a transient error, return immediately without retrying + if !isTransientError(err) { + logger.Debug("Non-transient error, not retrying", "attempt", attempt, "error", err) + return false, err + } + + // Transient error, retry + lastErr = err + logger.Info("Transient error encountered, retrying", "attempt", attempt, "max_attempts", r.backoff.Steps, "error", err) + return false, nil + }) + + // If wait.ExponentialBackoff returned an error, it means we exhausted retries + if err != nil { + if lastErr != nil { + logger.Warn("All retry attempts exhausted", "total_attempts", attempt, "error", lastErr) + return lastErr + } + logger.Warn("All retry attempts exhausted", "total_attempts", attempt, "error", err) + return err + } + + return nil +} + +// Create implements dynamic.ResourceInterface +func (r *retryResourceInterface) Create(ctx context.Context, obj *unstructured.Unstructured, options metav1.CreateOptions, subresources ...string) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.Create(ctx, obj, options, subresources...) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// Update implements dynamic.ResourceInterface +func (r *retryResourceInterface) Update(ctx context.Context, obj *unstructured.Unstructured, options metav1.UpdateOptions, subresources ...string) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.Update(ctx, obj, options, subresources...) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// UpdateStatus implements dynamic.ResourceInterface +func (r *retryResourceInterface) UpdateStatus(ctx context.Context, obj *unstructured.Unstructured, options metav1.UpdateOptions) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.UpdateStatus(ctx, obj, options) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// Delete implements dynamic.ResourceInterface +func (r *retryResourceInterface) Delete(ctx context.Context, name string, options metav1.DeleteOptions, subresources ...string) error { + return r.retryWithBackoff(ctx, func() error { + return r.client.Delete(ctx, name, options, subresources...) + }) +} + +// DeleteCollection implements dynamic.ResourceInterface +func (r *retryResourceInterface) DeleteCollection(ctx context.Context, options metav1.DeleteOptions, listOptions metav1.ListOptions) error { + return r.retryWithBackoff(ctx, func() error { + return r.client.DeleteCollection(ctx, options, listOptions) + }) +} + +// Get implements dynamic.ResourceInterface +func (r *retryResourceInterface) Get(ctx context.Context, name string, options metav1.GetOptions, subresources ...string) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.Get(ctx, name, options, subresources...) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// List implements dynamic.ResourceInterface +func (r *retryResourceInterface) List(ctx context.Context, opts metav1.ListOptions) (*unstructured.UnstructuredList, error) { + var result *unstructured.UnstructuredList + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.List(ctx, opts) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// Watch implements dynamic.ResourceInterface +func (r *retryResourceInterface) Watch(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { + // Watch operations are long-lived and shouldn't be retried in the same way + // Return the watch interface directly + return r.client.Watch(ctx, opts) +} + +// Patch implements dynamic.ResourceInterface +func (r *retryResourceInterface) Patch(ctx context.Context, name string, pt types.PatchType, data []byte, options metav1.PatchOptions, subresources ...string) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.Patch(ctx, name, pt, data, options, subresources...) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// Apply implements dynamic.ResourceInterface +func (r *retryResourceInterface) Apply(ctx context.Context, name string, obj *unstructured.Unstructured, options metav1.ApplyOptions, subresources ...string) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.Apply(ctx, name, obj, options, subresources...) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} + +// ApplyStatus implements dynamic.ResourceInterface +func (r *retryResourceInterface) ApplyStatus(ctx context.Context, name string, obj *unstructured.Unstructured, options metav1.ApplyOptions) (*unstructured.Unstructured, error) { + var result *unstructured.Unstructured + var err error + + retryErr := r.retryWithBackoff(ctx, func() error { + result, err = r.client.ApplyStatus(ctx, name, obj, options) + return err + }) + + if retryErr != nil { + return nil, retryErr + } + return result, nil +} diff --git a/pkg/registry/apis/provisioning/resources/retry_client_test.go b/pkg/registry/apis/provisioning/resources/retry_client_test.go new file mode 100644 index 00000000000..dab40bcb001 --- /dev/null +++ b/pkg/registry/apis/provisioning/resources/retry_client_test.go @@ -0,0 +1,507 @@ +package resources + +import ( + "context" + "errors" + "net" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/apimachinery/pkg/watch" +) + +func TestIsTransientError(t *testing.T) { + tests := []struct { + name string + err error + expected bool + }{ + { + name: "nil error", + err: nil, + expected: false, + }, + { + name: "service unavailable", + err: apierrors.NewServiceUnavailable("service unavailable"), + expected: true, + }, + { + name: "server timeout", + err: apierrors.NewServerTimeout(schema.GroupResource{}, "operation", 0), + expected: true, + }, + { + name: "too many requests", + err: apierrors.NewTooManyRequests("too many requests", 0), + expected: true, + }, + { + name: "internal error", + err: apierrors.NewInternalError(errors.New("internal error")), + expected: true, + }, + { + name: "network timeout error", + err: &net.DNSError{Err: "timeout", IsTimeout: true}, + expected: true, + }, + // Note: Temporary() is deprecated in Go 1.18+, so we no longer check for temporary errors + // Timeout errors are still checked and will be retried + { + name: "network op error", + err: &net.OpError{Op: "read", Err: errors.New("connection refused")}, + expected: true, + }, + { + name: "not found error", + err: apierrors.NewNotFound(schema.GroupResource{}, "resource"), + expected: false, + }, + { + name: "bad request error", + err: apierrors.NewBadRequest("bad request"), + expected: false, + }, + { + name: "forbidden error", + err: apierrors.NewForbidden(schema.GroupResource{}, "resource", errors.New("forbidden")), + expected: false, + }, + { + name: "generic error", + err: errors.New("generic error"), + expected: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + result := isTransientError(tt.err) + assert.Equal(t, tt.expected, result, "isTransientError(%v) = %v, want %v", tt.err, result, tt.expected) + }) + } +} + +func TestRetryResourceInterface_Create(t *testing.T) { + tests := []struct { + name string + setupMock func(*MockDynamicResourceInterface) + backoff wait.Backoff + expectedCalls int + expectError bool + }{ + { + name: "success on first attempt", + setupMock: func(m *MockDynamicResourceInterface) { + m.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(&unstructured.Unstructured{}, nil).Once() + }, + backoff: defaultRetryBackoff(), + expectedCalls: 1, + expectError: false, + }, + { + name: "success after transient errors", + setupMock: func(m *MockDynamicResourceInterface) { + m.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Twice() + m.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + }, + backoff: wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }, + expectedCalls: 3, + expectError: false, + }, + { + name: "max retries exceeded", + setupMock: func(m *MockDynamicResourceInterface) { + m.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")) + }, + backoff: wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 3, // Only 3 steps = 2 retries + 1 initial + Cap: 100 * time.Millisecond, + }, + expectedCalls: 3, + expectError: true, + }, + { + name: "non-transient error - no retry", + setupMock: func(m *MockDynamicResourceInterface) { + m.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewBadRequest("bad request")).Once() + }, + backoff: defaultRetryBackoff(), + expectedCalls: 1, + expectError: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + tt.setupMock(mockClient) + + retryClient := newRetryResourceInterface(mockClient, tt.backoff) + obj := &unstructured.Unstructured{} + _, err := retryClient.Create(context.Background(), obj, metav1.CreateOptions{}) + + if tt.expectError { + assert.Error(t, err) + } else { + assert.NoError(t, err) + } + mockClient.AssertNumberOfCalls(t, "Create", tt.expectedCalls) + }) + } +} + +func TestRetryResourceInterface_Update(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("Update", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Once() + mockClient.On("Update", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + obj := &unstructured.Unstructured{} + result, err := retryClient.Update(context.Background(), obj, metav1.UpdateOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "Update", 2) +} + +func TestRetryResourceInterface_Get(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("Get", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServerTimeout(schema.GroupResource{}, "operation", 0)).Twice() + mockClient.On("Get", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + result, err := retryClient.Get(context.Background(), "test-resource", metav1.GetOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "Get", 3) +} + +func TestRetryResourceInterface_Delete(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("Delete", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(apierrors.NewTooManyRequests("too many requests", 0)).Once() + mockClient.On("Delete", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + err := retryClient.Delete(context.Background(), "test-resource", metav1.DeleteOptions{}) + + assert.NoError(t, err) + mockClient.AssertNumberOfCalls(t, "Delete", 2) +} + +func TestRetryResourceInterface_List(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("List", mock.Anything, mock.Anything). + Return(nil, apierrors.NewInternalError(errors.New("internal error"))).Once() + mockClient.On("List", mock.Anything, mock.Anything). + Return(&unstructured.UnstructuredList{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + result, err := retryClient.List(context.Background(), metav1.ListOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "List", 2) +} + +func TestRetryResourceInterface_Patch(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("Patch", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, &net.OpError{Op: "read", Err: errors.New("connection refused")}).Once() + mockClient.On("Patch", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + result, err := retryClient.Patch(context.Background(), "test-resource", types.MergePatchType, []byte(`{}`), metav1.PatchOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "Patch", 2) +} + +func TestRetryResourceInterface_Apply(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("Apply", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Once() + mockClient.On("Apply", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + obj := &unstructured.Unstructured{} + result, err := retryClient.Apply(context.Background(), "test-resource", obj, metav1.ApplyOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "Apply", 2) +} + +func TestRetryResourceInterface_UpdateStatus(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("UpdateStatus", mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Once() + mockClient.On("UpdateStatus", mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + obj := &unstructured.Unstructured{} + result, err := retryClient.UpdateStatus(context.Background(), obj, metav1.UpdateOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "UpdateStatus", 2) +} + +func TestRetryResourceInterface_ApplyStatus(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("ApplyStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Once() + mockClient.On("ApplyStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(&unstructured.Unstructured{}, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + obj := &unstructured.Unstructured{} + result, err := retryClient.ApplyStatus(context.Background(), "test-resource", obj, metav1.ApplyOptions{}) + + assert.NoError(t, err) + assert.NotNil(t, result) + mockClient.AssertNumberOfCalls(t, "ApplyStatus", 2) +} + +func TestRetryResourceInterface_DeleteCollection(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("DeleteCollection", mock.Anything, mock.Anything, mock.Anything). + Return(apierrors.NewServiceUnavailable("service unavailable")).Once() + mockClient.On("DeleteCollection", mock.Anything, mock.Anything, mock.Anything). + Return(nil).Once() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 10 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 100 * time.Millisecond, + }) + + err := retryClient.DeleteCollection(context.Background(), metav1.DeleteOptions{}, metav1.ListOptions{}) + + assert.NoError(t, err) + mockClient.AssertNumberOfCalls(t, "DeleteCollection", 2) +} + +func TestRetryResourceInterface_Watch(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockWatch := &mockWatch{} + mockClient.On("Watch", mock.Anything, mock.Anything).Return(mockWatch, nil).Once() + + retryClient := newRetryResourceInterface(mockClient, defaultRetryBackoff()) + + watch, err := retryClient.Watch(context.Background(), metav1.ListOptions{}) + + assert.NoError(t, err) + assert.Equal(t, mockWatch, watch) + // Watch should not retry, so only one call + mockClient.AssertNumberOfCalls(t, "Watch", 1) +} + +func TestRetryResourceInterface_ContextCancellation(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + mockClient.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Maybe() + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 100 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 1 * time.Second, + }) + + ctx, cancel := context.WithCancel(context.Background()) + cancel() // Cancel immediately + + obj := &unstructured.Unstructured{} + _, err := retryClient.Create(ctx, obj, metav1.CreateOptions{}) + + assert.Error(t, err) + assert.Equal(t, context.Canceled, err) + // Should not retry after context cancellation + mockClient.AssertNumberOfCalls(t, "Create", 0) +} + +func TestRetryResourceInterface_ExponentialBackoff(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + callCount := 0 + + // Set up expectations for multiple calls + mockClient.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Run(func(args mock.Arguments) { + callCount++ + }). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")).Twice() + mockClient.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Run(func(args mock.Arguments) { + callCount++ + }). + Return(&unstructured.Unstructured{}, nil).Once() + + start := time.Now() + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 50 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, + Cap: 1 * time.Second, + }) + + obj := &unstructured.Unstructured{} + _, err := retryClient.Create(context.Background(), obj, metav1.CreateOptions{}) + + duration := time.Since(start) + + assert.NoError(t, err) + assert.Equal(t, 3, callCount) + // Should have waited at least initialDelay + (initialDelay * multiplier) = 50ms + 100ms = 150ms + assert.GreaterOrEqual(t, duration, 100*time.Millisecond) + // But not too long (with some buffer for jitter) + assert.Less(t, duration, 1*time.Second) +} + +func TestRetryResourceInterface_MaxDelayRespected(t *testing.T) { + mockClient := &MockDynamicResourceInterface{} + callCount := 0 + + // Set up expectations for multiple calls - always return transient error + mockClient.On("Create", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Run(func(args mock.Arguments) { + callCount++ + }). + Return(nil, apierrors.NewServiceUnavailable("service unavailable")). + Maybe() // Allow unlimited calls + + retryClient := newRetryResourceInterface(mockClient, wait.Backoff{ + Duration: 50 * time.Millisecond, + Factor: 2.0, + Jitter: 0.1, + Steps: 5, // 5 total attempts + Cap: 200 * time.Millisecond, // Max delay is 200ms (cap should prevent exponential growth beyond this) + }) + + start := time.Now() + obj := &unstructured.Unstructured{} + _, err := retryClient.Create(context.Background(), obj, metav1.CreateOptions{}) + duration := time.Since(start) + + assert.Error(t, err) + // Should have retried multiple times + assert.GreaterOrEqual(t, callCount, 2, "should have retried at least once") + // Verify that the delay was capped - if it wasn't capped, duration would be much longer + // With cap at 200ms and Steps=5, max total time should be reasonable + // The cap ensures delays don't grow exponentially beyond 200ms + assert.Less(t, duration, 2*time.Second, "duration should be reasonable due to cap") +} + +func TestDefaultRetryBackoff(t *testing.T) { + backoff := defaultRetryBackoff() + + assert.Equal(t, 100*time.Millisecond, backoff.Duration) + assert.Equal(t, 2.0, backoff.Factor) + assert.Equal(t, 0.1, backoff.Jitter) + assert.Equal(t, 8, backoff.Steps) // Updated to 8 steps for ~10s total retry window + assert.Equal(t, 5*time.Second, backoff.Cap) +} + +// mockWatch implements watch.Interface for testing +type mockWatch struct{} + +func (m *mockWatch) Stop() {} + +func (m *mockWatch) ResultChan() <-chan watch.Event { + return make(chan watch.Event) +} + +var _ watch.Interface = (*mockWatch)(nil)