Skip to content
Merged
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
6 changes: 5 additions & 1 deletion cmd/convertor/builder/builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,9 @@ type BuilderOptions struct {
// Push manifests with subject
Referrer bool

// Number of retries for registry upload operations when encountering 429 rate limiting
RetryCount int

// CustomResolver allows using a custom resolver instead of the default docker resolver
// Used for tar import/export functionality
CustomResolver remotes.Resolver
Expand Down Expand Up @@ -223,7 +226,7 @@ func (b *graphBuilder) process(ctx context.Context, src v1.Descriptor, tag bool)
} else {
pusher = b.pusher
}
if err := uploadBytes(ctx, pusher, expected, indexBytes); err != nil {
if err := uploadBytesWithRetry(ctx, pusher, expected, indexBytes, b.RetryCount); err != nil {
return v1.Descriptor{}, fmt.Errorf("failed to upload index: %w", err)
}
log.G(ctx).Infof("index uploaded, %s", expected.Digest)
Expand Down Expand Up @@ -296,6 +299,7 @@ func (b *graphBuilder) buildOne(ctx context.Context, src v1.Descriptor, tag bool
engineBase.reserve = b.Reserve
engineBase.noUpload = b.NoUpload
engineBase.dumpManifest = b.DumpManifest
engineBase.retryCount = b.RetryCount
if _, ok := b.Resolver.(*FileBasedResolver); ok {
engineBase.tarExport = true
}
Expand Down
5 changes: 3 additions & 2 deletions cmd/convertor/builder/builder_engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,7 @@ type builderEngineBase struct {
dumpManifest bool
referrer bool
tarExport bool
retryCount int
}

func (e *builderEngineBase) isGzipLayer(ctx context.Context, idx int) (bool, error) {
Expand Down Expand Up @@ -183,7 +184,7 @@ func (e *builderEngineBase) uploadManifestAndConfig(ctx context.Context) (specs.
Size: (int64)(len(cbuf)),
}
if shouldUploadBlob {
if err = uploadBytes(ctx, e.pusher, e.manifest.Config, cbuf); err != nil {
if err = uploadBytesWithRetry(ctx, e.pusher, e.manifest.Config, cbuf, e.retryCount); err != nil {
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload config")
}
log.G(ctx).Infof("config uploaded")
Expand All @@ -207,7 +208,7 @@ func (e *builderEngineBase) uploadManifestAndConfig(ctx context.Context) (specs.
Size: (int64)(len(cbuf)),
}
if shouldUploadBlob {
if err = uploadBytes(ctx, e.pusher, manifestDesc, cbuf); err != nil {
if err = uploadBytesWithRetry(ctx, e.pusher, manifestDesc, cbuf, e.retryCount); err != nil {
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload manifest")
}
e.outputDesc = manifestDesc
Expand Down
138 changes: 111 additions & 27 deletions cmd/convertor/builder/builder_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,12 +24,16 @@ import (
"encoding/json"
"fmt"
"io"
"math"
"math/rand"
"os"
"path"
"time"

"github.com/containerd/containerd/v2/core/content"
"github.com/containerd/containerd/v2/core/images"
"github.com/containerd/containerd/v2/core/remotes"
"github.com/containerd/containerd/v2/core/remotes/docker"
"github.com/containerd/containerd/v2/pkg/archive/compression"
"github.com/containerd/continuity"
"github.com/containerd/errdefs"
Expand All @@ -43,6 +47,70 @@ import (
t "github.com/containerd/accelerated-container-image/pkg/types"
)

// isRetryableError checks if the error is retryable (429 or 5xx errors)
func isRetryableError(err error) bool {
if err == nil {
return false
}

// Check for containerd docker error types
var dockerErr *docker.Error
if errors.As(err, &dockerErr) {
switch dockerErr.Code {
case docker.ErrorCodeTooManyRequests:
return true
case docker.ErrorCodeUnavailable:
return true
default:
return false
}
}

return false
}

// retryWithBackoff executes a function with exponential backoff on retryable errors
func retryWithBackoff(ctx context.Context, maxRetries int, operation func() error) error {
var lastErr error

for attempt := 0; attempt <= maxRetries; attempt++ {
lastErr = operation()

if lastErr == nil {
return nil
}

if !isRetryableError(lastErr) {
return lastErr
}

if attempt == maxRetries {
logrus.Warnf("max retries (%d) reached for retryable error: %v", maxRetries, lastErr)
return lastErr
}

// Exponential backoff with random jitter: base delay of 1s, max 30s
baseDelay := time.Duration(math.Min(float64(time.Second)*math.Pow(2, float64(attempt)), float64(30*time.Second)))
// Add random jitter: ±25% of the base delay
jitter := time.Duration(float64(baseDelay) * (rand.Float64() - 0.5) * 0.5)
backoffDelay := baseDelay + jitter
// Ensure minimum delay of 500ms
if backoffDelay < 500*time.Millisecond {
backoffDelay = 500 * time.Millisecond
}
logrus.Infof("received retryable error, retrying in %v (attempt %d/%d): %v", backoffDelay, attempt+1, maxRetries, lastErr)

select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoffDelay):
continue
}
}

return lastErr
}

func fetch(ctx context.Context, fetcher remotes.Fetcher, desc specs.Descriptor, target any) error {
rc, err := fetcher.Fetch(ctx, desc)
if err != nil {
Expand Down Expand Up @@ -200,41 +268,57 @@ func getFileDesc(filepath string, decompress bool) (specs.Descriptor, error) {
}

func uploadBlob(ctx context.Context, pusher remotes.Pusher, path string, desc specs.Descriptor) error {
cw, err := pusher.Push(ctx, desc)
if err != nil {
if errdefs.IsAlreadyExists(err) {
logrus.Infof("layer %s exists", desc.Digest.String())
return nil
return uploadBlobWithRetry(ctx, pusher, path, desc, 0)
}

func uploadBlobWithRetry(ctx context.Context, pusher remotes.Pusher, path string, desc specs.Descriptor, retryCount int) error {
return retryWithBackoff(ctx, retryCount, func() error {
cw, err := pusher.Push(ctx, desc)
if err != nil {
if errdefs.IsAlreadyExists(err) {
logrus.Infof("layer %s exists", desc.Digest.String())
return nil
}
return err
}
return err
}

defer cw.Close()
fobd, err := os.Open(path)
if err != nil {
return err
}
defer fobd.Close()
if err = content.Copy(ctx, cw, fobd, desc.Size, desc.Digest); err != nil {
return err
}
return nil
defer cw.Close()
fobd, err := os.Open(path)
if err != nil {
return err
}
defer fobd.Close()
if err = content.Copy(ctx, cw, fobd, desc.Size, desc.Digest); err != nil {
return err
}
return nil
})
}

func uploadBytes(ctx context.Context, pusher remotes.Pusher, desc specs.Descriptor, data []byte) error {
cw, err := pusher.Push(ctx, desc)
if err != nil {
if errdefs.IsAlreadyExists(err) {
logrus.Infof("content %s exists", desc.Digest.String())
return nil
return uploadBytesWithRetry(ctx, pusher, desc, data, 0)
}

func uploadBytesWithRetry(ctx context.Context, pusher remotes.Pusher, desc specs.Descriptor, data []byte, retryCount int) error {
return retryWithBackoff(ctx, retryCount, func() error {
cw, err := pusher.Push(ctx, desc)
if err != nil {
if errdefs.IsAlreadyExists(err) {
logrus.Infof("content %s exists", desc.Digest.String())
return nil
}
return err
}
return err
}
defer cw.Close()
return content.Copy(ctx, cw, bytes.NewReader(data), desc.Size, desc.Digest)
defer cw.Close()
return content.Copy(ctx, cw, bytes.NewReader(data), desc.Size, desc.Digest)
})
}

func tagPreviouslyConvertedManifest(ctx context.Context, pusher remotes.Pusher, fetcher remotes.Fetcher, desc specs.Descriptor) error {
return tagPreviouslyConvertedManifestWithRetry(ctx, pusher, fetcher, desc, 0)
}

func tagPreviouslyConvertedManifestWithRetry(ctx context.Context, pusher remotes.Pusher, fetcher remotes.Fetcher, desc specs.Descriptor, retryCount int) error {
manifest := specs.Manifest{}
if err := fetch(ctx, fetcher, desc, &manifest); err != nil {
return fmt.Errorf("failed to fetch converted manifest: %w", err)
Expand All @@ -243,7 +327,7 @@ func tagPreviouslyConvertedManifest(ctx context.Context, pusher remotes.Pusher,
if err != nil {
return err
}
if err := uploadBytes(ctx, pusher, desc, cbuf); err != nil {
if err := uploadBytesWithRetry(ctx, pusher, desc, cbuf, retryCount); err != nil {
return fmt.Errorf("failed to tag converted manifest: %w", err)
}
return nil
Expand Down
8 changes: 4 additions & 4 deletions cmd/convertor/builder/overlaybd_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -192,7 +192,7 @@ func (e *overlaybdBuilderEngine) UploadLayer(ctx context.Context, idx int) error
}
shouldUploadBlob := !e.noUpload || e.tarExport
if shouldUploadBlob {
if err := uploadBlob(ctx, e.pusher, path.Join(layerDir, commitFile), desc); err != nil {
if err := uploadBlobWithRetry(ctx, e.pusher, path.Join(layerDir, commitFile), desc, e.retryCount); err != nil {
return errors.Wrapf(err, "failed to upload layer %d", idx)
}
}
Expand Down Expand Up @@ -354,7 +354,7 @@ func (e *overlaybdBuilderEngine) CheckForConvertedManifest(ctx context.Context)

// If a converted manifest has been found we still need to tag it to match the expected output tag.
func (e *overlaybdBuilderEngine) TagPreviouslyConvertedManifest(ctx context.Context, desc specs.Descriptor) error {
return tagPreviouslyConvertedManifest(ctx, e.pusher, e.fetcher, desc)
return tagPreviouslyConvertedManifestWithRetry(ctx, e.pusher, e.fetcher, desc, e.retryCount)
}

// mountImage is responsible for mounting a specific manifest from a source repository, this includes
Expand Down Expand Up @@ -391,7 +391,7 @@ func (e *overlaybdBuilderEngine) mountImage(ctx context.Context, manifest specs.
if err != nil {
return err
}
return uploadBytes(ctx, e.pusher, desc, cbuf)
return uploadBytesWithRetry(ctx, e.pusher, desc, cbuf, e.retryCount)
}

func (e *overlaybdBuilderEngine) StoreConvertedManifestDetails(ctx context.Context) error {
Expand Down Expand Up @@ -483,7 +483,7 @@ func (e *overlaybdBuilderEngine) uploadBaseLayer(ctx context.Context) (specs.Des
}
shouldUploadBlob := !e.noUpload || e.tarExport
if shouldUploadBlob {
if err = uploadBlob(ctx, e.pusher, tarFile, baseDesc); err != nil {
if err = uploadBlobWithRetry(ctx, e.pusher, tarFile, baseDesc, e.retryCount); err != nil {
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload baselayer")
}
logrus.Infof("baselayer uploaded")
Expand Down
6 changes: 3 additions & 3 deletions cmd/convertor/builder/turboOCI_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,7 @@ func (e *turboOCIBuilderEngine) UploadLayer(ctx context.Context, idx int) error
}
}
desc.Annotations[label.TurboOCIMediaType] = targetMediaType
if err := uploadBlob(ctx, e.pusher, path.Join(layerDir, tociLayerTar), desc); err != nil {
if err := uploadBlobWithRetry(ctx, e.pusher, path.Join(layerDir, tociLayerTar), desc, e.retryCount); err != nil {
return errors.Wrapf(err, "failed to upload layer %d", idx)
}
e.tociLayers[idx] = desc
Expand Down Expand Up @@ -196,7 +196,7 @@ func (e *turboOCIBuilderEngine) UploadImage(ctx context.Context) (specs.Descript
},
}
if !e.mkfs {
if err := uploadBlob(ctx, e.pusher, overlaybdBaseLayer, baseDesc); err != nil {
if err := uploadBlobWithRetry(ctx, e.pusher, overlaybdBaseLayer, baseDesc, e.retryCount); err != nil {
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload baselayer %q", overlaybdBaseLayer)
}
e.manifest.Layers = append([]specs.Descriptor{baseDesc}, e.manifest.Layers...)
Expand All @@ -215,7 +215,7 @@ func (e *turboOCIBuilderEngine) UploadImage(ctx context.Context) (specs.Descript

// If a converted manifest has been found we still need to tag it to match the expected output tag.
func (e *turboOCIBuilderEngine) TagPreviouslyConvertedManifest(ctx context.Context, desc specs.Descriptor) error {
return tagPreviouslyConvertedManifest(ctx, e.pusher, e.fetcher, desc)
return tagPreviouslyConvertedManifestWithRetry(ctx, e.pusher, e.fetcher, desc, e.retryCount)
}

// Layer deduplication in FastOCI is not currently supported due to conversion not
Expand Down
4 changes: 4 additions & 0 deletions cmd/convertor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ var (
concurrencyLimit int
disableSparse bool
referrer bool
retryCount int

// tar import/export
importTar string
Expand Down Expand Up @@ -279,6 +280,7 @@ Version: ` + commitID,
ConcurrencyLimit: concurrencyLimit,
DisableSparse: disableSparse,
Referrer: referrer,
RetryCount: retryCount,
}
} else {
// Normal registry mode
Expand Down Expand Up @@ -307,6 +309,7 @@ Version: ` + commitID,
ConcurrencyLimit: concurrencyLimit,
DisableSparse: disableSparse,
Referrer: referrer,
RetryCount: retryCount,
}
}
if overlaybd != "" {
Expand Down Expand Up @@ -393,6 +396,7 @@ func init() {
rootCmd.Flags().IntVar(&concurrencyLimit, "concurrency-limit", 4, "the number of manifests that can be built at the same time, used for multi-arch images, 0 means no limit")
rootCmd.Flags().BoolVar(&disableSparse, "disable-sparse", false, "disable sparse file for overlaybd")
rootCmd.Flags().BoolVar(&referrer, "referrer", false, "push converted manifests with subject, note '--oci' will be enabled automatically if '--referrer' is set, cause the referrer must be in OCI format.")
rootCmd.Flags().IntVar(&retryCount, "retry-count", 5, "number of retries for registry upload operations when encountering 429 rate limiting")

// tar import/export
rootCmd.Flags().StringVar(&importTar, "import-tar", "", "import image from tar file (OCI layout format)")
Expand Down
Loading