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
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user