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
176 changes: 148 additions & 28 deletions internal/handler/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,11 @@ type Proxy struct {
fetchMu sync.Mutex
inFlight map[string]*inflightFetch

// inFlightMeta does the same for metadata misses. Keyed by
// metadataCoalesceKey.
metaMu sync.Mutex
inFlightMeta map[string]*inflightMetadata

// rewrites caches metadata documents after their handler rewrites them.
// Nil leaves every request to rewrite its own copy.
rewrites *rewriteCache
Expand Down Expand Up @@ -909,34 +914,20 @@ func (p *Proxy) FetchOrCacheMetadata(ctx context.Context, ecosystem, cacheKey, u
// replayed as sent. The ProxyCached path uses "identity" for signed indexes and
// "gzip" where both hops should stay compressed.
// validate, when supplied, runs before caching or serving a document. Validation
// failures follow the same stale-cache fallback path as upstream failures.
// failures follow the same stale-cache fallback path as upstream failures. It
// runs for every caller, including one that shares another caller's fetch, so
// it can also decode the document into request-local state; a caller that
// 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) {
if containsPathTraversal(cacheKey) {
return nil, "", "", fmt.Errorf("invalid cache key: %q", cacheKey)
}

storagePath := metadataStoragePath(ecosystem, cacheKey)

// Check for existing cache entry (for ETag revalidation and TTL)
var entry *database.MetadataCacheEntry
if p.CacheMetadata && p.DB != nil {
entry, _ = p.DB.GetMetadataCache(ecosystem, cacheKey)
}

// Serve from cache if within TTL (skip upstream entirely)
if entry != nil && p.MetadataTTL > 0 && entry.FetchedAt.Valid {
if time.Since(entry.FetchedAt.Time) < p.MetadataTTL {
data, ct, readErr := p.readCachedMetadata(ctx, entry, validate)
if readErr == nil {
metrics.RecordCacheHit(ecosystem)
return data, ct, entry.ContentEncoding.String, nil
}
if validate != nil {
// Do not revalidate an unusable cached body with its ETag.
entry = nil
}
// Cache file missing/unreadable, fall through to upstream
}
if _, hit := p.cachedMetadataState(ctx, ecosystem, cacheKey, validate); hit != nil {
metrics.RecordCacheHit(ecosystem)
return hit.body, hit.contentType, hit.contentEncoding, nil
}
p.recordMetadataCacheMiss(ecosystem)

Expand All @@ -945,7 +936,62 @@ func (p *Proxy) fetchOrCacheMetadata(ctx context.Context, ecosystem, cacheKey, u
accept = acceptHeaders[0]
}

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

// coalescedMetadataMiss handles a metadata cache miss, sharing one upstream
// fetch among concurrent callers with the same key.
func (p *Proxy) coalescedMetadataMiss(ctx context.Context, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding string, validate func([]byte) error) metadataResult {
key := metadataCoalesceKey(ecosystem, cacheKey, upstreamURL, accept, acceptEncoding, validate != nil)
res, shared := p.coalesceMetadata(ctx, key, func(fetchCtx context.Context) metadataResult {
// The caller's lookup ran before it took the key, so a fetch that
// finished in between has already refreshed the row. Recheck it
// rather than fetching again, and revalidate against the row as it
// is now.
entry, hit := p.cachedMetadataState(fetchCtx, ecosystem, cacheKey, validate)
if hit != nil {
return *hit
}
return p.fetchMetadataFromUpstream(fetchCtx, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding, validate, entry)
})
// The fetch ran the first caller's validate. A caller that joined it runs
// its own on the shared bytes, since validate may also decode them for it.
if shared && res.err == nil && validate != nil {
if err := validate(res.body); err != nil {
return metadataResult{err: err}
}
}
return res
}

// cachedMetadataState reads the cache row for a metadata lookup. hit is set
// when the row is within TTL and its bytes are usable. Otherwise entry is the
// row to revalidate against with its ETag: nil when there is none, or when
// validate rejected the cached body, since revalidating an unusable body would
// keep it.
func (p *Proxy) cachedMetadataState(ctx context.Context, ecosystem, cacheKey string, validate func([]byte) error) (*database.MetadataCacheEntry, *metadataResult) {
if !p.CacheMetadata || p.DB == nil {
return nil, nil
}
entry, _ := p.DB.GetMetadataCache(ecosystem, cacheKey)
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}
}
if validate != nil {
return nil, nil
}
// Cache file missing/unreadable, fall through to upstream
}
return entry, nil
}

// fetchMetadataFromUpstream fetches metadata after a cache miss, caches it,
// and falls back to the cached copy if upstream fails. entry is the row to
// revalidate against, or nil.
func (p *Proxy) fetchMetadataFromUpstream(ctx context.Context, ecosystem, cacheKey, upstreamURL, accept, acceptEncoding string, validate func([]byte) error, entry *database.MetadataCacheEntry) metadataResult {
meta, err := p.fetchUpstreamMetadata(ctx, upstreamURL, entry, accept, acceptEncoding)
if errors.Is(err, errStale304) {
// 304 but cached file is gone; retry without ETag
Expand All @@ -956,14 +1002,14 @@ func (p *Proxy) fetchOrCacheMetadata(ctx context.Context, ecosystem, cacheKey, u
}
if err == nil {
if p.CacheMetadata {
p.cacheMetadataBlob(ctx, ecosystem, cacheKey, storagePath, meta)
p.cacheMetadataBlob(ctx, ecosystem, cacheKey, metadataStoragePath(ecosystem, cacheKey), meta)
}
return meta.body, meta.contentType, meta.contentEncoding, nil
return metadataResult{body: meta.body, contentType: meta.contentType, contentEncoding: meta.contentEncoding}
}

// Upstream failed -- fall back to cache if available
if !p.CacheMetadata || entry == nil {
return nil, "", "", fmt.Errorf("upstream failed and no cached metadata: %w", err)
return metadataResult{err: fmt.Errorf("upstream failed and no cached metadata: %w", err)}
}

p.Logger.Warn("upstream metadata fetch failed, checking cache",
Expand All @@ -976,12 +1022,86 @@ func (p *Proxy) fetchOrCacheMetadata(ctx context.Context, ecosystem, cacheKey, u

data, ct, readErr := p.readCachedMetadata(ctx, entry, validate)
if readErr != nil {
return nil, "", "", fmt.Errorf("upstream failed and cached metadata unusable (%v): %w", readErr, err)
return metadataResult{err: fmt.Errorf("upstream failed and cached metadata unusable (%v): %w", readErr, err)}
}

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

// metadataResult is what a metadata lookup hands back. A shared fetch gives
// every waiter the same body, so callers must treat it as read-only.
type metadataResult struct {
body []byte
contentType string
contentEncoding string
err error
}

// inflightMetadata is one metadata fetch that concurrent callers share. res is
// written before done closes and read only after, so the close is the handoff.
type inflightMetadata struct {
done chan struct{}
res metadataResult

// waiters counts callers that joined the fetch, guarded by metaMu. Tests
// read it to know every caller has joined before the fetch finishes.
waiters int
}

// metadataCoalesceKey identifies metadata requests that can share one fetch.
// Accept and Accept-Encoding are part of it because they change the bytes
// upstream returns: npm serves an abbreviated or a full document for the same
// package, and the cached-proxy paths ask for identity or gzip. Whether the
// caller validates is part of it so an unvalidated caller is never handed a
// result that skipped validation, or the reverse.
func metadataCoalesceKey(ecosystem, cacheKey, upstreamURL, accept, acceptEncoding string, validated bool) string {
return strings.Join([]string{ecosystem, cacheKey, upstreamURL, accept, acceptEncoding, strconv.FormatBool(validated)}, "\x00")
}

// coalesceMetadata runs fetch at most once for concurrent callers sharing key
// and gives each the result, reporting whether this caller joined another's
// fetch rather than running it. It follows coalesceFetch with one difference:
// fetch runs on a context detached from the first caller's cancellation. A
// metadata fetch has no scan or mirror that depends on the caller aborting it,
// and CI jobs asking for the same Composer or npm metadata would otherwise all
// fail when the first of them disconnects. It stays bounded by the HTTP
// client's timeout. Waiters still leave when their own context ends.
func (p *Proxy) coalesceMetadata(ctx context.Context, key string, fetch func(context.Context) metadataResult) (metadataResult, bool) {
p.metaMu.Lock()
if p.inFlightMeta == nil {
p.inFlightMeta = make(map[string]*inflightMetadata)
}
f, joined := p.inFlightMeta[key]
if joined {
f.waiters++
} else {
f = &inflightMetadata{done: make(chan struct{})}
p.inFlightMeta[key] = f
}
p.metaMu.Unlock()

if joined {
select {
case <-ctx.Done():
return metadataResult{err: ctx.Err()}, true
case <-f.done:
return f.res, true
}
}

// Set before running so a panicking fetch leaves waiters with an error
// rather than an empty body.
f.res = metadataResult{err: errSharedFetchAbandoned}
defer func() {
p.metaMu.Lock()
delete(p.inFlightMeta, key)
p.metaMu.Unlock()
close(f.done)
}()
f.res = fetch(context.WithoutCancel(ctx))
return f.res, false
}

func (p *Proxy) readCachedMetadata(ctx context.Context, entry *database.MetadataCacheEntry, validate func([]byte) error) ([]byte, string, error) {
Expand Down
Loading
Loading