From 034b2cee9a05e4bd1e04946b36b7d7bd4da38342 Mon Sep 17 00:00:00 2001 From: PhantomPhoton <120879831+PhantomPhoton@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:14:00 +0000 Subject: [PATCH 1/4] feat(storage): expose efficient seek for local files Open file-backed artifacts with os.Open so local readers support efficient seeking. Leave cloud storage readers unchanged. Keep the legacy sidecar mapping test focused on fileblob's reader: Blob.Open now reads local files directly and no longer consults the sidecar. Continue checking sidecar cleanup and post-store reads. --- internal/storage/blob.go | 11 +++++++++++ internal/storage/blob_test.go | 15 ++++++++++++--- internal/storage/storage.go | 3 +++ 3 files changed, 26 insertions(+), 3 deletions(-) diff --git a/internal/storage/blob.go b/internal/storage/blob.go index 48574aa6..4cadf8af 100644 --- a/internal/storage/blob.go +++ b/internal/storage/blob.go @@ -228,6 +228,17 @@ func (b *Blob) Store(ctx context.Context, path string, r io.Reader) (int64, stri } func (b *Blob) Open(ctx context.Context, path string) (io.ReadCloser, error) { + if localPath := b.localPath(path); localPath != "" { + r, err := os.Open(localPath) + if err != nil { + if os.IsNotExist(err) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("opening local reader: %w", err) + } + return r, nil + } + r, err := b.bucket.NewReader(ctx, path, nil) if err != nil { if isNotExist(err) { diff --git a/internal/storage/blob_test.go b/internal/storage/blob_test.go index 57e85f59..efc1b765 100644 --- a/internal/storage/blob_test.go +++ b/internal/storage/blob_test.go @@ -67,13 +67,20 @@ func TestBlobOpen(t *testing.T) { t.Fatalf("Open failed: %v", err) } defer func() { _ = r.Close() }() + seeker, ok := r.(io.Seeker) + if !ok { + t.Fatal("local file reader does not implement io.Seeker") + } + if _, err := seeker.Seek(9, io.SeekStart); err != nil { + t.Fatalf("Seek failed: %v", err) + } data, err := io.ReadAll(r) if err != nil { t.Fatalf("ReadAll failed: %v", err) } - if string(data) != content { - t.Errorf("content = %q, want %q", string(data), content) + if string(data) != "content" { + t.Errorf("content after seek = %q, want %q", string(data), "content") } } @@ -529,7 +536,9 @@ func assertStoreClearsSidecar(t *testing.T, key string) { } b := openFileBlob(t, dir) - if _, err := b.Open(ctx, key); err == nil { + reader, err := b.bucket.NewReader(ctx, key, nil) + if err == nil { + _ = reader.Close() t.Fatal("corrupt sidecar did not fail the read, so it is not the file fileblob reads for this key") } diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 0f64ed77..3b27ef4b 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -46,6 +46,9 @@ type Storage interface { // Open returns a reader for the content at path. // The caller must close the reader when done. + // The reader may also implement io.Seeker when seeking is efficient for the + // storage backend. Backends must not expose io.Seeker when seeking requires + // reading and discarding the bytes before the requested position. // Returns ErrNotFound if the path does not exist. Open(ctx context.Context, path string) (io.ReadCloser, error) From a475ab42ab55f36777574818ca4a7b0227dcaf57 Mon Sep 17 00:00:00 2001 From: PhantomPhoton <120879831+PhantomPhoton@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:20:54 +0000 Subject: [PATCH 2/4] feat(handler): preserve efficient seek through integrity checks Keep the cache result's single Reader field and preserve io.Seeker when the underlying storage reader supports efficient seeking. Sequential reads continue through whole-object integrity verification. A successful seek switches subsequent reads to the source for range responses, where full-object verification is not possible. Test both modes and verify the source is closed once. --- internal/handler/integrity.go | 45 +++++++++++++++++++- internal/handler/integrity_test.go | 67 ++++++++++++++++++++++++++++++ 2 files changed, 111 insertions(+), 1 deletion(-) diff --git a/internal/handler/integrity.go b/internal/handler/integrity.go index 97f41d99..437252f7 100644 --- a/internal/handler/integrity.go +++ b/internal/handler/integrity.go @@ -40,7 +40,22 @@ func newIntegrityChecks(contentHash, native string) (integrityChecks, error) { } func (c integrityChecks) wrap(source io.ReadCloser, onMismatch func(string)) (io.ReadCloser, error) { - return c.newVerifyingReader(source, onMismatch, false) + if len(c.algorithms) == 0 { + return source, nil + } + verified, err := c.newVerifyingReader(source, onMismatch, false) + if err != nil { + return nil, err + } + seeker, ok := source.(io.Seeker) + if !ok { + return verified, nil + } + return &seekableVerifyingReader{ + source: source, + verified: verified, + seeker: seeker, + }, nil } // wrapFailOnMismatch is wrap for bytes that have not been checked anywhere @@ -80,6 +95,34 @@ type verifyingReader struct { mismatched bool } +// seekableVerifyingReader verifies ordinary reads until a successful seek, +// after which reads bypass whole-object verification for range responses. +type seekableVerifyingReader struct { + source io.ReadCloser + verified io.ReadCloser + seeker io.Seeker + bypass bool +} + +func (r *seekableVerifyingReader) Read(p []byte) (int, error) { + if r.bypass { + return r.source.Read(p) + } + return r.verified.Read(p) +} + +func (r *seekableVerifyingReader) Seek(offset int64, whence int) (int64, error) { + position, err := r.seeker.Seek(offset, whence) + if err == nil { + r.bypass = true + } + return position, err +} + +func (r *seekableVerifyingReader) Close() error { + return r.verified.Close() +} + func (r *verifyingReader) Read(p []byte) (int, error) { n, err := r.reader.Read(p) if err == io.EOF { diff --git a/internal/handler/integrity_test.go b/internal/handler/integrity_test.go index 95992c03..9aa99589 100644 --- a/internal/handler/integrity_test.go +++ b/internal/handler/integrity_test.go @@ -1,6 +1,7 @@ package handler import ( + "bytes" "crypto/sha256" "crypto/sha512" "encoding/base64" @@ -126,6 +127,72 @@ func TestVerifyingReader(t *testing.T) { } } +func TestVerifyingReaderPreservesSeekCapability(t *testing.T) { + const data = "hello world" + t.Run("ordinary reads stay verified", func(t *testing.T) { + source := &seekableCloseTrackingReader{Reader: bytes.NewReader([]byte(data))} + var calls int + reader := wrapIntegrityReader(t, source, sha256Hex("different"), "", func(string) { calls++ }) + + got, err := io.ReadAll(reader) + if err != nil { + t.Fatalf("ReadAll: %v", err) + } + if string(got) != data { + t.Errorf("data = %q, want %q", got, data) + } + if calls != 1 { + t.Errorf("onMismatch called %d times, want 1", calls) + } + if err := reader.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + if source.closeCount != 1 { + t.Errorf("source closed %d times, want 1", source.closeCount) + } + }) + + t.Run("seek bypasses whole-object verification", func(t *testing.T) { + source := &seekableCloseTrackingReader{Reader: bytes.NewReader([]byte(data))} + var calls int + reader := wrapIntegrityReader(t, source, sha256Hex("different"), "", func(string) { calls++ }) + seeker, ok := reader.(io.Seeker) + if !ok { + t.Fatal("seekable source did not retain io.Seeker") + } + if _, err := seeker.Seek(6, io.SeekStart); err != nil { + t.Fatalf("Seek: %v", err) + } + + got, err := io.ReadAll(reader) + if err != nil { + t.Fatalf("ReadAll: %v", err) + } + if string(got) != "world" { + t.Errorf("data after seek = %q, want %q", got, "world") + } + if calls != 0 { + t.Errorf("onMismatch called %d times for a partial read, want 0", calls) + } + if err := reader.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + if source.closeCount != 1 { + t.Errorf("source closed %d times, want 1", source.closeCount) + } + }) +} + +type seekableCloseTrackingReader struct { + *bytes.Reader + closeCount int +} + +func (r *seekableCloseTrackingReader) Close() error { + r.closeCount++ + return nil +} + func TestVerifyingReaderUsesStrongestNativeAlgorithm(t *testing.T) { const data = "artifact" tests := []struct { From e85fe143cf368348b01e92494bd0ead1556e598b Mon Sep 17 00:00:00 2001 From: PhantomPhoton <120879831+PhantomPhoton@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:31:13 +0000 Subject: [PATCH 3/4] feat(handler): add request-aware byte range responses Add a request-aware artifact serving helper for bounded, open-ended, and suffix byte ranges, including If-Range handling, 206 responses, and 416 responses for valid unsatisfiable ranges. Ignore malformed and multi-range requests by serving the full artifact. Advertise byte-range support only for seekable readers, preserving HEAD, redirect, and non-seekable behavior. Clarify that net/http recovers ErrAbortHandler per request and keeps the server running. --- internal/handler/handler.go | 176 ++++++++++++++++++++++++++++--- internal/handler/handler_test.go | 155 +++++++++++++++++++++++++++ 2 files changed, 318 insertions(+), 13 deletions(-) diff --git a/internal/handler/handler.go b/internal/handler/handler.go index b04f9419..df6a63c2 100644 --- a/internal/handler/handler.go +++ b/internal/handler/handler.go @@ -756,7 +756,21 @@ func ServeArtifact(w http.ResponseWriter, result *CacheResult) { serveArtifact(w, http.MethodGet, result) } +// ServeArtifactRequest writes a CacheResult to an HTTP response using request +// headers such as Range and If-Range. +func ServeArtifactRequest(w http.ResponseWriter, r *http.Request, result *CacheResult) { + if r == nil { + ServeArtifact(w, result) + return + } + serveArtifactResponse(w, r.Method, r, result) +} + func serveArtifact(w http.ResponseWriter, method string, result *CacheResult) { + serveArtifactResponse(w, method, nil, result) +} + +func serveArtifactResponse(w http.ResponseWriter, method string, request *http.Request, result *CacheResult) { contentHash := "" if result.Artifact.Digest != "" { contentHash = result.Artifact.Digest.Encoded() @@ -777,26 +791,162 @@ func serveArtifact(w http.ResponseWriter, method string, result *CacheResult) { if result.Artifact.MediaType != "" { w.Header().Set(headerContentType, result.Artifact.MediaType) } - if result.Artifact.Size > 0 || (method == http.MethodHead && result.Artifact.Size == 0) { - w.Header().Set(headerContentLength, strconv.FormatInt(result.Artifact.Size, 10)) - } if contentHash != "" { w.Header().Set(headerETag, `"`+contentHash+`"`) } + var seeker io.Seeker + if result.Reader != nil { + seeker, _ = result.Reader.(io.Seeker) + } + if seeker != nil && result.Artifact.Size >= 0 { + w.Header().Set("Accept-Ranges", "bytes") + if request != nil && serveArtifactRange(w, request, result, seeker) { + return + } + } + + if result.Artifact.Size > 0 || (method == http.MethodHead && result.Artifact.Size == 0) { + w.Header().Set(headerContentLength, strconv.FormatInt(result.Artifact.Size, 10)) + } + w.WriteHeader(http.StatusOK) if method != http.MethodHead && result.Reader != nil { - buffer := artifactCopyBufferPool.Get().(*[]byte) - defer artifactCopyBufferPool.Put(buffer) - // Hide optional ReaderFrom methods so io.CopyBuffer uses the pooled buffer. - written, err := io.CopyBuffer(struct{ io.Writer }{w}, result.Reader, *buffer) - if err != nil || (result.Artifact.Size > 0 && written != result.Artifact.Size) { - // Headers are already committed, so an error status is no longer - // possible. Aborting leaves the response unterminated and the - // client discards it instead of keeping a truncated or unverified - // artifact. - panic(http.ErrAbortHandler) + copyArtifactBody(w, result.Reader, result.Artifact.Size, false) + } +} + +type parsedByteRange struct { + start int64 + end int64 +} + +func serveArtifactRange(w http.ResponseWriter, request *http.Request, result *CacheResult, seeker io.Seeker) bool { + if request.Method != http.MethodGet { + return false + } + ranges := request.Header.Values("Range") + if len(ranges) != 1 || !ifRangeMatches(request, w.Header().Get(headerETag)) { + return false + } + + byteRange, valid, satisfiable := parseByteRange(ranges[0], result.Artifact.Size) + if !valid { + return false + } + if !satisfiable { + w.Header().Set("Content-Range", fmt.Sprintf("bytes */%d", result.Artifact.Size)) + w.Header().Set(headerContentLength, "0") + w.WriteHeader(http.StatusRequestedRangeNotSatisfiable) + return true + } + + if _, err := seeker.Seek(byteRange.start, io.SeekStart); err != nil { + w.Header().Del(headerContentType) + w.Header().Del(headerContentLength) + w.Header().Del(headerETag) + w.Header().Del("Accept-Ranges") + http.Error(w, "failed to seek cached artifact", http.StatusInternalServerError) + return true + } + + length := byteRange.end - byteRange.start + 1 + w.Header().Set("Content-Range", fmt.Sprintf("bytes %d-%d/%d", byteRange.start, byteRange.end, result.Artifact.Size)) + w.Header().Set(headerContentLength, strconv.FormatInt(length, 10)) + w.WriteHeader(http.StatusPartialContent) + copyArtifactBody(w, result.Reader, length, true) + return true +} + +func ifRangeMatches(request *http.Request, etag string) bool { + values := request.Header.Values("If-Range") + if len(values) == 0 { + return true + } + return len(values) == 1 && etag != "" && strings.TrimSpace(values[0]) == etag +} + +func parseByteRange(value string, size int64) (parsedByteRange, bool, bool) { + unit, spec, ok := strings.Cut(strings.TrimSpace(value), "=") + if !ok || !strings.EqualFold(strings.TrimSpace(unit), "bytes") { + return parsedByteRange{}, false, false + } + spec = strings.TrimSpace(spec) + if spec == "" || strings.Contains(spec, ",") { + return parsedByteRange{}, false, false + } + first, last, ok := strings.Cut(spec, "-") + if !ok { + return parsedByteRange{}, false, false + } + + if first == "" { + suffixLength, ok := parseRangeNumber(last) + if !ok { + return parsedByteRange{}, false, false } + if size == 0 || suffixLength == 0 { + return parsedByteRange{}, true, false + } + if suffixLength >= size { + return parsedByteRange{start: 0, end: size - 1}, true, true + } + return parsedByteRange{start: size - suffixLength, end: size - 1}, true, true + } + + start, ok := parseRangeNumber(first) + if !ok { + return parsedByteRange{}, false, false + } + if last == "" { + if size == 0 || start >= size { + return parsedByteRange{}, true, false + } + return parsedByteRange{start: start, end: size - 1}, true, true + } + end, ok := parseRangeNumber(last) + if !ok || start > end { + return parsedByteRange{}, false, false + } + if size == 0 || start >= size { + return parsedByteRange{}, true, false + } + if end >= size { + end = size - 1 + } + return parsedByteRange{start: start, end: end}, true, true +} + +func parseRangeNumber(value string) (int64, bool) { + if value == "" { + return 0, false + } + for _, digit := range value { + if digit < '0' || digit > '9' { + return 0, false + } + } + number, err := strconv.ParseInt(value, 10, 64) + return number, err == nil +} + +func copyArtifactBody(w http.ResponseWriter, reader io.Reader, expectedSize int64, bounded bool) { + buffer := artifactCopyBufferPool.Get().(*[]byte) + defer artifactCopyBufferPool.Put(buffer) + + source := reader + if bounded { + source = io.LimitReader(reader, expectedSize) + } + // Hide optional ReaderFrom methods so io.CopyBuffer uses the pooled buffer. + written, err := io.CopyBuffer(struct{ io.Writer }{w}, source, *buffer) + if err != nil || (bounded && written != expectedSize) || (!bounded && expectedSize > 0 && written != expectedSize) { + // Headers are already committed, so an error status is no longer + // possible. Aborting leaves the response unterminated so clients discard + // a truncated or unverified artifact. + // net/http recovers ErrAbortHandler for this request and keeps the server + // running; on HTTP/1.x it may close this connection as well. + panic(http.ErrAbortHandler) } } diff --git a/internal/handler/handler_test.go b/internal/handler/handler_test.go index c2714c30..cecf59d2 100644 --- a/internal/handler/handler_test.go +++ b/internal/handler/handler_test.go @@ -737,6 +737,161 @@ func TestServeArtifact_Stream(t *testing.T) { } } +func TestServeArtifactRequestRanges(t *testing.T) { + const payload = "hello world" + tests := []struct { + name string + rangeHeader string + secondRange string + wantStatus int + wantRange string + wantLength string + wantBody string + }{ + {name: "bounded", rangeHeader: "bytes=1-3", wantStatus: http.StatusPartialContent, wantRange: "bytes 1-3/11", wantLength: "3", wantBody: "ell"}, + {name: "open ended", rangeHeader: "bytes=6-", wantStatus: http.StatusPartialContent, wantRange: "bytes 6-10/11", wantLength: "5", wantBody: "world"}, + {name: "suffix", rangeHeader: "bytes=-4", wantStatus: http.StatusPartialContent, wantRange: "bytes 7-10/11", wantLength: "4", wantBody: "orld"}, + {name: "suffix larger than artifact", rangeHeader: "bytes=-99", wantStatus: http.StatusPartialContent, wantRange: "bytes 0-10/11", wantLength: "11", wantBody: payload}, + {name: "unsatisfiable", rangeHeader: "bytes=11-", wantStatus: http.StatusRequestedRangeNotSatisfiable, wantRange: "bytes */11", wantLength: "0"}, + {name: "malformed", rangeHeader: "bytes=invalid", wantStatus: http.StatusOK, wantLength: "11", wantBody: payload}, + {name: "reversed", rangeHeader: "bytes=4-2", wantStatus: http.StatusOK, wantLength: "11", wantBody: payload}, + {name: "multiple ranges", rangeHeader: "bytes=0-1,4-5", wantStatus: http.StatusOK, wantLength: "11", wantBody: payload}, + {name: "multiple range fields", rangeHeader: "bytes=0-1", secondRange: "bytes=4-5", wantStatus: http.StatusOK, wantLength: "11", wantBody: payload}, + {name: "no range", wantStatus: http.StatusOK, wantLength: "11", wantBody: payload}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + request := httptest.NewRequest(http.MethodGet, "/artifact", nil) + if test.rangeHeader != "" { + request.Header.Set("Range", test.rangeHeader) + } + if test.secondRange != "" { + request.Header.Add("Range", test.secondRange) + } + w := httptest.NewRecorder() + ServeArtifactRequest(w, request, newRangeTestResult(payload)) + + if w.Code != test.wantStatus { + t.Errorf("status = %d, want %d", w.Code, test.wantStatus) + } + if got := w.Header().Get("Content-Range"); got != test.wantRange { + t.Errorf("Content-Range = %q, want %q", got, test.wantRange) + } + if got := w.Header().Get(headerContentLength); got != test.wantLength { + t.Errorf("Content-Length = %q, want %q", got, test.wantLength) + } + if got := w.Body.String(); got != test.wantBody { + t.Errorf("body = %q, want %q", got, test.wantBody) + } + if got := w.Header().Get("Accept-Ranges"); got != "bytes" { + t.Errorf("Accept-Ranges = %q, want bytes", got) + } + }) + } +} + +func TestServeArtifactRequestIfRange(t *testing.T) { + const payload = "hello world" + tests := []struct { + name string + ifRange string + wantStatus int + wantRange string + wantBody string + }{ + {name: "matching etag", ifRange: `"sha256-matching"`, wantStatus: http.StatusPartialContent, wantRange: "bytes 0-4/11", wantBody: "hello"}, + {name: "mismatching etag", ifRange: `"sha256-stale"`, wantStatus: http.StatusOK, wantBody: payload}, + {name: "weak etag", ifRange: `W/"sha256-matching"`, wantStatus: http.StatusOK, wantBody: payload}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + result := newRangeTestResult(payload) + request := httptest.NewRequest(http.MethodGet, "/artifact", nil) + request.Header.Set("Range", "bytes=0-4") + if test.name == "matching etag" { + test.ifRange = `"` + result.Artifact.Digest.Encoded() + `"` + } + request.Header.Set("If-Range", test.ifRange) + w := httptest.NewRecorder() + ServeArtifactRequest(w, request, result) + + if w.Code != test.wantStatus { + t.Errorf("status = %d, want %d", w.Code, test.wantStatus) + } + if got := w.Header().Get("Content-Range"); got != test.wantRange { + t.Errorf("Content-Range = %q, want %q", got, test.wantRange) + } + if got := w.Body.String(); got != test.wantBody { + t.Errorf("body = %q, want %q", got, test.wantBody) + } + }) + } +} + +func TestServeArtifactRequestHeadIgnoresRange(t *testing.T) { + request := httptest.NewRequest(http.MethodHead, "/artifact", nil) + request.Header.Set("Range", "bytes=0-1") + w := httptest.NewRecorder() + ServeArtifactRequest(w, request, newRangeTestResult("hello world")) + + if w.Code != http.StatusOK { + t.Errorf("status = %d, want %d", w.Code, http.StatusOK) + } + if got := w.Header().Get(headerContentLength); got != "11" { + t.Errorf("Content-Length = %q, want 11", got) + } + if got := w.Header().Get("Content-Range"); got != "" { + t.Errorf("Content-Range = %q, want empty", got) + } + if w.Body.Len() != 0 { + t.Errorf("HEAD body = %q, want empty", w.Body.String()) + } +} + +func TestServeArtifactRequestWithoutSeekCapability(t *testing.T) { + request := httptest.NewRequest(http.MethodGet, "/artifact", nil) + request.Header.Set("Range", "bytes=1-2") + w := httptest.NewRecorder() + ServeArtifactRequest(w, request, &CacheResult{ + Reader: io.NopCloser(strings.NewReader("payload")), + Artifact: testArtifact("payload", "pkg:npm/example@1.0.0", "example.tgz", "application/gzip"), + }) + + if w.Code != http.StatusOK || w.Body.String() != "payload" { + t.Errorf("response = %d %q, want 200 with full payload", w.Code, w.Body.String()) + } + if got := w.Header().Get("Accept-Ranges"); got != "" { + t.Errorf("Accept-Ranges = %q, want empty", got) + } +} + +func TestServeArtifactRequestRedirectDoesNotAdvertiseRanges(t *testing.T) { + request := httptest.NewRequest(http.MethodGet, "/artifact", nil) + request.Header.Set("Range", "bytes=1-2") + w := httptest.NewRecorder() + ServeArtifactRequest(w, request, &CacheResult{RedirectURL: "https://storage.example/artifact"}) + + if w.Code != http.StatusFound { + t.Errorf("status = %d, want %d", w.Code, http.StatusFound) + } + if got := w.Header().Get("Accept-Ranges"); got != "" { + t.Errorf("Accept-Ranges = %q, want empty", got) + } +} + +type rangeTestReadSeeker struct { + *bytes.Reader +} + +func (r *rangeTestReadSeeker) Close() error { return nil } + +func newRangeTestResult(payload string) *CacheResult { + return &CacheResult{ + Reader: &rangeTestReadSeeker{Reader: bytes.NewReader([]byte(payload))}, + Artifact: testArtifact(payload, "pkg:npm/example@1.0.0", "example.tgz", "application/gzip"), + } +} + func TestGetOrFetchArtifactFromURL_CacheHit(t *testing.T) { proxy, db, store, fetcher := setupTestProxy(t) seedPackage(t, db, store, "pypi", "requests", "2.28.0", "requests-2.28.0.tar.gz", "pypi content") From f00e0e49437006f39b33cb3941732800f683fa69 Mon Sep 17 00:00:00 2001 From: PhantomPhoton <120879831+PhantomPhoton@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:48:02 +0000 Subject: [PATCH 4/4] feat(handler): enable byte ranges across artifact downloads Forward each artifact request to the shared range-aware response helper, including cached generic release assets and OCI blobs. This lets the handler honor ranges only when its reader supports efficient seeking, without changing direct-storage redirect behavior. Add warm-cache endpoint tests that assert partial response headers and bytes without an upstream refetch. Document supported range behavior, unsupported reader fallback, and the integrity limitation of partial reads. --- docs/architecture.md | 8 ++++++++ internal/handler/apk.go | 2 +- internal/handler/cargo.go | 2 +- internal/handler/composer.go | 2 +- internal/handler/conan.go | 4 ++-- internal/handler/conda.go | 2 +- internal/handler/container.go | 4 ++-- internal/handler/container_test.go | 13 +++++++++---- internal/handler/cran.go | 4 ++-- internal/handler/debian.go | 2 +- internal/handler/filename_download.go | 2 +- internal/handler/generic.go | 2 +- internal/handler/generic_test.go | 20 +++++++++++++++++++- internal/handler/go.go | 2 +- internal/handler/handler_test.go | 10 ++++++++++ internal/handler/helm.go | 8 ++++---- internal/handler/julia.go | 6 +++--- internal/handler/maven.go | 2 +- internal/handler/npm.go | 2 +- internal/handler/nuget.go | 2 +- internal/handler/pub.go | 2 +- internal/handler/pypi.go | 2 +- internal/handler/rpm.go | 2 +- internal/handler/swift.go | 4 ++-- 24 files changed, 75 insertions(+), 34 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 42122c5a..0d37a60a 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -62,6 +62,14 @@ Metadata is not cached - always fetched fresh. This ensures clients see new vers - Return reader to handler - Handler streams file to client +Artifact handlers honor single byte ranges when the storage reader supports +efficient seeking (currently the local filesystem reader). They advertise +`Accept-Ranges: bytes`, return `206` or `416` as appropriate, and respect +`If-Range`. Malformed and multi-range requests fall back to the full response. +Non-seekable storage readers do not advertise range support, and direct-storage +redirects rely on the destination's capabilities. Range responses cannot verify +the full artifact digest because they read only part of the artifact. + ``` ┌────────┐ GET /npm/lodash/-/lodash-4.17.21.tgz ┌─────────────┐ │ Client │ ──────────────────────────────────────▶│ NPMHandler │ diff --git a/internal/handler/apk.go b/internal/handler/apk.go index 9acc3cce..e0b4b31b 100644 --- a/internal/handler/apk.go +++ b/internal/handler/apk.go @@ -132,7 +132,7 @@ func (h *APKHandler) handlePackageDownload(w http.ResponseWriter, r *http.Reques if result.Artifact.MediaType == "" { result.Artifact.MediaType = "application/octet-stream" } - serveArtifact(w, r.Method, result) + ServeArtifactRequest(w, r, result) } // handleMetadata serves repository indexes and signatures through the diff --git a/internal/handler/cargo.go b/internal/handler/cargo.go index 01b62b9e..f87819c1 100644 --- a/internal/handler/cargo.go +++ b/internal/handler/cargo.go @@ -218,5 +218,5 @@ func (h *CargoHandler) handleDownload(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } diff --git a/internal/handler/composer.go b/internal/handler/composer.go index 75a6dc81..2f475ac8 100644 --- a/internal/handler/composer.go +++ b/internal/handler/composer.go @@ -350,7 +350,7 @@ func (h *ComposerHandler) handleDownload(w http.ResponseWriter, r *http.Request) return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // isDevVersion reports whether a Composer version string refers to a diff --git a/internal/handler/conan.go b/internal/handler/conan.go index 862f48c1..aa41668b 100644 --- a/internal/handler/conan.go +++ b/internal/handler/conan.go @@ -94,7 +94,7 @@ func (h *ConanHandler) handleRecipeFile(w http.ResponseWriter, r *http.Request) return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handlePackageFile serves a package file, fetching and caching from upstream if needed. @@ -131,7 +131,7 @@ func (h *ConanHandler) handlePackageFile(w http.ResponseWriter, r *http.Request) return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // shouldCacheFile returns true if the file should be cached. diff --git a/internal/handler/conda.go b/internal/handler/conda.go index ad01f167..07c2aea1 100644 --- a/internal/handler/conda.go +++ b/internal/handler/conda.go @@ -84,7 +84,7 @@ func (h *CondaHandler) handleDownload(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // isPackageFile returns true if the filename is a Conda package. diff --git a/internal/handler/container.go b/internal/handler/container.go index 5a47d795..f73f0f54 100644 --- a/internal/handler/container.go +++ b/internal/handler/container.go @@ -164,7 +164,7 @@ func (h *ContainerHandler) handleBlobDownload(w http.ResponseWriter, r *http.Req if cached.Artifact.MediaType == "" { cached.Artifact.MediaType = "application/octet-stream" } - serveArtifact(w, r.Method, cached) + ServeArtifactRequest(w, r, cached) return } @@ -208,7 +208,7 @@ func (h *ContainerHandler) handleBlobDownload(w http.ResponseWriter, r *http.Req if result.Artifact.MediaType == "" { result.Artifact.MediaType = "application/octet-stream" } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handleManifest serves immutable manifests from cache and revalidates mutable tags. diff --git a/internal/handler/container_test.go b/internal/handler/container_test.go index 5b7aa032..bd237cf7 100644 --- a/internal/handler/container_test.go +++ b/internal/handler/container_test.go @@ -665,6 +665,7 @@ func TestContainerHandler_BlobDownload_CacheHitSkipsAuth(t *testing.T) { proxy, db, store, fetcher := setupTestProxy(t) digest := "sha256:abc123def456abc123def456abc123def456abc123def456abc123def456abcd" seedPackage(t, db, store, "oci", "library/nginx", digest, digest, "cached blob") + store.seekable = true upstreamRequests := 0 upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { @@ -680,14 +681,18 @@ func TestContainerHandler_BlobDownload_CacheHitSkipsAuth(t *testing.T) { } req := httptest.NewRequest(http.MethodGet, "/library/nginx/blobs/"+digest, nil) + req.Header.Set("Range", "bytes=0-5") w := httptest.NewRecorder() h.Routes().ServeHTTP(w, req) - if w.Code != http.StatusOK { - t.Fatalf("status = %d, want %d; body: %s", w.Code, http.StatusOK, w.Body.String()) + if w.Code != http.StatusPartialContent { + t.Fatalf("status = %d, want %d; body: %s", w.Code, http.StatusPartialContent, w.Body.String()) + } + if got := w.Body.String(); got != "cached" { + t.Errorf("body = %q, want %q", got, "cached") } - if got := w.Body.String(); got != "cached blob" { - t.Errorf("body = %q, want %q", got, "cached blob") + if got := w.Header().Get("Content-Range"); got != "bytes 0-5/11" { + t.Errorf("Content-Range = %q, want %q", got, "bytes 0-5/11") } if upstreamRequests != 0 { t.Errorf("upstream requests = %d, want 0", upstreamRequests) diff --git a/internal/handler/cran.go b/internal/handler/cran.go index 4a6ded88..6a363a45 100644 --- a/internal/handler/cran.go +++ b/internal/handler/cran.go @@ -83,7 +83,7 @@ func (h *CRANHandler) handleSourceDownload(w http.ResponseWriter, r *http.Reques return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handleBinaryDownload serves a binary package, fetching and caching from upstream. @@ -117,7 +117,7 @@ func (h *CRANHandler) handleBinaryDownload(w http.ResponseWriter, r *http.Reques return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // parseSourceFilename extracts name and version from a CRAN source filename. diff --git a/internal/handler/debian.go b/internal/handler/debian.go index 513a6cc0..7b6635ef 100644 --- a/internal/handler/debian.go +++ b/internal/handler/debian.go @@ -156,7 +156,7 @@ func (h *DebianHandler) handlePackageDownload( } w.Header().Set(headerContentType, "application/vnd.debian.binary-package") - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handleMetadata serves repository metadata files through the metadata cache, diff --git a/internal/handler/filename_download.go b/internal/handler/filename_download.go index e3c1162e..fc96586a 100644 --- a/internal/handler/filename_download.go +++ b/internal/handler/filename_download.go @@ -39,5 +39,5 @@ func (p *Proxy) handleFilenameDownload(w http.ResponseWriter, r *http.Request, d return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } diff --git a/internal/handler/generic.go b/internal/handler/generic.go index 31f5dc39..a2458b28 100644 --- a/internal/handler/generic.go +++ b/internal/handler/generic.go @@ -133,7 +133,7 @@ func (h *GenericHandler) handleReleaseAsset(w http.ResponseWriter, r *http.Reque if result.Artifact.MediaType == "" { result.Artifact.MediaType = "application/octet-stream" } - serveArtifact(w, r.Method, result) + ServeArtifactRequest(w, r, result) } // handleMetadata serves any other path through the metadata cache. The query diff --git a/internal/handler/generic_test.go b/internal/handler/generic_test.go index 7d47cb94..e67a6b10 100644 --- a/internal/handler/generic_test.go +++ b/internal/handler/generic_test.go @@ -100,7 +100,8 @@ func TestGenericHandler_ReleaseAssetIsCachedAndServedWhenUpstreamDown(t *testing })) defer upstream.Close() - proxy, _, _, _ := setupTestProxy(t) + proxy, _, store, _ := setupTestProxy(t) + store.seekable = true fetcher := fetch.NewFetcher(fetch.WithHTTPClient(upstream.Client()), fetch.WithMaxRetries(0)) proxy.Fetcher = fetcher t.Cleanup(func() { _ = fetcher.Close() }) @@ -120,6 +121,23 @@ func TestGenericHandler_ReleaseAssetIsCachedAndServedWhenUpstreamDown(t *testing // Second request must be served from cache, even with the upstream down. available.Store(false) + rangeRequest := httptest.NewRequest(http.MethodGet, "/github"+testReleaseAssetPath, nil) + rangeRequest.Header.Set("Range", "bytes=0-2") + rangeResponse := httptest.NewRecorder() + h.Routes().ServeHTTP(rangeResponse, rangeRequest) + if rangeResponse.Code != http.StatusPartialContent { + t.Fatalf("range: status = %d, want 206: %s", rangeResponse.Code, rangeResponse.Body.String()) + } + if got := rangeResponse.Body.String(); got != string(asset[:3]) { + t.Errorf("range: body = %q, want %q", got, asset[:3]) + } + if got := rangeResponse.Header().Get("Content-Range"); got != "bytes 0-2/15" { + t.Errorf("range: Content-Range = %q, want %q", got, "bytes 0-2/15") + } + if got := upstreamRequests.Load(); got != 1 { + t.Errorf("upstream requests after range cache hit = %d, want 1", got) + } + w = serveGenericRequest(h, "/github"+testReleaseAssetPath) if w.Code != http.StatusOK { t.Fatalf("cached: status = %d, want 200: %s", w.Code, w.Body.String()) diff --git a/internal/handler/go.go b/internal/handler/go.go index 90b93b70..fb85587d 100644 --- a/internal/handler/go.go +++ b/internal/handler/go.go @@ -126,7 +126,7 @@ func (h *GoHandler) handleDownload(w http.ResponseWriter, r *http.Request, modul return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // proxyUpstream forwards a request to proxy.golang.org without caching. diff --git a/internal/handler/handler_test.go b/internal/handler/handler_test.go index cecf59d2..4c05cbd9 100644 --- a/internal/handler/handler_test.go +++ b/internal/handler/handler_test.go @@ -37,6 +37,7 @@ type mockStorage struct { openErr error signedURL string signErr error + seekable bool } func newMockStorage() *mockStorage { @@ -67,9 +68,18 @@ func (s *mockStorage) Open(_ context.Context, path string) (io.ReadCloser, error if !ok { return nil, storage.ErrNotFound } + if s.seekable { + return &mockSeekableReadCloser{Reader: bytes.NewReader(data)}, nil + } return io.NopCloser(bytes.NewReader(data)), nil } +type mockSeekableReadCloser struct { + *bytes.Reader +} + +func (r *mockSeekableReadCloser) Close() error { return nil } + func (s *mockStorage) Exists(_ context.Context, path string) (bool, error) { s.mu.Lock() defer s.mu.Unlock() diff --git a/internal/handler/helm.go b/internal/handler/helm.go index c0f86fde..31ef9015 100644 --- a/internal/handler/helm.go +++ b/internal/handler/helm.go @@ -94,7 +94,7 @@ func (h *HelmHandler) handleChart(w http.ResponseWriter, r *http.Request) { return } if cached != nil { - h.serveChart(w, repository, digest, filename, cached) + h.serveChart(w, r, repository, digest, filename, cached) return } @@ -127,10 +127,10 @@ func (h *HelmHandler) handleChart(w http.ResponseWriter, r *http.Request) { h.proxy.serveArtifactError(w, err, "failed to fetch chart") return } - h.serveChart(w, repository, digest, filename, result) + h.serveChart(w, r, repository, digest, filename, result) } -func (h *HelmHandler) serveChart(w http.ResponseWriter, repository, digest, filename string, result *CacheResult) { +func (h *HelmHandler) serveChart(w http.ResponseWriter, r *http.Request, repository, digest, filename string, result *CacheResult) { if !strings.EqualFold(result.Artifact.Digest.Encoded(), digest) { if result.Reader != nil { _ = result.Reader.Close() @@ -145,7 +145,7 @@ func (h *HelmHandler) serveChart(w http.ResponseWriter, repository, digest, file if result.Artifact.MediaType == "" { w.Header().Set(headerContentType, "application/gzip") } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } func (h *HelmHandler) repositoryForRequest(r *http.Request) (name, upstreamURL string, ok bool) { diff --git a/internal/handler/julia.go b/internal/handler/julia.go index 211d0922..66a22459 100644 --- a/internal/handler/julia.go +++ b/internal/handler/julia.go @@ -103,7 +103,7 @@ func (h *JuliaHandler) handleRegistry(w http.ResponseWriter, r *http.Request) { go h.refreshNamesFromRegistry(uuid, hash) - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handlePackage serves an immutable package source tarball. @@ -129,7 +129,7 @@ func (h *JuliaHandler) handlePackage(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handleArtifact serves an immutable binary artifact tarball. Artifacts are @@ -150,7 +150,7 @@ func (h *JuliaHandler) handleArtifact(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // proxyUpstream forwards a request to the upstream Pkg server without caching. diff --git a/internal/handler/maven.go b/internal/handler/maven.go index 10e551e0..3f1cd9d1 100644 --- a/internal/handler/maven.go +++ b/internal/handler/maven.go @@ -134,7 +134,7 @@ func (h *MavenHandler) handleDownload(w http.ResponseWriter, r *http.Request, ur return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // parsePath extracts Maven coordinates from a URL path. diff --git a/internal/handler/npm.go b/internal/handler/npm.go index 07dcc485..663c8c66 100644 --- a/internal/handler/npm.go +++ b/internal/handler/npm.go @@ -406,7 +406,7 @@ func (h *NPMHandler) handleDownload(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // versionInCooldown reports whether a version is still inside the cooldown diff --git a/internal/handler/nuget.go b/internal/handler/nuget.go index d37e3910..612f32ee 100644 --- a/internal/handler/nuget.go +++ b/internal/handler/nuget.go @@ -192,7 +192,7 @@ func (h *NuGetHandler) handleDownload(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // proxyUpstream forwards a request to NuGet without caching. diff --git a/internal/handler/pub.go b/internal/handler/pub.go index 9d449bae..6e967e35 100644 --- a/internal/handler/pub.go +++ b/internal/handler/pub.go @@ -81,7 +81,7 @@ func (h *PubHandler) handleDownload(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handlePackageMetadata proxies package metadata and rewrites archive URLs. diff --git a/internal/handler/pypi.go b/internal/handler/pypi.go index 21d98051..e37fddb2 100644 --- a/internal/handler/pypi.go +++ b/internal/handler/pypi.go @@ -626,7 +626,7 @@ func (h *PyPIHandler) handleDownload(w http.ResponseWriter, r *http.Request) { return } - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // archiveExtensions are sdist formats of the form {name}-{version}{ext}. They diff --git a/internal/handler/rpm.go b/internal/handler/rpm.go index e06a3cc1..6b138cca 100644 --- a/internal/handler/rpm.go +++ b/internal/handler/rpm.go @@ -95,7 +95,7 @@ func (h *RPMHandler) handlePackageDownload(w http.ResponseWriter, r *http.Reques } w.Header().Set(headerContentType, "application/x-rpm") - ServeArtifact(w, result) + ServeArtifactRequest(w, r, result) } // handleMetadata proxies repository metadata files (repomd.xml, primary.xml.gz, etc.). diff --git a/internal/handler/swift.go b/internal/handler/swift.go index 93bc9a71..6cc7d78a 100644 --- a/internal/handler/swift.go +++ b/internal/handler/swift.go @@ -189,7 +189,7 @@ func (h *SwiftHandler) handleSourceArchive(w http.ResponseWriter, r *http.Reques result.Artifact.MediaType = "application/zip" setSwiftArchiveHeaders(w.Header(), name, version, result.Artifact.Digest.Encoded(), archiveInfo) - serveArtifact(w, r.Method, result) + ServeArtifactRequest(w, r, result) } func (h *SwiftHandler) handleSourceArchiveHead( @@ -208,7 +208,7 @@ func (h *SwiftHandler) handleSourceArchiveHead( if result != nil { result.Artifact.MediaType = "application/zip" setSwiftArchiveHeaders(w.Header(), name, version, result.Artifact.Digest.Encoded(), archiveInfo) - serveArtifact(w, r.Method, result) + ServeArtifactRequest(w, r, result) return }