diff --git a/README.md b/README.md index 083f18b..a1217dd 100644 --- a/README.md +++ b/README.md @@ -718,6 +718,8 @@ base_url: "http://localhost:8080" storage: url: "file:///var/cache/proxy/artifacts" max_size: "10GB" # Optional: evict LRU when exceeded + retention: + default: "30d" # Optional: evict artifacts not downloaded for 30 days database: driver: "sqlite" @@ -1191,6 +1193,7 @@ The proxy exposes Prometheus metrics at `GET /metrics`. All metric names are pre | `proxy_cache_misses_total` | counter | `ecosystem` | Cache misses | | `proxy_cache_size_bytes` | gauge | | Total size of cached artifacts | | `proxy_cached_artifacts_total` | gauge | | Number of cached artifacts | +| `proxy_artifacts_evicted_total` | counter | `reason`, `ecosystem` | Artifacts evicted by the size limit (`lru`) or by age (`retention`) | | `proxy_upstream_fetch_duration_seconds` | histogram | `ecosystem` | Time spent fetching from upstream | | `proxy_upstream_errors_total` | counter | `ecosystem`, `error_type` | Upstream fetch failures | | `proxy_storage_operation_duration_seconds` | histogram | `operation` | Storage read/write latency | @@ -1253,7 +1256,7 @@ The same figures are available as JSON from `GET /stats`, which reports `downloa `ecosystem` means three slightly different things across `/metrics`, and queries that join across them need to know which. -**From the package record, normalized.** The six `proxy_ecosystem_*` gauges, `proxy_cache_hits_total`, `proxy_cache_misses_total`, `proxy_integrity_failures_total` and the scan metrics. Aliases collapse here: `gem` reads as `rubygems`, `composer` as `packagist`, `go` as `golang`. +**From the package record, normalized.** The six `proxy_ecosystem_*` gauges, `proxy_cache_hits_total`, `proxy_cache_misses_total`, `proxy_artifacts_evicted_total`, `proxy_integrity_failures_total` and the scan metrics. Aliases collapse here: `gem` reads as `rubygems`, `composer` as `packagist`, `go` as `golang`. **From the request path.** `proxy_requests_total`, `proxy_request_duration_seconds` and `proxy_response_bytes_total`. The names mostly coincide with the normalized ones -- these also report `rubygems`, `packagist` and `golang` -- but the Debian route reports `debian` where the package record says `deb`, and any path that is not a package endpoint reports `other`, which corresponds to no ecosystem at all. diff --git a/config.example.yaml b/config.example.yaml index a34462a..b647846 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -74,6 +74,21 @@ storage: # Empty or "0" means unlimited max_size: "" + # Evict artifacts that have not been served for longer than a duration + # ("30d", "72h"), independent of max_size. Resolution order: package rule, + # then ecosystem rule, then default. "0" or empty means never evict by age. + # The default only applies to ecosystems that support retention; naming an + # ecosystem or package that does not is a configuration error. See + # docs/configuration.md for the supported ecosystems. + retention: + default: "" + # ecosystems: + # npm: "30d" + # packages: + # "pkg:npm/lodash": "0" + # How often to look for expired artifacts; at least "1m". + sweep_interval: "10m" + # 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 diff --git a/docs/configuration.md b/docs/configuration.md index 331db95..96581a7 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -46,6 +46,8 @@ storage: | `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.retention.default` | `PROXY_STORAGE_RETENTION_DEFAULT` | - | Evict artifacts not served for this long (e.g. "30d"); see [Retention](#retention) | +| `storage.retention.sweep_interval` | `PROXY_STORAGE_RETENTION_SWEEP_INTERVAL` | - | How often expired artifacts are looked for (default: "10m", at least "1m") | `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. @@ -62,6 +64,57 @@ Every download is a fresh upstream fetch, and concurrent requests for the same a `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. +### Retention + +Retention evicts cached artifacts that nobody has downloaded for a while. It works alongside `storage.max_size`: retention removes what has gone unused, and the size limit still applies to whatever is left. + +The `ecosystems` and `packages` entries below only pass validation once the ecosystems they name support retention; see the table further down. + +```yaml +storage: + retention: + default: "30d" + ecosystems: + npm: "14d" + maven: "0" + packages: + "pkg:npm/lodash": "0" + sweep_interval: "10m" +``` + +Each artifact (each version of a package) ages on its own, counted from the last time the proxy served it, or from when it was fetched if it has not been served since. A newer release of the same package does not shorten the life of older versions. Durations take a `d` suffix for days, and `"0"` or an empty value means never evict by age. With no retention configured, artifacts stay until `max_size` evicts them. `sweep_interval` must be at least `1m`. + +A package rule wins over an ecosystem rule, which wins over `default`. Package rules are keyed by package PURL without a version. The ecosystems and package rules can only name ecosystems that support retention in the running version; anything else, including a typo, stops the proxy at startup with an error naming the setting. `default` applies only to supported ecosystems. Support is added one ecosystem at a time, so a later release can make an existing `default` cover more ecosystems. On startup the proxy logs which ecosystems retention covers (`covered`) and which ones `default` does not reach yet (`default_not_applied`). + +| Ecosystem key | Retention supported | Package rule form | +|---------------|---------------------|-------------------| +| `alpine` | no | - | +| `cargo` | no | - | +| `composer` | no | - | +| `conan` | no | - | +| `conda` | no | - | +| `cran` | no | - | +| `deb` | no | - | +| `gem` | no | - | +| `golang` | no | - | +| `helm` | no | - | +| `hex` | no | - | +| `julia` | no | - | +| `maven` | no | - | +| `npm` | no | - | +| `nuget` | no | - | +| `oci` | no | - | +| `pub` | no | - | +| `pypi` | no | - | +| `rpm` | no | - | +| `swift` | no | - | + +An evicted artifact is fetched again from upstream on the next request. Packages preloaded with `proxy mirror` age from the time they were mirrored, so an unused mirrored package is evicted like any other. Running `proxy mirror` again over a package that is still cached counts as a download and restarts its clock, so packages a scheduled mirror job keeps listing never age out. A registry that removes old releases (Linux distribution archives in particular) cannot serve an evicted version again, so set those ecosystems' rules with that in mind. + +Expired artifacts are not deleted right away. Their records are cleared and the files wait in the same queue a replaced artifact goes through, for at least an hour (or `storage.direct_serve_ttl` if longer), so downloads already in progress can finish. `storage.max_size` counts the space as free from the moment the record is cleared. Once that wait is over, each minute the queue is worked through for up to 30 seconds, pausing until the next minute after a delete fails. One sweep clears at most 10,000 artifacts, so the first sweep over a large, old cache spreads its evictions, and the downloads that bring popular artifacts back, over several intervals. + +With several proxies sharing one database and `database.hit_flush_interval` set, a proxy's sweep cannot see downloads another proxy has not written yet. Keep retention durations well above the flush interval. + ### Amazon S3 ```yaml @@ -490,7 +543,7 @@ gradle: |--------|-------------|-------------| | `gradle.build_cache.read_only` | `PROXY_GRADLE_BUILD_CACHE_READ_ONLY` | Disable PUT uploads and keep GET/HEAD read-only | | `gradle.build_cache.max_upload_size` | `PROXY_GRADLE_BUILD_CACHE_MAX_UPLOAD_SIZE` | Maximum accepted PUT body size (must be > 0) | -| `gradle.build_cache.max_age` | `PROXY_GRADLE_BUILD_CACHE_MAX_AGE` | Delete entries older than this duration (default `168h`, set `0` to disable) | +| `gradle.build_cache.max_age` | `PROXY_GRADLE_BUILD_CACHE_MAX_AGE` | Delete entries older than this duration (default `168h`, also accepts days such as `7d`; set `0` to disable) | | `gradle.build_cache.max_size` | `PROXY_GRADLE_BUILD_CACHE_MAX_SIZE` | Total size cap for `_gradle/http-build-cache`, deleting oldest first (`0` disables) | | `gradle.build_cache.sweep_interval` | `PROXY_GRADLE_BUILD_CACHE_SWEEP_INTERVAL` | Frequency for background eviction sweeps | diff --git a/internal/config/config.go b/internal/config/config.go index 2df21d8..ad67101 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -73,6 +73,7 @@ import ( "time" "github.com/git-pkgs/proxy/internal/denylist" + "github.com/git-pkgs/proxy/internal/retention" "github.com/git-pkgs/purl" "gopkg.in/yaml.v3" ) @@ -406,6 +407,33 @@ type StorageConfig struct { // 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"` + + // Retention evicts cached artifacts that have not been accessed for a + // configured time, per ecosystem and per package. + Retention RetentionConfig `json:"retention" yaml:"retention"` +} + +// RetentionConfig configures age-based eviction of cached artifacts. An +// artifact's age is the time since it was last served, or since it was +// fetched when it has not been served since. Durations accept a "d" suffix +// for days; "0" or empty means never evict by age. +type RetentionConfig struct { + // Default applies to every ecosystem that supports retention and has no + // rule of its own. + Default string `json:"default" yaml:"default"` + + // Ecosystems overrides the default per ecosystem key (e.g. "npm", "oci"). + // Configuring an ecosystem whose retention is not supported yet is an + // error. + Ecosystems map[string]string `json:"ecosystems" yaml:"ecosystems"` + + // Packages overrides the ecosystem rule for single packages, keyed by + // package PURL. + Packages map[string]string `json:"packages" yaml:"packages"` + + // SweepInterval is how often expired artifacts are looked for, at least + // one minute. Default: "10m". + SweepInterval string `json:"sweep_interval" yaml:"sweep_interval"` } // GradleConfig configures Gradle-specific features. @@ -861,6 +889,9 @@ func Default() *Config { Path: "./cache/artifacts", MaxSize: "", CacheArtifacts: true, + Retention: RetentionConfig{ + SweepInterval: defaultRetentionSweepIntervalStr, + }, }, Database: DatabaseConfig{ Driver: "sqlite", @@ -1000,6 +1031,8 @@ func (c *Config) LoadFromEnv() { 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.Storage.Retention.Default, "PROXY_STORAGE_RETENTION_DEFAULT") + setEnvString(&c.Storage.Retention.SweepInterval, "PROXY_STORAGE_RETENTION_SWEEP_INTERVAL") setEnvString(&c.Database.Driver, "PROXY_DATABASE_DRIVER") setEnvString(&c.Database.Path, "PROXY_DATABASE_PATH") setEnvString(&c.Database.URL, "PROXY_DATABASE_URL") @@ -1178,6 +1211,10 @@ func (c *Config) validateComponents() error { return err } + if _, err := c.RetentionRules(retention.Default); err != nil { + return err + } + if _, err := denylist.New(c.Denylist.Packages); err != nil { return err } @@ -1228,7 +1265,7 @@ func (g *GradleBuildCacheConfig) Validate() error { } if g.MaxAge != "" && g.MaxAge != "0" { - if _, err := time.ParseDuration(g.MaxAge); err != nil { + if _, err := parseDayDuration(g.MaxAge); err != nil { return fmt.Errorf("invalid gradle.build_cache.max_age %q: %w", g.MaxAge, err) } } @@ -1263,6 +1300,10 @@ const ( defaultGradleBuildCacheSweepInterval = 10 * time.Minute defaultGradleMaxUploadSizeStr = "100MB" defaultGradleSweepIntervalStr = "10m" + defaultRetentionSweepInterval = 10 * time.Minute + defaultRetentionSweepIntervalStr = "10m" + minRetentionSweepInterval = time.Minute + maxDurationDays = 36500 defaultScanningTimeoutStr = "30s" ) @@ -1431,7 +1472,7 @@ func (c *Config) ParseGradleBuildCacheMaxAge() time.Duration { if c.Gradle.BuildCache.MaxAge == "" || c.Gradle.BuildCache.MaxAge == "0" { return 0 } - d, err := time.ParseDuration(c.Gradle.BuildCache.MaxAge) + d, err := parseDayDuration(c.Gradle.BuildCache.MaxAge) if err != nil || d <= 0 { return 0 } diff --git a/internal/config/retention.go b/internal/config/retention.go new file mode 100644 index 0000000..753ddb6 --- /dev/null +++ b/internal/config/retention.go @@ -0,0 +1,154 @@ +package config + +import ( + "fmt" + "math" + "sort" + "strconv" + "strings" + "time" + + "github.com/git-pkgs/cooldown" + "github.com/git-pkgs/proxy/internal/retention" + "github.com/git-pkgs/purl" +) + +// parseDayDuration parses a duration that may use a "d" suffix for days, as +// cooldown values do. Day counts that are not finite or exceed +// maxDurationDays are rejected, since converting them to a time.Duration +// would overflow. +func parseDayDuration(s string) (time.Duration, error) { + if num, ok := strings.CutSuffix(strings.TrimSpace(s), "d"); ok { + days, err := strconv.ParseFloat(num, 64) + if err == nil && (math.IsNaN(days) || math.IsInf(days, 0) || math.Abs(days) > maxDurationDays) { + return 0, fmt.Errorf("invalid duration %q: day count out of range", s) + } + } + return cooldown.ParseDuration(s) +} + +// parseRetentionDuration parses one retention value. field names the +// setting in error messages. +func parseRetentionDuration(field, value string) (time.Duration, error) { + d, err := parseDayDuration(value) + if err != nil { + return 0, fmt.Errorf("invalid %s: %w", field, err) + } + if d < 0 { + return 0, fmt.Errorf("invalid %s %q: must not be negative", field, value) + } + return d, nil +} + +// RetentionRules validates storage.retention against the ecosystems +// registered with reg and resolves it. An ecosystem or package that reg +// does not support is an error, so a rule is never silently ignored. +func (c *Config) RetentionRules(reg *retention.Registry) (retention.Rules, error) { + r := c.Storage.Retention + rules := retention.Rules{ + Ecosystems: map[string]time.Duration{}, + Packages: map[string]time.Duration{}, + } + + var err error + if rules.Default, err = parseRetentionDuration("storage.retention.default", r.Default); err != nil { + return retention.Rules{}, err + } + + if r.SweepInterval != "" { + d, err := parseDayDuration(r.SweepInterval) + if err != nil { + return retention.Rules{}, fmt.Errorf("invalid storage.retention.sweep_interval: %w", err) + } + if d < minRetentionSweepInterval { + return retention.Rules{}, fmt.Errorf("invalid storage.retention.sweep_interval %q: must be at least %s", + r.SweepInterval, minRetentionSweepInterval) + } + } + + for _, key := range sortedKeys(r.Ecosystems) { + field := "storage.retention.ecosystems." + key + if !retention.IsKnown(key) { + return retention.Rules{}, fmt.Errorf("invalid %s: unknown ecosystem %q (valid: %s)", + field, key, strings.Join(retention.KnownKeys, ", ")) + } + if _, ok := reg.Lookup(key); !ok { + return retention.Rules{}, fmt.Errorf("invalid %s: retention is not supported for %q in this version", field, key) + } + d, err := parseRetentionDuration(field, r.Ecosystems[key]) + if err != nil { + return retention.Rules{}, err + } + rules.Ecosystems[key] = d + } + + for _, key := range sortedKeys(r.Packages) { + field := fmt.Sprintf("storage.retention.packages[%q]", key) + canonical, err := retentionPackageKey(reg, key) + if err != nil { + return retention.Rules{}, fmt.Errorf("invalid %s: %w", field, err) + } + d, err := parseRetentionDuration(field, r.Packages[key]) + if err != nil { + return retention.Rules{}, err + } + if prev, ok := rules.Packages[canonical]; ok && prev != d { + return retention.Rules{}, fmt.Errorf("invalid %s: conflicts with another key for package %s", field, canonical) + } + rules.Packages[canonical] = d + } + + return rules, nil +} + +// retentionPackageKey returns the stored package PURL a configured package +// key stands for. +func retentionPackageKey(reg *retention.Registry, key string) (string, error) { + // A julia PURL needs a uuid qualifier to parse at all, and the julia + // handler stores none, so no key could ever match one of its packages. + if strings.HasPrefix(strings.ToLower(key), "pkg:julia/") { + return "", fmt.Errorf("package rules are not supported for julia") + } + p, err := purl.Parse(key) + if err != nil { + return "", fmt.Errorf("not a valid package PURL: %w", err) + } + if p.Version != "" { + return "", fmt.Errorf("package PURL must not carry a version") + } + // A rule names a whole package; qualifiers and subpaths would be + // dropped on the way to the stored PURL, so refuse them rather than + // let "?type=pom" quietly cover every file of the package. + if len(p.Qualifiers) > 0 || p.Subpath != "" { + return "", fmt.Errorf("package PURL must not carry qualifiers or a subpath") + } + canonical, ok := reg.Canonical(p) + if ok { + return canonical, nil + } + if key, registered := reg.KeyForDBEcosystem(purl.PURLTypeToEcosystem(p.Type)); registered { + return "", fmt.Errorf("does not match how %s packages are stored; see the package rule forms in docs/configuration.md", key) + } + return "", fmt.Errorf("retention is not supported for %s packages in this version", p.Type) +} + +func sortedKeys(m map[string]string) []string { + keys := make([]string, 0, len(m)) + for k := range m { + keys = append(keys, k) + } + sort.Strings(keys) + return keys +} + +// ParseRetentionSweepInterval returns how often the retention sweep runs. +func (c *Config) ParseRetentionSweepInterval() time.Duration { + if c.Storage.Retention.SweepInterval == "" { + return defaultRetentionSweepInterval + } + d, err := parseDayDuration(c.Storage.Retention.SweepInterval) + if err != nil || d < minRetentionSweepInterval { + return defaultRetentionSweepInterval + } + return d +} diff --git a/internal/config/retention_test.go b/internal/config/retention_test.go new file mode 100644 index 0000000..d057ad9 --- /dev/null +++ b/internal/config/retention_test.go @@ -0,0 +1,164 @@ +package config + +import ( + "strings" + "testing" + "time" + + "github.com/git-pkgs/proxy/internal/retention" +) + +const retentionDay = 24 * time.Hour + +func retentionTestRegistry() *retention.Registry { + reg := retention.NewRegistry() + reg.Register(retention.Spec{Key: "npm", CanonicalPackage: retention.DefaultCanonical("npm")}) + reg.Register(retention.Spec{Key: "julia"}) + return reg +} + +func TestRetentionRulesUnsetIsDisabled(t *testing.T) { + rules, err := Default().RetentionRules(retentionTestRegistry()) + if err != nil { + t.Fatalf("RetentionRules() error = %v", err) + } + if got := rules.MinPositive(); got != 0 { + t.Errorf("MinPositive() = %v, want 0 for an unset retention block", got) + } +} + +func TestRetentionRulesResolves(t *testing.T) { + cfg := Default() + cfg.Storage.Retention = RetentionConfig{ + Default: "30d", + Ecosystems: map[string]string{"npm": "14d"}, + Packages: map[string]string{ + "pkg:npm/lodash": "0", + "pkg:npm/@babel/core": "90d", + }, + SweepInterval: "5m", + } + + rules, err := cfg.RetentionRules(retentionTestRegistry()) + if err != nil { + t.Fatalf("RetentionRules() error = %v", err) + } + if rules.Default != 30*retentionDay { + t.Errorf("Default = %v, want 30d", rules.Default) + } + if got := rules.For("npm", "pkg:npm/left-pad"); got != 14*retentionDay { + t.Errorf("npm rule = %v, want 14d", got) + } + if got, ok := rules.Packages["pkg:npm/lodash"]; !ok || got != 0 { + t.Errorf("lodash rule = %v, %v, want 0, true", got, ok) + } + if got := rules.For("npm", "pkg:npm/%40babel/core"); got != 90*retentionDay { + t.Errorf("babel rule = %v, want 90d under its canonical key", got) + } + if got := cfg.ParseRetentionSweepInterval(); got != 5*time.Minute { + t.Errorf("ParseRetentionSweepInterval() = %v, want 5m", got) + } +} + +func TestRetentionRulesErrors(t *testing.T) { + tests := []struct { + name string + modify func(*RetentionConfig) + wantErr string + }{ + {"bad default", func(r *RetentionConfig) { r.Default = "7x" }, "storage.retention.default"}, + {"negative default", func(r *RetentionConfig) { r.Default = "-1d" }, "must not be negative"}, + {"NaN days", func(r *RetentionConfig) { r.Default = "NaNd" }, "out of range"}, + {"infinite days", func(r *RetentionConfig) { r.Default = "Infd" }, "out of range"}, + {"overflowing days", func(r *RetentionConfig) { r.Default = "1e9d" }, "out of range"}, + {"zero sweep interval", func(r *RetentionConfig) { r.SweepInterval = "0s" }, "must be at least 1m0s"}, + {"sweep interval below a minute", func(r *RetentionConfig) { r.SweepInterval = "30s" }, "must be at least 1m0s"}, + {"bad sweep interval", func(r *RetentionConfig) { r.SweepInterval = "often" }, "sweep_interval"}, + {"unknown ecosystem", func(r *RetentionConfig) { r.Ecosystems = map[string]string{"nmp": "1d"} }, "unknown ecosystem"}, + {"generic is not an ecosystem key", func(r *RetentionConfig) { r.Ecosystems = map[string]string{"generic": "1d"} }, "unknown ecosystem"}, + {"ecosystem not registered", func(r *RetentionConfig) { r.Ecosystems = map[string]string{"oci": "14d"} }, + `storage.retention.ecosystems.oci: retention is not supported for "oci"`}, + {"bad ecosystem duration", func(r *RetentionConfig) { r.Ecosystems = map[string]string{"npm": "soon"} }, "storage.retention.ecosystems.npm"}, + {"package of unregistered ecosystem", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:cargo/serde": "1d"} }, + "not supported for cargo packages"}, + {"julia package", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:julia/Example": "1d"} }, + "not supported for julia"}, + {"invalid package PURL", func(r *RetentionConfig) { r.Packages = map[string]string{"lodash": "1d"} }, "not a valid package PURL"}, + {"npm package without its scope marker", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:npm/babel/core": "1d"} }, + "does not match how npm packages are stored"}, + {"npm package with an encoded scope separator", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:npm/%40babel%2Fcore": "1d"} }, + "does not match how npm packages are stored"}, + {"versioned package PURL", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:npm/lodash@4.17.21": "1d"} }, + "must not carry a version"}, + {"conflicting package keys", func(r *RetentionConfig) { + r.Packages = map[string]string{"pkg:npm/%40babel/core": "1d", "pkg:npm/@babel/core": "2d"} + }, "conflicts with another key"}, + {"package PURL with qualifiers", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:npm/lodash?type=tgz": "1d"} }, + "must not carry qualifiers or a subpath"}, + {"package PURL with subpath", func(r *RetentionConfig) { r.Packages = map[string]string{"pkg:npm/lodash#lib": "1d"} }, + "must not carry qualifiers or a subpath"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + cfg := Default() + tt.modify(&cfg.Storage.Retention) + _, err := cfg.RetentionRules(retentionTestRegistry()) + if err == nil { + t.Fatalf("RetentionRules() succeeded, want error containing %q", tt.wantErr) + } + if !strings.Contains(err.Error(), tt.wantErr) { + t.Errorf("RetentionRules() error = %q, want it to contain %q", err, tt.wantErr) + } + }) + } +} + +// Validate checks retention against the handlers' registry, which this +// package's tests leave empty, so any ecosystem rule fails while a bare +// default passes. +func TestValidateRetention(t *testing.T) { + cfg := Default() + cfg.Storage.Retention.Default = "30d" + if err := cfg.Validate(); err != nil { + t.Errorf("Validate() with only a default = %v, want nil", err) + } + + cfg.Storage.Retention.Ecosystems = map[string]string{"npm": "14d"} + err := cfg.Validate() + if err == nil || !strings.Contains(err.Error(), "not supported") { + t.Errorf("Validate() with an unsupported ecosystem = %v, want a not supported error", err) + } +} + +func TestRetentionEnvOverrides(t *testing.T) { + t.Setenv("PROXY_STORAGE_RETENTION_DEFAULT", "21d") + t.Setenv("PROXY_STORAGE_RETENTION_SWEEP_INTERVAL", "1h") + + cfg := Default() + cfg.LoadFromEnv() + + if cfg.Storage.Retention.Default != "21d" { + t.Errorf("Retention.Default = %q, want 21d", cfg.Storage.Retention.Default) + } + if got := cfg.ParseRetentionSweepInterval(); got != time.Hour { + t.Errorf("ParseRetentionSweepInterval() = %v, want 1h", got) + } +} + +func TestRetentionSweepIntervalDefault(t *testing.T) { + cfg := Default() + if got := cfg.ParseRetentionSweepInterval(); got != 10*time.Minute { + t.Errorf("ParseRetentionSweepInterval() = %v, want 10m", got) + } +} + +func TestGradleBuildCacheMaxAgeAcceptsDays(t *testing.T) { + cfg := Default() + cfg.Gradle.BuildCache.MaxAge = "7d" + if err := cfg.Validate(); err != nil { + t.Fatalf("Validate() with max_age 7d = %v, want nil", err) + } + if got := cfg.ParseGradleBuildCacheMaxAge(); got != 7*retentionDay { + t.Errorf("ParseGradleBuildCacheMaxAge() = %v, want 168h", got) + } +} diff --git a/internal/database/retention.go b/internal/database/retention.go new file mode 100644 index 0000000..25bd440 --- /dev/null +++ b/internal/database/retention.go @@ -0,0 +1,131 @@ +package database + +import ( + "database/sql" + "errors" + "time" +) + +// RetentionCandidate is a cached artifact that has not been accessed since a +// cutoff, together with the package record it belongs to. +type RetentionCandidate struct { + ID int64 `db:"id"` + VersionPURL string `db:"version_purl"` + Filename string `db:"filename"` + StoragePath string `db:"storage_path"` + Size sql.NullInt64 `db:"size"` + FetchedAt sql.NullTime `db:"fetched_at"` + LastAccessedAt sql.NullTime `db:"last_accessed_at"` + Ecosystem string `db:"ecosystem"` + PackagePURL string `db:"package_purl"` +} + +// ExpiredBefore reports whether the artifact was neither fetched nor +// accessed at or after cutoff. Both are checked: clearing a record keeps its +// last access time, so a record fetched again after an eviction has a fresh +// fetch time next to an old access time. +func (c RetentionCandidate) ExpiredBefore(cutoff time.Time) bool { + if !c.FetchedAt.Valid || !c.FetchedAt.Time.Before(cutoff) { + return false + } + return !c.LastAccessedAt.Valid || c.LastAccessedAt.Time.Before(cutoff) +} + +// expiredCondition selects cached artifacts fetched before the cutoff and not +// accessed since. It takes the cutoff twice. +const expiredCondition = ` + storage_path IS NOT NULL + AND fetched_at < ? + AND (last_accessed_at IS NULL OR last_accessed_at < ?)` + +// MaxArtifactID returns the highest artifact id, or zero for an empty table. +func (db *DB) MaxArtifactID() (int64, error) { + var maxID sql.NullInt64 + if err := db.Get(&maxID, `SELECT MAX(id) FROM artifacts`); err != nil { + return 0, err + } + return maxID.Int64, nil +} + +// GetRetentionCandidates returns the cached artifacts with fromID < id <= toID +// that expired before cutoff. Bounding the scan by id keeps each query short, +// which matters on SQLite, where every request shares one connection. +func (db *DB) GetRetentionCandidates(fromID, toID int64, cutoff time.Time) ([]RetentionCandidate, error) { + var candidates []RetentionCandidate + query := db.Rebind(` + SELECT a.id, a.version_purl, a.filename, a.storage_path, a.size, + a.fetched_at, a.last_accessed_at, p.ecosystem, p.purl AS package_purl + FROM ( + SELECT * FROM artifacts + WHERE id > ? AND id <= ? AND` + expiredCondition + ` + ) a + JOIN versions v ON v.purl = a.version_purl + JOIN packages p ON p.purl = v.package_purl + ORDER BY a.id + `) + if err := db.Select(&candidates, query, fromID, toID, cutoff, cutoff); err != nil { + return nil, err + } + if db.dialect == DialectPostgres { + for i := range candidates { + candidates[i].FetchedAt = asLocalWallClock(candidates[i].FetchedAt) + candidates[i].LastAccessedAt = asLocalWallClock(candidates[i].LastAccessedAt) + } + } + return candidates, nil +} + +// asLocalWallClock reinterprets a time read from a Postgres TIMESTAMP +// column. The column has no zone and holds the local wall clock time the +// proxy wrote, but lib/pq hands it back labelled UTC, which east of UTC +// makes it look hours newer than it is. +func asLocalWallClock(t sql.NullTime) sql.NullTime { + if !t.Valid { + return t + } + v := t.Time + t.Time = time.Date(v.Year(), v.Month(), v.Day(), v.Hour(), v.Minute(), v.Second(), v.Nanosecond(), time.Local) + return t +} + +// ClearExpiredArtifact marks an artifact uncached like ClearArtifactCache, +// but only while it is still expired before cutoff, so an artifact served +// between the scan and the clear stays cached. +func (db *DB) ClearExpiredArtifact(versionPURL, filename, storagePath string, cutoff time.Time) (bool, error) { + query := db.Rebind(` + UPDATE artifacts + SET storage_path = NULL, content_hash = NULL, size = NULL, + content_type = NULL, fetched_at = NULL, updated_at = ? + WHERE version_purl = ? AND filename = ? AND storage_path = ? AND` + expiredCondition) + res, err := db.Exec(query, time.Now(), versionPURL, filename, storagePath, cutoff, cutoff) + if err != nil { + return false, err + } + n, err := res.RowsAffected() + return n > 0, err +} + +// GetArtifactEcosystem returns the packages.ecosystem value of the package an +// artifact version belongs to, or "" when it has no package record. +func (db *DB) GetArtifactEcosystem(versionPURL string) (string, error) { + var ecosystem string + query := db.Rebind(` + SELECT p.ecosystem FROM versions v + JOIN packages p ON p.purl = v.package_purl + WHERE v.purl = ? + `) + err := db.Get(&ecosystem, query, versionPURL) + if errors.Is(err, sql.ErrNoRows) { + return "", nil + } + return ecosystem, err +} + +// FlushHits writes hits buffered by BatchHits now, so a caller about to read +// last_accessed_at sees them. +func (db *DB) FlushHits() error { + if db.hits == nil { + return nil + } + return db.flushHits() +} diff --git a/internal/database/retention_test.go b/internal/database/retention_test.go new file mode 100644 index 0000000..8f102d9 --- /dev/null +++ b/internal/database/retention_test.go @@ -0,0 +1,304 @@ +package database + +import ( + "database/sql" + "testing" + "time" +) + +// seedRetentionArtifact stores a cached artifact for pkg:{ecosystem}/{name} +// with the given fetch and last access times; a zero accessed leaves the +// artifact never served. +func seedRetentionArtifact(t *testing.T, db *DB, ecosystem, name string, fetched, accessed time.Time) { + t.Helper() + pkgPURL := "pkg:" + ecosystem + "/" + name + versionPURL := pkgPURL + "@1.0.0" + if err := db.UpsertPackage(&Package{PURL: pkgPURL, Ecosystem: ecosystem, Name: name}); err != nil { + t.Fatalf("UpsertPackage: %v", err) + } + if err := db.UpsertVersion(&Version{PURL: versionPURL, PackagePURL: pkgPURL}); err != nil { + t.Fatalf("UpsertVersion: %v", err) + } + art := &Artifact{ + VersionPURL: versionPURL, + Filename: name + ".tgz", + UpstreamURL: "https://example.com/" + name + ".tgz", + StoragePath: sql.NullString{String: ecosystem + "/" + name + "/1.0.0/" + name + ".tgz", Valid: true}, + Size: sql.NullInt64{Int64: 100, Valid: true}, + FetchedAt: sql.NullTime{Time: fetched, Valid: true}, + } + if !accessed.IsZero() { + art.LastAccessedAt = sql.NullTime{Time: accessed, Valid: true} + } + if err := db.UpsertArtifact(art); err != nil { + t.Fatalf("UpsertArtifact: %v", err) + } +} + +func candidateNames(candidates []RetentionCandidate) map[string]RetentionCandidate { + byName := make(map[string]RetentionCandidate, len(candidates)) + for _, c := range candidates { + byName[c.Filename] = c + } + return byName +} + +func TestGetRetentionCandidates(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + now := time.Now() + cutoff := now.Add(-24 * time.Hour) + old := now.Add(-48 * time.Hour) + + seedRetentionArtifact(t, db, "npm", "stale", old, old) + seedRetentionArtifact(t, db, "npm", "never-served", old, time.Time{}) + seedRetentionArtifact(t, db, "npm", "recently-served", old, now) + // Fetched again after an eviction: the clear kept the old access time. + seedRetentionArtifact(t, db, "npm", "refetched", now, old) + seedRetentionArtifact(t, db, "gem", "other-ecosystem", old, old) + + maxID, err := db.MaxArtifactID() + if err != nil { + t.Fatalf("MaxArtifactID: %v", err) + } + candidates, err := db.GetRetentionCandidates(0, maxID, cutoff) + if err != nil { + t.Fatalf("GetRetentionCandidates: %v", err) + } + got := candidateNames(candidates) + + for _, want := range []string{"stale.tgz", "never-served.tgz", "other-ecosystem.tgz"} { + if _, ok := got[want]; !ok { + t.Errorf("%s missing from candidates", want) + } + } + for _, unwanted := range []string{"recently-served.tgz", "refetched.tgz"} { + if _, ok := got[unwanted]; ok { + t.Errorf("%s returned as a candidate", unwanted) + } + } + + stale := got["stale.tgz"] + if stale.Ecosystem != "npm" || stale.PackagePURL != "pkg:npm/stale" { + t.Errorf("stale candidate ecosystem/package = %q/%q, want npm/pkg:npm/stale", stale.Ecosystem, stale.PackagePURL) + } + if !stale.ExpiredBefore(cutoff) { + t.Error("stale candidate not expired before the cutoff") + } + if got["other-ecosystem.tgz"].Ecosystem != "gem" { + t.Errorf("other-ecosystem candidate ecosystem = %q, want gem", got["other-ecosystem.tgz"].Ecosystem) + } + }) +} + +func TestGetRetentionCandidatesWindow(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + old := time.Now().Add(-48 * time.Hour) + for _, name := range []string{"a", "b", "c"} { + seedRetentionArtifact(t, db, "npm", name, old, old) + } + cutoff := time.Now().Add(-time.Hour) + + first, err := db.GetRetentionCandidates(0, 1, cutoff) + if err != nil { + t.Fatalf("GetRetentionCandidates: %v", err) + } + rest, err := db.GetRetentionCandidates(1, 3, cutoff) + if err != nil { + t.Fatalf("GetRetentionCandidates: %v", err) + } + if len(first) != 1 || len(rest) != 2 { + t.Fatalf("window sizes = %d and %d, want 1 and 2", len(first), len(rest)) + } + if first[0].ID != 1 || rest[0].ID != 2 || rest[1].ID != 3 { + t.Errorf("window ids = %d, %d, %d, want 1, 2, 3", first[0].ID, rest[0].ID, rest[1].ID) + } + }) +} + +func TestClearExpiredArtifact(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + now := time.Now() + old := now.Add(-48 * time.Hour) + cutoff := now.Add(-24 * time.Hour) + seedRetentionArtifact(t, db, "npm", "stale", old, old) + seedRetentionArtifact(t, db, "npm", "served", old, old) + const stalePath = "npm/stale/1.0.0/stale.tgz" + + // A hit lands between the scan and the clear. + if err := db.RecordArtifactHit("pkg:npm/served@1.0.0", "served.tgz"); err != nil { + t.Fatalf("RecordArtifactHit: %v", err) + } + cleared, err := db.ClearExpiredArtifact("pkg:npm/served@1.0.0", "served.tgz", "npm/served/1.0.0/served.tgz", cutoff) + if err != nil || cleared { + t.Errorf("ClearExpiredArtifact(served) = %v, %v, want false, nil", cleared, err) + } + + cleared, err = db.ClearExpiredArtifact("pkg:npm/stale@1.0.0", "stale.tgz", "elsewhere", cutoff) + if err != nil || cleared { + t.Errorf("ClearExpiredArtifact(moved record) = %v, %v, want false, nil", cleared, err) + } + + cleared, err = db.ClearExpiredArtifact("pkg:npm/stale@1.0.0", "stale.tgz", stalePath, cutoff) + if err != nil || !cleared { + t.Fatalf("ClearExpiredArtifact(stale) = %v, %v, want true, nil", cleared, err) + } + art, err := db.GetArtifact("pkg:npm/stale@1.0.0", "stale.tgz") + if err != nil { + t.Fatalf("GetArtifact: %v", err) + } + if art.StoragePath.Valid || art.FetchedAt.Valid { + t.Error("stale artifact still cached after the clear") + } + if !art.LastAccessedAt.Valid { + t.Error("clear dropped last_accessed_at") + } + }) +} + +func TestFlushHitsWritesBufferedHits(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + old := time.Now().Add(-48 * time.Hour) + seedRetentionArtifact(t, db, "npm", "buffered", old, old) + db.BatchHits(time.Hour, nil) + + if err := db.RecordArtifactHit("pkg:npm/buffered@1.0.0", "buffered.tgz"); err != nil { + t.Fatalf("RecordArtifactHit: %v", err) + } + if err := db.FlushHits(); err != nil { + t.Fatalf("FlushHits: %v", err) + } + art, err := db.GetArtifact("pkg:npm/buffered@1.0.0", "buffered.tgz") + if err != nil { + t.Fatalf("GetArtifact: %v", err) + } + if !art.LastAccessedAt.Time.After(old) { + t.Errorf("last_accessed_at = %v after FlushHits, want the buffered hit", art.LastAccessedAt.Time) + } + }) +} + +func TestFlushHitsWithoutBatching(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + if err := db.FlushHits(); err != nil { + t.Errorf("FlushHits without batching = %v, want nil", err) + } + }) +} + +func TestMaxArtifactIDEmpty(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + maxID, err := db.MaxArtifactID() + if err != nil || maxID != 0 { + t.Errorf("MaxArtifactID() = %d, %v, want 0, nil", maxID, err) + } + }) +} + +func TestGetArtifactEcosystem(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + seedRetentionArtifact(t, db, "rubygems", "rails", time.Now(), time.Time{}) + + got, err := db.GetArtifactEcosystem("pkg:rubygems/rails@1.0.0") + if err != nil || got != "rubygems" { + t.Errorf("GetArtifactEcosystem() = %q, %v, want rubygems, nil", got, err) + } + got, err = db.GetArtifactEcosystem("pkg:npm/missing@1.0.0") + if err != nil || got != "" { + t.Errorf("GetArtifactEcosystem(missing) = %q, %v, want empty, nil", got, err) + } + }) +} + +func TestPendingDeletesQueuedAtIndex(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + indexQuery := `SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name = 'idx_pending_deletes_queued_at'` + if db.Dialect() == DialectPostgres { + indexQuery = `SELECT COUNT(*) FROM pg_indexes WHERE indexname = 'idx_pending_deletes_queued_at'` + } + hasIndex := func() bool { + t.Helper() + var count int + if err := db.Get(&count, indexQuery); err != nil { + t.Fatalf("querying indexes: %v", err) + } + return count == 1 + } + + if !hasIndex() { + t.Error("fresh schema has no idx_pending_deletes_queued_at") + } + if _, err := db.Exec(`DROP INDEX idx_pending_deletes_queued_at`); err != nil { + t.Fatalf("dropping index: %v", err) + } + if err := migrateAddPendingDeletesQueuedAtIndex(db); err != nil { + t.Fatalf("migration: %v", err) + } + if err := migrateAddPendingDeletesQueuedAtIndex(db); err != nil { + t.Fatalf("migration is not idempotent: %v", err) + } + if !hasIndex() { + t.Error("migration did not create idx_pending_deletes_queued_at") + } + }) +} + +func TestRetentionCandidateExpiredBefore(t *testing.T) { + now := time.Now() + cutoff := now.Add(-24 * time.Hour) + old := sql.NullTime{Time: now.Add(-48 * time.Hour), Valid: true} + recent := sql.NullTime{Time: now, Valid: true} + + tests := []struct { + name string + fetched, accessed sql.NullTime + want bool + }{ + {"fetched and accessed before", old, old, true}, + {"never served", old, sql.NullTime{}, true}, + {"accessed since", old, recent, false}, + {"fetched again since", recent, old, false}, + {"not cached", sql.NullTime{}, old, false}, + } + for _, tt := range tests { + c := RetentionCandidate{FetchedAt: tt.fetched, LastAccessedAt: tt.accessed} + if got := c.ExpiredBefore(cutoff); got != tt.want { + t.Errorf("%s: ExpiredBefore() = %v, want %v", tt.name, got, tt.want) + } + } +} + +// A TIMESTAMP column on Postgres holds the local wall clock time without a +// zone. Read back, the candidate's times must still compare correctly +// against a cutoff east of UTC, where a time labelled UTC would look hours +// newer than it is. +func TestGetRetentionCandidatesTimezone(t *testing.T) { + tokyo, err := time.LoadLocation("Asia/Tokyo") + if err != nil { + t.Skipf("time zone data unavailable: %v", err) + } + prev := time.Local + time.Local = tokyo + t.Cleanup(func() { time.Local = prev }) + + runWithBothDatabases(t, func(t *testing.T, db *DB) { + now := time.Now() + fetched := now.Add(-90 * time.Minute) + seedRetentionArtifact(t, db, "npm", "tz", fetched, time.Time{}) + cutoff := now.Add(-60 * time.Minute) + + candidates, err := db.GetRetentionCandidates(0, 10, cutoff) + if err != nil { + t.Fatalf("GetRetentionCandidates: %v", err) + } + if len(candidates) != 1 { + t.Fatalf("got %d candidates, want 1", len(candidates)) + } + c := candidates[0] + if !c.ExpiredBefore(cutoff) { + t.Errorf("candidate fetched at %v not expired before %v", c.FetchedAt.Time, cutoff) + } + if diff := c.FetchedAt.Time.Sub(fetched).Abs(); diff > time.Second { + t.Errorf("fetched_at read back as %v, want %v", c.FetchedAt.Time, fetched) + } + }) +} diff --git a/internal/database/schema.go b/internal/database/schema.go index eb45789..4aef75c 100644 --- a/internal/database/schema.go +++ b/internal/database/schema.go @@ -82,6 +82,7 @@ CREATE TABLE IF NOT EXISTS pending_deletes ( path TEXT NOT NULL PRIMARY KEY, queued_at DATETIME NOT NULL ); +CREATE INDEX IF NOT EXISTS idx_pending_deletes_queued_at ON pending_deletes(queued_at); CREATE TABLE IF NOT EXISTS vulnerabilities ( id INTEGER PRIMARY KEY, @@ -190,6 +191,7 @@ CREATE TABLE IF NOT EXISTS pending_deletes ( path TEXT NOT NULL PRIMARY KEY, queued_at TIMESTAMP NOT NULL ); +CREATE INDEX IF NOT EXISTS idx_pending_deletes_queued_at ON pending_deletes(queued_at); CREATE TABLE IF NOT EXISTS vulnerabilities ( id SERIAL PRIMARY KEY, @@ -379,6 +381,7 @@ var migrations = []migration{ {"007_add_metadata_link", migrateAddMetadataLink}, {"008_add_metadata_content_encoding", migrateAddMetadataContentEncoding}, {"009_add_pending_deletes", migrateAddPendingDeletes}, + {"010_add_pending_deletes_queued_at_index", migrateAddPendingDeletesQueuedAtIndex}, } // isTableNotFound returns true if the error indicates a missing table. @@ -668,6 +671,16 @@ func migrateAddPendingDeletes(db *DB) error { return nil } +// migrateAddPendingDeletesQueuedAtIndex indexes the column reclaim orders +// the queue by, which grows large after a retention sweep. +func migrateAddPendingDeletesQueuedAtIndex(db *DB) error { + _, err := db.Exec(`CREATE INDEX IF NOT EXISTS idx_pending_deletes_queued_at ON pending_deletes(queued_at)`) + if err != nil { + return fmt.Errorf("creating pending_deletes queued_at index: %w", err) + } + return nil +} + // EnsureMetadataCacheTable creates the metadata_cache table if it doesn't exist. func (db *DB) EnsureMetadataCacheTable() error { has, err := db.HasTable("metadata_cache") diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 1958165..23dca2c 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -62,6 +62,14 @@ var ( }, ) + ArtifactsEvicted = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "proxy_artifacts_evicted_total", + Help: "Total number of cached artifacts evicted, by reason (lru or retention) and ecosystem", + }, + []string{"reason", "ecosystem"}, + ) + // Upstream metrics UpstreamFetchDuration = prometheus.NewHistogramVec( prometheus.HistogramOpts{ @@ -250,6 +258,7 @@ func init() { CacheMisses, CacheSize, CachedArtifacts, + ArtifactsEvicted, UpstreamFetchDuration, UpstreamErrors, CircuitBreakerState, @@ -309,6 +318,16 @@ func RecordCacheHit(ecosystem string) { CacheHits.WithLabelValues(purl.NormalizeEcosystem(ecosystem)).Inc() } +// RecordArtifactEvicted counts one evicted artifact. ecosystem is the +// package record's value, normalized like the cache hit counter; an artifact +// without a package record counts as "other". +func RecordArtifactEvicted(reason, ecosystem string) { + if ecosystem == "" { + ecosystem = "other" + } + ArtifactsEvicted.WithLabelValues(reason, purl.NormalizeEcosystem(ecosystem)).Inc() +} + // RecordCacheMiss increments cache miss counter. func RecordCacheMiss(ecosystem string) { CacheMisses.WithLabelValues(purl.NormalizeEcosystem(ecosystem)).Inc() diff --git a/internal/retention/retention.go b/internal/retention/retention.go new file mode 100644 index 0000000..c26ef4a --- /dev/null +++ b/internal/retention/retention.go @@ -0,0 +1,209 @@ +// Package retention holds the per-ecosystem rules for evicting cached +// artifacts that have not been accessed for a configured time. +// +// An ecosystem takes part only once its handler registers a Spec with +// Default, usually from an init function in the handler's own file. Until +// then configuring retention for it is a validation error, and the global +// default does not apply to it. +package retention + +import ( + "fmt" + "slices" + "sort" + "sync" + "time" + + "github.com/git-pkgs/purl" +) + +// KnownKeys lists the ecosystem keys storage.retention.ecosystems accepts. +// They are the literals the handlers pass when they fetch an artifact. A key +// outside this list is a typo; a key in it whose handler has not registered +// yet is not supported by this version. +var KnownKeys = []string{ + "alpine", "cargo", "composer", "conan", "conda", "cran", "deb", "gem", + "golang", "helm", "hex", "julia", "maven", "npm", "nuget", "oci", "pub", + "pypi", "rpm", "swift", +} + +// IsKnown reports whether key is one of KnownKeys. +func IsKnown(key string) bool { + return slices.Contains(KnownKeys, key) +} + +// Spec describes how one ecosystem takes part in retention. +type Spec struct { + // Key is the ecosystem's config key, one of KnownKeys. + Key string + + // ExtraDBEcosystems lists packages.ecosystem values besides the ones + // derived from Key that belong to this ecosystem. + ExtraDBEcosystems []string + + // CanonicalPackage turns a package PURL from the config into the + // package PURL the handler stores in the cache, reporting false when the + // PURL does not belong to this ecosystem. Nil means the ecosystem accepts + // no package overrides. + CanonicalPackage func(p *purl.PURL) (string, bool) +} + +// dbEcosystems returns every packages.ecosystem value that belongs to the +// spec: the key, its PURL type and the ecosystem name that type maps back +// to, the same spellings git-pkgs matches. The proxy's mirror command and +// git-pkgs store "rubygems" where the gem handler stores "gem", for example. +func (s Spec) dbEcosystems() []string { + purlType := purl.EcosystemToPURLType(s.Key) + values := []string{s.Key, purlType, purl.PURLTypeToEcosystem(purlType)} + values = append(values, s.ExtraDBEcosystems...) + slices.Sort(values) + return slices.Compact(values) +} + +// DefaultCanonical returns a CanonicalPackage for ecosystems that store the +// package PURL exactly as purl.MakePURLString builds it from the PURL's full +// name. It accepts a PURL only when building it back that way gives the same +// package, so for an ecosystem that stores a different form, such as one +// with a default namespace, it refuses the key instead of producing one +// that never matches. Such ecosystems need their own CanonicalPackage. +func DefaultCanonical(key string) func(p *purl.PURL) (string, bool) { + purlType := purl.EcosystemToPURLType(key) + return func(p *purl.PURL) (string, bool) { + if p == nil || p.Type != purlType { + return "", false + } + canonical := purl.MakePURLString(key, p.FullName(), "") + if canonical == "" { + return "", false + } + given := purl.New(p.Type, p.Namespace, p.Name, "", nil).String() + return canonical, canonical == given + } +} + +// Registry maps ecosystem keys to their specs. +type Registry struct { + mu sync.RWMutex + specs map[string]Spec + byDB map[string]string +} + +// NewRegistry returns an empty registry. +func NewRegistry() *Registry { + return &Registry{specs: map[string]Spec{}, byDB: map[string]string{}} +} + +// Default is the registry the handlers register with. +var Default = NewRegistry() + +// Register adds s. It panics when the key is not one of KnownKeys, is +// already registered, or claims a packages.ecosystem value another spec +// claims, since each of those is a programming error. +func (r *Registry) Register(s Spec) { + r.mu.Lock() + defer r.mu.Unlock() + + if !IsKnown(s.Key) { + panic(fmt.Sprintf("retention: unknown ecosystem key %q", s.Key)) + } + if _, ok := r.specs[s.Key]; ok { + panic(fmt.Sprintf("retention: ecosystem %q registered twice", s.Key)) + } + values := s.dbEcosystems() + for _, v := range values { + if owner, ok := r.byDB[v]; ok { + panic(fmt.Sprintf("retention: ecosystem value %q claimed by both %q and %q", v, owner, s.Key)) + } + } + r.specs[s.Key] = s + for _, v := range values { + r.byDB[v] = s.Key + } +} + +// Lookup returns the spec registered for key. +func (r *Registry) Lookup(key string) (Spec, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + s, ok := r.specs[key] + return s, ok +} + +// Keys returns the registered keys in sorted order. +func (r *Registry) Keys() []string { + r.mu.RLock() + defer r.mu.RUnlock() + keys := make([]string, 0, len(r.specs)) + for k := range r.specs { + keys = append(keys, k) + } + sort.Strings(keys) + return keys +} + +// KeyForDBEcosystem returns the registered key a packages.ecosystem value +// belongs to. +func (r *Registry) KeyForDBEcosystem(ecosystem string) (string, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + key, ok := r.byDB[ecosystem] + return key, ok +} + +// Canonical returns the stored package PURL a configured package PURL +// stands for, asking the registered specs in key order. +func (r *Registry) Canonical(p *purl.PURL) (string, bool) { + for _, key := range r.Keys() { + s, _ := r.Lookup(key) + if s.CanonicalPackage == nil { + continue + } + if canonical, ok := s.CanonicalPackage(p); ok { + return canonical, true + } + } + return "", false +} + +// Rules are resolved retention durations. A duration of zero means an +// artifact is never evicted by age. +type Rules struct { + // Default applies to registered ecosystems without their own rule. + Default time.Duration + // Ecosystems holds rules keyed by registered ecosystem key. + Ecosystems map[string]time.Duration + // Packages holds rules keyed by the package PURL stored in the cache. + Packages map[string]time.Duration +} + +// For returns the retention for an artifact of the given ecosystem key and +// stored package PURL: the package rule, else the ecosystem rule, else the +// default. +func (r Rules) For(key, packagePURL string) time.Duration { + if d, ok := r.Packages[packagePURL]; ok { + return d + } + if d, ok := r.Ecosystems[key]; ok { + return d + } + return r.Default +} + +// MinPositive returns the shortest non-zero duration among all rules, or +// zero when every rule is zero. +func (r Rules) MinPositive() time.Duration { + minimum := time.Duration(0) + consider := func(d time.Duration) { + if d > 0 && (minimum == 0 || d < minimum) { + minimum = d + } + } + consider(r.Default) + for _, d := range r.Ecosystems { + consider(d) + } + for _, d := range r.Packages { + consider(d) + } + return minimum +} diff --git a/internal/retention/retention_test.go b/internal/retention/retention_test.go new file mode 100644 index 0000000..9fa67b7 --- /dev/null +++ b/internal/retention/retention_test.go @@ -0,0 +1,163 @@ +package retention + +import ( + "slices" + "testing" + "time" + + "github.com/git-pkgs/purl" +) + +func TestSpecDBEcosystems(t *testing.T) { + tests := []struct { + key string + want []string + }{ + {"npm", []string{"npm"}}, + {"gem", []string{"gem", "rubygems"}}, + {"composer", []string{"composer", "packagist"}}, + {"alpine", []string{"alpine", "apk"}}, + {"golang", []string{"golang"}}, + {"deb", []string{"deb"}}, + } + for _, tt := range tests { + t.Run(tt.key, func(t *testing.T) { + got := Spec{Key: tt.key}.dbEcosystems() + if !slices.Equal(got, tt.want) { + t.Errorf("dbEcosystems() = %v, want %v", got, tt.want) + } + }) + } +} + +func TestRegistryRegister(t *testing.T) { + reg := NewRegistry() + reg.Register(Spec{Key: "gem"}) + reg.Register(Spec{Key: "npm"}) + + if got := reg.Keys(); !slices.Equal(got, []string{"gem", "npm"}) { + t.Errorf("Keys() = %v, want [gem npm]", got) + } + for _, value := range []string{"gem", "rubygems"} { + if key, ok := reg.KeyForDBEcosystem(value); !ok || key != "gem" { + t.Errorf("KeyForDBEcosystem(%q) = %q, %v, want gem, true", value, key, ok) + } + } + if _, ok := reg.KeyForDBEcosystem("pypi"); ok { + t.Error("KeyForDBEcosystem(pypi) found a key for an unregistered ecosystem") + } + if _, ok := reg.Lookup("pypi"); ok { + t.Error("Lookup(pypi) found an unregistered ecosystem") + } +} + +func TestRegistryRegisterPanics(t *testing.T) { + tests := []struct { + name string + specs []Spec + }{ + {"unknown key", []Spec{{Key: "generic"}}}, + {"duplicate key", []Spec{{Key: "npm"}, {Key: "npm"}}}, + {"overlapping ecosystem value", []Spec{{Key: "gem"}, {Key: "npm", ExtraDBEcosystems: []string{"rubygems"}}}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + reg := NewRegistry() + defer func() { + if recover() == nil { + t.Error("Register did not panic") + } + }() + for _, s := range tt.specs { + reg.Register(s) + } + }) + } +} + +func mustParse(t *testing.T, s string) *purl.PURL { + t.Helper() + p, err := purl.Parse(s) + if err != nil { + t.Fatalf("parsing %q: %v", s, err) + } + return p +} + +func TestDefaultCanonical(t *testing.T) { + canonical := DefaultCanonical("npm") + + got, ok := canonical(mustParse(t, "pkg:npm/%40babel/core")) + if !ok || got != "pkg:npm/%40babel/core" { + t.Errorf("scoped package = %q, %v, want pkg:npm/%%40babel/core, true", got, ok) + } + if got, ok := canonical(mustParse(t, "pkg:pypi/requests")); ok { + t.Errorf("other type accepted as %q", got) + } + if _, ok := canonical(nil); ok { + t.Error("nil PURL accepted") + } + + // alpine stores pkg:apk/alpine/{name} and deb stores pkg:deb/{name}, so + // building their keys back from the full name gives a different package. + // The default refuses them rather than produce keys that never match. + for _, tt := range []struct{ key, purl string }{ + {"alpine", "pkg:apk/alpine/curl"}, + {"alpine", "pkg:apk/curl"}, + {"deb", "pkg:deb/debian/curl"}, + } { + if got, ok := DefaultCanonical(tt.key)(mustParse(t, tt.purl)); ok { + t.Errorf("DefaultCanonical(%q) accepted %s as %q", tt.key, tt.purl, got) + } + } +} + +func TestRegistryCanonical(t *testing.T) { + reg := NewRegistry() + reg.Register(Spec{Key: "julia"}) + reg.Register(Spec{Key: "npm", CanonicalPackage: DefaultCanonical("npm")}) + + if got, ok := reg.Canonical(mustParse(t, "pkg:npm/lodash")); !ok || got != "pkg:npm/lodash" { + t.Errorf("Canonical(npm) = %q, %v, want pkg:npm/lodash, true", got, ok) + } + if got, ok := reg.Canonical(mustParse(t, "pkg:cargo/serde")); ok { + t.Errorf("Canonical(cargo) = %q for an unregistered ecosystem", got) + } +} + +func TestRulesFor(t *testing.T) { + rules := Rules{ + Default: 30 * 24 * time.Hour, + Ecosystems: map[string]time.Duration{"npm": 14 * 24 * time.Hour, "maven": 0}, + Packages: map[string]time.Duration{"pkg:npm/lodash": 0, "pkg:maven/a/b": 90 * 24 * time.Hour}, + } + tests := []struct { + key, pkg string + want time.Duration + }{ + {"npm", "pkg:npm/lodash", 0}, + {"npm", "pkg:npm/left-pad", 14 * 24 * time.Hour}, + {"maven", "pkg:maven/a/b", 90 * 24 * time.Hour}, + {"maven", "pkg:maven/c/d", 0}, + {"cargo", "pkg:cargo/serde", 30 * 24 * time.Hour}, + } + for _, tt := range tests { + if got := rules.For(tt.key, tt.pkg); got != tt.want { + t.Errorf("For(%q, %q) = %v, want %v", tt.key, tt.pkg, got, tt.want) + } + } +} + +func TestRulesMinPositive(t *testing.T) { + if got := (Rules{}).MinPositive(); got != 0 { + t.Errorf("empty rules MinPositive() = %v, want 0", got) + } + rules := Rules{ + Default: 0, + Ecosystems: map[string]time.Duration{"npm": 14 * time.Hour, "maven": 0}, + Packages: map[string]time.Duration{"pkg:npm/lodash": 2 * time.Hour}, + } + if got := rules.MinPositive(); got != 2*time.Hour { + t.Errorf("MinPositive() = %v, want 2h", got) + } +} diff --git a/internal/server/eviction.go b/internal/server/eviction.go index f9a706a..7fe108a 100644 --- a/internal/server/eviction.go +++ b/internal/server/eviction.go @@ -141,6 +141,7 @@ func evictBatch(ctx context.Context, db *database.DB, store storage.Storage, log freed += art.Size.Int64 } cleared++ + recordLRUEviction(db, logger, art.VersionPURL) } return cleared, freed diff --git a/internal/server/reclaim.go b/internal/server/reclaim.go index 72a4e8f..aa986a1 100644 --- a/internal/server/reclaim.go +++ b/internal/server/reclaim.go @@ -17,6 +17,10 @@ const ( reclaimInterval = 1 * time.Minute reclaimBatch = 100 reclaimMinGrace = 1 * time.Hour + // reclaimTickBudget bounds how long one tick keeps deleting batches, so + // a long queue, such as one left by a retention sweep, drains faster than + // one batch per interval without one tick running on indefinitely. + reclaimTickBudget = 30 * time.Second ) func (s *Server) startReclaimLoop(ctx context.Context) { @@ -30,29 +34,46 @@ func (s *Server) startReclaimLoop(ctx context.Context) { case <-ctx.Done(): return case <-ticker.C: - reclaimStorage(ctx, s.db, s.storage, s.logger, time.Now().Add(-grace)) + reclaimDue(ctx, s.db, s.storage, s.logger, time.Now().Add(-grace), reclaimTickBudget) } } } -// reclaimStorage deletes up to one batch of objects queued before cutoff. A -// delete that fails is queued again, behind the rest, so objects the backend -// keeps refusing cannot fill every batch. -func reclaimStorage(ctx context.Context, db *database.DB, store storage.Storage, logger *slog.Logger, cutoff time.Time) { +// reclaimDue deletes batches of objects queued before cutoff until a batch +// comes back short, a delete fails or budget has passed. Stopping on a +// failed delete keeps a storage outage to one batch of attempts per tick. +func reclaimDue(ctx context.Context, db *database.DB, store storage.Storage, logger *slog.Logger, + cutoff time.Time, budget time.Duration) { + start := time.Now() + for ctx.Err() == nil && time.Since(start) < budget { + listed, failed := reclaimStorage(ctx, db, store, logger, cutoff) + if listed < reclaimBatch || failed > 0 { + return + } + } +} + +// reclaimStorage deletes up to one batch of objects queued before cutoff and +// reports how many it listed and how many deletes failed. A delete that +// fails is queued again, behind the rest, so objects the backend keeps +// refusing cannot fill every batch. +func reclaimStorage(ctx context.Context, db *database.DB, store storage.Storage, logger *slog.Logger, + cutoff time.Time) (listed, failed int) { paths, err := db.GetDuePendingDeletes(cutoff, reclaimBatch) if err != nil { logger.Warn("reclaim: failed to list pending deletes", "error", err) - return + return 0, 0 } for _, path := range paths { if ctx.Err() != nil { - return + return len(paths), failed } if err := store.Delete(ctx, path); err != nil { if ctx.Err() != nil { - return + return len(paths), failed } + failed++ logger.Warn("reclaim: failed to delete object, will retry", "path", path, "error", err) if err := db.QueuePendingDelete(path); err != nil { logger.Warn("reclaim: failed to requeue object", "path", path, "error", err) @@ -63,4 +84,5 @@ func reclaimStorage(ctx context.Context, db *database.DB, store storage.Storage, logger.Warn("reclaim: failed to dequeue deleted object", "path", path, "error", err) } } + return len(paths), failed } diff --git a/internal/server/retention.go b/internal/server/retention.go new file mode 100644 index 0000000..69f7985 --- /dev/null +++ b/internal/server/retention.go @@ -0,0 +1,179 @@ +package server + +import ( + "context" + "log/slog" + "slices" + "time" + + "github.com/git-pkgs/proxy/internal/database" + "github.com/git-pkgs/proxy/internal/metrics" + "github.com/git-pkgs/proxy/internal/retention" +) + +// retentionWindow is how many artifact ids one query scans. +const retentionWindow = 5000 + +// retentionMaxClears caps how many artifacts one sweep clears, so a first +// sweep over a large, old cache spreads its deletes and the refetches that +// follow over several intervals. A variable so tests can lower it. +var retentionMaxClears = 10000 + +func (s *Server) startRetentionLoop(ctx context.Context) { + reg := retention.Default + rules, err := s.cfg.RetentionRules(reg) + if err != nil { + // Validate already rejected an invalid config at startup. + s.logger.Error("retention: invalid configuration, sweep disabled", "error", err) + return + } + if rules.MinPositive() == 0 { + return + } + + covered := reg.Keys() + if len(covered) == 0 { + // Only a default can be set without any registered ecosystem; + // validation rejects every other rule. + s.logger.Info("retention: no ecosystem supports retention in this version, sweep disabled", + "default_not_applied", notCovered(covered)) + return + } + s.logger.Info("retention enabled", + "covered", covered, + "default_not_applied", notCovered(covered), + "sweep_interval", s.cfg.ParseRetentionSweepInterval()) + + ticker := time.NewTicker(s.cfg.ParseRetentionSweepInterval()) + defer ticker.Stop() + + sweepRetention(ctx, s.db, s.logger, reg, rules, time.Now()) + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + sweepRetention(ctx, s.db, s.logger, reg, rules, time.Now()) + } + } +} + +// notCovered returns the known ecosystem keys that are not in covered. +func notCovered(covered []string) []string { + var missing []string + for _, key := range retention.KnownKeys { + if !slices.Contains(covered, key) { + missing = append(missing, key) + } + } + return missing +} + +// sweepRetention clears cached artifacts that have outlived their retention +// as of now and queues their objects for deletion. Objects are not deleted +// directly: a hit still buffered when the scan ran can make an artifact that +// is being served look expired, and the queue's grace period covers that. +func sweepRetention(ctx context.Context, db *database.DB, logger *slog.Logger, + reg *retention.Registry, rules retention.Rules, now time.Time) { + minimum := rules.MinPositive() + if minimum == 0 { + return + } + + // Without the buffered hits written, an artifact being served right now + // can look expired, so skip this sweep rather than clear it. + if err := db.FlushHits(); err != nil { + logger.Warn("retention: failed to write buffered hits, skipping this sweep", "error", err) + return + } + + maxID, err := db.MaxArtifactID() + if err != nil { + logger.Warn("retention: failed to read artifact ids", "error", err) + return + } + + // Rows that are not expired under the shortest rule cannot be expired + // under any, so the query uses that cutoff and each row is then checked + // against its own rule. + scanCutoff := now.Add(-minimum) + cleared := 0 + freed := int64(0) + capped := false + + for from := int64(0); from < maxID && !capped; from += retentionWindow { + if ctx.Err() != nil { + return + } + candidates, err := db.GetRetentionCandidates(from, from+retentionWindow, scanCutoff) + if err != nil { + logger.Warn("retention: failed to scan artifacts", "from_id", from, "error", err) + return + } + for _, c := range candidates { + if ctx.Err() != nil { + return + } + if cleared >= retentionMaxClears { + capped = true + break + } + if !clearIfExpired(db, logger, reg, rules, now, c) { + continue + } + metrics.RecordArtifactEvicted("retention", c.Ecosystem) + cleared++ + if c.Size.Valid { + freed += c.Size.Int64 + } + } + } + + if cleared > 0 { + logger.Info("retention: sweep completed", + "cleared", cleared, "freed_bytes", freed, "capped", capped) + } +} + +// clearIfExpired clears one candidate when its own rule says it has expired +// and queues its object for deletion, reporting whether it did. +func clearIfExpired(db *database.DB, logger *slog.Logger, reg *retention.Registry, + rules retention.Rules, now time.Time, c database.RetentionCandidate) bool { + key, ok := reg.KeyForDBEcosystem(c.Ecosystem) + if !ok { + return false + } + d := rules.For(key, c.PackagePURL) + if d <= 0 { + return false + } + cutoff := now.Add(-d) + if !c.ExpiredBefore(cutoff) { + return false + } + + recordCleared, err := db.ClearExpiredArtifact(c.VersionPURL, c.Filename, c.StoragePath, cutoff) + if err != nil { + logger.Warn("retention: failed to clear artifact record", + "version_purl", c.VersionPURL, "filename", c.Filename, "error", err) + return false + } + if !recordCleared { + return false + } + if err := db.QueuePendingDelete(c.StoragePath); err != nil { + logger.Warn("retention: failed to queue object for deletion", "path", c.StoragePath, "error", err) + } + return true +} + +// recordLRUEviction counts an artifact the size limit evicted under the +// ecosystem of its package record. +func recordLRUEviction(db *database.DB, logger *slog.Logger, versionPURL string) { + ecosystem, err := db.GetArtifactEcosystem(versionPURL) + if err != nil { + logger.Debug("eviction: failed to look up ecosystem for metrics", "version_purl", versionPURL, "error", err) + } + metrics.RecordArtifactEvicted("lru", ecosystem) +} diff --git a/internal/server/retention_test.go b/internal/server/retention_test.go new file mode 100644 index 0000000..9f1227a --- /dev/null +++ b/internal/server/retention_test.go @@ -0,0 +1,379 @@ +package server + +import ( + "context" + "database/sql" + "fmt" + "slices" + "strings" + "testing" + "time" + + "github.com/git-pkgs/proxy/internal/database" + "github.com/git-pkgs/proxy/internal/metrics" + "github.com/git-pkgs/proxy/internal/retention" + "github.com/git-pkgs/proxy/internal/storage" + "github.com/prometheus/client_golang/prometheus/testutil" +) + +const retentionDay = 24 * time.Hour + +// retentionRegistry registers npm and gem the way their handlers would, in a +// registry of the test's own so the handlers' registrations stay untouched. +func retentionRegistry() *retention.Registry { + reg := retention.NewRegistry() + reg.Register(retention.Spec{Key: "npm", CanonicalPackage: retention.DefaultCanonical("npm")}) + reg.Register(retention.Spec{Key: "gem", CanonicalPackage: retention.DefaultCanonical("gem")}) + return reg +} + +// seedAgedArtifact stores a cached artifact of pkg:{ecosystem}/{name}, with +// the package record's ecosystem set to dbEcosystem, fetched and last +// accessed at the given times. A zero accessed leaves it never served. +func seedAgedArtifact(t *testing.T, db *database.DB, store storage.Storage, dbEcosystem, purlType, name string, + fetched, accessed time.Time) string { + t.Helper() + pkgPURL := "pkg:" + purlType + "/" + name + versionPURL := pkgPURL + "@1.0.0" + filename := name + "-1.0.0.tgz" + if err := db.UpsertPackage(&database.Package{PURL: pkgPURL, Ecosystem: dbEcosystem, Name: name}); err != nil { + t.Fatalf("UpsertPackage: %v", err) + } + if err := db.UpsertVersion(&database.Version{PURL: versionPURL, PackagePURL: pkgPURL}); err != nil { + t.Fatalf("UpsertVersion: %v", err) + } + path := storage.ArtifactPath(dbEcosystem, "", name, "1.0.0", filename) + size, hash, err := store.Store(context.Background(), path, strings.NewReader("artifact")) + if err != nil { + t.Fatalf("Store: %v", err) + } + art := &database.Artifact{ + VersionPURL: versionPURL, + Filename: filename, + UpstreamURL: "https://example.com/" + filename, + StoragePath: sql.NullString{String: path, Valid: true}, + ContentHash: sql.NullString{String: hash, Valid: true}, + Size: sql.NullInt64{Int64: size, Valid: true}, + FetchedAt: sql.NullTime{Time: fetched, Valid: true}, + } + if !accessed.IsZero() { + art.LastAccessedAt = sql.NullTime{Time: accessed, Valid: true} + } + if err := db.UpsertArtifact(art); err != nil { + t.Fatalf("UpsertArtifact: %v", err) + } + return path +} + +func isCached(t *testing.T, db *database.DB, purlType, name string) bool { + t.Helper() + art, err := db.GetArtifact("pkg:"+purlType+"/"+name+"@1.0.0", name+"-1.0.0.tgz") + if err != nil { + t.Fatalf("GetArtifact: %v", err) + } + return art.StoragePath.Valid +} + +func TestSweepRetentionAppliesEachRule(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + ago := func(days int) time.Time { return now.Add(-time.Duration(days) * retentionDay) } + + rules := retention.Rules{ + Default: 30 * retentionDay, + Ecosystems: map[string]time.Duration{"npm": 14 * retentionDay}, + Packages: map[string]time.Duration{ + "pkg:npm/pinned": 0, + "pkg:npm/long-lived": 90 * retentionDay, + "pkg:npm/long-lived-expired": 90 * retentionDay, + "pkg:gem/short": 2 * retentionDay, + }, + } + + cases := []struct { + dbEco, purlType, name string + fetched, accessed time.Time + wantCached bool + }{ + {"npm", "npm", "stale", ago(60), ago(20), false}, + {"npm", "npm", "fresh", ago(60), ago(5), true}, + {"npm", "npm", "never-served", ago(20), time.Time{}, false}, + {"npm", "npm", "pinned", ago(400), ago(400), true}, + {"npm", "npm", "long-lived", ago(60), ago(60), true}, + {"npm", "npm", "long-lived-expired", ago(100), ago(100), false}, + {"gem", "gem", "default-applies", ago(40), ago(40), false}, + {"gem", "gem", "within-default", ago(40), ago(20), true}, + {"rubygems", "gem", "mirrored", ago(40), time.Time{}, false}, + {"gem", "gem", "short", ago(3), ago(3), false}, + {"pypi", "pypi", "unregistered", ago(400), ago(400), true}, + } + for _, c := range cases { + seedAgedArtifact(t, db, store, c.dbEco, c.purlType, c.name, c.fetched, c.accessed) + } + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now) + + for _, c := range cases { + if got := isCached(t, db, c.purlType, c.name); got != c.wantCached { + t.Errorf("%s cached = %v, want %v", c.name, got, c.wantCached) + } + } +} + +func TestSweepRetentionQueuesInsteadOfDeleting(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + old := now.Add(-30 * retentionDay) + path := seedAgedArtifact(t, db, store, "npm", "npm", "stale", old, old) + + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now) + + if isCached(t, db, "npm", "stale") { + t.Fatal("stale artifact still cached") + } + if !objectExists(t, store, path) { + t.Error("object deleted directly instead of waiting in the queue") + } + if got := queuedPaths(t, db); !slices.Equal(got, []string{path}) { + t.Errorf("queue = %v, want [%s]", got, path) + } +} + +// An artifact fetched again after an eviction keeps the access time of its +// earlier copy and must not be cleared again by the next sweep. +func TestSweepRetentionKeepsRefetchedArtifact(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + old := now.Add(-30 * retentionDay) + seedAgedArtifact(t, db, store, "npm", "npm", "popular", old, old) + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now) + if isCached(t, db, "npm", "popular") { + t.Fatal("first sweep kept the expired artifact") + } + + // Refetch: the upsert sets a new fetch time but leaves last_accessed_at. + seedAgedArtifact(t, db, store, "npm", "npm", "popular", now, time.Time{}) + art, err := db.GetArtifact("pkg:npm/popular@1.0.0", "popular-1.0.0.tgz") + if err != nil { + t.Fatalf("GetArtifact: %v", err) + } + if !art.LastAccessedAt.Valid || art.LastAccessedAt.Time.After(old.Add(time.Minute)) { + t.Fatalf("test setup: last_accessed_at = %v, want the old access kept", art.LastAccessedAt) + } + + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now.Add(time.Hour)) + if !isCached(t, db, "npm", "popular") { + t.Error("refetched artifact cleared again by the next sweep") + } +} + +func TestSweepRetentionFlushesBufferedHits(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + old := now.Add(-30 * retentionDay) + seedAgedArtifact(t, db, store, "npm", "npm", "busy", old, old) + db.BatchHits(time.Hour, nil) + if err := db.RecordArtifactHit("pkg:npm/busy@1.0.0", "busy-1.0.0.tgz"); err != nil { + t.Fatalf("RecordArtifactHit: %v", err) + } + + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, time.Now()) + + if !isCached(t, db, "npm", "busy") { + t.Error("artifact with a buffered hit was cleared") + } +} + +func TestSweepRetentionCapsClearsPerSweep(t *testing.T) { + db, store := setupEvictionTest(t) + prev := retentionMaxClears + retentionMaxClears = 2 + t.Cleanup(func() { retentionMaxClears = prev }) + + now := time.Now() + old := now.Add(-30 * retentionDay) + names := []string{"a", "b", "c"} + for _, name := range names { + seedAgedArtifact(t, db, store, "npm", "npm", name, old, old) + } + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + + countCached := func() int { + n := 0 + for _, name := range names { + if isCached(t, db, "npm", name) { + n++ + } + } + return n + } + + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now) + if got := countCached(); got != 1 { + t.Fatalf("cached after the first sweep = %d, want 1 (cap of 2)", got) + } + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now) + if got := countCached(); got != 0 { + t.Errorf("cached after the second sweep = %d, want 0", got) + } +} + +func TestSweepRetentionStopsOnCancel(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + old := now.Add(-30 * retentionDay) + seedAgedArtifact(t, db, store, "npm", "npm", "stale", old, old) + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + sweepRetention(ctx, db, discardLogger(), retentionRegistry(), rules, now) + + if !isCached(t, db, "npm", "stale") { + t.Error("cancelled sweep cleared an artifact") + } +} + +func TestSweepRetentionWithoutRulesDoesNothing(t *testing.T) { + db, store := setupEvictionTest(t) + old := time.Now().Add(-400 * retentionDay) + seedAgedArtifact(t, db, store, "npm", "npm", "ancient", old, old) + + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), retention.Rules{}, time.Now()) + + if !isCached(t, db, "npm", "ancient") { + t.Error("sweep without rules cleared an artifact") + } +} + +func TestNotCovered(t *testing.T) { + got := notCovered([]string{"npm", "pypi"}) + if slices.Contains(got, "npm") || slices.Contains(got, "pypi") { + t.Errorf("notCovered = %v, includes a covered key", got) + } + if len(got) != len(retention.KnownKeys)-2 { + t.Errorf("notCovered returned %d keys, want %d", len(got), len(retention.KnownKeys)-2) + } +} + +func TestReclaimDueDrainsSeveralBatches(t *testing.T) { + db, store := setupEvictionTest(t) + const queued = reclaimBatch*2 + 10 + for i := range queued { + storeQueued(t, db, store, fmt.Sprintf("npm/old/1.0.0/%d/old.tgz", i)) + } + + reclaimDue(context.Background(), db, store, discardLogger(), time.Now().Add(time.Hour), time.Minute) + + remaining, err := db.GetDuePendingDeletes(time.Now().Add(time.Hour), queued) + if err != nil { + t.Fatalf("GetDuePendingDeletes: %v", err) + } + if len(remaining) != 0 { + t.Errorf("%d objects still queued, want the whole queue reclaimed in one tick", len(remaining)) + } +} + +func TestReclaimDueRespectsBudget(t *testing.T) { + db, store := setupEvictionTest(t) + storeQueued(t, db, store, "npm/old/1.0.0/a/old.tgz") + + reclaimDue(context.Background(), db, store, discardLogger(), time.Now().Add(time.Hour), 0) + + if got := queuedPaths(t, db); len(got) != 1 { + t.Errorf("queue = %v, want the entry kept with no budget", got) + } +} + +// A hit that lands after the scan but before the clear keeps the artifact: +// the clear re-checks the artifact against its own rule's cutoff. +func TestClearIfExpiredRechecksAgainstItsRule(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + old := now.Add(-30 * retentionDay) + seedAgedArtifact(t, db, store, "npm", "npm", "racing", old, old) + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + + candidates, err := db.GetRetentionCandidates(0, 10, now.Add(-retentionDay)) + if err != nil || len(candidates) != 1 { + t.Fatalf("GetRetentionCandidates = %v, %v, want one candidate", candidates, err) + } + if err := db.RecordArtifactHit("pkg:npm/racing@1.0.0", "racing-1.0.0.tgz"); err != nil { + t.Fatalf("RecordArtifactHit: %v", err) + } + + if clearIfExpired(db, discardLogger(), retentionRegistry(), rules, now.Add(time.Minute), candidates[0]) { + t.Error("cleared an artifact served after the scan") + } + if !isCached(t, db, "npm", "racing") { + t.Error("artifact served after the scan is no longer cached") + } +} + +// When buffered hits cannot be written, the sweep must not run on the stale +// access times they would have replaced. +func TestSweepRetentionSkipsWhenHitsCannotBeWritten(t *testing.T) { + db, store := setupEvictionTest(t) + now := time.Now() + old := now.Add(-30 * retentionDay) + seedAgedArtifact(t, db, store, "npm", "npm", "busy", old, old) + db.BatchHits(time.Hour, nil) + if err := db.RecordArtifactHit("pkg:npm/busy@1.0.0", "busy-1.0.0.tgz"); err != nil { + t.Fatalf("RecordArtifactHit: %v", err) + } + // Refuse hit writes only; clearing a record leaves hit_count alone. + if _, err := db.Exec(`CREATE TRIGGER refuse_hits BEFORE UPDATE ON artifacts + WHEN NEW.hit_count > OLD.hit_count BEGIN SELECT RAISE(ABORT, 'hits refused'); END`); err != nil { + t.Fatalf("creating trigger: %v", err) + } + + rules := retention.Rules{Ecosystems: map[string]time.Duration{"npm": retentionDay}} + sweepRetention(context.Background(), db, discardLogger(), retentionRegistry(), rules, now) + + if !isCached(t, db, "npm", "busy") { + t.Error("sweep cleared an artifact whose buffered hit could not be written") + } + if _, err := db.Exec(`DROP TRIGGER refuse_hits`); err != nil { + t.Fatalf("dropping trigger: %v", err) + } +} + +func TestEvictionsAreCounted(t *testing.T) { + db, store := setupEvictionTest(t) + ctx := context.Background() + now := time.Now() + old := now.Add(-30 * retentionDay) + + retained := testutil.ToFloat64(metrics.ArtifactsEvicted.WithLabelValues("retention", "rubygems")) + seedAgedArtifact(t, db, store, "gem", "gem", "expired", old, old) + rules := retention.Rules{Ecosystems: map[string]time.Duration{"gem": retentionDay}} + sweepRetention(ctx, db, discardLogger(), retentionRegistry(), rules, now) + if got := testutil.ToFloat64(metrics.ArtifactsEvicted.WithLabelValues("retention", "rubygems")) - retained; got != 1 { + t.Errorf("retention evictions counted = %v, want 1", got) + } + + lru := testutil.ToFloat64(metrics.ArtifactsEvicted.WithLabelValues("lru", "npm")) + seedArtifact(t, ctx, db, store, "big", 500, now) + evictLRU(ctx, db, store, discardLogger(), 100) + if got := testutil.ToFloat64(metrics.ArtifactsEvicted.WithLabelValues("lru", "npm")) - lru; got != 1 { + t.Errorf("lru evictions counted = %v, want 1", got) + } +} + +func TestReclaimDueStopsAfterAFailedDelete(t *testing.T) { + db, store := setupEvictionTest(t) + queued := reclaimBatch * 2 + for i := range queued { + storeQueued(t, db, store, fmt.Sprintf("npm/old/1.0.0/%d/old.tgz", i)) + } + + undeletable := &undeletableStorage{Storage: store} + reclaimDue(context.Background(), db, undeletable, discardLogger(), time.Now().Add(time.Hour), time.Minute) + + if got := undeletable.deletes.Load(); got != reclaimBatch { + t.Errorf("delete attempts = %d, want one batch of %d during an outage", got, reclaimBatch) + } +} diff --git a/internal/server/server.go b/internal/server/server.go index 66ed14b..0026d25 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -374,6 +374,7 @@ func (s *Server) serve(listener net.Listener) error { "database", s.cfg.Database.String()) go s.updateCacheStatsMetrics() go s.startEvictionLoop(bgCtx) + go s.startRetentionLoop(bgCtx) go s.startReclaimLoop(bgCtx) if listener != nil {