Skip to content
Open
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
40 changes: 24 additions & 16 deletions internal/handler/composer.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,14 +104,13 @@ func (h *ComposerHandler) handlePackageMetadata(w http.ResponseWriter, r *http.R

upstreamURL := fmt.Sprintf("%s/p2/%s/%s.json", h.repoURL, vendor, pkg)

if rewritten, ok := h.proxy.storedRewrite("composer", packageName, h.proxyURL, packageName); ok {
w.Header().Set(headerContentType, "application/json")
_, _ = w.Write(rewritten)
if stored, ok := h.proxy.storedRewrite("composer", packageName, h.proxyURL, packageName); ok {
serveRewrittenMetadata(w, r, stored)
return
}

body, _, err := h.proxy.FetchOrCacheMetadata(r.Context(), "composer", packageName, upstreamURL)
if err != nil {
upstream := h.proxy.fetchMetadataDocument(r.Context(), "composer", packageName, upstreamURL, contentTypeJSON)
if err := upstream.err; err != nil {
if errors.Is(err, ErrUpstreamNotFound) {
http.Error(w, "not found", http.StatusNotFound)
return
Expand All @@ -121,19 +120,18 @@ func (h *ComposerHandler) handlePackageMetadata(w http.ResponseWriter, r *http.R
return
}

rewritten, err := h.proxy.cachedRewrite(r.Context(), "composer", h.proxyURL, packageName, body, h.rewriteMetadata)
rewritten, err := h.proxy.cachedRewrite(r.Context(), "composer", h.proxyURL, packageName, upstream.body, upstream.digest, h.rewriteMetadataKeeping)
if err != nil {
if r.Context().Err() != nil {
return // the client left while waiting on a shared rewrite
}
h.proxy.Logger.Warn("failed to rewrite metadata, proxying original", "error", err)
w.Header().Set(headerContentType, "application/json")
_, _ = w.Write(body)
_, _ = w.Write(upstream.body)
return
}

w.Header().Set(headerContentType, "application/json")
_, _ = w.Write(rewritten)
serveRewrittenMetadata(w, r, rewritten)
}

// rewriteMetadata rewrites dist URLs in Composer metadata to point at this proxy.
Expand All @@ -146,16 +144,24 @@ func (h *ComposerHandler) handlePackageMetadata(w http.ResponseWriter, r *http.R
// of the original bytes: only dist values are written fresh, and minified
// metadata stays minified.
func (h *ComposerHandler) rewriteMetadata(body []byte) ([]byte, error) {
out, _, err := h.rewriteMetadataKeeping(body)
return out, err
}

// rewriteMetadataKeeping is rewriteMetadata that also reports the versions
// it kept, each as its package name and version.
func (h *ComposerHandler) rewriteMetadataKeeping(body []byte) ([]byte, []string, error) {
if !json.Valid(body) {
return nil, errors.New("composer metadata is not valid JSON")
return nil, nil, errors.New("composer metadata is not valid JSON")
}
format, _, err := lookupJSONString(body, "minified")
if err != nil {
return nil, err
return nil, nil, err
}
minified := format == composerMinified

var out bytes.Buffer
var kept []string
out.Grow(len(body) + len(body)/8)
err = rewriteJSONMembers(&out, body, func(out *bytes.Buffer, m jsonMember) (bool, error) {
packages := m.value(body)
Expand All @@ -171,24 +177,25 @@ func (h *ComposerHandler) rewriteMetadata(body []byte) ([]byte, error) {
if err != nil {
return false, err
}
return true, h.writeVersions(out, packageName, versions, minified)
return true, h.writeVersions(out, packageName, versions, minified, &kept)
})
})
if err != nil {
return nil, err
return nil, nil, err
}
return out.Bytes(), nil
return out.Bytes(), kept, nil
}

// writeVersions writes one package's version list without the versions in
// cooldown and with each dist pointing at this proxy.
// cooldown and with each dist pointing at this proxy, and appends each version
// it writes to kept as its package name and version.
//
// Minified lists stay minified. Each version written carries the fields that
// differ from what the client has expanded so far, which is the upstream
// entry itself unless filtered versions were dropped before it; then it also
// carries the changes those versions made. Every version gets its own dist,
// since the proxied URL names the version.
func (h *ComposerHandler) writeVersions(out *bytes.Buffer, packageName string, versions []byte, minified bool) error {
func (h *ComposerHandler) writeVersions(out *bytes.Buffer, packageName string, versions []byte, minified bool, kept *[]string) error {
packagePURL := canonicalPackagePURL("composer", packageName)
upstream, written := newComposerFields(), newComposerFields()
devReset, first := false, true
Expand Down Expand Up @@ -226,6 +233,7 @@ func (h *ComposerHandler) writeVersions(out *bytes.Buffer, packageName string, v
devReset = false
}
h.writeVersion(out, packageName, version, upstream, written)
*kept = append(*kept, packageName+"\x00"+version)
return nil
})
out.WriteByte(']')
Expand Down
71 changes: 52 additions & 19 deletions internal/handler/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -1078,24 +1078,35 @@ func (p *Proxy) FetchOrCacheMetadata(ctx context.Context, ecosystem, cacheKey, u
// joined a fetch gets its own validation error directly. It must not modify the
// body, which joined callers share.
func (p *Proxy) fetchOrCacheMetadata(ctx context.Context, ecosystem, cacheKey, upstreamURL, acceptEncoding string, validate func([]byte) error, acceptHeaders ...string) ([]byte, string, string, error) {
accept := contentTypeJSON
if len(acceptHeaders) > 0 && acceptHeaders[0] != "" {
accept = acceptHeaders[0]
}
res := p.lookupMetadata(ctx, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding, validate)
return res.body, res.contentType, res.contentEncoding, res.err
}

// fetchMetadataDocument is FetchOrCacheMetadata for callers that rewrite the
// document and need to know which one they got: the result carries its
// digest along with the body.
func (p *Proxy) fetchMetadataDocument(ctx context.Context, ecosystem, cacheKey, upstreamURL, accept string) metadataResult {
return p.lookupMetadata(ctx, ecosystem, cacheKey, upstreamURL, accept, "", nil)
}

// lookupMetadata implements fetchOrCacheMetadata and fetchMetadataDocument.
func (p *Proxy) lookupMetadata(ctx context.Context, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding string, validate func([]byte) error) metadataResult {
if containsPathTraversal(cacheKey) {
return nil, "", "", fmt.Errorf("invalid cache key: %q", cacheKey)
return metadataResult{err: fmt.Errorf("invalid cache key: %q", cacheKey)}
}

// Serve from cache if within TTL (skip upstream entirely)
if _, hit := p.cachedMetadataState(ctx, ecosystem, cacheKey, validate); hit != nil {
metrics.RecordCacheHit(ecosystem)
return hit.body, hit.contentType, hit.contentEncoding, nil
return *hit
}
p.recordMetadataCacheMiss(ecosystem)

accept := contentTypeJSON
if len(acceptHeaders) > 0 && acceptHeaders[0] != "" {
accept = acceptHeaders[0]
}

res := p.coalescedMetadataMiss(ctx, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding, validate)
return res.body, res.contentType, res.contentEncoding, res.err
return p.coalescedMetadataMiss(ctx, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding, validate)
}

// coalescedMetadataMiss handles a metadata cache miss, sharing one upstream
Expand Down Expand Up @@ -1136,7 +1147,8 @@ func (p *Proxy) cachedMetadataState(ctx context.Context, ecosystem, cacheKey str
if entry != nil && p.MetadataTTL > 0 && entry.FetchedAt.Valid && time.Since(entry.FetchedAt.Time) < p.MetadataTTL {
data, ct, err := p.readCachedMetadata(ctx, entry, validate)
if err == nil {
return entry, &metadataResult{body: data, contentType: ct, contentEncoding: entry.ContentEncoding.String}
hit := metadataResult{body: data, contentType: ct, contentEncoding: entry.ContentEncoding.String}.describedBy(entry)
return entry, &hit
}
if validate != nil {
return nil, nil
Expand All @@ -1159,10 +1171,11 @@ func (p *Proxy) fetchMetadataFromUpstream(ctx context.Context, ecosystem, cacheK
err = validate(meta.body)
}
if err == nil {
var digest string
if p.CacheMetadata {
p.cacheMetadataBlob(ctx, ecosystem, cacheKey, metadataStoragePath(ecosystem, cacheKey), meta)
digest = p.cacheMetadataBlob(ctx, ecosystem, cacheKey, metadataStoragePath(ecosystem, cacheKey), meta)
}
return metadataResult{body: meta.body, contentType: meta.contentType, contentEncoding: meta.contentEncoding}
return metadataResult{body: meta.body, contentType: meta.contentType, contentEncoding: meta.contentEncoding, digest: digest}
}

// Upstream failed -- fall back to cache if available
Expand All @@ -1185,7 +1198,7 @@ func (p *Proxy) fetchMetadataFromUpstream(ctx context.Context, ecosystem, cacheK

p.Logger.Info("serving metadata from cache",
"ecosystem", ecosystem, "key", cacheKey)
return metadataResult{body: data, contentType: ct, contentEncoding: entry.ContentEncoding.String}
return metadataResult{body: data, contentType: ct, contentEncoding: entry.ContentEncoding.String}.describedBy(entry)
}

// metadataResult is what a metadata lookup hands back. A shared fetch gives
Expand All @@ -1194,7 +1207,19 @@ type metadataResult struct {
body []byte
contentType string
contentEncoding string
err error
// digest identifies body when the metadata cache stored it, and is
// empty otherwise.
digest string
err error
}

// describedBy returns res with the digest that entry, the cache row its body
// was read for, records.
func (res metadataResult) describedBy(entry *database.MetadataCacheEntry) metadataResult {
if entry.ContentDigest.Valid {
res.digest = entry.ContentDigest.String
}
return res
}

// inflightMetadata is one metadata fetch that concurrent callers share. res is
Expand Down Expand Up @@ -1377,16 +1402,22 @@ func (p *Proxy) fetchUpstreamMetadata(ctx context.Context, upstreamURL string, e
return meta, nil
}

// cacheMetadataBlob stores metadata bytes in storage and updates the database.
func (p *Proxy) cacheMetadataBlob(ctx context.Context, ecosystem, cacheKey, storagePath string, meta *upstreamMetadata) {
// cacheMetadataBlob stores metadata bytes in storage and updates the
// database. It returns the digest it recorded for the bytes, or "" when they
// were not cached.
func (p *Proxy) cacheMetadataBlob(ctx context.Context, ecosystem, cacheKey, storagePath string, meta *upstreamMetadata) string {
if p.DB == nil || p.Storage == nil {
return
return ""
}

size, hash, err := p.Storage.Store(ctx, storagePath, bytes.NewReader(meta.body))
if err != nil {
p.Logger.Warn("failed to cache metadata", "ecosystem", ecosystem, "key", cacheKey, "error", err)
return
return ""
}
var digest string
if hash != "" {
digest = "sha256:" + hash
}

err = p.DB.UpsertMetadataCache(&database.MetadataCacheEntry{
Expand All @@ -1398,7 +1429,7 @@ func (p *Proxy) cacheMetadataBlob(ctx context.Context, ecosystem, cacheKey, stor
ContentEncoding: sql.NullString{String: meta.contentEncoding, Valid: meta.contentEncoding != ""},
// The digest identifies the stored bytes, so a rewrite cached for them
// can be found without reading them back (see storedRewrite).
ContentDigest: sql.NullString{String: "sha256:" + hash, Valid: hash != ""},
ContentDigest: sql.NullString{String: digest, Valid: digest != ""},
Size: sql.NullInt64{Int64: size, Valid: true},
LastModified: sql.NullTime{Time: meta.lastModified, Valid: !meta.lastModified.IsZero()},
FetchedAt: sql.NullTime{Time: time.Now(), Valid: true},
Expand All @@ -1412,7 +1443,9 @@ func (p *Proxy) cacheMetadataBlob(ctx context.Context, ecosystem, cacheKey, stor
if delErr := p.Storage.Delete(ctx, storagePath); delErr != nil {
p.Logger.Warn("failed to discard metadata blob", "ecosystem", ecosystem, "key", cacheKey, "error", delErr)
}
return ""
}
return digest
}

// currentMetadataEntry re-reads the metadata cache row and returns it, or
Expand Down
Loading
Loading