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
4 changes: 4 additions & 0 deletions cmd/proxy/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -479,6 +479,10 @@ func runMirror() {
fmt.Fprintf(os.Stderr, "invalid configuration: %v\n", err)
os.Exit(1)
}
if !cfg.Storage.CacheArtifacts {
fmt.Fprintf(os.Stderr, "error: mirror is not available with storage.cache_artifacts: false: mirrored artifacts would never be served\n")
os.Exit(1)
}

logger := setupLogger("info", "text")

Expand Down
6 changes: 6 additions & 0 deletions config.example.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,12 @@ storage:
# Empty or "0" means unlimited
max_size: ""

# Store fetched artifacts. Set to false to stream every download from
# upstream without storing it; metadata filtering, cooldown and the
# denylist still apply. Useful when another cache sits in front of the
# proxy. false is incompatible with direct_serve, scanning and mirror_api.
cache_artifacts: true

# Redirect cached artifact downloads to presigned storage URLs (HTTP 302)
# instead of streaming through the proxy. Only effective for S3, GCS, and Azure.
# Leave disabled if clients reach the proxy through an authenticating gateway,
Expand Down
14 changes: 14 additions & 0 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -43,9 +43,23 @@ storage:
| `storage.url` | `PROXY_STORAGE_URL` | `-storage-url` | Storage URL (file:// or s3://) |
| `storage.path` | `PROXY_STORAGE_PATH` | `-storage-path` | Local path (deprecated, use url) |
| `storage.max_size` | `PROXY_STORAGE_MAX_SIZE` | - | Max cache size (e.g., "10GB") |
| `storage.cache_artifacts` | `PROXY_STORAGE_CACHE_ARTIFACTS` | - | Store fetched artifacts (default: true); `false` streams them from upstream |

`storage.max_size` counts cached artifacts only. An artifact replaced by a refetch stays in storage for at least an hour, or `storage.direct_serve_ttl` if longer, so requests already reading it can finish, and storage use can exceed the limit by what was replaced in that time.

### Serving artifacts without storing them

With `storage.cache_artifacts: false` the proxy streams every artifact download from upstream to the client and stores no artifacts. Metadata is still filtered and cached, so cooldown and the denylist apply as usual, and the database still holds the publish times cooldown needs. Use it when another caching layer, such as an Artifactory remote repository, sits in front of the proxy and caching artifacts twice only costs storage.

```yaml
storage:
cache_artifacts: false
```

Every download is a fresh upstream fetch, and concurrent requests for the same artifact are not combined. Artifacts with a digest known up front (OCI blobs, Swift archives, Helm charts) are verified while streaming: the response is sent chunked, and on a mismatch the connection is aborted before the response completes so the client never receives a tampered artifact as a good one. The same happens when the upstream connection fails mid-download.

`cache_artifacts: false` cannot be combined with `scanning.enabled`, `storage.direct_serve` or `mirror_api`, which all need stored artifacts, and the `mirror` command refuses to run with it.

### Amazon S3

```yaml
Expand Down
33 changes: 31 additions & 2 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -370,6 +370,14 @@ type StorageConfig struct {
// storage at an internal address (e.g. 127.0.0.1 or a Docker hostname)
// but clients must use a public one.
DirectServeBaseURL string `json:"direct_serve_base_url" yaml:"direct_serve_base_url"`

// CacheArtifacts stores fetched artifacts so later downloads are served
// from storage. When false, every download streams from upstream and
// nothing is stored; metadata is still cached, and cooldown and the
// denylist still apply. Useful when another caching layer sits in front
// of the proxy. False is incompatible with scanning, direct_serve and
// mirror_api, which all depend on stored artifacts. Default: true.
CacheArtifacts bool `json:"cache_artifacts" yaml:"cache_artifacts"`
}

// GradleConfig configures Gradle-specific features.
Expand Down Expand Up @@ -757,8 +765,9 @@ func Default() *Config {
Listen: ":8080",
BaseURL: "http://localhost:8080",
Storage: StorageConfig{
Path: "./cache/artifacts",
MaxSize: "",
Path: "./cache/artifacts",
MaxSize: "",
CacheArtifacts: true,
},
Database: DatabaseConfig{
Driver: "sqlite",
Expand Down Expand Up @@ -893,6 +902,7 @@ func (c *Config) LoadFromEnv() {
setEnvBool(&c.Storage.DirectServe, "PROXY_STORAGE_DIRECT_SERVE")
setEnvString(&c.Storage.DirectServeTTL, "PROXY_STORAGE_DIRECT_SERVE_TTL")
setEnvString(&c.Storage.DirectServeBaseURL, "PROXY_STORAGE_DIRECT_SERVE_BASE_URL")
setEnvBool(&c.Storage.CacheArtifacts, "PROXY_STORAGE_CACHE_ARTIFACTS")
setEnvString(&c.Database.Driver, "PROXY_DATABASE_DRIVER")
setEnvString(&c.Database.Path, "PROXY_DATABASE_PATH")
setEnvString(&c.Database.URL, "PROXY_DATABASE_URL")
Expand Down Expand Up @@ -1023,6 +1033,10 @@ func (c *Config) Validate() error {
}
}

if err := c.validateCacheArtifacts(); err != nil {
return err
}

// Validate metadata TTL if specified
if c.MetadataTTL != "" && c.MetadataTTL != "0" {
if _, err := time.ParseDuration(c.MetadataTTL); err != nil {
Expand All @@ -1041,6 +1055,21 @@ func (c *Config) Validate() error {
return c.validateComponents()
}

func (c *Config) validateCacheArtifacts() error {
if c.Storage.CacheArtifacts {
return nil
}
switch {
case c.Scanning.Enabled:
return fmt.Errorf("storage.cache_artifacts: false cannot be combined with scanning.enabled: scanning needs stored artifacts")
case c.Storage.DirectServe:
return fmt.Errorf("storage.cache_artifacts: false cannot be combined with storage.direct_serve: no artifacts are stored to redirect to")
case c.MirrorAPI:
return fmt.Errorf("storage.cache_artifacts: false cannot be combined with mirror_api: mirrored artifacts would never be served")
}
return nil
}

func (c *Config) validateComponents() error {
if _, err := denylist.New(c.Denylist.Packages); err != nil {
return err
Expand Down
58 changes: 58 additions & 0 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1303,3 +1303,61 @@ func TestValidateNamedUpstreams(t *testing.T) {
})
}
}

func TestValidateCacheArtifactsDisabled(t *testing.T) {
tests := []struct {
name string
modify func(*Config)
wantErr string
}{
{"alone", func(*Config) {}, ""},
{"with scanning", func(c *Config) { c.Scanning.Enabled = true }, "scanning.enabled"},
{"with direct_serve", func(c *Config) { c.Storage.DirectServe = true }, "storage.direct_serve"},
{"with mirror_api", func(c *Config) { c.MirrorAPI = true }, "mirror_api"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
cfg := Default()
cfg.Storage.CacheArtifacts = false
tt.modify(cfg)
err := cfg.Validate()
if tt.wantErr == "" {
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
return
}
if err == nil || !strings.Contains(err.Error(), tt.wantErr) {
t.Fatalf("Validate() = %v, want an error mentioning %q", err, tt.wantErr)
}
})
}
}

func TestCacheArtifactsDefaultsToTrue(t *testing.T) {
if !Default().Storage.CacheArtifacts {
t.Fatal("Default().Storage.CacheArtifacts = false, want true")
}

path := filepath.Join(t.TempDir(), "config.yaml")
if err := os.WriteFile(path, []byte("storage:\n url: \"file:///tmp/cache\"\n"), 0o600); err != nil {
t.Fatal(err)
}
cfg, err := Load(path)
if err != nil {
t.Fatalf("Load: %v", err)
}
if !cfg.Storage.CacheArtifacts {
t.Error("a config file without storage.cache_artifacts disabled artifact caching")
}
}

func TestLoadCacheArtifactsFromEnv(t *testing.T) {
cfg := Default()
t.Setenv("PROXY_STORAGE_CACHE_ARTIFACTS", "false")
cfg.LoadFromEnv()

if cfg.Storage.CacheArtifacts {
t.Error("Storage.CacheArtifacts should be false")
}
}
118 changes: 117 additions & 1 deletion internal/handler/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,11 @@ type Proxy struct {
HTTPClient *http.Client
AuthForURL func(string) (headerName, headerValue string)

// StreamArtifacts streams artifacts from upstream without storing them.
// Each request fetches its own copy: there is no cache to check and
// nothing for concurrent misses to share.
StreamArtifacts bool

// Scanners runs pre-cache artifact scanning (e.g. trivy, ClamAV, Wiz).
// Nil or disabled means artifacts are cached without scanning.
Scanners *scanner.Group
Expand Down Expand Up @@ -242,6 +247,17 @@ func (p *Proxy) GetOrFetchArtifact(ctx context.Context, ecosystem, name, version
return nil, errors.New("resolved artifact has no filename")
}
}
if p.StreamArtifacts {
if info == nil {
if info, err = p.resolveArtifact(ctx, ecosystem, name, version); err != nil {
return nil, err
}
}
return p.streamFromUpstream(ctx, ecosystem, name, version, filename, versionPURL, info.URL, "",
func(fetchCtx context.Context) (*fetch.Artifact, error) {
return p.Fetcher.Fetch(fetchCtx, info.URL)
})
}
if cached, err := p.checkCache(ctx, pkgPURL, versionPURL, filename); err != nil {
return nil, err
} else if cached != nil {
Expand Down Expand Up @@ -289,10 +305,15 @@ func (p *Proxy) ClearCachedArtifact(ecosystem, name, version, filename string) e
}

// checkCache looks up an artifact in the cache. Returns nil if not cached.
// With StreamArtifacts set it always reports a miss, so entries stored before the
// mode was enabled are never served.
func (p *Proxy) checkCache(ctx context.Context, pkgPURL, versionPURL, filename string) (*CacheResult, error) {
if p.Denylist.Denied(versionPURL) {
return nil, fmt.Errorf("%w: %s", ErrVersionDenied, versionPURL)
}
if p.StreamArtifacts {
return nil, nil
}
artifact, err := p.DB.GetCachedArtifact(pkgPURL, versionPURL, filename)
if err != nil {
return nil, fmt.Errorf("checking artifact cache: %w", err)
Expand Down Expand Up @@ -764,7 +785,14 @@ func serveArtifact(w http.ResponseWriter, method string, result *CacheResult) {
buffer := artifactCopyBufferPool.Get().(*[]byte)
defer artifactCopyBufferPool.Put(buffer)
// Hide optional ReaderFrom methods so io.CopyBuffer uses the pooled buffer.
_, _ = io.CopyBuffer(struct{ io.Writer }{w}, result.Reader, *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)
}
}
}

Expand Down Expand Up @@ -1320,6 +1348,12 @@ func (p *Proxy) getOrFetchArtifactFromURLWithCachePURLs(ctx context.Context, eco
if p.versionDenied(ecosystem, name, version) {
return nil, fmt.Errorf("%w: %s", ErrVersionDenied, canonicalVersionPURL(ecosystem, name, version))
}
if p.StreamArtifacts {
return p.streamFromUpstream(ctx, ecosystem, name, version, filename, versionPURL, downloadURL, upstreamHash,
func(fetchCtx context.Context) (*fetch.Artifact, error) {
return p.Fetcher.FetchWithHeaders(fetchCtx, downloadURL, headers)
})
}
if cached, err := p.getCachedArtifactWithUpstreamHash(ctx, pkgPURL, versionPURL, filename, upstreamHash); err != nil {
return nil, err
} else if cached != nil {
Expand Down Expand Up @@ -1402,6 +1436,88 @@ func (p *Proxy) fetchAndCacheFromURL(ctx context.Context, ecosystem, name, versi
return p.storeArtifact(ctx, ecosystem, name, version, filename, pkgPURL, versionPURL, downloadURL, upstreamHash, artifact)
}

// streamFromUpstream fetches an artifact when StreamArtifacts is set and hands its
// body to the caller without storing it.
//
// With upstreamHash set the body is verified as it streams. The size is left
// unknown so the response goes out chunked: a mismatch only shows at EOF,
// after the bytes have been sent, and an unterminated chunked response is the
// only way left to make the client reject them (see serveArtifact).
func (p *Proxy) streamFromUpstream(ctx context.Context, ecosystem, name, version, filename, versionPURL, upstreamURL, upstreamHash string, fetchArtifact func(context.Context) (*fetch.Artifact, error)) (*CacheResult, error) {
p.Logger.Info("streaming from upstream",
"ecosystem", ecosystem, "name", name, "version", version, "url", upstreamURL)

fetchStart := time.Now()
artifact, err := fetchArtifact(ctx)
metrics.RecordUpstreamFetch(ecosystem, time.Since(fetchStart))
if err != nil {
metrics.RecordUpstreamError(ecosystem, "fetch_failed")
if errors.Is(err, fetch.ErrNotFound) {
return nil, ErrUpstreamNotFound
}
return nil, fmt.Errorf("fetching from upstream: %w", err)
}

body := &streamErrorLogger{
ReadCloser: artifact.Body,
onError: func(read int64, err error) {
p.Logger.Warn("streaming artifact from upstream failed",
"purl", versionPURL, "filename", filename, "url", upstreamURL, "bytes", read, "error", err)
metrics.RecordUpstreamError(ecosystem, "stream_failed")
},
}
result := &CacheResult{
Reader: body,
Artifact: artifacts.Artifact{
PURL: versionPURL,
Size: artifact.Size,
Filename: filename,
MediaType: artifact.ContentType,
},
}
if upstreamHash == "" {
return result, nil
}

hash := strings.ToLower(upstreamHash)
checks, err := newIntegrityChecks(hash, "")
if err != nil {
_ = artifact.Body.Close()
return nil, fmt.Errorf("parsing upstream digest: %w", err)
}
result.Reader, err = checks.wrapFailOnMismatch(body, func(reason string) {
p.Logger.Error("streamed artifact failed integrity check",
"purl", versionPURL, "filename", filename, "url", upstreamURL, "reason", reason)
metrics.RecordIntegrityFailure(purl.NormalizeEcosystem(ecosystem))
})
if err != nil {
_ = artifact.Body.Close()
return nil, err
}
result.Artifact.Digest = digest.Digest("sha256:" + hash)
result.Artifact.Size = -1
return result, nil
}

// streamErrorLogger reports the first read error of a streamed upstream body.
// The error itself still reaches serveArtifact, which aborts the response.
type streamErrorLogger struct {
io.ReadCloser
onError func(read int64, err error)
read int64
logged bool
}

func (r *streamErrorLogger) Read(p []byte) (int, error) {
n, err := r.ReadCloser.Read(p)
r.read += int64(n)
if err != nil && err != io.EOF && !r.logged {
r.logged = true
r.onError(r.read, err)
}
return n, err
}

// ErrArtifactDigestMismatch indicates that fetched bytes did not match the
// checksum the upstream declared and were not recorded in the cache database.
var ErrArtifactDigestMismatch = errors.New("artifact digest mismatch")
Expand Down
10 changes: 8 additions & 2 deletions internal/handler/helm.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,8 +115,14 @@ func (h *HelmHandler) handleChart(w http.ResponseWriter, r *http.Request) {
return
}

result, err := h.proxy.GetOrFetchArtifactFromURL(
r.Context(), helmMetadataEcosystem, repository, digest, filename, downloadURL)
// A streamed fetch is never stored, so serveChart's digest check has
// nothing to compare against: verify the stream itself instead.
expectedDigest := ""
if h.proxy.StreamArtifacts {
expectedDigest = "sha256:" + digest
}
result, err := h.proxy.GetOrFetchArtifactFromURLWithDigest(
r.Context(), helmMetadataEcosystem, repository, digest, filename, downloadURL, expectedDigest)
if err != nil {
h.proxy.serveArtifactError(w, err, "failed to fetch chart")
return
Expand Down
Loading
Loading