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
8 changes: 8 additions & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 │
Expand Down
2 changes: 1 addition & 1 deletion internal/handler/apk.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion internal/handler/cargo.go
Original file line number Diff line number Diff line change
Expand Up @@ -218,5 +218,5 @@ func (h *CargoHandler) handleDownload(w http.ResponseWriter, r *http.Request) {
return
}

ServeArtifact(w, result)
ServeArtifactRequest(w, r, result)
}
2 changes: 1 addition & 1 deletion internal/handler/composer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions internal/handler/conan.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion internal/handler/conda.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 2 additions & 2 deletions internal/handler/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down Expand Up @@ -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.
Expand Down
13 changes: 9 additions & 4 deletions internal/handler/container_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions internal/handler/cran.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion internal/handler/debian.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion internal/handler/filename_download.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,5 +39,5 @@ func (p *Proxy) handleFilenameDownload(w http.ResponseWriter, r *http.Request, d
return
}

ServeArtifact(w, result)
ServeArtifactRequest(w, r, result)
}
2 changes: 1 addition & 1 deletion internal/handler/generic.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
20 changes: 19 additions & 1 deletion internal/handler/generic_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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() })
Expand All @@ -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())
Expand Down
2 changes: 1 addition & 1 deletion internal/handler/go.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
176 changes: 163 additions & 13 deletions internal/handler/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)
}
}

Expand Down
Loading
Loading