From c9818e3da09e743fdefa6c487de9cdfcc4bb82cd Mon Sep 17 00:00:00 2001 From: Vaibhav Patel Date: Wed, 16 Sep 2026 23:28:18 -0400 Subject: [PATCH] gantry: resume detached origin ingests --- cmd/gantry/main.go | 101 ++++++++-- cmd/gantry/origin_pull_test.go | 183 ++++++++++++++++++ designs/gantry-128gib-single-layer-pulls.md | 12 +- internal/gantry/containerdstore/store.go | 52 ++++- internal/gantry/containerdstore/store_test.go | 83 +++++++- internal/gantry/ifaces/ifaces.go | 12 ++ internal/gantry/origin/origin.go | 4 +- internal/gantry/origin/origin_test.go | 2 +- 8 files changed, 412 insertions(+), 37 deletions(-) diff --git a/cmd/gantry/main.go b/cmd/gantry/main.go index a2d933847..da5497525 100644 --- a/cmd/gantry/main.go +++ b/cmd/gantry/main.go @@ -1568,6 +1568,26 @@ type preIngestLeaseStore interface { CreateLease(ctx context.Context, d digest.Digest, registry, repository string) (*containerdstore.LeaseGuard, error) } +type resumableOriginStore interface { + ResumeWriter(ctx context.Context, d digest.Digest) (ifaces.ContentWriter, int64, error) +} + +type preservableOriginWriter interface { + Preserve() error +} + +func openOriginWriter(ctx context.Context, store ifaces.LocalContentStore, d digest.Digest, kind ifaces.OriginRefKind) (ifaces.ContentWriter, int64, error) { + if kind == ifaces.KindBlob { + if resumable, ok := store.(resumableOriginStore); ok { + return resumable.ResumeWriter(ctx, d) + } + } + + w, err := store.Writer(ctx, d) + + return w, 0, err +} + type pullerPumpGate struct { mu sync.Mutex accepting bool @@ -1852,23 +1872,7 @@ func runOriginPull(baseCtx context.Context, originClient ifaces.OriginPuller, cs Digest: d, Kind: kind, } - - rc, expectedSize, err := originClient.Pull(ctx, ref) - if err != nil { - // A delegated credential is requester-specific. Its origin failure - // must not poison the digest-wide cache for another requester. - recordOriginFailure(neg, d, err, lg, "origin pull failed", registry, repository, registryauth.Authorization(ctx) == "", - slog.String("pull_mode", "detached"), - slog.String("deadline_owner", originPullDeadlineOwner(ctx, err)), - slog.Duration("elapsed", time.Since(pullStartedAt)), - slog.Int64("expected_size", -1), - slog.Int64("written", 0), - ) - - return - } - - defer func() { _ = rc.Close() }() //nolint:errcheck // best-effort close + expectedSize := int64(-1) var leaseGuard *containerdstore.LeaseGuard @@ -1913,7 +1917,7 @@ func runOriginPull(baseCtx context.Context, originClient ifaces.OriginPuller, cs releaseCancel() } - w, err := cstore.Writer(ctx, d) + w, resumeOffset, err := openOriginWriter(ctx, cstore, d, kind) deadlineOwner := originPullDeadlineOwner(ctx, err) if err != nil { @@ -1923,6 +1927,7 @@ func runOriginPull(baseCtx context.Context, originClient ifaces.OriginPuller, cs slog.String("deadline_owner", deadlineOwner), slog.Duration("elapsed", time.Since(pullStartedAt)), slog.Int64("expected_size", expectedSize), + slog.Int64("resume_offset", resumeOffset), slog.Int64("written", 0), ) // Origin returned 2xx (we got past originClient.Pull above) @@ -1940,12 +1945,72 @@ func runOriginPull(baseCtx context.Context, originClient ifaces.OriginPuller, cs } defer func() { + if preservable, ok := w.(preservableOriginWriter); ok { + if err := preservable.Preserve(); err != nil { + lg.Warn("preserve partial origin ingest failed", slog.Any("err", err)) + } + + return + } + abortCtx, abortCancel := context.WithTimeout(context.Background(), 10*time.Second) defer abortCancel() _ = w.Abort(abortCtx) //nolint:errcheck // best-effort abort }() + ref.Offset = resumeOffset + + rc, expectedSize, err := originClient.Pull(ctx, ref) + if err != nil && resumeOffset > 0 { + var rangeUnsupported *ifaces.ErrRangeUnsupported + if errors.As(err, &rangeUnsupported) { + abortCtx, abortCancel := context.WithTimeout(context.Background(), 10*time.Second) + abortErr := w.Abort(abortCtx) + + abortCancel() + + if abortErr != nil { + err = fmt.Errorf("abort partial ingest before full retry: %w", abortErr) + } else { + replacement, replacementErr := cstore.Writer(ctx, d) + + err = replacementErr + if err == nil { + w = replacement + + lg.Info("origin does not support resume; restarting from byte zero", + slog.String("digest", d.String()), + slog.String("registry", registry), + slog.String("repository", repository), + slog.Int64("resume_offset", resumeOffset), + ) + + resumeOffset = 0 + ref.Offset = 0 + rc, expectedSize, err = originClient.Pull(ctx, ref) + } + } + } + } + + if err != nil { + // A delegated credential is requester-specific. Its origin failure + // must not poison the digest-wide cache for another requester. + recordOriginFailure(neg, d, err, lg, "origin pull failed", registry, repository, registryauth.Authorization(ctx) == "", + slog.String("pull_mode", "detached"), + slog.String("deadline_owner", originPullDeadlineOwner(ctx, err)), + slog.Duration("elapsed", time.Since(pullStartedAt)), + slog.Int64("expected_size", expectedSize), + slog.Int64("resume_offset", resumeOffset), + slog.Int64("written", 0), + ) + + return + } + + defer func() { _ = rc.Close() }() //nolint:errcheck // best-effort close + written, err := copyWithOriginProgressTimeout(ctx, cancel, w, rc, progressTimeout) if err != nil { releaseLeaseOnFailure() diff --git a/cmd/gantry/origin_pull_test.go b/cmd/gantry/origin_pull_test.go index 6b6b7fee9..38d84294f 100644 --- a/cmd/gantry/origin_pull_test.go +++ b/cmd/gantry/origin_pull_test.go @@ -4,10 +4,13 @@ package main import ( + "bytes" "context" "errors" + "fmt" "io" "log/slog" + "slices" "strings" "sync/atomic" "testing" @@ -184,6 +187,104 @@ func (originTimeoutError) Error() string { return "timed out" } func (originTimeoutError) Timeout() bool { return true } func (originTimeoutError) Temporary() bool { return true } +type offsetRecordingOrigin struct { + body []byte + offsets []int64 + rejectRange bool + fail bool +} + +func (o *offsetRecordingOrigin) Pull(_ context.Context, ref ifaces.OriginRef) (io.ReadCloser, int64, error) { + o.offsets = append(o.offsets, ref.Offset) + + if o.fail { + return nil, 0, errors.New("transient origin failure") + } + + if ref.Offset > 0 && o.rejectRange { + return nil, 0, &ifaces.OriginError{ + Ref: ref, + Class: ifaces.FailureTransient, + Err: &ifaces.ErrRangeUnsupported{Offset: ref.Offset, Reason: "status 200 OK"}, + } + } + + return io.NopCloser(bytes.NewReader(o.body[ref.Offset:])), int64(len(o.body)), nil +} + +func (o *offsetRecordingOrigin) Head(context.Context, ifaces.OriginRef) (int64, string, error) { + return int64(len(o.body)), "application/octet-stream", nil +} + +type resumableTestCache struct { + expected digest.Digest + partial []byte + committed []byte + aborts int +} + +func (c *resumableTestCache) Has(context.Context, digest.Digest) (bool, error) { + return c.committed != nil, nil +} + +func (c *resumableTestCache) Open(_ context.Context, d digest.Digest) (io.ReadCloser, int64, error) { + if c.committed == nil { + return nil, 0, &ifaces.ErrNotFound{Digest: d} + } + + return io.NopCloser(bytes.NewReader(c.committed)), int64(len(c.committed)), nil +} + +func (c *resumableTestCache) Writer(context.Context, digest.Digest) (ifaces.ContentWriter, error) { + return &resumableTestWriter{cache: c}, nil +} + +func (c *resumableTestCache) ResumeWriter(context.Context, digest.Digest) (ifaces.ContentWriter, int64, error) { + w := &resumableTestWriter{cache: c} + _, _ = w.body.Write(c.partial) + + return w, int64(len(c.partial)), nil +} + +type resumableTestWriter struct { + cache *resumableTestCache + body bytes.Buffer + finalized bool +} + +func (w *resumableTestWriter) Write(p []byte) (int, error) { return w.body.Write(p) } + +func (w *resumableTestWriter) Commit(context.Context) error { + if got := trackerDigestOf(w.body.Bytes()); got != w.cache.expected { + return fmt.Errorf("digest = %s; want %s", got, w.cache.expected) + } + + w.cache.committed = append([]byte(nil), w.body.Bytes()...) + w.cache.partial = nil + w.finalized = true + + return nil +} + +func (w *resumableTestWriter) Abort(context.Context) error { + if !w.finalized { + w.cache.partial = nil + w.cache.aborts++ + w.finalized = true + } + + return nil +} + +func (w *resumableTestWriter) Preserve() error { + if !w.finalized { + w.cache.partial = append([]byte(nil), w.body.Bytes()...) + w.finalized = true + } + + return nil +} + func (r *pacedReader) Read(p []byte) (int, error) { if r.remaining == 0 { return 0, io.EOF @@ -288,6 +389,88 @@ func TestOriginPullDeadlineOwner(t *testing.T) { } } +func TestRunOriginPullResumesPartialIngest(t *testing.T) { + body := []byte("partial-then-completed") + d := trackerDigestOf(body) + originPuller := &offsetRecordingOrigin{body: body} + cache := &resumableTestCache{expected: d, partial: append([]byte(nil), body[:8]...)} + h, _, _ := inflight.New(inflight.DefaultStalls(), nil).Start(d, ifaces.KindBlob, 0) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + + var successes int + + runOriginPull(context.Background(), originPuller, cache, nil, logger, h, "registry.example.com", "library/test", d, ifaces.KindBlob, 0, + func(context.Context, digest.Digest) bool { return true }, + func(string, int64) { successes++ }, + func(string, string) {}, + leaseMetricHooks{}, + ) + + if !slices.Equal(originPuller.offsets, []int64{8}) { + t.Fatalf("origin offsets = %v; want [8]", originPuller.offsets) + } + + if !bytes.Equal(cache.committed, body) { + t.Fatalf("committed body = %q; want %q", cache.committed, body) + } + + if successes != 1 { + t.Fatalf("successes = %d; want 1", successes) + } +} + +func TestRunOriginPullRestartsWhenRangeUnsupported(t *testing.T) { + body := []byte("partial-restarted-from-zero") + d := trackerDigestOf(body) + originPuller := &offsetRecordingOrigin{body: body, rejectRange: true} + cache := &resumableTestCache{expected: d, partial: append([]byte(nil), body[:8]...)} + h, _, _ := inflight.New(inflight.DefaultStalls(), nil).Start(d, ifaces.KindBlob, 0) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + + runOriginPull(context.Background(), originPuller, cache, nil, logger, h, "registry.example.com", "library/test", d, ifaces.KindBlob, 0, + func(context.Context, digest.Digest) bool { return true }, + func(string, int64) {}, + func(string, string) {}, + leaseMetricHooks{}, + ) + + if !slices.Equal(originPuller.offsets, []int64{8, 0}) { + t.Fatalf("origin offsets = %v; want [8 0]", originPuller.offsets) + } + + if cache.aborts != 1 { + t.Fatalf("aborts = %d; want 1", cache.aborts) + } + + if !bytes.Equal(cache.committed, body) { + t.Fatalf("committed body = %q; want %q", cache.committed, body) + } +} + +func TestRunOriginPullPreservesPartialOnTransientFailure(t *testing.T) { + body := []byte("partial-preserved") + d := trackerDigestOf(body) + originPuller := &offsetRecordingOrigin{body: body, fail: true} + cache := &resumableTestCache{expected: d, partial: append([]byte(nil), body[:8]...)} + h, _, _ := inflight.New(inflight.DefaultStalls(), nil).Start(d, ifaces.KindBlob, 0) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + + runOriginPull(context.Background(), originPuller, cache, nil, logger, h, "registry.example.com", "library/test", d, ifaces.KindBlob, 0, + func(context.Context, digest.Digest) bool { return true }, + func(string, int64) {}, + func(string, string) {}, + leaseMetricHooks{}, + ) + + if !bytes.Equal(cache.partial, body[:8]) { + t.Fatalf("partial body = %q; want %q", cache.partial, body[:8]) + } + + if cache.aborts != 0 { + t.Fatalf("aborts = %d; want 0", cache.aborts) + } +} + func TestRunOriginPull_ReopenFailurePreventsAdvertiseAndSuccess(t *testing.T) { body := []byte("committed-but-not-reopenable") d := trackerDigestOf(body) diff --git a/designs/gantry-128gib-single-layer-pulls.md b/designs/gantry-128gib-single-layer-pulls.md index 9a1cf96d5..0884a2c97 100644 --- a/designs/gantry-128gib-single-layer-pulls.md +++ b/designs/gantry-128gib-single-layer-pulls.md @@ -111,7 +111,7 @@ resumable. | ID | Issue | Evidence | Owning PR | Status | |---|---|---|---|---| | L4 | Gantry origin pulls do not send `Range` and cannot request an origin body from an offset. | `internal/gantry/origin/origin.go:25`, `internal/gantry/origin/origin.go:477-542` | PR 4 | Addressed | -| L5 | A failed background ingest is aborted. A later writer with a stale nonzero offset is also aborted because callers restart at byte zero. | `internal/gantry/containerdstore/store.go:328-383`, `internal/gantry/containerdstore/store.go:502-519` | PR 5 | Open | +| L5 | A failed background ingest is aborted. A later writer with a stale nonzero offset is also aborted because callers restart at byte zero. | `internal/gantry/containerdstore/store.go:328-383`, `internal/gantry/containerdstore/store.go:502-519` | PR 5 | Addressed | | L6 | Gantry does not handle the local containerd client's inbound `Range` header on the origin path. | `internal/gantry/mirror/mirror.go:685-930` | PR 4 | Addressed | L4-L6 do not cause the fixed five- or 30-minute failures. They determine @@ -310,14 +310,14 @@ resume failures easier to distinguish during review. **Purpose:** Reuse bytes already staged in containerd when a detached chair pull is interrupted. -- [ ] Expose a verified writer offset through the local content-store +- [x] Expose a verified writer offset through an optional local content-store abstraction. -- [ ] Reopen the origin at that offset using PR 4's origin range support. -- [ ] Continue the existing containerd ingest and rely on commit-time digest +- [x] Reopen the origin at that offset using PR 4's origin range support. +- [x] Continue the existing containerd ingest and rely on commit-time digest verification over the complete layer. -- [ ] Abort and restart from byte zero when the origin cannot honor the range, +- [x] Abort and restart from byte zero when the origin cannot honor the range, and log that decision. -- [ ] Test process-local retry, stale partial state, unsupported ranges, digest +- [x] Test process-local retry, stale partial state, unsupported ranges, digest mismatch, and delegated authorization. **Addresses:** L5. diff --git a/internal/gantry/containerdstore/store.go b/internal/gantry/containerdstore/store.go index b44af95d8..d3d06154c 100644 --- a/internal/gantry/containerdstore/store.go +++ b/internal/gantry/containerdstore/store.go @@ -338,6 +338,19 @@ func (s *Store) Descriptor(ctx context.Context, d gdigest.Digest) (ocispec.Descr // ErrAlreadyExists wrapped - callers who want "treat-as-committed" // semantics should call Has first. func (s *Store) Writer(ctx context.Context, d gdigest.Digest) (ifaces.ContentWriter, error) { + w, _, err := s.openWriter(ctx, d, false) + + return w, err +} + +// ResumeWriter returns the digest's existing ingest writer and its stored +// offset. Unlike Writer, it preserves partial data so callers that can issue a +// validated origin Range request may continue the ingest. +func (s *Store) ResumeWriter(ctx context.Context, d gdigest.Digest) (ifaces.ContentWriter, int64, error) { + return s.openWriter(ctx, d, true) +} + +func (s *Store) openWriter(ctx context.Context, d gdigest.Digest, resume bool) (ifaces.ContentWriter, int64, error) { ref := s.refPrefix + d.String() expected := godigest.Digest(d.String()) desc := ocispec.Descriptor{Digest: expected} @@ -350,26 +363,35 @@ func (s *Store) Writer(ctx context.Context, d gdigest.Digest) (ifaces.ContentWri if errors.Is(err, cerrdefs.ErrAlreadyExists) { ok, hasErr := s.Has(ctx, d) if hasErr != nil { - return nil, hasErr + return nil, 0, hasErr } if ok { - return alreadyCommittedContentWriter{}, nil + return alreadyCommittedContentWriter{}, 0, nil } } - return nil, &ifaces.ErrUnavailable{Op: "Writer", Cause: err} + return nil, 0, &ifaces.ErrUnavailable{Op: "Writer", Cause: err} } // If a previous crashed/failed pull left a partial ingest, the - // writer may have a non-zero offset. Callers always stream from - // byte 0, so appending to stale data would produce a corrupt - // commit. Abort and re-acquire a clean writer. + // writer may have a non-zero offset. Ordinary callers stream from byte 0, + // so Writer aborts it. ResumeWriter returns the offset to a caller that can + // request the matching origin suffix. if st, stErr := w.Status(); stErr == nil && st.Offset > 0 { + if resume { + return &contentWriter{ + inner: w, + expected: expected, + ref: ref, + store: s, + }, st.Offset, nil + } + _ = w.Close() //nolint:errcheck // closing stale writer if abErr := s.cs.Abort(s.withNS(ctx), ref); abErr != nil && !errors.Is(abErr, cerrdefs.ErrNotFound) { - return nil, &ifaces.ErrUnavailable{Op: "Writer(abort-stale)", Cause: abErr} + return nil, 0, &ifaces.ErrUnavailable{Op: "Writer(abort-stale)", Cause: abErr} } w, err = s.cs.Writer(s.withNS(ctx), @@ -377,7 +399,7 @@ func (s *Store) Writer(ctx context.Context, d gdigest.Digest) (ifaces.ContentWri content.WithDescriptor(desc), ) if err != nil { - return nil, &ifaces.ErrUnavailable{Op: "Writer(retry)", Cause: err} + return nil, 0, &ifaces.ErrUnavailable{Op: "Writer(retry)", Cause: err} } } @@ -386,7 +408,7 @@ func (s *Store) Writer(ctx context.Context, d gdigest.Digest) (ifaces.ContentWri expected: expected, ref: ref, store: s, - }, nil + }, 0, nil } // Inventory enumerates every sha256 digest currently present and @@ -519,6 +541,18 @@ func (w *contentWriter) Abort(ctx context.Context) error { return nil } +// Preserve closes the active writer without aborting its ingest so a later +// ResumeWriter call can continue from the stored offset. +func (w *contentWriter) Preserve() error { + if w.committedOrAborted { + return nil + } + + w.committedOrAborted = true + + return w.inner.Close() +} + // readerAtCloser bundles a SectionReader (Read+Seek) with the // underlying content.ReaderAt's Close so the transfer endpoint can // type-assert io.ReadSeeker for Range serving while still releasing diff --git a/internal/gantry/containerdstore/store_test.go b/internal/gantry/containerdstore/store_test.go index d96866c3d..17eae8435 100644 --- a/internal/gantry/containerdstore/store_test.go +++ b/internal/gantry/containerdstore/store_test.go @@ -164,6 +164,10 @@ func (f *fakeStore) Writer(_ context.Context, opts ...content.WriterOpt) (conten f.mu.Lock() defer f.mu.Unlock() + if existing := f.writers[wopts.Ref]; existing != nil { + return existing, nil + } + w := &fakeWriter{ref: wopts.Ref, expected: wopts.Desc.Digest, store: f} f.writers[wopts.Ref] = w @@ -182,7 +186,7 @@ func (w *fakeWriter) Write(p []byte) (int, error) { return w.buf.Write(p) } func (w *fakeWriter) Close() error { w.closed = true; return nil } func (w *fakeWriter) Digest() godigest.Digest { return godigest.FromBytes(w.buf.Bytes()) } func (w *fakeWriter) Status() (content.Status, error) { - return content.Status{Ref: w.ref}, nil + return content.Status{Ref: w.ref, Offset: int64(w.buf.Len())}, nil } func (w *fakeWriter) Truncate(_ int64) error { return nil } @@ -345,6 +349,83 @@ func TestStore_WriterCommit(t *testing.T) { } } +func TestStore_ResumeWriterPreservesPartialIngest(t *testing.T) { + cs := newFake() + s := New(cs) + payload := []byte("partial-then-resumed") + d := mustDigest(t, payload) + + partial, _, err := s.ResumeWriter(context.Background(), d) + if err != nil { + t.Fatalf("ResumeWriter initial: %v", err) + } + + const offset = 7 + if _, err := partial.Write(payload[:offset]); err != nil { + t.Fatalf("Write partial: %v", err) + } + + preservable, ok := partial.(interface{ Preserve() error }) + if !ok { + t.Fatalf("writer type = %T; want Preserve", partial) + } + + if err := preservable.Preserve(); err != nil { + t.Fatalf("Preserve: %v", err) + } + + resumed, gotOffset, err := s.ResumeWriter(context.Background(), d) + if err != nil { + t.Fatalf("ResumeWriter existing: %v", err) + } + + if gotOffset != offset { + t.Fatalf("offset = %d; want %d", gotOffset, offset) + } + + if _, err := resumed.Write(payload[offset:]); err != nil { + t.Fatalf("Write suffix: %v", err) + } + + if err := resumed.Commit(context.Background()); err != nil { + t.Fatalf("Commit: %v", err) + } + + has, err := s.Has(context.Background(), d) + if err != nil || !has { + t.Fatalf("Has after resume = %v, %v; want true, nil", has, err) + } +} + +func TestStore_WriterStillDiscardsPartialIngest(t *testing.T) { + cs := newFake() + s := New(cs) + payload := []byte("partial-discarded") + d := mustDigest(t, payload) + + partial, _, err := s.ResumeWriter(context.Background(), d) + if err != nil { + t.Fatalf("ResumeWriter: %v", err) + } + + if _, err := partial.Write([]byte("stale")); err != nil { + t.Fatalf("Write partial: %v", err) + } + + w, err := s.Writer(context.Background(), d) + if err != nil { + t.Fatalf("Writer: %v", err) + } + + if _, err := w.Write(payload); err != nil { + t.Fatalf("Write full: %v", err) + } + + if err := w.Commit(context.Background()); err != nil { + t.Fatalf("Commit: %v", err) + } +} + func TestStore_WriterAlreadyExistsRechecksPresence(t *testing.T) { cs := newFake() payload := []byte("already-committed-by-containerd") diff --git a/internal/gantry/ifaces/ifaces.go b/internal/gantry/ifaces/ifaces.go index bcd6c2c5c..8d797b5ed 100644 --- a/internal/gantry/ifaces/ifaces.go +++ b/internal/gantry/ifaces/ifaces.go @@ -487,6 +487,18 @@ func (e *ErrUnavailable) Error() string { func (e *ErrUnavailable) Unwrap() error { return e.Cause } +// ErrRangeUnsupported reports that an origin did not honor a requested blob +// offset. Callers may safely restart from byte zero; other transport and auth +// errors must preserve the partial ingest for a later retry. +type ErrRangeUnsupported struct { + Offset int64 + Reason string +} + +func (e *ErrRangeUnsupported) Error() string { + return fmt.Sprintf("origin range offset %d unsupported: %s", e.Offset, e.Reason) +} + // ErrPeerHTTPStatus is returned by PeerDialer implementations when a peer // transfer endpoint responds with an unexpected HTTP status. Callers can use // StatusCode to classify failures (auth/config, server error, protocol error) diff --git a/internal/gantry/origin/origin.go b/internal/gantry/origin/origin.go index cd3b5871a..8205b8bb8 100644 --- a/internal/gantry/origin/origin.go +++ b/internal/gantry/origin/origin.go @@ -553,7 +553,7 @@ func (r *registry) pull(ctx context.Context, ref ifaces.OriginRef) (io.ReadClose if resp.StatusCode != http.StatusPartialContent { defer func() { _ = resp.Body.Close() }() //nolint:errcheck // best-effort body close - return nil, 0, &ifaces.OriginError{Ref: ref, Class: ifaces.FailureTransient, Err: fmt.Errorf("origin ignored range offset %d: status %s", ref.Offset, resp.Status)} + return nil, 0, &ifaces.OriginError{Ref: ref, Class: ifaces.FailureTransient, Err: &ifaces.ErrRangeUnsupported{Offset: ref.Offset, Reason: "status " + resp.Status}} } start, end, size, ok := parseOriginContentRange(resp.Header.Get("Content-Range")) @@ -561,7 +561,7 @@ func (r *registry) pull(ctx context.Context, ref ifaces.OriginRef) (io.ReadClose (resp.ContentLength >= 0 && resp.ContentLength != end-start+1) { defer func() { _ = resp.Body.Close() }() //nolint:errcheck // best-effort body close - return nil, 0, &ifaces.OriginError{Ref: ref, Class: ifaces.FailureTransient, Err: fmt.Errorf("origin returned invalid Content-Range %q for offset %d", resp.Header.Get("Content-Range"), ref.Offset)} + return nil, 0, &ifaces.OriginError{Ref: ref, Class: ifaces.FailureTransient, Err: &ifaces.ErrRangeUnsupported{Offset: ref.Offset, Reason: fmt.Sprintf("invalid Content-Range %q", resp.Header.Get("Content-Range"))}} } return resp.Body, size, nil diff --git a/internal/gantry/origin/origin_test.go b/internal/gantry/origin/origin_test.go index 57cbad550..bea6e5433 100644 --- a/internal/gantry/origin/origin_test.go +++ b/internal/gantry/origin/origin_test.go @@ -177,7 +177,7 @@ func TestPullBlobRangeRejectsIgnoredOrInvalidResponse(t *testing.T) { contentRange string want string }{ - {name: "ignored", status: http.StatusOK, want: "ignored range"}, + {name: "ignored", status: http.StatusOK, want: "unsupported"}, {name: "invalid", status: http.StatusPartialContent, contentRange: "bytes 0-5/10", want: "invalid Content-Range"}, }