diff --git a/pkg/steps/pod.go b/pkg/steps/pod.go index 3c86c5a995..f7d61dc706 100644 --- a/pkg/steps/pod.go +++ b/pkg/steps/pod.go @@ -85,6 +85,18 @@ type podStep struct { resolvedGSMCredentials []api.CredentialReference } +// PodStepError reports a failed PodStep together with the last pod state +// observed by WaitForPodCompletion. Callers that need failure-specific +// lifecycle handling can use errors.As without racing a replacement pod via a +// second API lookup. +type PodStepError struct { + Pod *coreapi.Pod + err error +} + +func (e *PodStepError) Error() string { return e.err.Error() } +func (e *PodStepError) Unwrap() error { return e.err } + func (s *podStep) Inputs() (api.InputDefinition, error) { return nil, nil } @@ -133,14 +145,6 @@ func (s *podStep) run(ctx context.Context) error { pod.OwnerReferences = append(pod.OwnerReferences, *owner) } - go func() { - <-ctx.Done() - logrus.Infof("cleanup: Deleting %s pod %s", s.name, s.config.As) - if err := s.client.Delete(CleanupCtx, &coreapi.Pod{ObjectMeta: meta.ObjectMeta{Namespace: s.jobSpec.Namespace(), Name: s.config.As}}); err != nil && !kerrors.IsNotFound(err) { - logrus.WithError(err).Warnf("Could not delete %s pod.", s.name) - } - }() - s.client.MetricsAgent().StoreMachinesSnapshot(pod) pod, err = util.CreateOrRestartPod(ctx, s.client, pod) @@ -148,11 +152,22 @@ func (s *podStep) run(ctx context.Context) error { return fmt.Errorf("failed to create or restart %s pod: %w", s.name, err) } + // Register cleanup only after creation, when the API-assigned UID is known. + // A later attempt may reuse the name, so deleting by name alone is unsafe. + cleanupPod := pod.DeepCopy() + go func() { + <-ctx.Done() + logrus.Infof("cleanup: Deleting %s pod %s", s.name, s.config.As) + if err := util.DeletePodWithUID(CleanupCtx, s.client, cleanupPod); err != nil { + logrus.WithError(err).Warnf("Could not delete %s pod.", s.name) + } + }() + defer func() { s.subTests = testCaseNotifier.SubTests(s.Description() + " - ") }() - if _, err := util.WaitForPodCompletion(ctx, s.client, pod.Namespace, pod.Name, testCaseNotifier, s.config.WaitFlags); err != nil { - return fmt.Errorf("%s %q failed: %w", s.name, pod.Name, err) + if observedPod, err := util.WaitForPodCompletion(ctx, s.client, pod.Namespace, pod.Name, testCaseNotifier, s.config.WaitFlags); err != nil { + return &PodStepError{Pod: observedPod, err: fmt.Errorf("%s %q failed: %w", s.name, pod.Name, err)} } return nil } diff --git a/pkg/steps/pod_error_test.go b/pkg/steps/pod_error_test.go new file mode 100644 index 0000000000..6061af552c --- /dev/null +++ b/pkg/steps/pod_error_test.go @@ -0,0 +1,27 @@ +package steps + +import ( + "errors" + "fmt" + "testing" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func TestPodStepErrorContract(t *testing.T) { + cause := errors.New("pod failed") + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Namespace: "test", Name: "failed-pod"}} + err := fmt.Errorf("outer context: %w", &PodStepError{Pod: pod, err: cause}) + + var podStepErr *PodStepError + if !errors.As(err, &podStepErr) { + t.Fatalf("expected errors.As to find PodStepError in %v", err) + } + if podStepErr.Pod != pod { + t.Fatalf("expected observed pod to be retained, got %#v", podStepErr.Pod) + } + if !errors.Is(err, cause) { + t.Fatalf("expected errors.Is to find wrapped cause in %v", err) + } +} diff --git a/pkg/steps/release/import_release.go b/pkg/steps/release/import_release.go index c28f3e7ff3..88c9d298ae 100644 --- a/pkg/steps/release/import_release.go +++ b/pkg/steps/release/import_release.go @@ -66,6 +66,126 @@ type importReleaseStep struct { // originalPullSpec stores the original value before resolving the release to pass as an env var for multi-stage steps to utilize originalPullSpec string + + releaseImportRetryDelays func() []time.Duration + importTagWithRetryDelays releaseTagImporter +} + +type releaseTagImporter func(context.Context, ctrlruntimeclient.Client, string, string, string, string, []time.Duration, *metrics.MetricsAgent) (string, error) + +const ( + releaseExtractionContainerName = "release" + transientReleaseExtractionExitCode = 75 + transientReleaseExtractionErrorPattern = `too many requests|toomanyrequests|status:?( code)? (429|5[[:digit:]]{2})|unexpected http status|internal server error|bad gateway|service unavailable|gateway timeout|connection (refused|reset)|http2: client connection lost|i/o timeout|TLS handshake timeout|context deadline exceeded|temporary failure|unexpected EOF` +) + +type releaseImportSleep func(context.Context, time.Duration) error + +func sleepForReleaseImportRetry(ctx context.Context, duration time.Duration) error { + timer := time.NewTimer(duration) + defer timer.Stop() + select { + case <-ctx.Done(): + return ctx.Err() + case <-timer.C: + return nil + } +} + +type transientReleaseExtractionError struct { + err error +} + +func (e *transientReleaseExtractionError) Error() string { return e.err.Error() } +func (e *transientReleaseExtractionError) Unwrap() error { return e.err } + +func retryReleaseExtraction(ctx context.Context, name string, retryDelays []time.Duration, sleep releaseImportSleep, run func(context.Context) error) error { + start := time.Now() + var retryBudget time.Duration + for _, delay := range retryDelays { + if delay < 0 { + return fmt.Errorf("invalid release extraction retry delay %s", delay) + } + retryBudget += delay + } + for attempt := 0; ; attempt++ { + if err := ctx.Err(); err != nil { + return fmt.Errorf("release extraction pod %s canceled before attempt %d: %w", name, attempt+1, err) + } + err := run(ctx) + if err == nil { + if attempt > 0 { + logrus.WithFields(logrus.Fields{ + "attempts": attempt + 1, + "elapsed": time.Since(start), + }).Info("Release extraction recovered after retry") + } + return nil + } + if ctxErr := ctx.Err(); ctxErr != nil { + return fmt.Errorf("release extraction pod %s canceled after attempt %d: %w", name, attempt+1, ctxErr) + } + var transientErr *transientReleaseExtractionError + if !errors.As(err, &transientErr) { + return fmt.Errorf("release extraction pod %s failed on attempt %d: %w", name, attempt+1, err) + } + if attempt == len(retryDelays) { + logrus.WithFields(logrus.Fields{ + "attempts": attempt + 1, + "elapsed": time.Since(start), + "error_class": "transient_extraction", + "retry_budget": retryBudget, + }).Error("Release extraction retry attempts exhausted") + return fmt.Errorf("release extraction pod %s exhausted %d attempts with retry budget %s after %s: %w", name, attempt+1, retryBudget, time.Since(start).Truncate(time.Millisecond), err) + } + delay := retryDelays[attempt] + logrus.WithFields(logrus.Fields{ + "attempt": attempt + 1, + "delay": delay, + "error_class": "transient_extraction", + }).Warn("Release extraction failed, retrying") + if err := sleep(ctx, delay); err != nil { + return fmt.Errorf("release extraction pod %s canceled while waiting to retry attempt %d: %w", name, attempt+2, err) + } + } +} + +func transientReleaseExtractionPodError(pod *coreapi.Pod, runErr error) error { + if pod == nil { + return runErr + } + statuses := append([]coreapi.ContainerStatus{}, pod.Status.InitContainerStatuses...) + statuses = append(statuses, pod.Status.ContainerStatuses...) + for _, status := range statuses { + if status.Name == releaseExtractionContainerName && status.State.Terminated != nil && status.State.Terminated.ExitCode == transientReleaseExtractionExitCode { + return &transientReleaseExtractionError{err: runErr} + } + } + return runErr +} + +func runReleaseExtractionWithRetries(ctx context.Context, name string, step api.Step, client ctrlruntimeclient.Client, retryDelays []time.Duration, sleep releaseImportSleep) error { + return retryReleaseExtraction(ctx, name, retryDelays, sleep, func(ctx context.Context) error { + err := step.Run(ctx) + if err == nil { + return nil + } + var podStepErr *steps.PodStepError + if !errors.As(err, &podStepErr) { + return err + } + classifiedErr := transientReleaseExtractionPodError(podStepErr.Pod, err) + var transientErr *transientReleaseExtractionError + if !errors.As(classifiedErr, &transientErr) { + return classifiedErr + } + if err := util.DeletePodWithUID(ctx, client, podStepErr.Pod); err != nil { + // Preserve the extraction cause without its transient marker: a new + // attempt is unsafe until cleanup confirms that this pod UID is gone. + return fmt.Errorf("failed to confirm transient release extraction pod cleanup, cannot retry safely: %w", errors.Join(err, transientErr.err)) + } + return classifiedErr + }) } func (s *importReleaseStep) Inputs() (api.InputDefinition, error) { @@ -113,9 +233,10 @@ func (s *importReleaseStep) run(ctx context.Context) error { } logrus.WithField("name", s.name).Debugf("setting originalPullSpec to: %s for multi-stage steps to reference", pullSpec) s.originalPullSpec = pullSpec - // retry importing the image a few times because we might race against establishing credentials/roles - // and be unable to import images on the same cluster - if newPullSpec, err := utils.ImportTagWithRetries(ctx, s.client, s.jobSpec.Namespace(), "release", s.name, pullSpec, api.ImageStreamImportRetries, s.client.MetricsAgent()); err != nil { + // Retry importing for several minutes because we might race against establishing + // credentials/roles or a transient registry outage. + retryDelays := s.releaseImportRetryDelays() + if newPullSpec, err := s.importTagWithRetryDelays(ctx, s.client, s.jobSpec.Namespace(), "release", s.name, pullSpec, retryDelays, s.client.MetricsAgent()); err != nil { return fmt.Errorf("unable to import %s release image: %w", s.name, err) } else { logrus.WithField("pullSpec", pullSpec).WithField("newPullSpec", newPullSpec).WithField("name", s.name). @@ -149,7 +270,21 @@ if [[ -d /pull ]]; then cp /pull/.dockerconfigjson $HOME/.docker/config.json fi oc registry login --to $HOME/.docker/config.json -oc adm release extract --from=%q --file=image-references > ${ARTIFACT_DIR}/%s +extract_error=${ARTIFACT_DIR}/%s-extract-error.log +set +e +oc adm release extract --from=%q --file=image-references > ${ARTIFACT_DIR}/%s 2> "${extract_error}" +extract_status=$? +set -e +cat "${extract_error}" >&2 +if (( extract_status != 0 )); then + # Exit 75 is reserved for failures that are safe to retry in a fresh pod. + # Authorization, invalid references, and malformed payloads retain the + # original exit status and fail immediately. + if grep -Eqi '(%s)' "${extract_error}"; then + exit %d + fi + exit "${extract_status}" +fi # while release creation may happen more than once in the lifetime of a test # namespace, only one release creation Pod will ever run at once. Therefore, # while actions editing the output ConfigMap may race if done from ci-operator @@ -160,7 +295,7 @@ if oc get configmap release-%s; then oc delete configmap release-%s fi oc create configmap release-%s --from-file=%s.yaml=${ARTIFACT_DIR}/%s -`, pullSpec, target, target, target, target, target, target) +`, target, pullSpec, target, transientReleaseExtractionErrorPattern, transientReleaseExtractionExitCode, target, target, target, target, target) // run adm release extract and grab the raw image-references from the payload podConfig := steps.PodStepConfiguration{ @@ -185,9 +320,9 @@ oc create configmap release-%s --from-file=%s.yaml=${ARTIFACT_DIR}/%s copied[podConfig.As] = api.ResourceRequirements{Requests: api.ResourceList{"cpu": "50m", "memory": "400Mi"}} resources = copied } - step := steps.PodStep("release", podConfig, resources, s.client, s.jobSpec, nil) - if err := step.Run(ctx); err != nil { - return err + step := steps.PodStep(releaseExtractionContainerName, podConfig, resources, s.client, s.jobSpec, nil) + if err := runReleaseExtractionWithRetries(ctx, target, step, s.client, retryDelays, sleepForReleaseImportRetry); err != nil { + return fmt.Errorf("failed to extract release image %s: %w", s.name, err) } // read the contents from the configmap we created @@ -360,6 +495,8 @@ func ImportReleaseStep( jobSpec: jobSpec, pullSecret: pullSecret, overrideCLIReleaseExtractImage: overrideCLIReleaseExtractImage, + releaseImportRetryDelays: utils.ReleaseImportRetryDelays, + importTagWithRetryDelays: utils.ImportTagWithRetryDelays, } } diff --git a/pkg/steps/release/import_release_test.go b/pkg/steps/release/import_release_test.go index a188e70a51..a88951034a 100644 --- a/pkg/steps/release/import_release_test.go +++ b/pkg/steps/release/import_release_test.go @@ -1,26 +1,718 @@ package release import ( + "bytes" "context" + "errors" + "fmt" + "os/exec" + "strings" + "sync" + "syscall" "testing" "time" + "github.com/sirupsen/logrus" + corev1 "k8s.io/api/core/v1" + rbacv1 "k8s.io/api/rbac/v1" + kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/watch" + utilpointer "k8s.io/utils/pointer" ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" fakectrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client/fake" + prowapi "sigs.k8s.io/prow/pkg/apis/prowjobs/v1" + "sigs.k8s.io/prow/pkg/entrypoint" + "sigs.k8s.io/prow/pkg/pod-utils/downwardapi" imagev1 "github.com/openshift/api/image/v1" "github.com/openshift/ci-tools/pkg/api" + "github.com/openshift/ci-tools/pkg/metrics" + "github.com/openshift/ci-tools/pkg/steps" "github.com/openshift/ci-tools/pkg/steps/loggingclient" + steputils "github.com/openshift/ci-tools/pkg/steps/utils" testhelperkube "github.com/openshift/ci-tools/pkg/testhelper/kubernetes" + ciutil "github.com/openshift/ci-tools/pkg/util" ) const testCLIImage = "quay.io/test/cli:latest" +func deterministicReleaseImportRetryDelays() []time.Duration { + return []time.Duration{ + time.Second, + 2 * time.Second, + 4 * time.Second, + 8 * time.Second, + 16 * time.Second, + 32 * time.Second, + 64 * time.Second, + 128 * time.Second, + } +} + +func TestRetryReleaseExtractionRecoversAfterVirtualTwoMinuteOutage(t *testing.T) { + retryDelays := deterministicReleaseImportRetryDelays() + var elapsed time.Duration + var attempts int + var logs bytes.Buffer + logger := logrus.StandardLogger() + originalOutput := logger.Out + logger.SetOutput(&logs) + defer logger.SetOutput(originalOutput) + + err := retryReleaseExtraction(context.Background(), "release-images-latest", retryDelays, func(_ context.Context, delay time.Duration) error { + elapsed += delay + return nil + }, func(context.Context) error { + attempts++ + if elapsed < 2*time.Minute { + return &transientReleaseExtractionError{err: errors.New("registry unavailable")} + } + return nil + }) + if err != nil { + t.Fatalf("expected extraction to recover: %v", err) + } + if attempts != 8 { + t.Fatalf("expected 8 extraction attempts, got %d", attempts) + } + if elapsed != 127*time.Second { + t.Fatalf("expected virtual recovery after 127s, got %s", elapsed) + } + if !strings.Contains(logs.String(), "Release extraction failed, retrying") || strings.Contains(logs.String(), "release-images-latest") { + t.Fatalf("expected retry log evidence, got %q", logs.String()) + } + if !strings.Contains(logs.String(), "Release extraction recovered after retry") || !strings.Contains(logs.String(), "attempts=8") { + t.Fatalf("expected recovery log with attempt count, got %q", logs.String()) + } +} + +func TestRetryReleaseExtractionIsBoundedAndContextAware(t *testing.T) { + t.Run("bounded", func(t *testing.T) { + var elapsed time.Duration + var attempts int + var logs bytes.Buffer + logger := logrus.StandardLogger() + originalOutput := logger.Out + logger.SetOutput(&logs) + defer logger.SetOutput(originalOutput) + expectedErr := errors.New("extract failed") + err := retryReleaseExtraction(context.Background(), "release-images-latest", deterministicReleaseImportRetryDelays(), func(_ context.Context, delay time.Duration) error { + elapsed += delay + return nil + }, func(context.Context) error { + attempts++ + return &transientReleaseExtractionError{err: expectedErr} + }) + if !errors.Is(err, expectedErr) { + t.Fatalf("expected final extraction error, got %v", err) + } + if attempts != 9 { + t.Fatalf("expected 9 bounded attempts, got %d", attempts) + } + if elapsed != 255*time.Second { + t.Fatalf("expected a 255s retry budget, got %s", elapsed) + } + if !strings.Contains(logs.String(), "Release extraction retry attempts exhausted") || + !strings.Contains(logs.String(), "attempts=9") || + !strings.Contains(logs.String(), "retry_budget=4m15s") || + !strings.Contains(err.Error(), "retry budget 4m15s") { + t.Fatalf("expected terminal exhaustion log with attempt count, got %q", logs.String()) + } + }) + + t.Run("permanent failure", func(t *testing.T) { + var attempts int + expectedErr := errors.New("malformed release payload") + err := retryReleaseExtraction(context.Background(), "release-images-latest", deterministicReleaseImportRetryDelays(), func(context.Context, time.Duration) error { + t.Fatal("permanent failure must not sleep") + return nil + }, func(context.Context) error { + attempts++ + return expectedErr + }) + if !errors.Is(err, expectedErr) || attempts != 1 { + t.Fatalf("expected one permanent attempt, got err=%v attempts=%d", err, attempts) + } + }) + + t.Run("pre-canceled", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + var attempts int + err := retryReleaseExtraction(ctx, "release-images-latest", deterministicReleaseImportRetryDelays(), sleepForReleaseImportRetry, func(context.Context) error { + attempts++ + return nil + }) + if !errors.Is(err, context.Canceled) || attempts != 0 || !strings.Contains(err.Error(), "release extraction pod") { + t.Fatalf("expected contextual pre-cancellation, got err=%v attempts=%d", err, attempts) + } + }) + + t.Run("cancellation during backoff", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + var attempts int + err := retryReleaseExtraction(ctx, "release-images-latest", deterministicReleaseImportRetryDelays(), func(ctx context.Context, _ time.Duration) error { + cancel() + return ctx.Err() + }, func(context.Context) error { + attempts++ + return &transientReleaseExtractionError{err: errors.New("extract failed")} + }) + if !errors.Is(err, context.Canceled) { + t.Fatalf("expected cancellation, got %v", err) + } + if attempts != 1 { + t.Fatalf("expected one attempt before cancellation, got %d", attempts) + } + }) + + t.Run("cancellation during final attempt", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + var attempts int + err := retryReleaseExtraction(ctx, "release-images-latest", []time.Duration{0}, func(context.Context, time.Duration) error { return nil }, func(context.Context) error { + attempts++ + if attempts == 2 { + cancel() + } + return &transientReleaseExtractionError{err: errors.New("extract failed")} + }) + if !errors.Is(err, context.Canceled) || attempts != 2 { + t.Fatalf("expected final-attempt cancellation, got err=%v attempts=%d", err, attempts) + } + }) + + t.Run("attempt retains parent deadline", func(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), time.Hour) + defer cancel() + err := retryReleaseExtraction(ctx, "release-images-latest", nil, sleepForReleaseImportRetry, func(attemptCtx context.Context) error { + deadline, ok := attemptCtx.Deadline() + if !ok || time.Until(deadline) < 30*time.Minute { + return fmt.Errorf("attempt context was unexpectedly capped: deadline=%v present=%t", deadline, ok) + } + return nil + }) + if err != nil { + t.Fatalf("expected parent deadline to govern the attempt: %v", err) + } + }) + + t.Run("expired deadline", func(t *testing.T) { + ctx, cancel := context.WithDeadline(context.Background(), time.Now().Add(-time.Second)) + defer cancel() + var attempts int + err := retryReleaseExtraction(ctx, "release-images-latest", deterministicReleaseImportRetryDelays(), sleepForReleaseImportRetry, func(context.Context) error { + attempts++ + return nil + }) + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("expected deadline error, got %v", err) + } + if attempts != 0 { + t.Fatalf("expected no attempt after deadline, got %d", attempts) + } + }) +} + +func TestTransientReleaseExtractionErrorPattern(t *testing.T) { + testCases := []struct { + name string + output string + transient bool + }{ + {name: "too many requests", output: "error: too many requests", transient: true}, + {name: "docker hub rate limit", output: "toomanyrequests: You have reached your pull rate limit", transient: true}, + {name: "server status", output: "registry returned status code 503", transient: true}, + {name: "canonical internal server error", output: "received unexpected HTTP status: 500 Internal Server Error", transient: true}, + {name: "canonical bad gateway", output: "received unexpected HTTP status: 502 Bad Gateway", transient: true}, + {name: "connection reset", output: "read: connection reset by peer", transient: true}, + {name: "HTTP2 connection lost", output: "http2: client connection lost", transient: true}, + {name: "timeout", output: "TLS handshake timeout", transient: true}, + {name: "unauthorized", output: "unauthorized: authentication required"}, + {name: "forbidden", output: "forbidden: access denied"}, + {name: "manifest unknown", output: "manifest unknown: manifest unknown"}, + {name: "invalid reference", output: "invalid reference format"}, + {name: "malformed payload", output: "image-references is malformed"}, + } + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + command := exec.Command("grep", "-Eqi", "("+transientReleaseExtractionErrorPattern+")") + command.Stdin = strings.NewReader(testCase.output) + got := command.Run() == nil + if got != testCase.transient { + t.Fatalf("classification = %t, want %t for %q", got, testCase.transient, testCase.output) + } + }) + } +} + +type releasePodLifecycleClient struct { + *testhelperkube.FakePodClient + lock sync.Mutex + statuses []corev1.PodStatus + createCount int + deletedUIDs []types.UID + staleDeleteCount int + deleteErr error + deleteResponseErr error + deleteCalls chan struct{} +} + +func (c *releasePodLifecycleClient) Create(ctx context.Context, obj ctrlruntimeclient.Object, opts ...ctrlruntimeclient.CreateOption) error { + pod, ok := obj.(*corev1.Pod) + if !ok { + return c.FakePodExecutor.LoggingClient.Create(ctx, obj, opts...) + } + c.lock.Lock() + c.createCount++ + attempt := c.createCount + pod.UID = types.UID(fmt.Sprintf("release-pod-%d", attempt)) + pod.CreationTimestamp = metav1.NewTime(time.Now().Add(-time.Minute)) + pod.Status = c.statuses[attempt-1] + c.CreatedPods = append(c.CreatedPods, pod.DeepCopy()) + c.lock.Unlock() + return c.FakePodExecutor.LoggingClient.Create(ctx, pod, opts...) +} + +func (c *releasePodLifecycleClient) Watch(ctx context.Context, list ctrlruntimeclient.ObjectList, _ ...ctrlruntimeclient.ListOption) (watch.Interface, error) { + if err := c.FakePodExecutor.LoggingClient.List(ctx, list); err != nil { + return nil, err + } + items := list.(*corev1.PodList).Items + ch := make(chan watch.Event, len(items)) + for i := range items { + ch <- watch.Event{Type: watch.Modified, Object: items[i].DeepCopy()} + } + return watch.NewProxyWatcher(ch), nil +} + +func (c *releasePodLifecycleClient) Delete(ctx context.Context, obj ctrlruntimeclient.Object, opts ...ctrlruntimeclient.DeleteOption) error { + c.deleteCalls <- struct{}{} + pod, ok := obj.(*corev1.Pod) + if !ok { + return c.FakePodExecutor.LoggingClient.Delete(ctx, obj, opts...) + } + c.lock.Lock() + deleteErr := c.deleteErr + deleteResponseErr := c.deleteResponseErr + c.lock.Unlock() + if deleteErr != nil { + return deleteErr + } + deleteOptions := &ctrlruntimeclient.DeleteOptions{} + deleteOptions.ApplyOptions(opts) + var expectedUID types.UID + if deleteOptions.Preconditions != nil && deleteOptions.Preconditions.UID != nil { + expectedUID = *deleteOptions.Preconditions.UID + } + current := &corev1.Pod{} + if err := c.FakePodExecutor.LoggingClient.Get(ctx, ctrlruntimeclient.ObjectKeyFromObject(pod), current); err != nil { + return err + } + if expectedUID == "" || current.UID != expectedUID { + c.lock.Lock() + c.staleDeleteCount++ + c.lock.Unlock() + return kerrors.NewConflict(corev1.Resource("pods"), pod.Name, errors.New("UID precondition mismatch")) + } + c.lock.Lock() + c.deletedUIDs = append(c.deletedUIDs, expectedUID) + c.DeletedPods = append(c.DeletedPods, current.DeepCopy()) + c.lock.Unlock() + if err := c.FakePodExecutor.LoggingClient.Delete(ctx, current, opts...); err != nil { + return err + } + return deleteResponseErr +} + +func newReleasePodLifecycleClient(t *testing.T, statuses ...corev1.PodStatus) *releasePodLifecycleClient { + t.Helper() + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatalf("add core API to scheme: %v", err) + } + loggingClient := loggingclient.New(fakectrlruntimeclient.NewClientBuilder(). + WithScheme(scheme). + WithIndex(&corev1.Pod{}, "metadata.name", func(obj ctrlruntimeclient.Object) []string { return []string{obj.GetName()} }). + WithIndex(&corev1.Event{}, "involvedObject.uid", func(obj ctrlruntimeclient.Object) []string { + return []string{string(obj.(*corev1.Event).InvolvedObject.UID)} + }). + Build(), nil) + return &releasePodLifecycleClient{ + FakePodClient: &testhelperkube.FakePodClient{ + FakePodExecutor: &testhelperkube.FakePodExecutor{LoggingClient: loggingClient}, + PendingTimeout: 0, + }, + statuses: statuses, + deleteCalls: make(chan struct{}, 20), + } +} + +func waitForReleasePodDeleteCalls(t *testing.T, client *releasePodLifecycleClient, count int) { + t.Helper() + for i := 0; i < count; i++ { + select { + case <-client.deleteCalls: + case <-time.After(time.Second): + t.Fatalf("timed out waiting for delete call %d of %d", i+1, count) + } + } +} + +func releaseExtractionPodStep(client *releasePodLifecycleClient, namespace, podName string) api.Step { + jobSpec := &api.JobSpec{JobSpec: downwardapi.JobSpec{ + Job: "release-import-test", + Type: prowapi.PresubmitJob, + Refs: &prowapi.Refs{Org: "openshift", Repo: "ci-tools", BaseRef: "main", Pulls: []prowapi.Pull{{Number: 5376, SHA: "test-sha"}}}, + DecorationConfig: &prowapi.DecorationConfig{ + Timeout: &prowapi.Duration{Duration: time.Hour}, + GracePeriod: &prowapi.Duration{Duration: time.Second}, + UtilityImages: &prowapi.UtilityImages{ + Sidecar: "sidecar", + Entrypoint: "entrypoint", + }, + SkipCloning: utilpointer.Bool(true), + }, + }} + jobSpec.SetNamespace(namespace) + return steps.PodStep(releaseExtractionContainerName, steps.PodStepConfiguration{ + WaitFlags: ciutil.SkipLogs, + As: podName, + From: api.ImageStreamTagReference{Name: api.PipelineImageStream, Tag: "cli"}, + Commands: "oc adm release extract", + }, nil, client, jobSpec, nil) +} + +func transientRunningReleasePodStatus() corev1.PodStatus { + return corev1.PodStatus{ + Phase: corev1.PodRunning, + ContainerStatuses: []corev1.ContainerStatus{ + {Name: releaseExtractionContainerName, State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{ExitCode: transientReleaseExtractionExitCode}}}, + {Name: "sidecar", State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}}, + }, + } +} + +type importReleaseOrchestrationClient struct { + *releasePodLifecycleClient + imageImportAttempts int + extractionCommands string +} + +func (c *importReleaseOrchestrationClient) Create(ctx context.Context, obj ctrlruntimeclient.Object, opts ...ctrlruntimeclient.CreateOption) error { + switch typed := obj.(type) { + case *imagev1.ImageStreamImport: + c.imageImportAttempts++ + if c.imageImportAttempts == 1 { + return kerrors.NewServiceUnavailable("registry temporarily unavailable") + } + typed.Status.Images = []imagev1.ImageImportStatus{{Image: &imagev1.Image{DockerImageReference: "registry.test/release@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}}} + return nil + case *corev1.Pod: + if typed.Name == "release-images-latest" && len(typed.Spec.Containers) > 0 { + for _, env := range typed.Spec.Containers[0].Env { + if env.Name == entrypoint.JSONConfigEnvVar { + var options entrypoint.Options + if options.LoadConfig(env.Value) == nil { + c.extractionCommands = strings.Join(options.Args, "\n") + } + break + } + } + } + } + return c.releasePodLifecycleClient.Create(ctx, obj, opts...) +} + +func (c *importReleaseOrchestrationClient) Watch(ctx context.Context, list ctrlruntimeclient.ObjectList, opts ...ctrlruntimeclient.ListOption) (watch.Interface, error) { + if _, ok := list.(*corev1.PodList); ok { + return c.releasePodLifecycleClient.Watch(ctx, list, opts...) + } + return c.FakePodExecutor.LoggingClient.Watch(ctx, list, opts...) +} + +func newImportReleaseOrchestrationClient(t *testing.T, namespace string) *importReleaseOrchestrationClient { + t.Helper() + scheme := runtime.NewScheme() + for name, addToScheme := range map[string]func(*runtime.Scheme) error{ + "core": corev1.AddToScheme, + "image": imagev1.AddToScheme, + "rbac": rbacv1.AddToScheme, + } { + if err := addToScheme(scheme); err != nil { + t.Fatalf("add %s API to scheme: %v", name, err) + } + } + stableStream := &imagev1.ImageStream{ + ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: api.ReleaseStreamFor(api.LatestReleaseName)}, + Spec: imagev1.ImageStreamSpec{Tags: []imagev1.TagReference{{ + Name: "cli", + From: &corev1.ObjectReference{Kind: "DockerImage", Name: testCLIImage}, + }}}, + Status: imagev1.ImageStreamStatus{Tags: []imagev1.NamedTagEventList{{ + Tag: "cli", + Items: []imagev1.TagEvent{{DockerImageReference: "quay.io/test/cli@sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"}}, + }}}, + } + stableCLI := &imagev1.ImageStreamTag{ + ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: stableStream.Name + ":cli"}, + Tag: &imagev1.TagReference{From: &corev1.ObjectReference{Kind: "DockerImage", Name: testCLIImage}}, + } + serviceAccount := &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: "ci-operator"}, + ImagePullSecrets: []corev1.LocalObjectReference{{Name: api.RegistryPullCredentialsSecret}, {Name: "generated-dockercfg"}}, + } + loggingClient := loggingclient.New(fakectrlruntimeclient.NewClientBuilder(). + WithScheme(scheme). + WithObjects(stableStream, stableCLI, serviceAccount). + WithIndex(&corev1.Pod{}, "metadata.name", func(obj ctrlruntimeclient.Object) []string { return []string{obj.GetName()} }). + WithIndex(&corev1.Event{}, "involvedObject.uid", func(obj ctrlruntimeclient.Object) []string { + return []string{string(obj.(*corev1.Event).InvolvedObject.UID)} + }). + Build(), nil) + base := &releasePodLifecycleClient{ + FakePodClient: &testhelperkube.FakePodClient{ + FakePodExecutor: &testhelperkube.FakePodExecutor{LoggingClient: loggingClient}, + PendingTimeout: 0, + }, + statuses: []corev1.PodStatus{ + transientRunningReleasePodStatus(), + {Phase: corev1.PodFailed, ContainerStatuses: []corev1.ContainerStatus{{ + Name: releaseExtractionContainerName, + State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{ExitCode: 1}}, + }}}, + }, + deleteCalls: make(chan struct{}, 20), + } + return &importReleaseOrchestrationClient{releasePodLifecycleClient: base} +} + +func TestImportReleaseRunOrchestratesRetriesAndExtractionLifecycle(t *testing.T) { + const namespace = "test-namespace" + client := newImportReleaseOrchestrationClient(t, namespace) + jobSpec := &api.JobSpec{JobSpec: downwardapi.JobSpec{ + Job: "release-import-test", + Type: prowapi.PresubmitJob, + Refs: &prowapi.Refs{Org: "openshift", Repo: "ci-tools", BaseRef: "main", Pulls: []prowapi.Pull{{Number: 5376, SHA: "test-sha"}}}, + DecorationConfig: &prowapi.DecorationConfig{Timeout: &prowapi.Duration{Duration: time.Hour}, GracePeriod: &prowapi.Duration{Duration: time.Second}, UtilityImages: &prowapi.UtilityImages{Sidecar: "sidecar", Entrypoint: "entrypoint"}, SkipCloning: utilpointer.Bool(true)}, + }} + jobSpec.SetNamespace(namespace) + step := ImportReleaseStep( + api.LatestReleaseName, + "", + "release:latest", + imagev1.SourceTagReferencePolicy, + NewReleaseSourceFromPullSpec("registry.test/source:latest"), + false, + nil, + client, + jobSpec, + nil, + nil, + ).(*importReleaseStep) + expectedDelays := []time.Duration{0} + step.releaseImportRetryDelays = func() []time.Duration { return append([]time.Duration(nil), expectedDelays...) } + importCalled := false + step.importTagWithRetryDelays = func(ctx context.Context, gotClient ctrlruntimeclient.Client, ns, name, tag, pullSpec string, delays []time.Duration, metricsAgent *metrics.MetricsAgent) (string, error) { + importCalled = true + if gotClient != client || ns != namespace || name != "release" || tag != api.LatestReleaseName || pullSpec != "registry.test/source:latest" { + t.Fatalf("unexpected extended import arguments: client_match=%t ns=%q name=%q tag=%q pull_spec=%q", gotClient == client, ns, name, tag, pullSpec) + } + if len(delays) != 1 || delays[0] != 0 { + t.Fatalf("unexpected release import retry delays: %v", delays) + } + return steputils.ImportTagWithRetryDelays(ctx, gotClient, ns, name, tag, pullSpec, delays, metricsAgent) + } + + ctx, cancel := context.WithCancel(context.Background()) + err := step.Run(ctx) + if err == nil || !strings.Contains(err.Error(), "failed to extract release image") { + t.Fatalf("expected permanent extraction failure from real Run wiring, got %v", err) + } + if !importCalled || client.imageImportAttempts != 2 { + t.Fatalf("expected extended import API to retry once, called=%t attempts=%d", importCalled, client.imageImportAttempts) + } + if client.createCount != 2 || len(client.deletedUIDs) != 1 || client.deletedUIDs[0] != "release-pod-1" { + t.Fatalf("expected exit-75 pod recreation followed by one permanent failure, creates=%d deleted_uids=%v", client.createCount, client.deletedUIDs) + } + for _, expected := range []string{ + "extract_error=${ARTIFACT_DIR}/release-images-latest-extract-error.log", + "grep -Eqi '(" + transientReleaseExtractionErrorPattern + ")'", + "exit 75", + "oc adm release extract --from=\"registry.test/release@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\"", + } { + if !strings.Contains(client.extractionCommands, expected) { + t.Fatalf("generated extraction command missing %q:\n%s", expected, client.extractionCommands) + } + } + var transientErr *transientReleaseExtractionError + if errors.As(err, &transientErr) { + t.Fatalf("permanent second attempt retained transient classification: %v", err) + } + cancel() + waitForReleasePodDeleteCalls(t, client.releasePodLifecycleClient, 3) +} + +func TestReleaseExtractionUsesRealPodStepLifecycle(t *testing.T) { + const ( + namespace = "test-namespace" + podName = "release-images-latest" + ) + ctx, cancel := context.WithCancel(context.Background()) + client := newReleasePodLifecycleClient(t, transientRunningReleasePodStatus(), corev1.PodStatus{Phase: corev1.PodSucceeded}) + err := runReleaseExtractionWithRetries(ctx, podName, releaseExtractionPodStep(client, namespace, podName), client, []time.Duration{0}, func(context.Context, time.Duration) error { return nil }) + if err != nil { + t.Fatalf("expected recreated extraction pod to succeed: %v", err) + } + if client.createCount != 2 { + t.Fatalf("expected two PodStep pod creations, got %d", client.createCount) + } + if len(client.deletedUIDs) != 1 || client.deletedUIDs[0] != "release-pod-1" { + t.Fatalf("expected UID-safe deletion of the transient pod, got %v", client.deletedUIDs) + } + replacement := &corev1.Pod{} + if err := client.Get(ctx, ctrlruntimeclient.ObjectKey{Namespace: namespace, Name: podName}, replacement); err != nil || replacement.UID != "release-pod-2" { + t.Fatalf("expected recreated pod with new UID, got pod=%#v err=%v", replacement, err) + } + if err := ciutil.DeletePodWithUID(ctx, client, client.CreatedPods[0]); err != nil { + t.Fatalf("stale cleanup should be harmless: %v", err) + } + if err := client.Get(ctx, ctrlruntimeclient.ObjectKey{Namespace: namespace, Name: podName}, replacement); err != nil || replacement.UID != "release-pod-2" || client.staleDeleteCount != 1 { + t.Fatalf("stale cleanup affected replacement: pod=%#v err=%v stale_deletes=%d", replacement, err, client.staleDeleteCount) + } + cancel() + waitForReleasePodDeleteCalls(t, client, 4) +} + +func TestReleaseExtractionRecreatesAfterCommittedDeleteLostResponse(t *testing.T) { + const ( + namespace = "test-namespace" + podName = "release-images-latest" + ) + ctx, cancel := context.WithCancel(context.Background()) + client := newReleasePodLifecycleClient(t, transientRunningReleasePodStatus(), corev1.PodStatus{Phase: corev1.PodSucceeded}) + client.deleteResponseErr = syscall.ECONNRESET + + err := runReleaseExtractionWithRetries(ctx, podName, releaseExtractionPodStep(client, namespace, podName), client, []time.Duration{0}, func(context.Context, time.Duration) error { return nil }) + if err != nil { + t.Fatalf("expected lost deletion response to be reconciled before retry: %v", err) + } + if client.createCount != 2 { + t.Fatalf("expected extraction pod recreation, got %d creations", client.createCount) + } + if len(client.deletedUIDs) != 1 || client.deletedUIDs[0] != "release-pod-1" { + t.Fatalf("expected confirmed deletion of the first pod UID, got %v", client.deletedUIDs) + } + replacement := &corev1.Pod{} + if err := client.Get(ctx, ctrlruntimeclient.ObjectKey{Namespace: namespace, Name: podName}, replacement); err != nil || replacement.UID != "release-pod-2" { + t.Fatalf("expected successful replacement pod, got pod=%#v err=%v", replacement, err) + } + + cancel() + waitForReleasePodDeleteCalls(t, client, 3) +} + +func TestReleaseExtractionRealPodStepTransientExhaustion(t *testing.T) { + client := newReleasePodLifecycleClient(t, transientRunningReleasePodStatus(), transientRunningReleasePodStatus()) + ctx, cancel := context.WithCancel(context.Background()) + err := runReleaseExtractionWithRetries(ctx, "extract", releaseExtractionPodStep(client, "ns", "extract"), client, []time.Duration{0}, func(context.Context, time.Duration) error { return nil }) + var transientErr *transientReleaseExtractionError + if !errors.As(err, &transientErr) { + t.Fatalf("expected typed transient exhaustion cause, got %v", err) + } + if client.createCount != 2 || len(client.deletedUIDs) != 2 { + t.Fatalf("expected two real attempts and UID-safe deletions, got creates=%d deletes=%v", client.createCount, client.deletedUIDs) + } + cancel() + waitForReleasePodDeleteCalls(t, client, 4) +} + +func TestReleaseExtractionCleanupFailureStopsRetry(t *testing.T) { + client := newReleasePodLifecycleClient(t, transientRunningReleasePodStatus()) + cleanupErr := errors.New("pod deletion failed") + client.deleteErr = cleanupErr + ctx, cancel := context.WithCancel(context.Background()) + err := runReleaseExtractionWithRetries(ctx, "extract", releaseExtractionPodStep(client, "ns", "extract"), client, []time.Duration{0}, func(context.Context, time.Duration) error { + t.Fatal("cleanup failure must stop before retry backoff") + return nil + }) + if !errors.Is(err, cleanupErr) || !strings.Contains(err.Error(), "cannot retry") { + t.Fatalf("expected permanent cleanup failure, got %v", err) + } + var transientErr *transientReleaseExtractionError + if errors.As(err, &transientErr) { + t.Fatalf("cleanup failure retained transient marker: %v", err) + } + if client.createCount != 1 { + t.Fatalf("cleanup failure retried the PodStep, got %d creations", client.createCount) + } + cancel() + waitForReleasePodDeleteCalls(t, client, 2) +} + +func TestReleaseExtractionPendingImagePullIsNotRetried(t *testing.T) { + pending := corev1.PodStatus{Phase: corev1.PodPending, ContainerStatuses: []corev1.ContainerStatus{{ + Name: releaseExtractionContainerName, + State: corev1.ContainerState{Waiting: &corev1.ContainerStateWaiting{ + Reason: "ImagePullBackOff", + Message: "authentication required", + }}, + }}} + client := newReleasePodLifecycleClient(t, pending) + ctx, cancel := context.WithCancel(context.Background()) + err := runReleaseExtractionWithRetries(ctx, "extract", releaseExtractionPodStep(client, "ns", "extract"), client, []time.Duration{0}, func(context.Context, time.Duration) error { + t.Fatal("pending image-pull failure must not retry") + return nil + }) + var transientErr *transientReleaseExtractionError + if err == nil || errors.As(err, &transientErr) || client.createCount != 1 { + t.Fatalf("expected one permanent pending failure, got err=%v creates=%d", err, client.createCount) + } + cancel() + waitForReleasePodDeleteCalls(t, client, 1) +} + +func TestTransientReleaseExtractionPodErrorClassification(t *testing.T) { + testCases := []struct { + name string + phase corev1.PodPhase + exitCode int32 + transient bool + }{ + {name: "explicit transient exit", phase: corev1.PodFailed, exitCode: transientReleaseExtractionExitCode, transient: true}, + {name: "decorated running pod with transient release exit", phase: corev1.PodRunning, exitCode: transientReleaseExtractionExitCode, transient: true}, + {name: "permanent container failure", phase: corev1.PodFailed, exitCode: 1}, + {name: "successful pod", phase: corev1.PodSucceeded}, + } + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + pod := &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{Namespace: "ns", Name: "extract"}, + Status: corev1.PodStatus{ + Phase: testCase.phase, + ContainerStatuses: []corev1.ContainerStatus{{ + Name: releaseExtractionContainerName, + State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{ExitCode: testCase.exitCode}}, + }}, + }, + } + original := errors.New("pod failed") + err := transientReleaseExtractionPodError(pod, original) + var transientErr *transientReleaseExtractionError + if got := errors.As(err, &transientErr); got != testCase.transient { + t.Fatalf("transient classification = %t, want %t (err=%v)", got, testCase.transient, err) + } + if !errors.Is(err, original) { + t.Fatalf("classification must preserve original error, got %v", err) + } + }) + } +} + func TestResolveCLIImageFromStreamWaitsForSpecVisibility(t *testing.T) { const ( namespace = "test-namespace" diff --git a/pkg/steps/utils/image.go b/pkg/steps/utils/image.go index de97617705..edb389ea77 100644 --- a/pkg/steps/utils/image.go +++ b/pkg/steps/utils/image.go @@ -2,6 +2,7 @@ package utils import ( "context" + "errors" "fmt" "sort" "strings" @@ -13,6 +14,7 @@ import ( kerrors "k8s.io/apimachinery/pkg/api/errors" meta "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + utilnet "k8s.io/apimachinery/pkg/util/net" "k8s.io/apimachinery/pkg/util/sets" "k8s.io/apimachinery/pkg/util/wait" ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" @@ -178,9 +180,18 @@ func FindStatusTag(is *imagev1.ImageStream, tag string) (*coreapi.ObjectReferenc return nil, "" } -const DefaultImageImportTimeout = 45 * time.Minute +const ( + DefaultImageImportTimeout = 45 * time.Minute + maxImageImportRetryDelay = 5 * time.Minute +) + +type imageTagImporter func(context.Context, ctrlruntimeclient.Client, string, string, string, string, int, *metrics.MetricsAgent) (string, error) func getEvaluator(ctx context.Context, client ctrlruntimeclient.Client, ns, name string, tags sets.Set[string], waitForSpecTags bool, metricsAgent *metrics.MetricsAgent) func(obj runtime.Object) (bool, error) { + return getEvaluatorWithImporter(ctx, client, ns, name, tags, waitForSpecTags, metricsAgent, ImportTagWithRetries) +} + +func getEvaluatorWithImporter(ctx context.Context, client ctrlruntimeclient.Client, ns, name string, tags sets.Set[string], waitForSpecTags bool, metricsAgent *metrics.MetricsAgent, importer imageTagImporter) func(obj runtime.Object) (bool, error) { return func(obj runtime.Object) (bool, error) { switch stream := obj.(type) { case *imagev1.ImageStream: @@ -209,7 +220,11 @@ func getEvaluator(ctx context.Context, client ctrlruntimeclient.Client, ns, name // should never happen return false, fmt.Errorf("failed to import tag %s/%s:%s from an empty source", stream.Namespace, stream.Name, tag.Name) } - if _, err := ImportTagWithRetries(ctx, client, ns, name, tag.Name, tag.From.Name, api.ImageStreamImportRetries, metricsAgent); err != nil { + if _, err := importer(ctx, client, ns, name, tag.Name, tag.From.Name, api.ImageStreamImportRetries, metricsAgent); err != nil { + if isTransientImageImportError(err) { + logrus.WithField("error_class", "transient_import_exhausted").Warnf("Failed to reimport tag %s/%s:%s after a transient registry error, continuing to wait", stream.Namespace, stream.Name, tag.Name) + return false, nil + } return false, fmt.Errorf("failed to reimport the tag %s/%s:%s: %w", stream.Namespace, stream.Name, tag.Name, err) } } @@ -286,17 +301,155 @@ func WaitForImportingISTag(ctx context.Context, client ctrlruntimeclient.WithWat return err } -// ImportTagWithRetries imports image with retries +type importRetrySleep func(context.Context, time.Duration) error + +func sleepForImageImportRetry(ctx context.Context, duration time.Duration) error { + timer := time.NewTimer(duration) + defer timer.Stop() + select { + case <-ctx.Done(): + return ctx.Err() + case <-timer.C: + return nil + } +} + +func exponentialImageImportRetryDelays(attempts int) []time.Duration { + if attempts < 2 { + return nil + } + delays := make([]time.Duration, 0, attempts-1) + delay := time.Second + for len(delays) < attempts-1 { + delays = append(delays, delay) + delay *= 2 + } + return delays +} + +// ReleaseImportRetryDelays returns the extended, jittered retry schedule used +// by both phases of release import. Jitter prevents concurrent jobs from +// synchronizing their requests during a registry or API outage. +func ReleaseImportRetryDelays() []time.Duration { + delays := exponentialImageImportRetryDelays(9) + for i := range delays { + delays[i] = wait.Jitter(delays[i], 0.1) + } + return delays +} + +type transientImageImportError struct { + err error +} + +func (e *transientImageImportError) Error() string { return e.err.Error() } +func (e *transientImageImportError) Unwrap() error { return e.err } + +func isTransientImageImportError(err error) bool { + var transientErr *transientImageImportError + return errors.As(err, &transientErr) +} + +func isRetryableImageImportAPIError(err error) bool { + return utilnet.IsConnectionReset(err) || + utilnet.IsConnectionRefused(err) || + utilnet.IsHTTP2ConnectionLost(err) || + utilnet.IsProbableEOF(err) || + utilnet.IsTimeout(err) || + kerrors.IsConflict(err) || + kerrors.IsTooManyRequests(err) || + kerrors.IsServerTimeout(err) || + kerrors.IsTimeout(err) || + kerrors.IsInternalError(err) || + kerrors.IsServiceUnavailable(err) || + kerrors.IsUnexpectedServerError(err) +} + +func isRetryableImageImportCreateError(err error) bool { + return isRetryableImageImportAPIError(err) || kerrors.IsForbidden(err) +} + +func isRetryableImageImportStatusError(err error) bool { + return isRetryableImageImportAPIError(err) || kerrors.IsUnauthorized(err) +} + +func imageImportRetryErrorClass(err error) string { + switch { + case err == nil: + return "status_not_ready" + case utilnet.IsConnectionReset(err): + return "connection_reset" + case utilnet.IsConnectionRefused(err): + return "connection_refused" + case utilnet.IsHTTP2ConnectionLost(err): + return "http2_connection_lost" + case utilnet.IsProbableEOF(err): + return "connection_closed" + case utilnet.IsTimeout(err): + return "network_timeout" + default: + reason := kerrors.ReasonForError(err) + if reason == meta.StatusReasonUnknown { + return "api_error" + } + return string(reason) + } +} + +func imageImportRetryDelay(err error, configured time.Duration) time.Duration { + if seconds, suggested := kerrors.SuggestsClientDelay(err); suggested { + serverDelay := time.Duration(seconds) * time.Second + if serverDelay > maxImageImportRetryDelay { + serverDelay = maxImageImportRetryDelay + } + if serverDelay > configured { + return serverDelay + } + } + return configured +} + +// ImportTagWithRetries imports an image with bounded retries for transient API +// and registry failures. Create-level Forbidden and status-level Unauthorized +// are retried because they can reflect propagation delays; status-level +// NotFound remains permanent because it normally identifies a missing image. func ImportTagWithRetries(ctx context.Context, client ctrlruntimeclient.Client, ns, name, tag, sourcePullSpec string, retries int, metricsAgent *metrics.MetricsAgent) (string, error) { + if retries < 1 { + return importTagWithRetryDelays(ctx, client, ns, name, tag, sourcePullSpec, nil, sleepForImageImportRetry, false, metricsAgent, 0) + } + return importTagWithRetryDelays(ctx, client, ns, name, tag, sourcePullSpec, exponentialImageImportRetryDelays(retries), sleepForImageImportRetry, false, metricsAgent, retries) +} + +// ImportTagWithRetryDelays imports an image using the provided delays between attempts. +func ImportTagWithRetryDelays(ctx context.Context, client ctrlruntimeclient.Client, ns, name, tag, sourcePullSpec string, retryDelays []time.Duration, metricsAgent *metrics.MetricsAgent) (string, error) { + return importTagWithRetryDelays(ctx, client, ns, name, tag, sourcePullSpec, retryDelays, sleepForImageImportRetry, true, metricsAgent, len(retryDelays)+1) +} + +func importTagWithRetryDelays(ctx context.Context, client ctrlruntimeclient.Client, ns, name, tag, sourcePullSpec string, retryDelays []time.Duration, sleep importRetrySleep, logRetries bool, metricsAgent *metrics.MetricsAgent, attempts int) (string, error) { if sourcePullSpec == "" { return "", fmt.Errorf("sourcePullSpec cannot be empty") } + if attempts != len(retryDelays)+1 && attempts != 0 { + return "", fmt.Errorf("invalid image import retry policy: %d attempts require %d delays, got %d", attempts, attempts-1, len(retryDelays)) + } + for _, delay := range retryDelays { + if delay < 0 { + return "", fmt.Errorf("invalid image import retry policy: delay %s must not be negative", delay) + } + } startTime := time.Now() var pullSpec string step := 0 retryCount := 0 - logger := logrus.WithField("tag", fmt.Sprintf(" %s/%s:%s", ns, name, tag)).WithField("sourcePullSpec", sourcePullSpec) - if err := wait.ExponentialBackoff(wait.Backoff{Steps: retries, Duration: 1 * time.Second, Factor: 2}, func() (bool, error) { + logger := logrus.WithField("tag", fmt.Sprintf(" %s/%s:%s", ns, name, tag)) + var importErr error + if attempts < 1 { + importErr = wait.ErrWaitTimeout + } + for step < attempts { + if err := ctx.Err(); err != nil { + return "", fmt.Errorf("unable to import tag %s/%s:%s before import (%d): %w", ns, name, tag, step, err) + } logger.WithField("step", step).Debug("Retrying importing tag ...") retryCount = step streamImport := &imagev1.ImageStreamImport{ @@ -322,52 +475,74 @@ func ImportTagWithRetries(ctx context.Context, client ctrlruntimeclient.Client, }, } step = step + 1 - if err := client.Create(ctx, streamImport); err != nil { - if kerrors.IsConflict(err) { - logger.WithField("step", step-1).Debug("Unable to create image stream import up to conflicts") - return false, nil + attemptErr := client.Create(ctx, streamImport) + if ctxErr := ctx.Err(); ctxErr != nil { + return "", fmt.Errorf("unable to import tag %s/%s:%s at import (%d): %w", ns, name, tag, step-1, ctxErr) + } + if attemptErr != nil { + if !isRetryableImageImportCreateError(attemptErr) { + return "", fmt.Errorf("unable to import tag %s/%s:%s at import (%d): %w", ns, name, tag, step-1, attemptErr) } - if kerrors.IsForbidden(err) { - logger.WithField("step", step-1).Debug("Unable to create image stream import up to permissions") - return false, nil + logger.WithFields(logrus.Fields{"error_class": imageImportRetryErrorClass(attemptErr), "step": step - 1}).Debug("Transient image stream import API error") + } else if len(streamImport.Status.Images) == 0 { + logger.WithField("step", step-1).Debug("Imports' status has no images") + } else { + image := streamImport.Status.Images[0] + if image.Image != nil { + pullSpec = image.Image.DockerImageReference + logrus.Debugf("Imported tag %s/%s:%s at import (%d)", ns, name, tag, step-1) + importErr = nil + break + } + if image.Status.Reason != "" || image.Status.Status == meta.StatusFailure { + statusErr := &kerrors.StatusError{ErrStatus: image.Status} + if !isRetryableImageImportStatusError(statusErr) { + return "", fmt.Errorf("unable to import tag %s/%s:%s at import (%d): %w", ns, name, tag, step-1, statusErr) + } + attemptErr = statusErr } - return false, err + logger.WithField("step", step-1).Debug("Imports' status' image is nil") } - if len(streamImport.Status.Images) == 0 { - logger.WithField("step", step-1).Debug("Imports' status has no images") - return false, nil + + if ctxErr := ctx.Err(); ctxErr != nil { + return "", fmt.Errorf("unable to import tag %s/%s:%s after import (%d): %w", ns, name, tag, step-1, ctxErr) } - image := streamImport.Status.Images[0] - if image.Image == nil { - logger.WithField("step", step-1).Debug("Imports' status' image is nil") - return false, nil + if step == attempts { + exhaustionErr := errors.Join(wait.ErrWaitTimeout, attemptErr) + importErr = &transientImageImportError{err: exhaustionErr} + logger.WithFields(logrus.Fields{"attempts": step, "error_class": imageImportRetryErrorClass(attemptErr)}).Error("Image stream import retry attempts exhausted") + break } - pullSpec = image.Image.DockerImageReference - logrus.Debugf("Imported tag %s/%s:%s at import (%d)", ns, name, tag, step-1) - return true, nil - }); err != nil { - if err == wait.ErrorInterrupted(err) { - var conditionMsg string - imagestream := imagev1.ImageStream{} - if err := client.Get(ctx, ctrlruntimeclient.ObjectKey{Namespace: ns, Name: name}, &imagestream); err != nil { - logger.WithError(err).Debug("Failed to get image stream for the tag") - } else { - for _, t := range imagestream.Status.Tags { - if t.Tag == tag { - if len(t.Conditions) > 0 { - conditionMsg = t.Conditions[0].Message - } - break + delay := imageImportRetryDelay(attemptErr, retryDelays[step-1]) + if logRetries { + logger.WithFields(logrus.Fields{"attempt": step, "delay": delay, "error_class": imageImportRetryErrorClass(attemptErr)}).Warn("Image stream import did not succeed, retrying") + } + if err := sleep(ctx, delay); err != nil { + return "", fmt.Errorf("unable to import tag %s/%s:%s while waiting to retry import (%d): %w", ns, name, tag, step-1, err) + } + } + if importErr != nil { + var conditionMsg string + imagestream := imagev1.ImageStream{} + if err := client.Get(ctx, ctrlruntimeclient.ObjectKey{Namespace: ns, Name: name}, &imagestream); err != nil { + logger.WithError(err).Debug("Failed to get image stream for the tag") + } else { + for _, t := range imagestream.Status.Tags { + if t.Tag == tag { + if len(t.Conditions) > 0 { + conditionMsg = t.Conditions[0].Message } + break } } - if conditionMsg == "" { - return "", fmt.Errorf("unable to import tag %s/%s:%s even after (%d) imports: %w", ns, name, tag, step, err) - } else { - return "", fmt.Errorf("unable to import tag %s/%s:%s with message %s on the image stream even after (%d) imports: %w", ns, name, tag, conditionMsg, step, err) - } } - return "", fmt.Errorf("unable to import tag %s/%s:%s at import (%d): %w", ns, name, tag, step-1, err) + if ctxErr := ctx.Err(); ctxErr != nil { + return "", fmt.Errorf("unable to import tag %s/%s:%s while collecting terminal status after %d imports: %w", ns, name, tag, step, ctxErr) + } + if conditionMsg == "" { + return "", fmt.Errorf("unable to import tag %s/%s:%s even after (%d) imports: %w", ns, name, tag, step, importErr) + } + return "", fmt.Errorf("unable to import tag %s/%s:%s with message %s on the image stream even after (%d) imports: %w", ns, name, tag, conditionMsg, step, importErr) } completionTime := time.Now() diff --git a/pkg/steps/utils/image_test.go b/pkg/steps/utils/image_test.go index 77463964ce..5ca047aed0 100644 --- a/pkg/steps/utils/image_test.go +++ b/pkg/steps/utils/image_test.go @@ -1,19 +1,26 @@ package utils import ( + "bytes" "context" "errors" "fmt" + "io" "strings" + "syscall" "testing" "time" "github.com/google/go-cmp/cmp" + "github.com/sirupsen/logrus" coreapi "k8s.io/api/core/v1" + kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/util/sets" + "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/kubernetes/scheme" ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" fakectrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client/fake" @@ -21,6 +28,7 @@ import ( imagev1 "github.com/openshift/api/image/v1" "github.com/openshift/ci-tools/pkg/api" + "github.com/openshift/ci-tools/pkg/metrics" "github.com/openshift/ci-tools/pkg/testhelper" ) @@ -263,6 +271,311 @@ func TestReimportTag(t *testing.T) { } } +func TestImportTagWithRetryDelaysRecoversAfterVirtualTwoMinuteOutage(t *testing.T) { + const namespace = "release-import" + var elapsed time.Duration + client := &outageImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + now: func() time.Duration { return elapsed }, + } + delays := exponentialImageImportRetryDelays(9) + var logs bytes.Buffer + logger := logrus.StandardLogger() + originalOutput := logger.Out + logger.SetOutput(&logs) + defer logger.SetOutput(originalOutput) + + pullSpec, err := importTagWithRetryDelays(context.Background(), client, namespace, "release", "latest", "quay.io/openshift/release:latest", delays, func(_ context.Context, delay time.Duration) error { + elapsed += delay + return nil + }, true, nil, len(delays)+1) + if err != nil { + t.Fatalf("expected import to recover: %v", err) + } + if pullSpec != "quay.io/openshift/release@sha256:resolved" { + t.Fatalf("unexpected resolved pull spec %q", pullSpec) + } + if client.attempts != 8 { + t.Fatalf("expected 8 import attempts, got %d", client.attempts) + } + if elapsed != 127*time.Second { + t.Fatalf("expected virtual recovery after 127s, got %s", elapsed) + } + if !strings.Contains(logs.String(), "Image stream import did not succeed, retrying") { + t.Fatalf("expected retry log evidence, got %q", logs.String()) + } +} + +func TestImportTagWithRetryDelaysDoesNotRetryPermanentErrors(t *testing.T) { + permanentErr := errors.New("malformed source") + client := &outageImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + permanentErr: permanentErr, + } + _, err := ImportTagWithRetryDelays(context.Background(), client, "ns", "release", "latest", "not a valid source", exponentialImageImportRetryDelays(9), nil) + if !errors.Is(err, permanentErr) { + t.Fatalf("expected permanent error, got %v", err) + } + if client.attempts != 1 { + t.Fatalf("expected permanent error after one attempt, got %d", client.attempts) + } +} + +func TestImportTagWithRetryDelaysPreservesTransientExhaustionCause(t *testing.T) { + lastCause := kerrors.NewTooManyRequests("registry overloaded", 0) + client := &scriptedImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + create: func(context.Context, int, *imagev1.ImageStreamImport) error { + return lastCause + }, + } + + _, err := ImportTagWithRetryDelays(context.Background(), client, "ns", "release", "latest", "registry/release:latest", []time.Duration{0, 0}, nil) + var transientErr *transientImageImportError + if !errors.As(err, &transientErr) { + t.Fatalf("expected typed transient exhaustion, got %v", err) + } + if !errors.Is(err, wait.ErrWaitTimeout) || !errors.Is(err, lastCause) { + t.Fatalf("expected exhaustion and final API cause to be preserved, got %v", err) + } + if client.attempts != 3 { + t.Fatalf("expected three bounded attempts, got %d", client.attempts) + } +} + +type outageImageImportClient struct { + ctrlruntimeclient.Client + now func() time.Duration + permanentErr error + attempts int +} + +type scriptedImageImportClient struct { + ctrlruntimeclient.Client + create func(context.Context, int, *imagev1.ImageStreamImport) error + attempts int +} + +func (c *scriptedImageImportClient) Create(ctx context.Context, obj ctrlruntimeclient.Object, _ ...ctrlruntimeclient.CreateOption) error { + streamImport, ok := obj.(*imagev1.ImageStreamImport) + if !ok { + return c.Client.Create(ctx, obj) + } + c.attempts++ + return c.create(ctx, c.attempts, streamImport) +} + +func TestImportTagWithRetryDelaysClassifiesTypedAPIErrors(t *testing.T) { + resource := schema.GroupResource{Group: "image.openshift.io", Resource: "imagestreamimports"} + testCases := []struct { + name string + err error + retryable bool + wantErrorClass string + }{ + {name: "conflict", err: kerrors.NewConflict(resource, "release", errors.New("conflict")), retryable: true}, + {name: "too many requests", err: kerrors.NewTooManyRequests("busy", 0), retryable: true}, + {name: "server timeout", err: kerrors.NewServerTimeout(resource, "create", 0), retryable: true}, + {name: "timeout", err: kerrors.NewTimeoutError("timeout", 0), retryable: true}, + {name: "internal error", err: kerrors.NewInternalError(errors.New("internal")), retryable: true}, + {name: "service unavailable", err: kerrors.NewServiceUnavailable("unavailable"), retryable: true}, + {name: "connection reset", err: syscall.ECONNRESET, retryable: true}, + {name: "connection refused", err: fmt.Errorf("dial API server: %w", syscall.ECONNREFUSED), retryable: true, wantErrorClass: "connection_refused"}, + {name: "HTTP2 connection lost", err: errors.New("Post API request: http2: client connection lost"), retryable: true, wantErrorClass: "http2_connection_lost"}, + {name: "unexpected EOF", err: io.ErrUnexpectedEOF, retryable: true}, + {name: "network timeout", err: context.DeadlineExceeded, retryable: true}, + {name: "forbidden", err: kerrors.NewForbidden(resource, "release", errors.New("denied")), retryable: true}, + {name: "unauthorized", err: kerrors.NewUnauthorized("unauthorized")}, + {name: "bad request", err: kerrors.NewBadRequest("malformed source")}, + {name: "generic", err: errors.New("generic client failure")}, + } + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + if testCase.wantErrorClass != "" { + if got := imageImportRetryErrorClass(testCase.err); got != testCase.wantErrorClass { + t.Fatalf("error class = %q, want %q", got, testCase.wantErrorClass) + } + } + client := &scriptedImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + create: func(_ context.Context, attempt int, streamImport *imagev1.ImageStreamImport) error { + if attempt == 1 { + return testCase.err + } + streamImport.Status.Images = []imagev1.ImageImportStatus{{Image: &imagev1.Image{DockerImageReference: "registry/release@sha256:resolved"}}} + return nil + }, + } + pullSpec, err := ImportTagWithRetryDelays(context.Background(), client, "ns", "release", "latest", "registry/release:latest", []time.Duration{0}, nil) + if testCase.retryable { + if err != nil { + t.Fatalf("expected retryable error to recover: %v", err) + } + if pullSpec != "registry/release@sha256:resolved" || client.attempts != 2 { + t.Fatalf("exported retry wiring got pull spec %q after %d attempts", pullSpec, client.attempts) + } + return + } + if !errors.Is(err, testCase.err) { + t.Fatalf("expected permanent error %v, got %v", testCase.err, err) + } + if client.attempts != 1 { + t.Fatalf("expected permanent error after one attempt, got %d", client.attempts) + } + }) + } +} + +func TestImportTagWithRetryDelaysClassifiesImportStatus(t *testing.T) { + testCases := []struct { + name string + reason metav1.StatusReason + retryable bool + }{ + {name: "service unavailable", reason: metav1.StatusReasonServiceUnavailable, retryable: true}, + {name: "too many requests", reason: metav1.StatusReasonTooManyRequests, retryable: true}, + {name: "unauthorized propagation delay", reason: metav1.StatusReasonUnauthorized, retryable: true}, + {name: "forbidden", reason: metav1.StatusReasonForbidden}, + {name: "not found missing image", reason: metav1.StatusReasonNotFound}, + {name: "invalid malformed source", reason: metav1.StatusReasonInvalid}, + } + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + client := &scriptedImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + create: func(_ context.Context, attempt int, streamImport *imagev1.ImageStreamImport) error { + if attempt == 1 { + streamImport.Status.Images = []imagev1.ImageImportStatus{{Status: metav1.Status{Status: metav1.StatusFailure, Reason: testCase.reason, Message: "import failed"}}} + return nil + } + streamImport.Status.Images = []imagev1.ImageImportStatus{{Image: &imagev1.Image{DockerImageReference: "registry/release@sha256:resolved"}}} + return nil + }, + } + _, err := ImportTagWithRetryDelays(context.Background(), client, "ns", "release", "latest", "registry/release:latest", []time.Duration{0}, nil) + if testCase.retryable { + if err != nil || client.attempts != 2 { + t.Fatalf("expected transient status to recover on attempt 2, got err=%v attempts=%d", err, client.attempts) + } + return + } + if err == nil || client.attempts != 1 { + t.Fatalf("expected permanent status to fail once, got err=%v attempts=%d", err, client.attempts) + } + }) + } +} + +func TestImportTagWithRetryDelaysCancellation(t *testing.T) { + t.Run("pre-canceled", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + client := &scriptedImageImportClient{Client: fakectrlruntimeclient.NewClientBuilder().Build(), create: func(context.Context, int, *imagev1.ImageStreamImport) error { return nil }} + _, err := ImportTagWithRetryDelays(ctx, client, "ns", "release", "latest", "registry/release:latest", nil, nil) + if !errors.Is(err, context.Canceled) || client.attempts != 0 { + t.Fatalf("expected pre-cancellation before any attempt, got err=%v attempts=%d", err, client.attempts) + } + }) + + t.Run("during backoff", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + client := &scriptedImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + create: func(context.Context, int, *imagev1.ImageStreamImport) error { + return kerrors.NewTooManyRequests("busy", 0) + }, + } + _, err := importTagWithRetryDelays(ctx, client, "ns", "release", "latest", "registry/release:latest", []time.Duration{time.Second}, func(ctx context.Context, _ time.Duration) error { + cancel() + return ctx.Err() + }, true, nil, 2) + if !errors.Is(err, context.Canceled) || client.attempts != 1 { + t.Fatalf("expected cancellation during backoff, got err=%v attempts=%d", err, client.attempts) + } + }) + + t.Run("on final attempt", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + client := &scriptedImageImportClient{ + Client: fakectrlruntimeclient.NewClientBuilder().Build(), + create: func(_ context.Context, attempt int, _ *imagev1.ImageStreamImport) error { + if attempt == 2 { + cancel() + } + return nil + }, + } + _, err := importTagWithRetryDelays(ctx, client, "ns", "release", "latest", "registry/release:latest", []time.Duration{0}, func(context.Context, time.Duration) error { return nil }, true, nil, 2) + if !errors.Is(err, context.Canceled) || errors.Is(err, wait.ErrWaitTimeout) { + t.Fatalf("expected final-attempt cancellation, got %v", err) + } + if client.attempts != 2 { + t.Fatalf("expected cancellation on second attempt, got %d attempts", client.attempts) + } + }) +} + +func TestImageImportRetryPolicyValidationAndJitter(t *testing.T) { + client := &scriptedImageImportClient{Client: fakectrlruntimeclient.NewClientBuilder().Build(), create: func(context.Context, int, *imagev1.ImageStreamImport) error { return nil }} + _, err := importTagWithRetryDelays(context.Background(), client, "ns", "release", "latest", "registry/release:latest", []time.Duration{time.Second}, func(context.Context, time.Duration) error { return nil }, false, nil, 3) + if err == nil || !strings.Contains(err.Error(), "invalid image import retry policy") || client.attempts != 0 { + t.Fatalf("expected inconsistent policy to fail before attempts, got err=%v attempts=%d", err, client.attempts) + } + _, err = ImportTagWithRetryDelays(context.Background(), client, "ns", "release", "latest", "registry/release:latest", []time.Duration{-time.Second}, nil) + if err == nil || !strings.Contains(err.Error(), "must not be negative") || client.attempts != 0 { + t.Fatalf("expected negative delay to fail before attempts, got err=%v attempts=%d", err, client.attempts) + } + + got := ReleaseImportRetryDelays() + wantBase := exponentialImageImportRetryDelays(9) + if len(got) != len(wantBase) { + t.Fatalf("release retry schedule has %d delays, want %d", len(got), len(wantBase)) + } + for i := range got { + if got[i] < wantBase[i] || got[i] > wantBase[i]+wantBase[i]/10 { + t.Fatalf("jittered delay %d = %s, want within [%s,%s]", i, got[i], wantBase[i], wantBase[i]+wantBase[i]/10) + } + } +} + +func TestImageImportRetryDelay(t *testing.T) { + testCases := []struct { + name string + err error + configured time.Duration + want time.Duration + }{ + {name: "no suggestion", err: errors.New("transient"), configured: 5 * time.Second, want: 5 * time.Second}, + {name: "shorter suggestion", err: kerrors.NewTooManyRequests("busy", 3), configured: 5 * time.Second, want: 5 * time.Second}, + {name: "equal suggestion", err: kerrors.NewTooManyRequests("busy", 5), configured: 5 * time.Second, want: 5 * time.Second}, + {name: "longer suggestion", err: kerrors.NewTooManyRequests("busy", 7), configured: 5 * time.Second, want: 7 * time.Second}, + {name: "suggestion is capped", err: kerrors.NewTooManyRequests("busy", 600), configured: 5 * time.Second, want: maxImageImportRetryDelay}, + } + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + if got := imageImportRetryDelay(testCase.err, testCase.configured); got != testCase.want { + t.Fatalf("delay = %s, want %s", got, testCase.want) + } + }) + } +} + +func (c *outageImageImportClient) Create(ctx context.Context, obj ctrlruntimeclient.Object, opts ...ctrlruntimeclient.CreateOption) error { + streamImport, ok := obj.(*imagev1.ImageStreamImport) + if !ok { + return c.Client.Create(ctx, obj, opts...) + } + c.attempts++ + if c.permanentErr != nil { + return c.permanentErr + } + if c.now() >= 2*time.Minute { + streamImport.Status.Images = []imagev1.ImageImportStatus{{Image: &imagev1.Image{DockerImageReference: "quay.io/openshift/release@sha256:resolved"}}} + } + return nil +} + func bcc(upstream ctrlruntimeclient.Client) ctrlruntimeclient.Client { c := &imageStreamImportStatusSettingClient{ Client: upstream, @@ -379,7 +692,7 @@ func TestGetEvaluator(t *testing.T) { expected: false, }, { - name: "reimport with error", + name: "permanent reimport failure is returned", client: bcc(fakectrlruntimeclient.NewClientBuilder().Build()), obj: &imagev1.ImageStream{ ObjectMeta: metav1.ObjectMeta{ @@ -682,6 +995,115 @@ func TestWaitForImportingISTagSpecTimeout(t *testing.T) { } } +func TestImportEvaluatorContinuesAfterTransientReimportFailure(t *testing.T) { + const ( + namespace = "test-namespace" + streamName = "stable" + ) + from := &coreapi.ObjectReference{Kind: "DockerImage", Name: "quay.io/openshift/release:latest"} + failing := &imagev1.ImageStream{ + ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: streamName}, + Spec: imagev1.ImageStreamSpec{Tags: []imagev1.TagReference{{ + Name: "cli", + From: from, + }}}, + Status: imagev1.ImageStreamStatus{Tags: []imagev1.NamedTagEventList{{ + Tag: "cli", + Conditions: []imagev1.TagEventCondition{{ + Message: "Internal error occurred: registry unavailable", + }}, + }}}, + } + recovered := failing.DeepCopy() + recovered.Status.PublicDockerImageRepository = "registry.example.com/test-namespace/stable" + recovered.Status.Tags = []imagev1.NamedTagEventList{{ + Tag: "cli", + Items: []imagev1.TagEvent{{ + Image: "sha256:resolved", + DockerImageReference: "quay.io/openshift/release@sha256:resolved", + }}, + }} + client := fakectrlruntimeclient.NewClientBuilder().WithScheme(scheme.Scheme).WithRuntimeObjects(failing).Build() + var importAttempts int + var logs bytes.Buffer + logger := logrus.StandardLogger() + originalOutput := logger.Out + logger.SetOutput(&logs) + defer logger.SetOutput(originalOutput) + + evaluate := getEvaluatorWithImporter(context.Background(), client, namespace, streamName, sets.New("cli"), false, nil, func(context.Context, ctrlruntimeclient.Client, string, string, string, string, int, *metrics.MetricsAgent) (string, error) { + importAttempts++ + return "", &transientImageImportError{err: wait.ErrWaitTimeout} + }) + done, err := evaluate(failing) + if err != nil { + t.Fatalf("expected transient failure to remain retryable: %v", err) + } + if done { + t.Fatal("expected evaluator to continue polling after transient reimport failure") + } + if importAttempts != 1 { + t.Fatalf("expected one failed reimport attempt, got %d", importAttempts) + } + if !strings.Contains(logs.String(), "continuing to wait") { + t.Fatalf("expected transient reimport warning, got %q", logs.String()) + } + done, err = evaluate(recovered) + if err != nil { + t.Fatalf("expected recovered import to succeed: %v", err) + } + if !done { + t.Fatal("expected evaluator to finish after the import recovered") + } +} + +func TestImportEvaluatorPreservesPermanentAndContextErrors(t *testing.T) { + resource := schema.GroupResource{Group: "image.openshift.io", Resource: "imagestreamimports"} + testCases := []struct { + name string + importErr error + transient bool + }{ + {name: "transient exhaustion", importErr: &transientImageImportError{err: wait.ErrWaitTimeout}, transient: true}, + {name: "forbidden", importErr: kerrors.NewForbidden(resource, "release", errors.New("denied"))}, + {name: "unauthorized", importErr: kerrors.NewUnauthorized("unauthorized")}, + {name: "generic", importErr: errors.New("client failure")}, + {name: "canceled", importErr: context.Canceled}, + } + stream := &imagev1.ImageStream{ + ObjectMeta: metav1.ObjectMeta{Namespace: "ns", Name: "release"}, + Spec: imagev1.ImageStreamSpec{Tags: []imagev1.TagReference{{ + Name: "latest", + From: &coreapi.ObjectReference{Kind: "DockerImage", Name: "registry/release:latest"}, + }}}, + Status: imagev1.ImageStreamStatus{Tags: []imagev1.NamedTagEventList{{ + Tag: "latest", + Conditions: []imagev1.TagEventCondition{{Message: "Internal error occurred: registry unavailable"}}, + }}}, + } + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + evaluate := getEvaluatorWithImporter(context.Background(), fakectrlruntimeclient.NewClientBuilder().Build(), "ns", "release", sets.New("latest"), false, nil, func(context.Context, ctrlruntimeclient.Client, string, string, string, string, int, *metrics.MetricsAgent) (string, error) { + return "", testCase.importErr + }) + done, err := evaluate(stream) + if done { + t.Fatal("expected evaluator not to complete") + } + if testCase.transient { + if err != nil { + t.Fatalf("expected transient error to continue polling: %v", err) + } + return + } + if !errors.Is(err, testCase.importErr) { + t.Fatalf("expected permanent/context error %v, got %v", testCase.importErr, err) + } + }) + } +} + func TestImageDigestForSpecTagWithoutFrom(t *testing.T) { is := &imagev1.ImageStream{ ObjectMeta: metav1.ObjectMeta{Namespace: "ns", Name: "pipeline"}, diff --git a/pkg/util/pods.go b/pkg/util/pods.go index a7717b9128..c0a8a9a786 100644 --- a/pkg/util/pods.go +++ b/pkg/util/pods.go @@ -21,6 +21,7 @@ import ( "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" + utilnet "k8s.io/apimachinery/pkg/util/net" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/utils/ptr" ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" @@ -88,22 +89,80 @@ func waitForCompletedPodDeletion(ctx context.Context, podClient ctrlruntimeclien return nil } - // delete the pod we expect, otherwise another user has relaunched this pod - uid := pod.UID - err := podClient.Delete(ctx, pod, ctrlruntimeclient.Preconditions(metav1.Preconditions{UID: &uid})) - if kerrors.IsNotFound(err) { - return nil + return DeletePodWithUID(ctx, podClient, pod) +} + +// DeletePodWithUID deletes exactly the observed pod and waits until that UID is +// gone. A pod recreated under the same namespace and name is left untouched. +func DeletePodWithUID(ctx context.Context, podClient ctrlruntimeclient.Client, pod *corev1.Pod) error { + return deletePodWithUID(ctx, podClient, pod, wait.Backoff{Duration: 2 * time.Second, Factor: 2, Steps: 10}) +} + +func deletePodWithUID(ctx context.Context, podClient ctrlruntimeclient.Client, pod *corev1.Pod, backoff wait.Backoff) error { + if pod == nil { + return errors.New("cannot delete a nil pod") + } + if pod.UID == "" { + return fmt.Errorf("cannot safely delete pod %s/%s without a UID", pod.Namespace, pod.Name) } - if kerrors.IsConflict(err) { - // UID precondition mismatch: the pod was replaced between our GET and - // DELETE. The completed pod we intended to remove is already gone. + + uid := pod.UID + key := ctrlruntimeclient.ObjectKey{Namespace: pod.Namespace, Name: pod.Name} + deleteAccepted := false + var lastRetryableErr error + err := wait.ExponentialBackoffWithContext(ctx, backoff, func(ctx context.Context) (bool, error) { + var deleteErr error + if !deleteAccepted { + deleteErr = podClient.Delete(ctx, &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Namespace: pod.Namespace, Name: pod.Name}}, ctrlruntimeclient.Preconditions(metav1.Preconditions{UID: &uid})) + switch { + case deleteErr == nil: + deleteAccepted = true + case kerrors.IsNotFound(deleteErr): + return true, nil + case kerrors.IsConflict(deleteErr), isRetryablePodRequestError(deleteErr): + lastRetryableErr = deleteErr + default: + return false, fmt.Errorf("could not delete completed pod: %w", deleteErr) + } + } + + current := &corev1.Pod{} + getErr := podClient.Get(ctx, key, current) + switch { + case kerrors.IsNotFound(getErr): + return true, nil + case getErr != nil && isRetryablePodRequestError(getErr): + lastRetryableErr = getErr + return false, nil + case getErr != nil: + return false, fmt.Errorf("could not retrieve deleting pod: %w", getErr) + case current.UID != uid: + return true, nil + case kerrors.IsConflict(deleteErr): + return false, fmt.Errorf("could not delete completed pod with matching UID: %w", deleteErr) + default: + if deleteAccepted { + lastRetryableErr = nil + } + logrus.Debugf("Waiting for pod %s to be deleted ...", pod.Name) + return false, nil + } + }) + if err == nil { return nil } - if err != nil { - return fmt.Errorf("could not delete completed pod: %w", err) + if wait.Interrupted(err) && lastRetryableErr != nil { + err = errors.Join(err, lastRetryableErr) } + return fmt.Errorf("could not confirm completed pod deletion: %w", err) +} - return WaitForPodDeletion(ctx, podClient, namespace, name, uid) +func isRetryablePodRequestError(err error) bool { + return utilnet.IsConnectionReset(err) || + utilnet.IsConnectionRefused(err) || + utilnet.IsHTTP2ConnectionLost(err) || + utilnet.IsProbableEOF(err) || + utilnet.IsTimeout(err) } func WaitForPodDeletion(ctx context.Context, podClient ctrlruntimeclient.Client, namespace, name string, uid types.UID) error { diff --git a/pkg/util/pods_test.go b/pkg/util/pods_test.go index ad90f0b221..7cbc6f6c01 100644 --- a/pkg/util/pods_test.go +++ b/pkg/util/pods_test.go @@ -3,6 +3,8 @@ package util import ( "context" "errors" + "fmt" + "syscall" "testing" "time" @@ -10,6 +12,7 @@ import ( apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/wait" ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" fakectrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client/fake" "sigs.k8s.io/controller-runtime/pkg/client/interceptor" @@ -80,8 +83,9 @@ func TestWaitForCompletedPodDeletion(t *testing.T) { }, }, { - name: "delete returns Conflict (UID mismatch) is treated as success", - objects: []ctrlruntimeclient.Object{succeededPod}, + name: "delete returns Conflict while the same UID remains", + objects: []ctrlruntimeclient.Object{succeededPod}, + expectErr: true, interceptorsFunc: func() interceptor.Funcs { return interceptor.Funcs{ Delete: func(_ context.Context, _ ctrlruntimeclient.WithWatch, _ ctrlruntimeclient.Object, _ ...ctrlruntimeclient.DeleteOption) error { @@ -124,6 +128,181 @@ func TestWaitForCompletedPodDeletion(t *testing.T) { } } +func TestDeletePodWithUIDReconcilesAmbiguousDeletion(t *testing.T) { + const ( + namespace = "test-ns" + name = "test-pod" + ) + observed := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: name, UID: types.UID("original-uid")}} + + t.Run("committed delete with lost response", func(t *testing.T) { + var deleteCalls int + client := fakectrlruntimeclient.NewClientBuilder().WithObjects(observed.DeepCopy()).WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(ctx context.Context, client ctrlruntimeclient.WithWatch, obj ctrlruntimeclient.Object, opts ...ctrlruntimeclient.DeleteOption) error { + deleteCalls++ + if err := client.Delete(ctx, obj, opts...); err != nil { + return err + } + return syscall.ECONNRESET + }, + }).Build() + + if err := deletePodWithUID(context.Background(), client, observed, wait.Backoff{Steps: 3}); err != nil { + t.Fatalf("expected committed deletion to be reconciled: %v", err) + } + if deleteCalls != 1 { + t.Fatalf("delete calls = %d, want 1", deleteCalls) + } + if err := client.Get(context.Background(), ctrlruntimeclient.ObjectKeyFromObject(observed), &corev1.Pod{}); !apierrors.IsNotFound(err) { + t.Fatalf("expected original pod to be gone, got %v", err) + } + }) + + t.Run("transient confirmation GET is retried without another DELETE", func(t *testing.T) { + var deleteCalls, getCalls int + client := fakectrlruntimeclient.NewClientBuilder().WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(context.Context, ctrlruntimeclient.WithWatch, ctrlruntimeclient.Object, ...ctrlruntimeclient.DeleteOption) error { + deleteCalls++ + return nil + }, + Get: func(_ context.Context, _ ctrlruntimeclient.WithWatch, _ ctrlruntimeclient.ObjectKey, _ ctrlruntimeclient.Object, _ ...ctrlruntimeclient.GetOption) error { + getCalls++ + if getCalls == 1 { + return fmt.Errorf("confirm deletion: %w", syscall.ECONNREFUSED) + } + return apierrors.NewNotFound(corev1.Resource("pods"), name) + }, + }).Build() + + if err := deletePodWithUID(context.Background(), client, observed, wait.Backoff{Steps: 3}); err != nil { + t.Fatalf("expected transient confirmation failure to recover: %v", err) + } + if deleteCalls != 1 || getCalls != 2 { + t.Fatalf("calls delete=%d get=%d, want delete=1 get=2", deleteCalls, getCalls) + } + }) + + t.Run("permanent confirmation GET fails", func(t *testing.T) { + forbidden := apierrors.NewForbidden(corev1.Resource("pods"), name, errors.New("denied")) + client := fakectrlruntimeclient.NewClientBuilder().WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(context.Context, ctrlruntimeclient.WithWatch, ctrlruntimeclient.Object, ...ctrlruntimeclient.DeleteOption) error { + return nil + }, + Get: func(context.Context, ctrlruntimeclient.WithWatch, ctrlruntimeclient.ObjectKey, ctrlruntimeclient.Object, ...ctrlruntimeclient.GetOption) error { + return forbidden + }, + }).Build() + + err := deletePodWithUID(context.Background(), client, observed, wait.Backoff{Steps: 3}) + if !apierrors.IsForbidden(err) { + t.Fatalf("expected permanent GET error, got %v", err) + } + }) + + t.Run("same UID after transient DELETE is retried boundedly", func(t *testing.T) { + var deleteCalls, getCalls int + client := fakectrlruntimeclient.NewClientBuilder().WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(context.Context, ctrlruntimeclient.WithWatch, ctrlruntimeclient.Object, ...ctrlruntimeclient.DeleteOption) error { + deleteCalls++ + return syscall.ECONNRESET + }, + Get: func(_ context.Context, _ ctrlruntimeclient.WithWatch, _ ctrlruntimeclient.ObjectKey, obj ctrlruntimeclient.Object, _ ...ctrlruntimeclient.GetOption) error { + getCalls++ + observed.DeepCopyInto(obj.(*corev1.Pod)) + return nil + }, + }).Build() + + err := deletePodWithUID(context.Background(), client, observed, wait.Backoff{Steps: 3}) + if !wait.Interrupted(err) || !errors.Is(err, syscall.ECONNRESET) { + t.Fatalf("expected bounded transient deletion failure, got %v", err) + } + if deleteCalls != 3 || getCalls != 3 { + t.Fatalf("calls delete=%d get=%d, want 3 each", deleteCalls, getCalls) + } + }) + + t.Run("accepted DELETE polls a remaining UID without redundant DELETEs", func(t *testing.T) { + var deleteCalls, getCalls int + client := fakectrlruntimeclient.NewClientBuilder().WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(context.Context, ctrlruntimeclient.WithWatch, ctrlruntimeclient.Object, ...ctrlruntimeclient.DeleteOption) error { + deleteCalls++ + return nil + }, + Get: func(_ context.Context, _ ctrlruntimeclient.WithWatch, _ ctrlruntimeclient.ObjectKey, obj ctrlruntimeclient.Object, _ ...ctrlruntimeclient.GetOption) error { + getCalls++ + observed.DeepCopyInto(obj.(*corev1.Pod)) + return nil + }, + }).Build() + + err := deletePodWithUID(context.Background(), client, observed, wait.Backoff{Steps: 3}) + if !wait.Interrupted(err) { + t.Fatalf("expected bounded confirmation failure, got %v", err) + } + if deleteCalls != 1 || getCalls != 3 { + t.Fatalf("calls delete=%d get=%d, want delete=1 get=3", deleteCalls, getCalls) + } + }) +} + +func TestDeletePodWithUIDSameUIDConflictFails(t *testing.T) { + const ( + namespace = "test-ns" + name = "test-pod" + ) + observed := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: name, UID: types.UID("original-uid")}} + var deleteCalls int + client := fakectrlruntimeclient.NewClientBuilder().WithObjects(observed.DeepCopy()).WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(context.Context, ctrlruntimeclient.WithWatch, ctrlruntimeclient.Object, ...ctrlruntimeclient.DeleteOption) error { + deleteCalls++ + return apierrors.NewConflict(corev1.Resource("pods"), name, errors.New("conflict")) + }, + }).Build() + + err := deletePodWithUID(context.Background(), client, observed, wait.Backoff{Steps: 3}) + if !apierrors.IsConflict(err) { + t.Fatalf("expected same-UID conflict to fail permanently, got %v", err) + } + if deleteCalls != 1 { + t.Fatalf("delete calls = %d, want 1", deleteCalls) + } +} + +func TestDeletePodWithUIDDoesNotDeleteReplacement(t *testing.T) { + const ( + namespace = "test-ns" + name = "test-pod" + ) + observed := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: name, UID: types.UID("old-uid")}} + replacement := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: name, UID: types.UID("replacement-uid")}} + var preconditionUID types.UID + client := fakectrlruntimeclient.NewClientBuilder().WithObjects(replacement).WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(_ context.Context, _ ctrlruntimeclient.WithWatch, _ ctrlruntimeclient.Object, opts ...ctrlruntimeclient.DeleteOption) error { + deleteOptions := &ctrlruntimeclient.DeleteOptions{} + deleteOptions.ApplyOptions(opts) + if deleteOptions.Preconditions != nil && deleteOptions.Preconditions.UID != nil { + preconditionUID = *deleteOptions.Preconditions.UID + } + return apierrors.NewConflict(corev1.Resource("pods"), name, errors.New("UID precondition mismatch")) + }, + }).Build() + + if err := DeletePodWithUID(context.Background(), client, observed); err != nil { + t.Fatalf("stale deletion should be harmless: %v", err) + } + if preconditionUID != observed.UID { + t.Fatalf("delete precondition UID = %q, want %q", preconditionUID, observed.UID) + } + current := &corev1.Pod{} + if err := client.Get(context.Background(), ctrlruntimeclient.ObjectKey{Namespace: namespace, Name: name}, current); err != nil { + t.Fatalf("replacement pod was deleted: %v", err) + } + if current.UID != replacement.UID { + t.Fatalf("got pod UID %q, want replacement UID %q", current.UID, replacement.UID) + } +} + func TestCheckPending(t *testing.T) { timeout, now := 30*time.Minute, time.Time{} withinLimit := metav1.Time{Time: now.Add(-time.Minute)}