Skip to content

Commit 095e48d

Browse files
authored
Merge pull request #12 from runloopai/gautam-retry-upload
convertor: Retry on docker registry upload failures
2 parents 3a38a2b + 5c2bc21 commit 095e48d

6 files changed

Lines changed: 130 additions & 37 deletions

File tree

cmd/convertor/builder/builder.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,9 @@ type BuilderOptions struct {
7575
// Push manifests with subject
7676
Referrer bool
7777

78+
// Number of retries for registry upload operations when encountering 429 rate limiting
79+
RetryCount int
80+
7881
// CustomResolver allows using a custom resolver instead of the default docker resolver
7982
// Used for tar import/export functionality
8083
CustomResolver remotes.Resolver
@@ -223,7 +226,7 @@ func (b *graphBuilder) process(ctx context.Context, src v1.Descriptor, tag bool)
223226
} else {
224227
pusher = b.pusher
225228
}
226-
if err := uploadBytes(ctx, pusher, expected, indexBytes); err != nil {
229+
if err := uploadBytesWithRetry(ctx, pusher, expected, indexBytes, b.RetryCount); err != nil {
227230
return v1.Descriptor{}, fmt.Errorf("failed to upload index: %w", err)
228231
}
229232
log.G(ctx).Infof("index uploaded, %s", expected.Digest)
@@ -296,6 +299,7 @@ func (b *graphBuilder) buildOne(ctx context.Context, src v1.Descriptor, tag bool
296299
engineBase.reserve = b.Reserve
297300
engineBase.noUpload = b.NoUpload
298301
engineBase.dumpManifest = b.DumpManifest
302+
engineBase.retryCount = b.RetryCount
299303
if _, ok := b.Resolver.(*FileBasedResolver); ok {
300304
engineBase.tarExport = true
301305
}

cmd/convertor/builder/builder_engine.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -117,6 +117,7 @@ type builderEngineBase struct {
117117
dumpManifest bool
118118
referrer bool
119119
tarExport bool
120+
retryCount int
120121
}
121122

122123
func (e *builderEngineBase) isGzipLayer(ctx context.Context, idx int) (bool, error) {
@@ -183,7 +184,7 @@ func (e *builderEngineBase) uploadManifestAndConfig(ctx context.Context) (specs.
183184
Size: (int64)(len(cbuf)),
184185
}
185186
if shouldUploadBlob {
186-
if err = uploadBytes(ctx, e.pusher, e.manifest.Config, cbuf); err != nil {
187+
if err = uploadBytesWithRetry(ctx, e.pusher, e.manifest.Config, cbuf, e.retryCount); err != nil {
187188
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload config")
188189
}
189190
log.G(ctx).Infof("config uploaded")
@@ -207,7 +208,7 @@ func (e *builderEngineBase) uploadManifestAndConfig(ctx context.Context) (specs.
207208
Size: (int64)(len(cbuf)),
208209
}
209210
if shouldUploadBlob {
210-
if err = uploadBytes(ctx, e.pusher, manifestDesc, cbuf); err != nil {
211+
if err = uploadBytesWithRetry(ctx, e.pusher, manifestDesc, cbuf, e.retryCount); err != nil {
211212
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload manifest")
212213
}
213214
e.outputDesc = manifestDesc

cmd/convertor/builder/builder_utils.go

Lines changed: 111 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -24,12 +24,16 @@ import (
2424
"encoding/json"
2525
"fmt"
2626
"io"
27+
"math"
28+
"math/rand"
2729
"os"
2830
"path"
31+
"time"
2932

3033
"github.com/containerd/containerd/v2/core/content"
3134
"github.com/containerd/containerd/v2/core/images"
3235
"github.com/containerd/containerd/v2/core/remotes"
36+
"github.com/containerd/containerd/v2/core/remotes/docker"
3337
"github.com/containerd/containerd/v2/pkg/archive/compression"
3438
"github.com/containerd/continuity"
3539
"github.com/containerd/errdefs"
@@ -43,6 +47,70 @@ import (
4347
t "github.com/containerd/accelerated-container-image/pkg/types"
4448
)
4549

50+
// isRetryableError checks if the error is retryable (429 or 5xx errors)
51+
func isRetryableError(err error) bool {
52+
if err == nil {
53+
return false
54+
}
55+
56+
// Check for containerd docker error types
57+
var dockerErr *docker.Error
58+
if errors.As(err, &dockerErr) {
59+
switch dockerErr.Code {
60+
case docker.ErrorCodeTooManyRequests:
61+
return true
62+
case docker.ErrorCodeUnavailable:
63+
return true
64+
default:
65+
return false
66+
}
67+
}
68+
69+
return false
70+
}
71+
72+
// retryWithBackoff executes a function with exponential backoff on retryable errors
73+
func retryWithBackoff(ctx context.Context, maxRetries int, operation func() error) error {
74+
var lastErr error
75+
76+
for attempt := 0; attempt <= maxRetries; attempt++ {
77+
lastErr = operation()
78+
79+
if lastErr == nil {
80+
return nil
81+
}
82+
83+
if !isRetryableError(lastErr) {
84+
return lastErr
85+
}
86+
87+
if attempt == maxRetries {
88+
logrus.Warnf("max retries (%d) reached for retryable error: %v", maxRetries, lastErr)
89+
return lastErr
90+
}
91+
92+
// Exponential backoff with random jitter: base delay of 1s, max 30s
93+
baseDelay := time.Duration(math.Min(float64(time.Second)*math.Pow(2, float64(attempt)), float64(30*time.Second)))
94+
// Add random jitter: ±25% of the base delay
95+
jitter := time.Duration(float64(baseDelay) * (rand.Float64() - 0.5) * 0.5)
96+
backoffDelay := baseDelay + jitter
97+
// Ensure minimum delay of 500ms
98+
if backoffDelay < 500*time.Millisecond {
99+
backoffDelay = 500 * time.Millisecond
100+
}
101+
logrus.Infof("received retryable error, retrying in %v (attempt %d/%d): %v", backoffDelay, attempt+1, maxRetries, lastErr)
102+
103+
select {
104+
case <-ctx.Done():
105+
return ctx.Err()
106+
case <-time.After(backoffDelay):
107+
continue
108+
}
109+
}
110+
111+
return lastErr
112+
}
113+
46114
func fetch(ctx context.Context, fetcher remotes.Fetcher, desc specs.Descriptor, target any) error {
47115
rc, err := fetcher.Fetch(ctx, desc)
48116
if err != nil {
@@ -200,41 +268,57 @@ func getFileDesc(filepath string, decompress bool) (specs.Descriptor, error) {
200268
}
201269

202270
func uploadBlob(ctx context.Context, pusher remotes.Pusher, path string, desc specs.Descriptor) error {
203-
cw, err := pusher.Push(ctx, desc)
204-
if err != nil {
205-
if errdefs.IsAlreadyExists(err) {
206-
logrus.Infof("layer %s exists", desc.Digest.String())
207-
return nil
271+
return uploadBlobWithRetry(ctx, pusher, path, desc, 0)
272+
}
273+
274+
func uploadBlobWithRetry(ctx context.Context, pusher remotes.Pusher, path string, desc specs.Descriptor, retryCount int) error {
275+
return retryWithBackoff(ctx, retryCount, func() error {
276+
cw, err := pusher.Push(ctx, desc)
277+
if err != nil {
278+
if errdefs.IsAlreadyExists(err) {
279+
logrus.Infof("layer %s exists", desc.Digest.String())
280+
return nil
281+
}
282+
return err
208283
}
209-
return err
210-
}
211284

212-
defer cw.Close()
213-
fobd, err := os.Open(path)
214-
if err != nil {
215-
return err
216-
}
217-
defer fobd.Close()
218-
if err = content.Copy(ctx, cw, fobd, desc.Size, desc.Digest); err != nil {
219-
return err
220-
}
221-
return nil
285+
defer cw.Close()
286+
fobd, err := os.Open(path)
287+
if err != nil {
288+
return err
289+
}
290+
defer fobd.Close()
291+
if err = content.Copy(ctx, cw, fobd, desc.Size, desc.Digest); err != nil {
292+
return err
293+
}
294+
return nil
295+
})
222296
}
223297

224298
func uploadBytes(ctx context.Context, pusher remotes.Pusher, desc specs.Descriptor, data []byte) error {
225-
cw, err := pusher.Push(ctx, desc)
226-
if err != nil {
227-
if errdefs.IsAlreadyExists(err) {
228-
logrus.Infof("content %s exists", desc.Digest.String())
229-
return nil
299+
return uploadBytesWithRetry(ctx, pusher, desc, data, 0)
300+
}
301+
302+
func uploadBytesWithRetry(ctx context.Context, pusher remotes.Pusher, desc specs.Descriptor, data []byte, retryCount int) error {
303+
return retryWithBackoff(ctx, retryCount, func() error {
304+
cw, err := pusher.Push(ctx, desc)
305+
if err != nil {
306+
if errdefs.IsAlreadyExists(err) {
307+
logrus.Infof("content %s exists", desc.Digest.String())
308+
return nil
309+
}
310+
return err
230311
}
231-
return err
232-
}
233-
defer cw.Close()
234-
return content.Copy(ctx, cw, bytes.NewReader(data), desc.Size, desc.Digest)
312+
defer cw.Close()
313+
return content.Copy(ctx, cw, bytes.NewReader(data), desc.Size, desc.Digest)
314+
})
235315
}
236316

237317
func tagPreviouslyConvertedManifest(ctx context.Context, pusher remotes.Pusher, fetcher remotes.Fetcher, desc specs.Descriptor) error {
318+
return tagPreviouslyConvertedManifestWithRetry(ctx, pusher, fetcher, desc, 0)
319+
}
320+
321+
func tagPreviouslyConvertedManifestWithRetry(ctx context.Context, pusher remotes.Pusher, fetcher remotes.Fetcher, desc specs.Descriptor, retryCount int) error {
238322
manifest := specs.Manifest{}
239323
if err := fetch(ctx, fetcher, desc, &manifest); err != nil {
240324
return fmt.Errorf("failed to fetch converted manifest: %w", err)
@@ -243,7 +327,7 @@ func tagPreviouslyConvertedManifest(ctx context.Context, pusher remotes.Pusher,
243327
if err != nil {
244328
return err
245329
}
246-
if err := uploadBytes(ctx, pusher, desc, cbuf); err != nil {
330+
if err := uploadBytesWithRetry(ctx, pusher, desc, cbuf, retryCount); err != nil {
247331
return fmt.Errorf("failed to tag converted manifest: %w", err)
248332
}
249333
return nil

cmd/convertor/builder/overlaybd_builder.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -192,7 +192,7 @@ func (e *overlaybdBuilderEngine) UploadLayer(ctx context.Context, idx int) error
192192
}
193193
shouldUploadBlob := !e.noUpload || e.tarExport
194194
if shouldUploadBlob {
195-
if err := uploadBlob(ctx, e.pusher, path.Join(layerDir, commitFile), desc); err != nil {
195+
if err := uploadBlobWithRetry(ctx, e.pusher, path.Join(layerDir, commitFile), desc, e.retryCount); err != nil {
196196
return errors.Wrapf(err, "failed to upload layer %d", idx)
197197
}
198198
}
@@ -354,7 +354,7 @@ func (e *overlaybdBuilderEngine) CheckForConvertedManifest(ctx context.Context)
354354

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

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

397397
func (e *overlaybdBuilderEngine) StoreConvertedManifestDetails(ctx context.Context) error {
@@ -483,7 +483,7 @@ func (e *overlaybdBuilderEngine) uploadBaseLayer(ctx context.Context) (specs.Des
483483
}
484484
shouldUploadBlob := !e.noUpload || e.tarExport
485485
if shouldUploadBlob {
486-
if err = uploadBlob(ctx, e.pusher, tarFile, baseDesc); err != nil {
486+
if err = uploadBlobWithRetry(ctx, e.pusher, tarFile, baseDesc, e.retryCount); err != nil {
487487
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload baselayer")
488488
}
489489
logrus.Infof("baselayer uploaded")

cmd/convertor/builder/turboOCI_builder.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -168,7 +168,7 @@ func (e *turboOCIBuilderEngine) UploadLayer(ctx context.Context, idx int) error
168168
}
169169
}
170170
desc.Annotations[label.TurboOCIMediaType] = targetMediaType
171-
if err := uploadBlob(ctx, e.pusher, path.Join(layerDir, tociLayerTar), desc); err != nil {
171+
if err := uploadBlobWithRetry(ctx, e.pusher, path.Join(layerDir, tociLayerTar), desc, e.retryCount); err != nil {
172172
return errors.Wrapf(err, "failed to upload layer %d", idx)
173173
}
174174
e.tociLayers[idx] = desc
@@ -196,7 +196,7 @@ func (e *turboOCIBuilderEngine) UploadImage(ctx context.Context) (specs.Descript
196196
},
197197
}
198198
if !e.mkfs {
199-
if err := uploadBlob(ctx, e.pusher, overlaybdBaseLayer, baseDesc); err != nil {
199+
if err := uploadBlobWithRetry(ctx, e.pusher, overlaybdBaseLayer, baseDesc, e.retryCount); err != nil {
200200
return specs.Descriptor{}, errors.Wrapf(err, "failed to upload baselayer %q", overlaybdBaseLayer)
201201
}
202202
e.manifest.Layers = append([]specs.Descriptor{baseDesc}, e.manifest.Layers...)
@@ -215,7 +215,7 @@ func (e *turboOCIBuilderEngine) UploadImage(ctx context.Context) (specs.Descript
215215

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

221221
// Layer deduplication in FastOCI is not currently supported due to conversion not

cmd/convertor/main.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,7 @@ var (
5959
concurrencyLimit int
6060
disableSparse bool
6161
referrer bool
62+
retryCount int
6263

6364
// tar import/export
6465
importTar string
@@ -279,6 +280,7 @@ Version: ` + commitID,
279280
ConcurrencyLimit: concurrencyLimit,
280281
DisableSparse: disableSparse,
281282
Referrer: referrer,
283+
RetryCount: retryCount,
282284
}
283285
} else {
284286
// Normal registry mode
@@ -307,6 +309,7 @@ Version: ` + commitID,
307309
ConcurrencyLimit: concurrencyLimit,
308310
DisableSparse: disableSparse,
309311
Referrer: referrer,
312+
RetryCount: retryCount,
310313
}
311314
}
312315
if overlaybd != "" {
@@ -393,6 +396,7 @@ func init() {
393396
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")
394397
rootCmd.Flags().BoolVar(&disableSparse, "disable-sparse", false, "disable sparse file for overlaybd")
395398
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.")
399+
rootCmd.Flags().IntVar(&retryCount, "retry-count", 5, "number of retries for registry upload operations when encountering 429 rate limiting")
396400

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

0 commit comments

Comments
 (0)