Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 25 additions & 10 deletions pkg/steps/pod.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -133,26 +145,29 @@ 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)
if err != nil {
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
}
Expand Down
27 changes: 27 additions & 0 deletions pkg/steps/pod_error_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
153 changes: 145 additions & 8 deletions pkg/steps/release/import_release.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

func (s *importReleaseStep) Inputs() (api.InputDefinition, error) {
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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
Expand All @@ -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{
Expand All @@ -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)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

// read the contents from the configmap we created
Expand Down Expand Up @@ -360,6 +495,8 @@ func ImportReleaseStep(
jobSpec: jobSpec,
pullSecret: pullSecret,
overrideCLIReleaseExtractImage: overrideCLIReleaseExtractImage,
releaseImportRetryDelays: utils.ReleaseImportRetryDelays,
importTagWithRetryDelays: utils.ImportTagWithRetryDelays,
}
}

Expand Down
Loading