diff --git a/README.md b/README.md index 491c6a81..38a588be 100644 --- a/README.md +++ b/README.md @@ -1139,10 +1139,16 @@ The proxy exposes Prometheus metrics at `GET /metrics`. All metric names are pre | `proxy_health_probe_failures_total` | counter | `step` | Storage health probe failures by failing step (`write`, `size`, `read`, `verify`, `delete`). | | `proxy_circuit_breaker_state` | gauge | `registry` | Artifact-fetch circuit breaker state per upstream registry (0 closed, 2 open). Published once that registry's breaker has tripped. | | `proxy_circuit_breaker_trips_total` | counter | `registry` | Circuit breaker trips per upstream registry. | +| `proxy_ecosystem_downloaded_bytes` | gauge | `ecosystem` | Accumulated bytes served from cache: cache hits multiplied by the artifact size they served. | +| `proxy_ecosystem_artifact_downloads` | gauge | `ecosystem` | Accumulated artifact downloads served from cache. | +| `proxy_ecosystem_cache_size_bytes` | gauge | `ecosystem` | Size of cached artifacts per ecosystem. | +| `proxy_ecosystem_cached_artifacts` | gauge | `ecosystem` | Number of cached artifacts per ecosystem. | +| `proxy_ecosystem_packages` | gauge | `ecosystem` | Known packages per ecosystem. | +| `proxy_ecosystem_versions` | gauge | `ecosystem` | Known package versions per ecosystem. | The `ecosystem` label on `proxy_requests_total` and `proxy_request_duration_seconds` is the mounted route a request arrived on, not the ecosystem recorded against the package it served: `/gem` reports as `rubygems`, `/go` as `golang`, `/composer` as `packagist`, `/apk` as `alpine` and `/v2` as `oci`. Anything outside a package route -- the UI, `/health`, `/metrics`, `/stats` -- reports as `other`, and so did `/apk`, `/helm`, `/homebrew`, `/generic` and `/swift` before they were listed; traffic on those five routes now appears under its own name instead. -Cache size and artifact count are refreshed every 60 seconds. Circuit breaker state is read from the fetcher on each scrape of `/metrics` and each `/health` request, so `proxy_circuit_breaker_trips_total` counts the trips visible between those reads — a breaker that opens and recovers entirely between two scrapes is not counted. The remaining metrics update on each request. +Cache size, artifact count and the per-ecosystem gauges are refreshed every 60 seconds, from a single pass over the database. Circuit breaker state is read from the fetcher on each scrape of `/metrics` and each `/health` request, so `proxy_circuit_breaker_trips_total` counts the trips visible between those reads — a breaker that opens and recovers entirely between two scrapes is not counted. The remaining metrics update on each request. The breaker metrics carry one series per upstream host, but only for hosts whose breaker has tripped at least once since startup. A breaker is created per host the proxy fetches artifacts from, and for some ecosystems that host comes from upstream metadata rather than from configuration (composer takes it from a package's `dist.url`, helm from the chart URLs in `index.yaml`), so publishing every host would let upstream content grow the series count for the lifetime of the process. Once a host has tripped it keeps reporting, so a recovery still shows up as a transition to 0 rather than as a series that vanishes. `/health` is not a persistent time series and lists every breaker, tripped or not. @@ -1150,6 +1156,18 @@ The `registry` label is the host of the URL the artifact was fetched from. Becau Alert on `proxy_circuit_breaker_state == 2` sustained for more than a few minutes: while a breaker is open, artifact downloads for that upstream fail with HTTP 502 on every cache miss, and only a single probe request per backoff interval reaches the upstream. Cached artifacts keep serving, and so does metadata for the same ecosystem (metadata does not go through the circuit breaker), so installs fail in a way that looks like a partial upstream outage. +#### Accumulated download size + +`proxy_ecosystem_downloaded_bytes` is, for every cached artifact, the number of times it was served multiplied by its size. It answers "how much traffic has this proxy actually carried", which is the number that matters when sizing egress or justifying the cache. Two properties are worth knowing before alerting on it. + +**It counts cache hits, not upstream fetches.** The request that first pulls an artifact through the proxy is a miss and is not counted; only later hits are. So the accumulated total is also the upstream bandwidth the cache has saved, not the total bytes the proxy has ever sent. + +**Eviction removes history.** Evicting an artifact clears its size, so its past hits drop out of the total. That is why these are gauges rather than counters, and why the figure can step downwards. Chart them with `max_over_time` rather than `increase`, and read a drop after an eviction sweep as expected rather than as data loss. + +Nothing at the schema level ties `artifacts.version_purl` to a version row, so a cached artifact can end up with no ecosystem to attribute it to. Those are reported under the ecosystem `unattributed` rather than dropped, which keeps the per-ecosystem figures adding up to `proxy_cache_size_bytes` and `proxy_cached_artifacts_total`. A non-zero `unattributed` means the database holds artifact rows whose version or package rows have gone missing. + +The `ecosystem` label on these six is taken from the package record and normalized, so aliases collapse: a database carrying both `gem` and `rubygems` rows -- the proxy writes the former, git-pkgs the latter -- reports one `rubygems` series with the two summed. + ### Health Check `/health` returns a structured JSON report of subsystem health. HTTP 200 if all checks pass; 503 if any fail. diff --git a/internal/database/analytics_test.go b/internal/database/analytics_test.go new file mode 100644 index 00000000..269b9d2f --- /dev/null +++ b/internal/database/analytics_test.go @@ -0,0 +1,375 @@ +package database + +import ( + "database/sql" + "path/filepath" + "testing" + "time" + + "github.com/git-pkgs/purl" +) + +// seedArtifact inserts a package, a version and a cached artifact so that the +// ecosystem aggregation has a complete join path to walk. +func seedArtifact(t *testing.T, db *DB, ecosystem, name, version string, size, hits int64) { + t.Helper() + + pkgPURL := "pkg:" + ecosystem + "/" + name + versionPURL := pkgPURL + "@" + version + + if err := db.UpsertPackage(&Package{PURL: pkgPURL, Ecosystem: ecosystem, Name: name}); err != nil { + t.Fatalf("UpsertPackage(%s): %v", pkgPURL, err) + } + if err := db.UpsertVersion(&Version{PURL: versionPURL, PackagePURL: pkgPURL}); err != nil { + t.Fatalf("UpsertVersion(%s): %v", versionPURL, err) + } + if err := db.UpsertArtifact(&Artifact{ + VersionPURL: versionPURL, + Filename: name + "-" + version + ".tgz", + UpstreamURL: "https://example.test/" + name, + StoragePath: sql.NullString{String: "objects/" + name, Valid: true}, + Size: sql.NullInt64{Int64: size, Valid: true}, + FetchedAt: sql.NullTime{Time: time.Now(), Valid: true}, + HitCount: hits, + }); err != nil { + t.Fatalf("UpsertArtifact(%s): %v", versionPURL, err) + } +} + +func newTestDB(t *testing.T) *DB { + t.Helper() + db, err := Create(filepath.Join(t.TempDir(), "analytics.db")) + if err != nil { + t.Fatalf("Create: %v", err) + } + t.Cleanup(func() { _ = db.Close() }) + return db +} + +func TestGetEcosystemStats(t *testing.T) { + db := newTestDB(t) + + // npm: 2 artifacts across 2 packages, 1000*3 + 500*1 = 3500 bytes served. + seedArtifact(t, db, "npm", "lodash", "4.17.21", 1000, 3) + seedArtifact(t, db, "npm", "express", "4.18.2", 500, 1) + // cargo: a single heavily-hit artifact, 200*10 = 2000 bytes served. + seedArtifact(t, db, "cargo", "serde", "1.0.0", 200, 10) + + stats, err := db.GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + if len(stats) != 2 { + t.Fatalf("expected 2 ecosystems, got %d: %+v", len(stats), stats) + } + + // Ordered by accumulated download volume, so npm (3500) precedes cargo (2000). + npm, cargo := stats[0], stats[1] + if npm.Ecosystem != "npm" || cargo.Ecosystem != "cargo" { + t.Fatalf("expected npm then cargo, got %q then %q", npm.Ecosystem, cargo.Ecosystem) + } + + if npm.DownloadedBytes != 3500 { + t.Errorf("npm DownloadedBytes = %d, want 3500", npm.DownloadedBytes) + } + if npm.Downloads != 4 { + t.Errorf("npm Downloads = %d, want 4", npm.Downloads) + } + if npm.CacheSize != 1500 { + t.Errorf("npm CacheSize = %d, want 1500", npm.CacheSize) + } + if npm.Artifacts != 2 { + t.Errorf("npm Artifacts = %d, want 2", npm.Artifacts) + } + if npm.Packages != 2 { + t.Errorf("npm Packages = %d, want 2", npm.Packages) + } + if npm.Versions != 2 { + t.Errorf("npm Versions = %d, want 2", npm.Versions) + } + + if cargo.DownloadedBytes != 2000 { + t.Errorf("cargo DownloadedBytes = %d, want 2000", cargo.DownloadedBytes) + } + if cargo.Downloads != 10 { + t.Errorf("cargo Downloads = %d, want 10", cargo.Downloads) + } +} + +func TestGetEcosystemStatsEmptyDatabase(t *testing.T) { + stats, err := newTestDB(t).GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + if len(stats) != 0 { + t.Errorf("expected no ecosystems, got %+v", stats) + } +} + +// An ecosystem known from metadata but with nothing cached should still appear, +// so the analytics table can show it as idle rather than omitting it. +func TestGetEcosystemStatsIncludesEcosystemsWithoutArtifacts(t *testing.T) { + db := newTestDB(t) + seedArtifact(t, db, "npm", "lodash", "4.17.21", 1000, 2) + + if err := db.UpsertPackage(&Package{PURL: "pkg:gem/rails", Ecosystem: "gem", Name: "rails"}); err != nil { + t.Fatalf("UpsertPackage: %v", err) + } + + stats, err := db.GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + if len(stats) != 2 { + t.Fatalf("expected 2 ecosystems, got %d: %+v", len(stats), stats) + } + + // gem has no download volume, so it sorts last. + gem := stats[1] + if gem.Ecosystem != "gem" { + t.Fatalf("expected gem last, got %q", gem.Ecosystem) + } + if gem.Packages != 1 { + t.Errorf("gem Packages = %d, want 1", gem.Packages) + } + if gem.Artifacts != 0 || gem.CacheSize != 0 || gem.DownloadedBytes != 0 { + t.Errorf("expected gem to report no cached artifacts, got %+v", gem) + } +} + +// Eviction clears storage_path and size, which drops the artifact's historical +// hits out of the accumulated total. Pinning that here so the documented +// behaviour of GetEcosystemStats does not drift. +func TestGetEcosystemStatsExcludesEvictedArtifacts(t *testing.T) { + db := newTestDB(t) + seedArtifact(t, db, "npm", "lodash", "4.17.21", 1000, 3) + seedArtifact(t, db, "npm", "express", "4.18.2", 500, 2) + + cleared, err := db.ClearArtifactCache("pkg:npm/lodash@4.17.21", "lodash-4.17.21.tgz", "objects/lodash") + if err != nil { + t.Fatalf("ClearArtifactCache: %v", err) + } + if !cleared { + t.Fatal("ClearArtifactCache cleared nothing") + } + + stats, err := db.GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + if len(stats) != 1 { + t.Fatalf("expected 1 ecosystem, got %d: %+v", len(stats), stats) + } + + // Only express remains cached: 500 * 2 = 1000. + if got := stats[0].DownloadedBytes; got != 1000 { + t.Errorf("DownloadedBytes = %d, want 1000 (evicted artifact must not count)", got) + } + if got := stats[0].Artifacts; got != 1 { + t.Errorf("Artifacts = %d, want 1", got) + } + // The package row survives eviction, so both packages still count. + if got := stats[0].Packages; got != 2 { + t.Errorf("Packages = %d, want 2", got) + } +} + +func TestSortEcosystemStatsTieBreaks(t *testing.T) { + stats := []EcosystemStats{ + {Ecosystem: "zzz"}, + {Ecosystem: "aaa"}, + {Ecosystem: "mid", CacheSize: 10}, + {Ecosystem: "top", DownloadedBytes: 5}, + } + sortEcosystemStats(stats) + + want := []string{"top", "mid", "aaa", "zzz"} + for i, name := range want { + if stats[i].Ecosystem != name { + t.Errorf("position %d = %q, want %q (full order: %+v)", i, stats[i].Ecosystem, name, stats) + } + } +} + +// A cached artifact whose version row is missing must still be counted, so the +// per-ecosystem totals reconcile with GetTotalCacheSize rather than quietly +// coming up short. Nothing at the schema level enforces the link. +func TestGetEcosystemStatsCountsUnattributedArtifacts(t *testing.T) { + db := newTestDB(t) + seedArtifact(t, db, "npm", "lodash", "4.17.21", 1000, 2) + + // An artifact pointing at a version that was never recorded. + if err := db.UpsertArtifact(&Artifact{ + VersionPURL: "pkg:npm/orphan@9.9.9", + Filename: "orphan-9.9.9.tgz", + UpstreamURL: "https://example.test/orphan", + StoragePath: sql.NullString{String: "objects/orphan", Valid: true}, + Size: sql.NullInt64{Int64: 500, Valid: true}, + FetchedAt: sql.NullTime{Time: time.Now(), Valid: true}, + HitCount: 4, + }); err != nil { + t.Fatalf("UpsertArtifact: %v", err) + } + + stats, err := db.GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + + var cacheSize, artifacts, downloaded int64 + var sawUnattributed bool + for _, e := range stats { + cacheSize += e.CacheSize + artifacts += e.Artifacts + downloaded += e.DownloadedBytes + if e.Ecosystem == unattributedEcosystem { + sawUnattributed = true + } + } + + if !sawUnattributed { + t.Errorf("the orphaned artifact was dropped; got %+v", stats) + } + + // These must match what the unjoined queries report, or the UI shows two + // totals that do not add up. + wantSize, err := db.GetTotalCacheSize() + if err != nil { + t.Fatalf("GetTotalCacheSize: %v", err) + } + if cacheSize != wantSize { + t.Errorf("per-ecosystem cache size sums to %d, but GetTotalCacheSize reports %d", cacheSize, wantSize) + } + + wantCount, err := db.GetCachedArtifactCount() + if err != nil { + t.Fatalf("GetCachedArtifactCount: %v", err) + } + if artifacts != wantCount { + t.Errorf("per-ecosystem artifacts sum to %d, but GetCachedArtifactCount reports %d", artifacts, wantCount) + } + + // 1000*2 from lodash plus 500*4 from the orphan. + if downloaded != 4000 { + t.Errorf("downloaded bytes = %d, want 4000", downloaded) + } +} + +// The proxy writes "gem" and git-pkgs writes "rubygems". A database that has +// seen both must report one ecosystem, not two rows splitting its share. +func TestGetEcosystemStatsMergesAliasedEcosystems(t *testing.T) { + db := newTestDB(t) + seedArtifact(t, db, "gem", "colorize", "1.1.0", 1000, 3) + seedArtifact(t, db, "rubygems", "rails", "7.1.0", 4000, 2) + seedArtifact(t, db, "npm", "lodash", "4.17.21", 500, 1) + + stats, err := db.GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + if len(stats) != 2 { + t.Fatalf("expected gem and rubygems merged into one row alongside npm, got %+v", stats) + } + + var ruby *EcosystemStats + for i := range stats { + if purl.NormalizeEcosystem(stats[i].Ecosystem) == "rubygems" { + ruby = &stats[i] + } + } + if ruby == nil { + t.Fatal("no row for the rubygems ecosystem") + } + + // 1000*3 + 4000*2 = 11000 + if ruby.DownloadedBytes != 11000 { + t.Errorf("DownloadedBytes = %d, want 11000", ruby.DownloadedBytes) + } + if ruby.CacheSize != 5000 { + t.Errorf("CacheSize = %d, want 5000", ruby.CacheSize) + } + if ruby.Artifacts != 2 || ruby.Packages != 2 { + t.Errorf("Artifacts=%d Packages=%d, want 2 and 2", ruby.Artifacts, ruby.Packages) + } + + // The surviving row keeps the spelling with the most cached bytes, because + // the UI filters and links by it and the packages table stores it raw. + if ruby.Ecosystem != "rubygems" { + t.Errorf("Ecosystem = %q, want the dominant raw spelling %q", ruby.Ecosystem, "rubygems") + } +} + +// A version whose package row is missing must be bucketed like an orphaned +// artifact, so the per-ecosystem figures reconcile with GetCacheStats. +func TestGetEcosystemStatsCountsUnattributedVersions(t *testing.T) { + db := newTestDB(t) + seedArtifact(t, db, "npm", "lodash", "4.17.21", 1000, 1) + if err := db.UpsertVersion(&Version{PURL: "pkg:npm/ghost@1.0.0", PackagePURL: "pkg:npm/ghost"}); err != nil { + t.Fatalf("UpsertVersion: %v", err) + } + + stats, err := db.GetEcosystemStats() + if err != nil { + t.Fatalf("GetEcosystemStats: %v", err) + } + + var versions int64 + for _, e := range stats { + versions += e.Versions + } + + cacheStats, err := db.GetCacheStats() + if err != nil { + t.Fatalf("GetCacheStats: %v", err) + } + if versions != cacheStats.TotalVersions { + t.Errorf("per-ecosystem versions sum to %d, but GetCacheStats reports %d", + versions, cacheStats.TotalVersions) + } +} + +// The surviving spelling must be the one with the most cached bytes whatever +// order the rows arrive in. GetEcosystemStats builds its input by ranging a +// map, so an implementation that compared against the running total instead of +// each row's own size would pick a different name between refreshes — flipping +// the analytics row's badge and its /ui/packages filter link. +func TestMergeAliasedEcosystemsPicksLargestWhateverTheOrder(t *testing.T) { + rows := []EcosystemStats{ + {Ecosystem: "gem", CacheSize: 5, Artifacts: 1}, + {Ecosystem: "rubygems", CacheSize: 4, Artifacts: 1}, + {Ecosystem: "RubyGems", CacheSize: 6, Artifacts: 1}, + } + + for _, order := range [][]int{ + {0, 1, 2}, {0, 2, 1}, {1, 0, 2}, {1, 2, 0}, {2, 0, 1}, {2, 1, 0}, + } { + in := make([]EcosystemStats, 0, len(order)) + for _, i := range order { + in = append(in, rows[i]) + } + + got := mergeAliasedEcosystems(in) + if len(got) != 1 { + t.Fatalf("order %v: got %d rows, want 1", order, len(got)) + } + if got[0].Ecosystem != "RubyGems" { + t.Errorf("order %v: Ecosystem = %q, want the largest spelling %q", + order, got[0].Ecosystem, "RubyGems") + } + if got[0].CacheSize != 15 { + t.Errorf("order %v: CacheSize = %d, want 15", order, got[0].CacheSize) + } + } +} + +// Equal sizes still have to resolve to one answer, for the same reason. +func TestMergeAliasedEcosystemsBreaksSizeTiesByName(t *testing.T) { + a := []EcosystemStats{{Ecosystem: "gem", CacheSize: 5}, {Ecosystem: "rubygems", CacheSize: 5}} + b := []EcosystemStats{{Ecosystem: "rubygems", CacheSize: 5}, {Ecosystem: "gem", CacheSize: 5}} + + first := mergeAliasedEcosystems(a)[0].Ecosystem + second := mergeAliasedEcosystems(b)[0].Ecosystem + if first != second { + t.Errorf("input order changed the surviving spelling: %q vs %q", first, second) + } +} diff --git a/internal/database/queries.go b/internal/database/queries.go index 4e568d69..49461330 100644 --- a/internal/database/queries.go +++ b/internal/database/queries.go @@ -4,9 +4,11 @@ import ( "database/sql" "errors" "fmt" + "sort" "time" "github.com/git-pkgs/artifacts" + "github.com/git-pkgs/purl" "github.com/opencontainers/go-digest" ) @@ -1132,3 +1134,202 @@ func (db *DB) UpsertMetadataCache(entry *MetadataCacheEntry) error { } return nil } + +// Analytics queries + +// EcosystemStats aggregates cache and download activity for one ecosystem. +// +// DownloadedBytes is the accumulated download volume: every cache hit on an +// artifact served its full size, so the sum of hit_count * size is the number +// of bytes the proxy has handed to clients from cache for this ecosystem. +type EcosystemStats struct { + Ecosystem string `db:"ecosystem"` + Packages int64 `db:"packages"` + Versions int64 `db:"versions"` + Artifacts int64 `db:"artifacts"` + CacheSize int64 `db:"cache_size"` + Downloads int64 `db:"downloads"` + DownloadedBytes int64 `db:"downloaded_bytes"` +} + +// GetEcosystemStats returns per-ecosystem cache and download totals, ordered by +// accumulated download volume descending. Ecosystems with rows in packages but +// nothing cached are included with zeroed artifact counters. +// +// Artifacts evicted from the cache no longer contribute: eviction clears the +// size column, so their historical hits drop out of the accumulated total. +func (db *DB) GetEcosystemStats() ([]EcosystemStats, error) { + byEcosystem := make(map[string]*EcosystemStats) + + get := func(ecosystem string) *EcosystemStats { + if s, ok := byEcosystem[ecosystem]; ok { + return s + } + s := &EcosystemStats{Ecosystem: ecosystem} + byEcosystem[ecosystem] = s + return s + } + + if err := db.eachCount(`SELECT ecosystem, COUNT(*) FROM packages GROUP BY ecosystem`, + func(ecosystem string, n int64) { get(ecosystem).Packages = n }); err != nil { + return nil, err + } + + // Left joined and bucketed for the same reason as artifacts below: a version + // whose package row is missing would otherwise vanish here while still + // counting in GetCacheStats' COUNT(*), leaving the two unable to reconcile. + if err := db.eachCount(` + SELECT COALESCE(p.ecosystem, '`+unattributedEcosystem+`'), COUNT(*) + FROM versions v + LEFT JOIN packages p ON p.purl = v.package_purl + GROUP BY COALESCE(p.ecosystem, '`+unattributedEcosystem+`') + `, func(ecosystem string, n int64) { get(ecosystem).Versions = n }); err != nil { + return nil, err + } + + // The artifacts table is proxy-specific: a database inherited from + // git-pkgs carries packages and versions without it. + hasArtifacts, err := db.HasTable("artifacts") + if err != nil { + return nil, err + } + if hasArtifacts { + if err := db.eachArtifactStat(get); err != nil { + return nil, err + } + } + + stats := make([]EcosystemStats, 0, len(byEcosystem)) + for _, s := range byEcosystem { + stats = append(stats, *s) + } + stats = mergeAliasedEcosystems(stats) + sortEcosystemStats(stats) + return stats, nil +} + +// mergeAliasedEcosystems combines rows whose ecosystem names normalize to the +// same canonical name. The proxy writes "gem" and git-pkgs writes "rubygems", +// so a database that has seen both carries two rows for one ecosystem; left +// split they would render as two table rows and two chart slices, each with +// half the real share. +// +// The surviving row keeps the raw spelling of whichever input held the most +// cached bytes, because that string is what the UI filters and links by: the +// packages table stores the raw value, so substituting the canonical name would +// produce links that match nothing. That comparison is against each input's own +// size, not the running total, which would otherwise let the first spelling win +// simply by being merged into first. +func mergeAliasedEcosystems(stats []EcosystemStats) []EcosystemStats { + merged := make(map[string]*EcosystemStats, len(stats)) + largest := make(map[string]int64, len(stats)) + order := make([]string, 0, len(stats)) + + for i := range stats { + key := purl.NormalizeEcosystem(stats[i].Ecosystem) + into, ok := merged[key] + if !ok { + row := stats[i] + merged[key] = &row + largest[key] = stats[i].CacheSize + order = append(order, key) + continue + } + + // The name tiebreak matters: GetEcosystemStats builds its input by + // ranging a map, so without it two spellings of equal size would swap + // between refreshes and flip the row's badge and filter link. + if stats[i].CacheSize > largest[key] || + (stats[i].CacheSize == largest[key] && stats[i].Ecosystem < into.Ecosystem) { + largest[key] = stats[i].CacheSize + into.Ecosystem = stats[i].Ecosystem + } + into.Packages += stats[i].Packages + into.Versions += stats[i].Versions + into.Artifacts += stats[i].Artifacts + into.CacheSize += stats[i].CacheSize + into.Downloads += stats[i].Downloads + into.DownloadedBytes += stats[i].DownloadedBytes + } + + out := make([]EcosystemStats, 0, len(order)) + for _, key := range order { + out = append(out, *merged[key]) + } + return out +} + +// unattributedEcosystem collects cached artifacts whose version or package row +// is missing. Nothing enforces that link at the schema level, so an inner join +// would silently drop such rows and leave the per-ecosystem totals short of +// GetTotalCacheSize — two numbers that sit side by side in the UI and in +// Grafana. Bucketing them keeps the two reconcilable. +const unattributedEcosystem = "unattributed" + +func (db *DB) eachArtifactStat(get func(string) *EcosystemStats) error { + rows, err := db.Query(` + SELECT COALESCE(p.ecosystem, '` + unattributedEcosystem + `'), + COUNT(*), + COALESCE(SUM(a.size), 0), + COALESCE(SUM(a.hit_count), 0), + COALESCE(SUM(a.hit_count * a.size), 0) + FROM artifacts a + LEFT JOIN versions v ON v.purl = a.version_purl + LEFT JOIN packages p ON p.purl = v.package_purl + WHERE a.storage_path IS NOT NULL + GROUP BY COALESCE(p.ecosystem, '` + unattributedEcosystem + `') + `) + if err != nil { + return err + } + defer func() { _ = rows.Close() }() + + for rows.Next() { + var ecosystem string + var artifacts, cacheSize, downloads, downloadedBytes int64 + if err := rows.Scan(&ecosystem, &artifacts, &cacheSize, &downloads, &downloadedBytes); err != nil { + return err + } + s := get(ecosystem) + s.Artifacts = artifacts + s.CacheSize = cacheSize + s.Downloads = downloads + s.DownloadedBytes = downloadedBytes + } + return rows.Err() +} + +// eachCount runs a two-column "group by" query and hands each (key, count) pair to fn. +func (db *DB) eachCount(query string, fn func(key string, n int64)) error { + rows, err := db.Query(query) + if err != nil { + return err + } + defer func() { _ = rows.Close() }() + + for rows.Next() { + var key string + var n int64 + if err := rows.Scan(&key, &n); err != nil { + return err + } + fn(key, n) + } + return rows.Err() +} + +// sortEcosystemStats orders by accumulated download volume, then by cache size, +// then by name, so that ecosystems with no traffic yet still sort predictably. +func sortEcosystemStats(stats []EcosystemStats) { + sort.Slice(stats, func(i, j int) bool { + a, b := stats[i], stats[j] + switch { + case a.DownloadedBytes != b.DownloadedBytes: + return a.DownloadedBytes > b.DownloadedBytes + case a.CacheSize != b.CacheSize: + return a.CacheSize > b.CacheSize + default: + return a.Ecosystem < b.Ecosystem + } + }) +} diff --git a/internal/metrics/analytics_test.go b/internal/metrics/analytics_test.go new file mode 100644 index 00000000..6d391bc6 --- /dev/null +++ b/internal/metrics/analytics_test.go @@ -0,0 +1,99 @@ +package metrics + +import ( + "testing" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/testutil" +) + +func TestUpdateEcosystemStats(t *testing.T) { + UpdateEcosystemStats([]EcosystemStats{ + {Ecosystem: "npm", Packages: 10, Versions: 25, Artifacts: 40, CacheSize: 1500, Downloads: 4, DownloadedBytes: 3500}, + {Ecosystem: "cargo", Packages: 2, Versions: 3, Artifacts: 3, CacheSize: 200, Downloads: 10, DownloadedBytes: 2000}, + }) + + if got := testutil.ToFloat64(EcosystemDownloadedBytes.WithLabelValues("npm")); got != 3500 { + t.Errorf("npm downloaded bytes = %v, want 3500", got) + } + if got := testutil.ToFloat64(EcosystemDownloads.WithLabelValues("npm")); got != 4 { + t.Errorf("npm downloads = %v, want 4", got) + } + if got := testutil.ToFloat64(EcosystemCacheSize.WithLabelValues("cargo")); got != 200 { + t.Errorf("cargo cache size = %v, want 200", got) + } + if got := testutil.ToFloat64(EcosystemCachedArtifacts.WithLabelValues("npm")); got != 40 { + t.Errorf("npm cached artifacts = %v, want 40", got) + } + if got := testutil.ToFloat64(EcosystemPackages.WithLabelValues("npm")); got != 10 { + t.Errorf("npm packages = %v, want 10", got) + } + if got := testutil.ToFloat64(EcosystemVersions.WithLabelValues("npm")); got != 25 { + t.Errorf("npm versions = %v, want 25", got) + } +} + +// A refresh that no longer mentions an ecosystem must drop its series rather +// than leave it frozen at the last observed value. +func TestUpdateEcosystemStatsDropsStaleSeries(t *testing.T) { + UpdateEcosystemStats([]EcosystemStats{ + {Ecosystem: "npm", DownloadedBytes: 3500}, + {Ecosystem: "gem", DownloadedBytes: 900}, + }) + if got := testutil.CollectAndCount(EcosystemDownloadedBytes); got != 2 { + t.Fatalf("expected 2 series after first refresh, got %d", got) + } + + UpdateEcosystemStats([]EcosystemStats{{Ecosystem: "npm", DownloadedBytes: 4000}}) + + if got := testutil.CollectAndCount(EcosystemDownloadedBytes); got != 1 { + t.Errorf("expected 1 series after gem disappeared, got %d", got) + } + if got := testutil.ToFloat64(EcosystemDownloadedBytes.WithLabelValues("npm")); got != 4000 { + t.Errorf("npm downloaded bytes = %v, want 4000", got) + } +} + +// Ecosystem labels are normalized so these gauges join against the other +// metrics that take their ecosystem from a package record. +func TestUpdateEcosystemStatsNormalizesLabels(t *testing.T) { + UpdateEcosystemStats([]EcosystemStats{{Ecosystem: "NPM", DownloadedBytes: 12}}) + + if got := testutil.ToFloat64(EcosystemDownloadedBytes.WithLabelValues("npm")); got != 12 { + t.Errorf("normalized npm gauge = %v, want 12", got) + } +} + +// GetEcosystemStats groups by the raw packages.ecosystem column, so a database +// carrying both spellings of an aliased ecosystem yields two rows that +// normalize to one label. They must sum rather than overwrite each other. +func TestUpdateEcosystemStatsSumsAliasedRows(t *testing.T) { + reg := prometheus.NewRegistry() + reg.MustRegister(EcosystemDownloadedBytes, EcosystemCacheSize, EcosystemPackages) + + UpdateEcosystemStats([]EcosystemStats{ + {Ecosystem: "gem", DownloadedBytes: 100, CacheSize: 10, Packages: 1}, + {Ecosystem: "rubygems", DownloadedBytes: 200, CacheSize: 20, Packages: 2}, + {Ecosystem: "npm", DownloadedBytes: 50, CacheSize: 5, Packages: 3}, + }) + + if got := testutil.ToFloat64(EcosystemDownloadedBytes.WithLabelValues("rubygems")); got != 300 { + t.Errorf("rubygems downloaded bytes = %v, want 300", got) + } + if got := testutil.ToFloat64(EcosystemCacheSize.WithLabelValues("rubygems")); got != 30 { + t.Errorf("rubygems cache size = %v, want 30", got) + } + if got := testutil.ToFloat64(EcosystemPackages.WithLabelValues("rubygems")); got != 3 { + t.Errorf("rubygems packages = %v, want 3", got) + } + // An unaliased ecosystem is unaffected. + if got := testutil.ToFloat64(EcosystemDownloadedBytes.WithLabelValues("npm")); got != 50 { + t.Errorf("npm downloaded bytes = %v, want 50", got) + } + + // A later refresh replaces the value rather than adding to it. + UpdateEcosystemStats([]EcosystemStats{{Ecosystem: "npm", DownloadedBytes: 50}}) + if got := testutil.ToFloat64(EcosystemDownloadedBytes.WithLabelValues("npm")); got != 50 { + t.Errorf("npm downloaded bytes after a second refresh = %v, want 50", got) + } +} diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index aff57682..2f20fb5a 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -4,6 +4,7 @@ package metrics import ( "net/http" "strconv" + "sync" "time" "github.com/git-pkgs/purl" @@ -138,6 +139,57 @@ var ( []string{"step"}, ) + // Per-ecosystem gauges, derived from the database rather than incremented + // in the request path, and refreshed on the same tick as the cache gauges + // above. + EcosystemDownloadedBytes = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "proxy_ecosystem_downloaded_bytes", + Help: "Accumulated bytes served from cache per ecosystem (cache hits x artifact size)", + }, + []string{"ecosystem"}, + ) + + EcosystemDownloads = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "proxy_ecosystem_artifact_downloads", + Help: "Accumulated artifact downloads served from cache per ecosystem", + }, + []string{"ecosystem"}, + ) + + EcosystemCacheSize = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "proxy_ecosystem_cache_size_bytes", + Help: "Size of cached artifacts per ecosystem in bytes", + }, + []string{"ecosystem"}, + ) + + EcosystemCachedArtifacts = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "proxy_ecosystem_cached_artifacts", + Help: "Number of cached artifacts per ecosystem", + }, + []string{"ecosystem"}, + ) + + EcosystemPackages = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "proxy_ecosystem_packages", + Help: "Number of known packages per ecosystem", + }, + []string{"ecosystem"}, + ) + + EcosystemVersions = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "proxy_ecosystem_versions", + Help: "Number of known package versions per ecosystem", + }, + []string{"ecosystem"}, + ) + // Scanning metrics ScanDuration = prometheus.NewHistogramVec( prometheus.HistogramOpts{ @@ -183,6 +235,12 @@ func init() { ActiveRequests, IntegrityFailures, HealthProbeFailures, + EcosystemDownloadedBytes, + EcosystemDownloads, + EcosystemCacheSize, + EcosystemCachedArtifacts, + EcosystemPackages, + EcosystemVersions, ScanDuration, ScanBlocked, ScanErrors, @@ -263,6 +321,79 @@ func UpdateCacheStats(sizeBytes, artifactCount int64) { CachedArtifacts.Set(float64(artifactCount)) } +// EcosystemStats is one ecosystem's row of the snapshot published as gauges. +type EcosystemStats struct { + Ecosystem string + Packages int64 + Versions int64 + Artifacts int64 + CacheSize int64 + Downloads int64 + DownloadedBytes int64 +} + +// publishedEcosystems tracks which labels the per-ecosystem gauges currently +// carry, so a label that disappears can be deleted individually. +var ( + publishedMu sync.Mutex + publishedEcosystems = map[string]bool{} +) + +// UpdateEcosystemStats republishes the per-ecosystem gauges from a fresh snapshot. +// +// Rows are summed by label before anything is published, because normalizing +// collapses aliases: a database holding both "gem" and "rubygems" rows -- the +// proxy writes the former, git-pkgs the latter -- arrives as two rows belonging +// to one label, and publishing them one at a time would leave only the last. +// +// The vectors are not Reset() first. Reset followed by a repopulating loop +// leaves a window in which a scrape sees the families empty or half filled, +// which renders as a spurious gap on any panel built from them. Each series is +// Set instead, and only labels that have actually disappeared are deleted. +func UpdateEcosystemStats(stats []EcosystemStats) { + totals := make(map[string]EcosystemStats, len(stats)) + for _, s := range stats { + ecosystem := purl.NormalizeEcosystem(s.Ecosystem) + t := totals[ecosystem] + t.Packages += s.Packages + t.Versions += s.Versions + t.Artifacts += s.Artifacts + t.CacheSize += s.CacheSize + t.Downloads += s.Downloads + t.DownloadedBytes += s.DownloadedBytes + totals[ecosystem] = t + } + + for ecosystem, t := range totals { + EcosystemDownloadedBytes.WithLabelValues(ecosystem).Set(float64(t.DownloadedBytes)) + EcosystemDownloads.WithLabelValues(ecosystem).Set(float64(t.Downloads)) + EcosystemCacheSize.WithLabelValues(ecosystem).Set(float64(t.CacheSize)) + EcosystemCachedArtifacts.WithLabelValues(ecosystem).Set(float64(t.Artifacts)) + EcosystemPackages.WithLabelValues(ecosystem).Set(float64(t.Packages)) + EcosystemVersions.WithLabelValues(ecosystem).Set(float64(t.Versions)) + } + + publishedMu.Lock() + defer publishedMu.Unlock() + + for ecosystem := range publishedEcosystems { + if _, still := totals[ecosystem]; still { + continue + } + EcosystemDownloadedBytes.DeleteLabelValues(ecosystem) + EcosystemDownloads.DeleteLabelValues(ecosystem) + EcosystemCacheSize.DeleteLabelValues(ecosystem) + EcosystemCachedArtifacts.DeleteLabelValues(ecosystem) + EcosystemPackages.DeleteLabelValues(ecosystem) + EcosystemVersions.DeleteLabelValues(ecosystem) + } + + publishedEcosystems = make(map[string]bool, len(totals)) + for ecosystem := range totals { + publishedEcosystems[ecosystem] = true + } +} + // UpdateCircuitBreakerState updates circuit breaker state gauge. // state: 0=closed, 1=half-open, 2=open func UpdateCircuitBreakerState(registry string, state int) { diff --git a/internal/server/server.go b/internal/server/server.go index 010c5d62..213cb3f2 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -516,6 +516,31 @@ func (s *Server) updateCacheStats() { return } metrics.UpdateCacheStats(stats.TotalSize, stats.TotalArtifacts) + + ecosystems, err := s.db.GetEcosystemStats() + if err != nil { + s.logger.Warn("failed to get ecosystem stats for metrics", "error", err) + return + } + metrics.UpdateEcosystemStats(ecosystemMetrics(ecosystems)) +} + +// ecosystemMetrics converts database rows into the metrics package's own +// snapshot type, so that package keeps no dependency on the database schema. +func ecosystemMetrics(stats []database.EcosystemStats) []metrics.EcosystemStats { + out := make([]metrics.EcosystemStats, 0, len(stats)) + for _, e := range stats { + out = append(out, metrics.EcosystemStats{ + Ecosystem: e.Ecosystem, + Packages: e.Packages, + Versions: e.Versions, + Artifacts: e.Artifacts, + CacheSize: e.CacheSize, + Downloads: e.Downloads, + DownloadedBytes: e.DownloadedBytes, + }) + } + return out } // Shutdown gracefully shuts down the server.