From 4ce6c37bd9b7872c2061d86bed407b56d2c4b4dc Mon Sep 17 00:00:00 2001 From: montehurd Date: Wed, 23 Sep 2026 14:40:13 -0700 Subject: [PATCH 1/2] Give each fetch its own storage path Every fetch of an artifact wrote one path, so fetches with different URLs or digests, which do not share a coalescing key, could overwrite each other's bytes, serve the wrong ones, or delete them on a digest mismatch while the other was about to open them. storeArtifact now writes {ecosystem}/{name}/{version}/{fetch id}/{filename}, with a random id per fetch. Most ecosystems learn the digest only after the fetch, and storage has no rename, so the id is the one rule that fits all of them. Existing records keep their paths and stay readable. On a file:// bucket, Delete now removes an emptied fetch directory, which fileblob leaves behind, and no other directory, since another fetch may be creating its own inside it. A path a record stops pointing at is queued in pending_deletes rather than deleted, since a request that read the record may still open it. A loop deletes queued paths after max(1h, direct_serve_ttl), whether or not max_size is set. UpsertArtifact and clears apply only if the record still holds the path the caller read, so racing commits each queue the path they replaced and a clear never orphans a newer commit. Eviction still deletes inline, and counts space as freed only when the record still pointed at what it deleted. --- docs/architecture.md | 20 ++- internal/database/database_test.go | 6 +- internal/database/pending_deletes_test.go | 199 ++++++++++++++++++++++ internal/database/queries.go | 113 +++++++++++- internal/database/schema.go | 28 +++ internal/handler/fetch_path_test.go | 163 ++++++++++++++++++ internal/handler/handler.go | 48 +++--- internal/handler/handler_test.go | 49 +++++- internal/handler/helm_test.go | 23 ++- internal/handler/stale_cache_test.go | 7 +- internal/server/eviction.go | 8 +- internal/server/eviction_test.go | 37 ++++ internal/server/reclaim.go | 59 +++++++ internal/server/reclaim_test.go | 97 +++++++++++ internal/server/server.go | 1 + internal/storage/blob.go | 31 +++- internal/storage/blob_test.go | 38 +++++ internal/storage/storage.go | 38 ++++- internal/storage/storage_test.go | 13 ++ 19 files changed, 918 insertions(+), 60 deletions(-) create mode 100644 internal/database/pending_deletes_test.go create mode 100644 internal/handler/fetch_path_test.go create mode 100644 internal/server/reclaim.go create mode 100644 internal/server/reclaim_test.go diff --git a/docs/architecture.md b/docs/architecture.md index 9e656ef4..42122c5a 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -146,6 +146,11 @@ artifacts ( ) -- indexes: (version_purl, filename) unique, storage_path, last_accessed_at +pending_deletes ( + path TEXT NOT NULL PRIMARY KEY, -- storage path no record points at + queued_at DATETIME NOT NULL +) + vulnerabilities ( id INTEGER PRIMARY KEY, vuln_id TEXT NOT NULL, -- e.g. CVE-2021-1234 @@ -209,10 +214,10 @@ type Storage interface { ``` **Filesystem implementation:** -- Stores files in nested directories: `{ecosystem}/{name}/{version}/{filename}` +- Stores files in nested directories: `{ecosystem}/{name}/{version}/{fetch id}/{filename}`, a new id per fetch, so fetches of one artifact never overwrite or delete each other's object (artifacts cached before this sit at `{ecosystem}/{name}/{version}/{filename}`) - Atomic writes using temp file + rename - Computes SHA256 hash during write -- Cleans up empty parent directories on delete +- Removes a fetch's directory once its object is deleted **Path structure:** @@ -221,15 +226,18 @@ cache/artifacts/ ├── npm/ │ ├── lodash/ │ │ └── 4.17.21/ -│ │ └── lodash-4.17.21.tgz +│ │ └── 3f9a0c1d2e4b5a67/ +│ │ └── lodash-4.17.21.tgz │ └── @babel/ │ └── core/ │ └── 7.23.0/ -│ └── core-7.23.0.tgz +│ └── 8c2e41f09a7d3b15/ +│ └── core-7.23.0.tgz └── cargo/ └── serde/ └── 1.0.193/ - └── serde-1.0.193.crate + └── d05b7e9c14a2f863/ + └── serde-1.0.193.crate ``` ### `internal/upstream` @@ -341,6 +349,8 @@ Eviction can be implemented as: 2. When over limit, get LRU artifacts 3. Delete from storage and clear database records +An object a record stops pointing at, because a refetch replaced it or its entry was discarded, is not deleted at once: a request that read the record may still be opening it. Its path goes into `pending_deletes`, and a background loop deletes it after a grace period of at least an hour, or `direct_serve_ttl` if longer, so signed URLs to it stay valid. This runs whether or not `max_size` is set. + ## Design Decisions **Why SQLite?** diff --git a/internal/database/database_test.go b/internal/database/database_test.go index bb73bc8a..13f6828c 100644 --- a/internal/database/database_test.go +++ b/internal/database/database_test.go @@ -613,7 +613,7 @@ func createTestPostgresDB(t *testing.T) *DB { // Drop and recreate every table CreateSchema creates for clean test state; // leftover migration records make the next CreateSchema fail on the // migrations primary key. - tables := []string{"artifacts", "versions", "packages", "vulnerabilities", "metadata_cache", "migrations", "schema_info"} + tables := []string{"artifacts", "pending_deletes", "versions", "packages", "vulnerabilities", "metadata_cache", "migrations", "schema_info"} for _, table := range tables { _, _ = db.Exec("DROP TABLE IF EXISTS " + table + " CASCADE") } @@ -781,6 +781,10 @@ func TestMigrationFromOldSchema(t *testing.T) { t.Errorf("expected package name test-package, got %s", pkg.Name) } + if has, err := db.HasTable("pending_deletes"); err != nil || !has { + t.Errorf("pending_deletes table missing after migration (err %v)", err) + } + // Verify migrations were recorded applied, err := db.appliedMigrations() if err != nil { diff --git a/internal/database/pending_deletes_test.go b/internal/database/pending_deletes_test.go new file mode 100644 index 00000000..9af6c48d --- /dev/null +++ b/internal/database/pending_deletes_test.go @@ -0,0 +1,199 @@ +package database + +import ( + "database/sql" + "fmt" + "slices" + "sync" + "testing" + "time" +) + +const ( + pendingVersionPURL = "pkg:npm/pending@1.0.0" + pendingFilename = "pending-1.0.0.tgz" +) + +func upsertPendingArtifact(t *testing.T, db *DB, storagePath string) { + t.Helper() + a := &Artifact{ + VersionPURL: pendingVersionPURL, + Filename: pendingFilename, + UpstreamURL: "https://example.com/" + pendingFilename, + StoragePath: sql.NullString{String: storagePath, Valid: true}, + Size: sql.NullInt64{Int64: 1, Valid: true}, + FetchedAt: sql.NullTime{Time: time.Now(), Valid: true}, + } + if err := db.UpsertArtifact(a); err != nil { + t.Fatalf("UpsertArtifact(%q): %v", storagePath, err) + } +} + +func recordedPath(t *testing.T, db *DB) sql.NullString { + t.Helper() + a, err := db.GetArtifact(pendingVersionPURL, pendingFilename) + if err != nil || a == nil { + t.Fatalf("GetArtifact: %v, %v", a, err) + } + return a.StoragePath +} + +// allPendingDeletes returns every queued path, due or not. +func allPendingDeletes(t *testing.T, db *DB) []string { + t.Helper() + paths, err := db.GetDuePendingDeletes(time.Now().Add(time.Hour), 100) + if err != nil { + t.Fatalf("GetDuePendingDeletes: %v", err) + } + return paths +} + +func TestUpsertArtifactQueuesThePathItReplaces(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + upsertPendingArtifact(t, db, "npm/pending/1.0.0/a/pending-1.0.0.tgz") + upsertPendingArtifact(t, db, "npm/pending/1.0.0/a/pending-1.0.0.tgz") + if got := allPendingDeletes(t, db); len(got) != 0 { + t.Fatalf("queued %v, want nothing while the record keeps its path", got) + } + + upsertPendingArtifact(t, db, "npm/pending/1.0.0/b/pending-1.0.0.tgz") + want := []string{"npm/pending/1.0.0/a/pending-1.0.0.tgz"} + if got := allPendingDeletes(t, db); !slices.Equal(got, want) { + t.Errorf("queued %v, want %v", got, want) + } + }) +} + +// TestUpsertArtifactSkipsRecordThatMoved is the race UpsertArtifact retries: +// another commit moved the record after this one read it, so the write must +// not apply, or the path that commit stored would never be queued. +func TestUpsertArtifactSkipsRecordThatMoved(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + upsertPendingArtifact(t, db, "npm/pending/1.0.0/a/pending-1.0.0.tgz") + read := recordedPath(t, db) + upsertPendingArtifact(t, db, "npm/pending/1.0.0/b/pending-1.0.0.tgz") + + late := &Artifact{ + VersionPURL: pendingVersionPURL, + Filename: pendingFilename, + UpstreamURL: "https://example.com/" + pendingFilename, + StoragePath: sql.NullString{String: "npm/pending/1.0.0/c/pending-1.0.0.tgz", Valid: true}, + } + applied, err := db.upsertArtifactFrom(late, read) + if err != nil { + t.Fatalf("upsertArtifactFrom: %v", err) + } + if applied { + t.Error("write applied over a record that moved since it was read") + } + if got := recordedPath(t, db).String; got != "npm/pending/1.0.0/b/pending-1.0.0.tgz" { + t.Errorf("record points at %q, want the newer commit kept", got) + } + + missing := &Artifact{VersionPURL: "pkg:npm/pending@2.0.0", Filename: pendingFilename, UpstreamURL: "u"} + applied, err = db.upsertArtifactFrom(missing, sql.NullString{}) + if err != nil || !applied { + t.Errorf("insert of a new record: applied=%v err=%v", applied, err) + } + }) +} + +// TestConcurrentUpsertsQueueEveryReplacedPath commits one artifact from several +// fetches at once: whichever commit ends up recorded, every other path must be +// queued, or its object is never deleted. +func TestConcurrentUpsertsQueueEveryReplacedPath(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + // A commit only retries after another commit succeeded, so this many + // always fit within upsertArtifactAttempts. + const commits = upsertArtifactAttempts - 1 + paths := make([]string, commits) + errs := make(chan error, commits) + var wg sync.WaitGroup + for i := range paths { + paths[i] = fmt.Sprintf("npm/pending/1.0.0/%d/pending-1.0.0.tgz", i) + wg.Go(func() { + errs <- db.UpsertArtifact(&Artifact{ + VersionPURL: pendingVersionPURL, + Filename: pendingFilename, + UpstreamURL: "https://example.com/" + pendingFilename, + StoragePath: sql.NullString{String: paths[i], Valid: true}, + }) + }) + } + wg.Wait() + close(errs) + for err := range errs { + if err != nil { + t.Fatalf("UpsertArtifact: %v", err) + } + } + + recorded := recordedPath(t, db).String + want := slices.DeleteFunc(slices.Clone(paths), func(p string) bool { return p == recorded }) + got := allPendingDeletes(t, db) + slices.Sort(want) + slices.Sort(got) + if !slices.Equal(got, want) { + t.Errorf("recorded %q, queued %v, want %v", recorded, got, want) + } + }) +} + +func TestClearArtifactCacheLeavesRecordThatMoved(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + upsertPendingArtifact(t, db, "npm/pending/1.0.0/a/pending-1.0.0.tgz") + upsertPendingArtifact(t, db, "npm/pending/1.0.0/b/pending-1.0.0.tgz") + + if cleared, err := db.ClearArtifactCache(pendingVersionPURL, pendingFilename, "npm/pending/1.0.0/a/pending-1.0.0.tgz"); err != nil || cleared { + t.Fatalf("ClearArtifactCache: cleared=%v err=%v, want nothing cleared", cleared, err) + } + if err := db.DiscardArtifact(pendingVersionPURL, pendingFilename, "npm/pending/1.0.0/a/pending-1.0.0.tgz"); err != nil { + t.Fatalf("DiscardArtifact: %v", err) + } + if got := recordedPath(t, db).String; got != "npm/pending/1.0.0/b/pending-1.0.0.tgz" { + t.Errorf("record points at %q, want the newer commit kept", got) + } + }) +} + +func TestDiscardArtifactClearsAndQueues(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + upsertPendingArtifact(t, db, "npm/pending/1.0.0/a/pending-1.0.0.tgz") + + if err := db.DiscardArtifact(pendingVersionPURL, pendingFilename, "npm/pending/1.0.0/a/pending-1.0.0.tgz"); err != nil { + t.Fatalf("DiscardArtifact: %v", err) + } + if got := recordedPath(t, db); got.Valid { + t.Errorf("record still points at %q", got.String) + } + want := []string{"npm/pending/1.0.0/a/pending-1.0.0.tgz"} + if got := allPendingDeletes(t, db); !slices.Equal(got, want) { + t.Errorf("queued %v, want %v", got, want) + } + }) +} + +func TestGetDuePendingDeletes(t *testing.T) { + runWithBothDatabases(t, func(t *testing.T, db *DB) { + for _, path := range []string{"old", "referenced"} { + if err := db.QueuePendingDelete(path); err != nil { + t.Fatalf("QueuePendingDelete: %v", err) + } + } + upsertPendingArtifact(t, db, "referenced") + + if got, err := db.GetDuePendingDeletes(time.Now().Add(-time.Hour), 100); err != nil || len(got) != 0 { + t.Errorf("before the grace period: got %v (err %v), want nothing", got, err) + } + if got := allPendingDeletes(t, db); !slices.Equal(got, []string{"old"}) { + t.Errorf("after the grace period: got %v, want [old] and not the path a record points at", got) + } + + if err := db.RemovePendingDelete("old"); err != nil { + t.Fatalf("RemovePendingDelete: %v", err) + } + if got := allPendingDeletes(t, db); len(got) != 0 { + t.Errorf("after removal: got %v", got) + } + }) +} diff --git a/internal/database/queries.go b/internal/database/queries.go index a7c4c7a4..4417c171 100644 --- a/internal/database/queries.go +++ b/internal/database/queries.go @@ -2,6 +2,7 @@ package database import ( "database/sql" + "errors" "fmt" "time" @@ -310,7 +311,42 @@ func (db *DB) GetArtifactsByVersionPURL(versionPURL string) ([]Artifact, error) return artifacts, nil } +// upsertArtifactAttempts bounds how often UpsertArtifact retries a record that +// concurrent commits keep moving. +const upsertArtifactAttempts = 5 + +// UpsertArtifact records a. A storage path the record stops pointing at is +// queued for deletion, since each fetch stores its own object and nothing else +// would remove it. The write applies only if the record still holds the path +// read before it, so commits racing on one artifact each queue the path they +// replaced. func (db *DB) UpsertArtifact(a *Artifact) error { + for range upsertArtifactAttempts { + var previous sql.NullString + query := db.Rebind(`SELECT storage_path FROM artifacts WHERE version_purl = ? AND filename = ?`) + if err := db.Get(&previous, query, a.VersionPURL, a.Filename); err != nil && !errors.Is(err, sql.ErrNoRows) { + return fmt.Errorf("reading artifact: %w", err) + } + applied, err := db.upsertArtifactFrom(a, previous) + if err != nil { + return err + } + if !applied { + continue + } + if previous.Valid && previous.String != a.StoragePath.String { + if err := db.QueuePendingDelete(previous.String); err != nil { + return fmt.Errorf("queueing replaced artifact: %w", err) + } + } + return nil + } + return errors.New("upserting artifact: record kept changing") +} + +// upsertArtifactFrom writes a if the record is absent or still holds +// previous, and reports whether it did. +func (db *DB) upsertArtifactFrom(a *Artifact, previous sql.NullString) (bool, error) { now := time.Now() var query string @@ -327,6 +363,7 @@ func (db *DB) UpsertArtifact(a *Artifact) error { content_type = EXCLUDED.content_type, fetched_at = EXCLUDED.fetched_at, updated_at = EXCLUDED.updated_at + WHERE artifacts.storage_path IS NOT DISTINCT FROM $13 ` } else { query = ` @@ -341,17 +378,22 @@ func (db *DB) UpsertArtifact(a *Artifact) error { content_type = excluded.content_type, fetched_at = excluded.fetched_at, updated_at = excluded.updated_at + WHERE artifacts.storage_path IS ? ` } - _, err := db.Exec(query, + res, err := db.Exec(query, a.VersionPURL, a.Filename, a.UpstreamURL, a.StoragePath, a.ContentHash, - a.Size, a.ContentType, a.FetchedAt, a.HitCount, a.LastAccessedAt, now, now, + a.Size, a.ContentType, a.FetchedAt, a.HitCount, a.LastAccessedAt, now, now, previous, ) if err != nil { - return fmt.Errorf("upserting artifact: %w", err) + return false, fmt.Errorf("upserting artifact: %w", err) } - return nil + n, err := res.RowsAffected() + if err != nil { + return false, fmt.Errorf("upserting artifact: %w", err) + } + return n > 0, nil } func (db *DB) RecordArtifactHit(versionPURL, filename string) error { @@ -415,14 +457,71 @@ func (db *DB) GetCachedArtifactCount() (int64, error) { return count, err } -func (db *DB) ClearArtifactCache(versionPURL, filename string) error { +// ClearArtifactCache marks an artifact uncached if its record still points at +// storagePath, and reports whether it did. A record a newer fetch has moved +// elsewhere is left alone, so a clear never orphans the object that fetch +// committed. +func (db *DB) ClearArtifactCache(versionPURL, filename, storagePath string) (bool, error) { + return db.clearArtifactCache(versionPURL, filename, storagePath) +} + +// DiscardArtifact clears the record as ClearArtifactCache does and queues +// storagePath for deletion, for callers that must not delete an object another +// request may be reading. +func (db *DB) DiscardArtifact(versionPURL, filename, storagePath string) error { + cleared, err := db.clearArtifactCache(versionPURL, filename, storagePath) + if err != nil || !cleared { + return err + } + return db.QueuePendingDelete(storagePath) +} + +func (db *DB) clearArtifactCache(versionPURL, filename, storagePath string) (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 = ? + WHERE version_purl = ? AND filename = ? AND storage_path = ? `) - _, err := db.Exec(query, time.Now(), versionPURL, filename) + res, err := db.Exec(query, time.Now(), versionPURL, filename, storagePath) + if err != nil { + return false, err + } + n, err := res.RowsAffected() + return n > 0, err +} + +// QueuePendingDelete queues a storage path no record points at any more. +// Queueing a path again restarts its grace period. +func (db *DB) QueuePendingDelete(path string) error { + query := db.Rebind(` + INSERT INTO pending_deletes (path, queued_at) VALUES (?, ?) + ON CONFLICT(path) DO UPDATE SET queued_at = excluded.queued_at + `) + _, err := db.Exec(query, path, time.Now().UTC()) + return err +} + +// GetDuePendingDeletes returns up to limit paths queued before cutoff, oldest +// first. A path a record points at again is skipped. +func (db *DB) GetDuePendingDeletes(cutoff time.Time, limit int) ([]string, error) { + var paths []string + query := db.Rebind(` + SELECT path FROM pending_deletes + WHERE queued_at < ? + AND NOT EXISTS (SELECT 1 FROM artifacts WHERE artifacts.storage_path = pending_deletes.path) + ORDER BY queued_at + LIMIT ? + `) + if err := db.Select(&paths, query, cutoff.UTC(), limit); err != nil { + return nil, err + } + return paths, nil +} + +// RemovePendingDelete drops a deleted path from the queue. +func (db *DB) RemovePendingDelete(path string) error { + _, err := db.Exec(db.Rebind(`DELETE FROM pending_deletes WHERE path = ?`), path) return err } diff --git a/internal/database/schema.go b/internal/database/schema.go index 564745fb..eb457899 100644 --- a/internal/database/schema.go +++ b/internal/database/schema.go @@ -78,6 +78,11 @@ CREATE UNIQUE INDEX IF NOT EXISTS idx_artifacts_version_filename ON artifacts(ve CREATE INDEX IF NOT EXISTS idx_artifacts_storage_path ON artifacts(storage_path); CREATE INDEX IF NOT EXISTS idx_artifacts_last_accessed ON artifacts(last_accessed_at); +CREATE TABLE IF NOT EXISTS pending_deletes ( + path TEXT NOT NULL PRIMARY KEY, + queued_at DATETIME NOT NULL +); + CREATE TABLE IF NOT EXISTS vulnerabilities ( id INTEGER PRIMARY KEY, vuln_id TEXT NOT NULL, @@ -181,6 +186,11 @@ CREATE UNIQUE INDEX IF NOT EXISTS idx_artifacts_version_filename ON artifacts(ve CREATE INDEX IF NOT EXISTS idx_artifacts_storage_path ON artifacts(storage_path); CREATE INDEX IF NOT EXISTS idx_artifacts_last_accessed ON artifacts(last_accessed_at); +CREATE TABLE IF NOT EXISTS pending_deletes ( + path TEXT NOT NULL PRIMARY KEY, + queued_at TIMESTAMP NOT NULL +); + CREATE TABLE IF NOT EXISTS vulnerabilities ( id SERIAL PRIMARY KEY, vuln_id TEXT NOT NULL, @@ -368,6 +378,7 @@ var migrations = []migration{ {"006_add_metadata_content_digest", migrateAddMetadataContentDigest}, {"007_add_metadata_link", migrateAddMetadataLink}, {"008_add_metadata_content_encoding", migrateAddMetadataContentEncoding}, + {"009_add_pending_deletes", migrateAddPendingDeletes}, } // isTableNotFound returns true if the error indicates a missing table. @@ -640,6 +651,23 @@ func migrateAddMetadataContentEncoding(db *DB) error { return nil } +// migrateAddPendingDeletes creates the queue of storage paths that no record +// points at any more, which the server deletes after a grace period. +func migrateAddPendingDeletes(db *DB) error { + ts := sqliteDatetime + if db.dialect == DialectPostgres { + ts = postgresTimestamp + } + query := fmt.Sprintf(`CREATE TABLE IF NOT EXISTS pending_deletes ( + path TEXT NOT NULL PRIMARY KEY, + queued_at %s NOT NULL + )`, ts) + if _, err := db.Exec(query); err != nil { + return fmt.Errorf("creating pending_deletes table: %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/handler/fetch_path_test.go b/internal/handler/fetch_path_test.go new file mode 100644 index 00000000..4f7602ac --- /dev/null +++ b/internal/handler/fetch_path_test.go @@ -0,0 +1,163 @@ +package handler + +import ( + "context" + "errors" + "io" + "net/http" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/git-pkgs/registries/fetch" +) + +// These tests cover fetches of one artifact that do not share a coalescing +// key because their URLs differ. Each must keep its own stored object. + +const ( + fetchPathURLA = "https://mirror-a.example/pkg-1.0.0.tgz" + fetchPathURLB = "https://mirror-b.example/pkg-1.0.0.tgz" +) + +// barrierFetcher serves a body per URL and holds each fetch until both have +// started, so both requests are past the cache check before either commits. +type barrierFetcher struct { + bodies map[string]string + arrived sync.WaitGroup +} + +func newBarrierFetcher(bodies map[string]string) *barrierFetcher { + f := &barrierFetcher{bodies: bodies} + f.arrived.Add(len(bodies)) + return f +} + +func (f *barrierFetcher) Fetch(ctx context.Context, url string) (*fetch.Artifact, error) { + return f.FetchWithHeaders(ctx, url, nil) +} + +func (f *barrierFetcher) FetchWithHeaders(_ context.Context, url string, _ http.Header) (*fetch.Artifact, error) { + f.arrived.Done() + f.arrived.Wait() + return artifactBody(f.bodies[url]), nil +} + +func (f *barrierFetcher) Head(_ context.Context, _ string) (int64, string, error) { + return 0, "", nil +} + +// gatedStorage holds the first Open until ready reports true, to order it +// after the other request's store or delete. +type gatedStorage struct { + *mockStorage + ready func(stores, deletes int32) bool + stores atomic.Int32 + deletes atomic.Int32 + opened atomic.Bool + release chan struct{} + releaseMu sync.Once +} + +func newGatedStorage(store *mockStorage, ready func(stores, deletes int32) bool) *gatedStorage { + return &gatedStorage{mockStorage: store, ready: ready, release: make(chan struct{})} +} + +func (g *gatedStorage) check() { + if g.ready(g.stores.Load(), g.deletes.Load()) { + g.releaseMu.Do(func() { close(g.release) }) + } +} + +func (g *gatedStorage) Store(ctx context.Context, path string, r io.Reader) (int64, string, error) { + size, hash, err := g.mockStorage.Store(ctx, path, r) + g.stores.Add(1) + g.check() + return size, hash, err +} + +func (g *gatedStorage) Delete(ctx context.Context, path string) error { + err := g.mockStorage.Delete(ctx, path) + g.deletes.Add(1) + g.check() + return err +} + +func (g *gatedStorage) Open(ctx context.Context, path string) (io.ReadCloser, error) { + if g.opened.CompareAndSwap(false, true) { + select { + case <-g.release: + case <-ctx.Done(): + return nil, ctx.Err() + } + } + return g.mockStorage.Open(ctx, path) +} + +func readResult(res *CacheResult) string { + b, _ := io.ReadAll(res.Reader) + _ = res.Reader.Close() + return string(b) +} + +// TestFetchesWithDifferentURLsServeTheirOwnBytes has both fetches store before +// either opens. Sharing one object, one of them would serve the other's bytes. +func TestFetchesWithDifferentURLsServeTheirOwnBytes(t *testing.T) { + proxy, _, store, _ := setupTestProxy(t) + bodies := map[string]string{fetchPathURLA: "bytes from a", fetchPathURLB: "bytes from b"} + proxy.Fetcher = newBarrierFetcher(bodies) + proxy.Storage = newGatedStorage(store, func(stores, _ int32) bool { return stores == 2 }) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + var wg sync.WaitGroup + for url, want := range bodies { + wg.Go(func() { + res, err := proxy.GetOrFetchArtifactFromURL(ctx, "npm", "pkg", "1.0.0", staleFilename, url) + if err != nil { + t.Errorf("fetch of %s: %v", url, err) + return + } + if got := readResult(res); got != want { + t.Errorf("fetch of %s served %q, want %q", url, got, want) + } + }) + } + wg.Wait() +} + +// TestDigestMismatchLeavesAnotherFetchsObject has one fetch discard bytes that +// fail the declared digest before another fetch, which stored good bytes, +// opens them. Sharing one object, the discard would delete the good bytes. +func TestDigestMismatchLeavesAnotherFetchsObject(t *testing.T) { + proxy, _, store, _ := setupTestProxy(t) + proxy.Fetcher = newBarrierFetcher(map[string]string{fetchPathURLA: "good bytes", fetchPathURLB: "tampered bytes"}) + proxy.Storage = newGatedStorage(store, func(_, deletes int32) bool { return deletes == 1 }) + declared := "sha256:" + sha256Hex("good bytes") + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + var wg sync.WaitGroup + wg.Go(func() { + res, err := proxy.GetOrFetchArtifactFromURLWithDigest(ctx, "npm", "pkg", "1.0.0", staleFilename, fetchPathURLA, declared) + if err != nil { + t.Errorf("good fetch: %v", err) + return + } + if got := readResult(res); got != "good bytes" { + t.Errorf("good fetch served %q", got) + } + }) + wg.Go(func() { + _, err := proxy.GetOrFetchArtifactFromURLWithDigest(ctx, "npm", "pkg", "1.0.0", staleFilename, fetchPathURLB, declared) + if !errors.Is(err, ErrArtifactDigestMismatch) { + t.Errorf("tampered fetch: got %v, want a digest mismatch", err) + } + }) + wg.Wait() + + if _, ok := store.files[recordedStoragePath(t, proxy.DB, staleVersionPURL, staleFilename)]; !ok { + t.Error("the recorded object was deleted") + } +} diff --git a/internal/handler/handler.go b/internal/handler/handler.go index 67d5d121..88e08bd9 100644 --- a/internal/handler/handler.go +++ b/internal/handler/handler.go @@ -267,9 +267,9 @@ func (p *Proxy) GetCachedArtifact(ctx context.Context, ecosystem, name, version, return p.checkCache(ctx, pkgPURL, versionPURL, filename) } -// ClearCachedArtifact removes both an artifact cache record and its stored -// bytes after an external integrity check fails. -func (p *Proxy) ClearCachedArtifact(ctx context.Context, ecosystem, name, version, filename string) error { +// ClearCachedArtifact clears an artifact cache record after an external +// integrity check fails, and queues its stored bytes for deletion. +func (p *Proxy) ClearCachedArtifact(_ context.Context, ecosystem, name, version, filename string) error { if p.DB == nil || p.Storage == nil { return nil } @@ -284,10 +284,7 @@ func (p *Proxy) ClearCachedArtifact(ctx context.Context, ecosystem, name, versio if cached == nil { return nil } - if err := p.Storage.Delete(ctx, cached.StoragePath); err != nil { - return fmt.Errorf("deleting cached artifact: %w", err) - } - return p.DB.ClearArtifactCache(versionPURL, filename) + return p.DB.DiscardArtifact(versionPURL, filename, cached.StoragePath) } // checkCache looks up an artifact in the cache. Returns nil if not cached. @@ -343,7 +340,7 @@ func (p *Proxy) checkCache(ctx context.Context, pkgPURL, versionPURL, filename s "purl", versionPURL, "filename", filename, "path", artifact.StoragePath, "reason", reason) metrics.RecordIntegrityFailure(purl.NormalizeEcosystem(artifact.Ecosystem)) - if err := p.DB.ClearArtifactCache(versionPURL, filename); err != nil { + if err := p.DB.DiscardArtifact(versionPURL, filename, artifact.StoragePath); err != nil { p.Logger.Warn("failed to clear corrupt artifact from cache", "error", err) } }) @@ -386,7 +383,7 @@ func (p *Proxy) rejectUnusableCacheRecord(artifact *database.CachedArtifact, ver "purl", versionPURL, "filename", filename, "path", artifact.StoragePath, "error", cause) metrics.RecordIntegrityFailure(purl.NormalizeEcosystem(artifact.Ecosystem)) - if err := p.DB.ClearArtifactCache(versionPURL, filename); err != nil { + if err := p.DB.DiscardArtifact(versionPURL, filename, artifact.StoragePath); err != nil { p.Logger.Warn("failed to clear unusable artifact from cache", "error", err) } } @@ -448,7 +445,7 @@ func (p *Proxy) fetchAndCache(ctx context.Context, ecosystem, name, version, fil // It returns the artifact and its storage path, not a reader; callers get one // from openStoredArtifact. func (p *Proxy) storeArtifact(ctx context.Context, ecosystem, name, version, filename, pkgPURL, versionPURL, upstreamURL, upstreamHash string, artifact *fetch.Artifact) (artifacts.Artifact, string, error) { - storagePath := storage.ArtifactPath(ecosystem, "", name, version, filename) + storagePath := storage.FetchPath(ecosystem, name, version, storage.NewFetchID(), filename) storeStart := time.Now() size, hash, err := p.Storage.Store(ctx, storagePath, artifact.Body) @@ -490,7 +487,11 @@ func (p *Proxy) storeArtifact(ctx context.Context, ecosystem, name, version, fil // Update database if err := p.updateCacheDB(ecosystem, name, pkgPURL, upstreamURL, storagePath, sharedArtifact); err != nil { p.Logger.Warn("failed to update cache database", "error", err) - // Continue anyway - we have the file + // Continue anyway - we have the file. Queue it for deletion in case + // no record points at it; the queue skips it while one does. + if qErr := p.DB.QueuePendingDelete(storagePath); qErr != nil { + p.Logger.Warn("failed to queue unrecorded artifact for deletion", "path", storagePath, "error", qErr) + } } return sharedArtifact, storagePath, nil @@ -1352,7 +1353,7 @@ func (p *Proxy) coalescedFetchFromURL(ctx context.Context, ecosystem, name, vers return p.cachedArtifactRecord(pkgPURL, versionPURL, filename, upstreamHash) } return p.coalesceFetch(ctx, key, recheck, func(fetchCtx context.Context) (artifacts.Artifact, string, error) { - p.discardStaleArtifact(fetchCtx, pkgPURL, versionPURL, filename, upstreamHash) + p.discardStaleArtifact(pkgPURL, versionPURL, filename, upstreamHash) return p.fetchAndCacheFromURL(fetchCtx, ecosystem, name, version, filename, pkgPURL, versionPURL, downloadURL, headers, upstreamHash) }) } @@ -1378,10 +1379,10 @@ func (p *Proxy) getCachedArtifactWithUpstreamHash(ctx context.Context, pkgPURL, return nil, nil } -// discardStaleArtifact removes the cached entry when its digest disagrees -// with upstreamHash. It runs under the coalescing key, after the recheck, so -// an entry a previous fetch refreshed is kept. -func (p *Proxy) discardStaleArtifact(ctx context.Context, pkgPURL, versionPURL, filename, upstreamHash string) { +// discardStaleArtifact clears the cached entry when its digest disagrees with +// upstreamHash, queueing its object for deletion. It runs under the coalescing +// key, after the recheck, so an entry a previous fetch refreshed is kept. +func (p *Proxy) discardStaleArtifact(pkgPURL, versionPURL, filename, upstreamHash string) { record, err := p.DB.GetCachedArtifact(pkgPURL, versionPURL, filename) if err != nil { p.Logger.Warn("failed to read cache record before refetch", @@ -1393,7 +1394,9 @@ func (p *Proxy) discardStaleArtifact(ctx context.Context, pkgPURL, versionPURL, } p.Logger.Warn("cached artifact hash disagrees with upstream metadata, discarding", "purl", versionPURL, "filename", filename, "cached", record.Artifact.Digest.Encoded(), "upstream", upstreamHash) - p.discardCachedArtifact(ctx, versionPURL, filename, record.StoragePath) + if err := p.DB.DiscardArtifact(versionPURL, filename, record.StoragePath); err != nil { + p.Logger.Warn("failed to clear artifact cache record", "purl", versionPURL, "filename", filename, "error", err) + } } func (p *Proxy) fetchAndCacheFromURL(ctx context.Context, ecosystem, name, version, filename, pkgPURL, versionPURL, downloadURL string, headers http.Header, upstreamHash string) (artifacts.Artifact, string, error) { @@ -1421,14 +1424,3 @@ var ErrArtifactDigestMismatch = errors.New("artifact digest mismatch") func artifactHashMatches(got, expected string) bool { return expected == "" || strings.EqualFold(got, expected) } - -func (p *Proxy) discardCachedArtifact(ctx context.Context, versionPURL, filename, storagePath string) { - if storagePath != "" { - if err := p.Storage.Delete(ctx, storagePath); err != nil { - p.Logger.Warn("failed to discard cached artifact", "path", storagePath, "error", err) - } - } - if err := p.DB.ClearArtifactCache(versionPURL, filename); err != nil { - p.Logger.Warn("failed to clear artifact cache record", "purl", versionPURL, "filename", filename, "error", err) - } -} diff --git a/internal/handler/handler_test.go b/internal/handler/handler_test.go index b632294d..c2714c30 100644 --- a/internal/handler/handler_test.go +++ b/internal/handler/handler_test.go @@ -10,6 +10,7 @@ import ( "log/slog" "net/http" "net/http/httptest" + "slices" "strings" "sync" "testing" @@ -246,6 +247,16 @@ func seedPackage(t testing.TB, db *database.DB, store *mockStorage, ecosystem, n } } +// recordedStoragePath returns the storage path the artifact's record points at. +func recordedStoragePath(t testing.TB, db *database.DB, versionPURL, filename string) string { + t.Helper() + art, err := db.GetArtifact(versionPURL, filename) + if err != nil || art == nil || !art.StoragePath.Valid { + t.Fatalf("no storage path recorded for %s %s (err %v)", versionPURL, filename, err) + } + return art.StoragePath.String +} + // pathParseCase holds a single test case for path parsing functions that return // (name, version, arch). type pathParseCase struct { @@ -394,6 +405,38 @@ func assertMalformedCacheRejected(t *testing.T, malformedHash, malformedIntegrit if artifact.StoragePath.Valid { t.Error("unusable cache record retained its storage path") } + assertQueuedForDeletion(t, db, storage.ArtifactPath("npm", "", packageName, version, filename)) +} + +// TestCorruptCachedArtifactIsQueuedForDeletion streams cached bytes that no +// longer match their recorded digest, which clears the record once the stream +// ends. With each fetch writing its own path, no later fetch overwrites them, +// so they must be queued for deletion. +func TestCorruptCachedArtifactIsQueuedForDeletion(t *testing.T) { + proxy, db, store, _ := setupTestProxy(t) + seedPackage(t, db, store, "npm", "corrupt", "1.0.0", "corrupt-1.0.0.tgz", "cached content") + storagePath := storage.ArtifactPath("npm", "", "corrupt", "1.0.0", "corrupt-1.0.0.tgz") + store.files[storagePath] = []byte("tampered bytes") + + result, err := proxy.GetCachedArtifact(context.Background(), "npm", "corrupt", "1.0.0", "corrupt-1.0.0.tgz") + if err != nil || result == nil { + t.Fatalf("GetCachedArtifact = %v, %v", result, err) + } + drain(result) + artifact, err := db.GetArtifact("pkg:npm/corrupt@1.0.0", "corrupt-1.0.0.tgz") + if err != nil || artifact == nil || artifact.StoragePath.Valid { + t.Errorf("record = %+v (err %v), want it cleared", artifact, err) + } + assertQueuedForDeletion(t, db, storagePath) +} + +// assertQueuedForDeletion fails unless path waits in the pending delete queue. +func assertQueuedForDeletion(t *testing.T, db *database.DB, path string) { + t.Helper() + queued, err := db.GetDuePendingDeletes(time.Now().Add(time.Hour), 100) + if err != nil || !slices.Contains(queued, path) { + t.Errorf("queued %v (err %v), want %q queued for deletion", queued, err, path) + } } func TestGetOrFetchArtifact_CacheMiss_NoPackage(t *testing.T) { @@ -454,7 +497,7 @@ func TestGetOrFetchArtifactFromURL_CacheMiss_StorageMissing(t *testing.T) { } // Verify the new content was stored - storagePath := storage.ArtifactPath("npm", "", "missing", "1.0.0", "missing-1.0.0.tgz") + storagePath := recordedStoragePath(t, db, "pkg:npm/missing@1.0.0", "missing-1.0.0.tgz") if _, ok := store.files[storagePath]; !ok { t.Error("refetched artifact should be stored") } @@ -718,7 +761,7 @@ func TestGetOrFetchArtifactFromURL_CacheHit(t *testing.T) { } func TestGetOrFetchArtifactFromURL_CacheMiss(t *testing.T) { - proxy, _, store, fetcher := setupTestProxy(t) + proxy, db, store, fetcher := setupTestProxy(t) missesBefore := testutil.ToFloat64(metrics.CacheMisses.WithLabelValues("pypi")) fetchesBefore := histogramSampleCount(t, metrics.UpstreamFetchDuration.WithLabelValues("pypi")) writesBefore := histogramSampleCount(t, metrics.StorageOperationDuration.WithLabelValues("write")) @@ -763,7 +806,7 @@ func TestGetOrFetchArtifactFromURL_CacheMiss(t *testing.T) { } // Verify it was stored - storagePath := storage.ArtifactPath("pypi", "", "newpkg", "1.0.0", "newpkg-1.0.0.tar.gz") + storagePath := recordedStoragePath(t, db, "pkg:pypi/newpkg@1.0.0", "newpkg-1.0.0.tar.gz") if _, ok := store.files[storagePath]; !ok { t.Error("artifact was not stored in storage") } diff --git a/internal/handler/helm_test.go b/internal/handler/helm_test.go index ec8d4a5a..c4a60fe4 100644 --- a/internal/handler/helm_test.go +++ b/internal/handler/helm_test.go @@ -6,6 +6,7 @@ import ( "fmt" "net/http" "net/http/httptest" + "slices" "strings" "sync/atomic" "testing" @@ -164,11 +165,23 @@ func TestHelmHandler_RejectsChartDigestMismatch(t *testing.T) { if requests != 2 { t.Errorf("chart requests = %d, want 2 after invalid cache entry is cleared", requests) } - storagePath := storage.ArtifactPath(helmMetadataEcosystem, "", "test", digest, "demo.tgz") - if exists, err := store.Exists(t.Context(), storagePath); err != nil { - t.Fatalf("checking rejected chart storage: %v", err) - } else if exists { - t.Errorf("rejected chart remains in storage at %q", storagePath) + queued, err := proxy.DB.GetDuePendingDeletes(time.Now().Add(time.Hour), 10) + if err != nil { + t.Fatalf("listing queued deletes: %v", err) + } + chartDir := storage.ArtifactPath(helmMetadataEcosystem, "", "test", digest, "") + stored := 0 + for path := range store.files { + if !strings.HasPrefix(path, chartDir) { + continue + } + stored++ + if !slices.Contains(queued, path) { + t.Errorf("rejected chart at %q is not queued for deletion", path) + } + } + if stored != requests { + t.Errorf("found %d stored charts, want one per request (%d)", stored, requests) } } diff --git a/internal/handler/stale_cache_test.go b/internal/handler/stale_cache_test.go index 3d520a43..c11edd0b 100644 --- a/internal/handler/stale_cache_test.go +++ b/internal/handler/stale_cache_test.go @@ -89,9 +89,12 @@ func TestStaleCacheIsDiscardedBeforeTheFetch(t *testing.T) { if got := cachedDigest(t, proxy); got != "" { t.Errorf("stale record survived a failed refresh, digest = %q", got) } - if bytesPresent(store) { - t.Error("stale bytes survived a failed refresh") + // A request that read the stale record may still be opening its bytes, + // so they are queued for deletion rather than deleted. + if !bytesPresent(store) { + t.Error("stale bytes were deleted while a reader may still open them") } + assertQueuedForDeletion(t, proxy.DB, staleStoragePath) } func TestStaleCacheIsReplacedByTheFetch(t *testing.T) { diff --git a/internal/server/eviction.go b/internal/server/eviction.go index ea29e525..c4639f7a 100644 --- a/internal/server/eviction.go +++ b/internal/server/eviction.go @@ -122,11 +122,17 @@ func evictBatch(ctx context.Context, db *database.DB, store storage.Storage, log continue } - if err := db.ClearArtifactCache(art.VersionPURL, art.Filename); err != nil { + recordCleared, err := db.ClearArtifactCache(art.VersionPURL, art.Filename, art.StoragePath.String) + if err != nil { logger.Warn("eviction: failed to clear artifact record", "version_purl", art.VersionPURL, "filename", art.Filename, "error", err) continue } + if !recordCleared { + // The record no longer points here, so this delete freed nothing + // the recorded size counts. + continue + } if art.Size.Valid { freed += art.Size.Int64 diff --git a/internal/server/eviction_test.go b/internal/server/eviction_test.go index efb7a2e9..58cd1dce 100644 --- a/internal/server/eviction_test.go +++ b/internal/server/eviction_test.go @@ -403,3 +403,40 @@ func clearRecordedSize(t *testing.T, db *database.DB, versionPURL, filename stri t.Fatalf("clearing recorded size: %v", err) } } + +// TestEvictBatch_SkipsRecordThatMoved evicts from a row read before a newer +// fetch moved the record. Only the old object goes, and nothing is counted as +// freed, since the size in use is unchanged and counting it would end the pass +// while the cache is still over its limit. +func TestEvictBatch_SkipsRecordThatMoved(t *testing.T) { + db, store := setupEvictionTest(t) + ctx := context.Background() + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + seedArtifact(t, ctx, db, store, "moved", 1000, time.Now().Add(-time.Hour)) + + stale, err := db.GetLeastRecentlyUsedArtifacts(evictionBatch) + if err != nil || len(stale) != 1 { + t.Fatalf("reading LRU rows: %v, %v", stale, err) + } + newer := storage.FetchPath("npm", "moved", "1.0.0", storage.NewFetchID(), "moved-1.0.0.tgz") + if _, _, err := store.Store(ctx, newer, strings.NewReader("refetched")); err != nil { + t.Fatalf("storing refetch: %v", err) + } + moved := stale[0] + moved.StoragePath = sql.NullString{String: newer, Valid: true} + if err := db.UpsertArtifact(&moved); err != nil { + t.Fatalf("moving record: %v", err) + } + + cleared, freed := evictBatch(ctx, db, store, logger, stale, 0, 1000) + if cleared != 0 || freed != 0 { + t.Errorf("cleared %d records freeing %d bytes, want nothing counted", cleared, freed) + } + record, err := db.GetArtifact(moved.VersionPURL, moved.Filename) + if err != nil || record == nil || record.StoragePath.String != newer { + t.Fatalf("record = %+v (err %v), want it kept at %q", record, err, newer) + } + if ok, _ := store.Exists(ctx, newer); !ok { + t.Error("the newer fetch's object was deleted") + } +} diff --git a/internal/server/reclaim.go b/internal/server/reclaim.go new file mode 100644 index 00000000..deaeaa8a --- /dev/null +++ b/internal/server/reclaim.go @@ -0,0 +1,59 @@ +package server + +import ( + "context" + "log/slog" + "time" + + "github.com/git-pkgs/proxy/internal/database" + "github.com/git-pkgs/proxy/internal/storage" +) + +// An object no record points at any more waits in the pending_deletes queue +// for a grace period before it is deleted, so a request that read the record +// before it changed can still open the object, and a signed URL to it stays +// valid. Reclaim runs whether or not a cache size limit is set. +const ( + reclaimInterval = 1 * time.Minute + reclaimBatch = 100 + reclaimMinGrace = 1 * time.Hour +) + +func (s *Server) startReclaimLoop(ctx context.Context) { + grace := max(reclaimMinGrace, s.cfg.ParseDirectServeTTL()) + + ticker := time.NewTicker(reclaimInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + reclaimStorage(ctx, s.db, s.storage, s.logger, time.Now().Add(-grace)) + } + } +} + +// reclaimStorage deletes up to one batch of objects queued before cutoff. A +// delete that fails stays queued for the next pass to retry. +func reclaimStorage(ctx context.Context, db *database.DB, store storage.Storage, logger *slog.Logger, cutoff time.Time) { + paths, err := db.GetDuePendingDeletes(cutoff, reclaimBatch) + if err != nil { + logger.Warn("reclaim: failed to list pending deletes", "error", err) + return + } + + for _, path := range paths { + if ctx.Err() != nil { + return + } + if err := store.Delete(ctx, path); err != nil { + logger.Warn("reclaim: failed to delete object, will retry", "path", path, "error", err) + continue + } + if err := db.RemovePendingDelete(path); err != nil { + logger.Warn("reclaim: failed to dequeue deleted object", "path", path, "error", err) + } + } +} diff --git a/internal/server/reclaim_test.go b/internal/server/reclaim_test.go new file mode 100644 index 00000000..cda8a316 --- /dev/null +++ b/internal/server/reclaim_test.go @@ -0,0 +1,97 @@ +package server + +import ( + "context" + "io" + "log/slog" + "slices" + "strings" + "testing" + "time" + + "github.com/git-pkgs/proxy/internal/database" + "github.com/git-pkgs/proxy/internal/storage" +) + +func storeQueued(t *testing.T, db *database.DB, store storage.Storage, path string) { + t.Helper() + if _, _, err := store.Store(context.Background(), path, strings.NewReader("superseded")); err != nil { + t.Fatalf("storing %s: %v", path, err) + } + if err := db.QueuePendingDelete(path); err != nil { + t.Fatalf("queueing %s: %v", path, err) + } +} + +func queuedPaths(t *testing.T, db *database.DB) []string { + t.Helper() + paths, err := db.GetDuePendingDeletes(time.Now().Add(time.Hour), 100) + if err != nil { + t.Fatalf("listing queue: %v", err) + } + return paths +} + +func objectExists(t *testing.T, store storage.Storage, path string) bool { + t.Helper() + ok, err := store.Exists(context.Background(), path) + if err != nil { + t.Fatalf("checking %s: %v", path, err) + } + return ok +} + +func TestReclaimStorageWaitsForGracePeriod(t *testing.T) { + db, store := setupEvictionTest(t) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + const path = "npm/old/1.0.0/a/old-1.0.0.tgz" + storeQueued(t, db, store, path) + + reclaimStorage(context.Background(), db, store, logger, time.Now().Add(-time.Hour)) + if !objectExists(t, store, path) { + t.Fatal("object deleted before its grace period ended") + } + + reclaimStorage(context.Background(), db, store, logger, time.Now().Add(time.Hour)) + if objectExists(t, store, path) { + t.Error("object survived after its grace period ended") + } + if got := queuedPaths(t, db); len(got) != 0 { + t.Errorf("queue = %v after reclaim, want empty", got) + } +} + +func TestReclaimStorageKeepsFailedDeletesQueued(t *testing.T) { + db, store := setupEvictionTest(t) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + const path = "npm/old/1.0.0/a/old-1.0.0.tgz" + storeQueued(t, db, store, path) + + undeletable := &undeletableStorage{Storage: store} + reclaimStorage(context.Background(), db, undeletable, logger, time.Now().Add(time.Hour)) + + if got := undeletable.deletes.Load(); got != 1 { + t.Errorf("delete attempts = %d, want 1", got) + } + if got := queuedPaths(t, db); !slices.Equal(got, []string{path}) { + t.Errorf("queue = %v, want the failed path kept for retry", got) + } +} + +func TestReclaimStorageStopsWhenContextCanceled(t *testing.T) { + db, store := setupEvictionTest(t) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + const path = "npm/old/1.0.0/a/old-1.0.0.tgz" + storeQueued(t, db, store, path) + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + reclaimStorage(ctx, db, store, logger, time.Now().Add(time.Hour)) + + if !objectExists(t, store, path) { + t.Error("object deleted after the context was canceled") + } + if got := queuedPaths(t, db); !slices.Equal(got, []string{path}) { + t.Errorf("queue = %v, want the path kept", got) + } +} diff --git a/internal/server/server.go b/internal/server/server.go index 6d2a2cb3..7fc326fc 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -358,6 +358,7 @@ func (s *Server) serve(listener net.Listener) error { "database", s.cfg.Database.String()) go s.updateCacheStatsMetrics() go s.startEvictionLoop(bgCtx) + go s.startReclaimLoop(bgCtx) if listener != nil { return s.http.Serve(listener) diff --git a/internal/storage/blob.go b/internal/storage/blob.go index cdc7aa0d..48574aa6 100644 --- a/internal/storage/blob.go +++ b/internal/storage/blob.go @@ -119,13 +119,21 @@ func OpenBucket(ctx context.Context, urlStr string) (Storage, error) { // legacySidecarPath gives the ".attrs" path an earlier version wrote for key, // or "" when that path would not be a file inside fileRoot. -// -// The key is escaped the way fileblob escapes it on the way to disk, so the -// sidecar is looked for where fileblob wrote it. filepath.Localize then -// validates the escaped form: it rejects an empty, absolute or ".." path, and -// "." would name fileRoot itself. What it declines are keys the proxy never -// produces. func (b *Blob) legacySidecarPath(key string) string { + if p := b.localPath(key); p != "" { + return p + attrsExt + } + return "" +} + +// localPath gives the file fileblob keeps key in, or "" when that would not +// be a file inside fileRoot. +// +// The key is escaped the way fileblob escapes it on the way to disk. +// filepath.Localize then validates the escaped form: it rejects an empty, +// absolute or ".." path, and "." would name fileRoot itself. What it declines +// are keys the proxy never produces. +func (b *Blob) localPath(key string) string { if b.fileRoot == "" { return "" } @@ -133,7 +141,7 @@ func (b *Blob) legacySidecarPath(key string) string { if err != nil || rel == "." { return "" } - return filepath.Join(b.fileRoot, rel) + attrsExt + return filepath.Join(b.fileRoot, rel) } // escapeKey mirrors fileblob's unexported escapeKey, which hex-escapes a rune @@ -238,11 +246,20 @@ func (b *Blob) Exists(ctx context.Context, path string) (bool, error) { return exists, nil } +// Delete removes the object at path. On a file:// bucket it also removes the +// object's fetch directory once empty, since fileblob leaves directories +// behind. Any other directory, such as a version directory another fetch may +// be creating its own directory in, is left alone. func (b *Blob) Delete(ctx context.Context, path string) error { err := b.bucket.Delete(ctx, path) if err != nil && !isNotExist(err) { return fmt.Errorf("deleting object: %w", err) } + if p := b.localPath(path); p != "" { + if dir := filepath.Dir(p); dir != b.fileRoot && isFetchDir(filepath.Base(dir)) { + _ = os.Remove(dir) // fails, harmlessly, while the directory holds anything + } + } return nil } diff --git a/internal/storage/blob_test.go b/internal/storage/blob_test.go index 3a97704f..57e85f59 100644 --- a/internal/storage/blob_test.go +++ b/internal/storage/blob_test.go @@ -127,6 +127,44 @@ func TestBlobDelete(t *testing.T) { } } +// A fetch stores its object in a directory of its own, so Delete removes that +// directory once it is empty rather than leave one behind per deleted object. +// A version directory from the layout before fetch directories is left alone, +// since a new fetch may be creating its directory inside it. +func TestBlobDeleteRemovesEmptyFetchDirectory(t *testing.T) { + dir := t.TempDir() + b := openFileBlob(t, dir) + ctx := context.Background() + const ( + emptied = "npm/pkg/1.0.0/0123456789abcdef/pkg.tgz" + shared = "npm/pkg/1.0.0/fedcba9876543210/pkg.tgz" + kept = "npm/pkg/1.0.0/fedcba9876543210/other.tgz" + legacy = "npm/pkg/2.0.0/pkg.tgz" + ) + for _, key := range []string{emptied, shared, kept, legacy} { + if _, _, err := b.Store(ctx, key, strings.NewReader("content")); err != nil { + t.Fatalf("Store(%q): %v", key, err) + } + } + + for _, key := range []string{emptied, shared, legacy} { + if err := b.Delete(ctx, key); err != nil { + t.Fatalf("Delete(%q): %v", key, err) + } + } + + pkg := filepath.Join(dir, "npm", "pkg") + if _, err := os.Stat(filepath.Join(pkg, "1.0.0", "0123456789abcdef")); !os.IsNotExist(err) { + t.Errorf("emptied fetch directory left behind (stat err %v)", err) + } + for _, d := range []string{filepath.Join("1.0.0", "fedcba9876543210"), "1.0.0", "2.0.0"} { + if _, err := os.Stat(filepath.Join(pkg, d)); err != nil { + t.Errorf("directory %s removed: %v", d, err) + } + } + assertReadsBack(t, b, kept, "content") +} + func TestBlobDeleteNotFound(t *testing.T) { b := createTestBlob(t) ctx := context.Background() diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 5ff86f29..cb807f07 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -14,6 +14,8 @@ package storage import ( "context" + "crypto/rand" + "encoding/hex" "errors" "io" "time" @@ -72,7 +74,8 @@ type Storage interface { Close() error } -// ArtifactPath builds a storage path for an artifact. +// ArtifactPath builds the storage path artifacts were cached under before each +// fetch got its own; records from then still point at such paths. // Format: {ecosystem}/{namespace}/{name}/{version}/{filename} // For packages without namespace: {ecosystem}/{name}/{version}/{filename} func ArtifactPath(ecosystem, namespace, name, version, filename string) string { @@ -81,3 +84,36 @@ func ArtifactPath(ecosystem, namespace, name, version, filename string) string { } return ecosystem + "/" + name + "/" + version + "/" + filename } + +// FetchPath builds the storage path for one fetch of an artifact: +// {ecosystem}/{name}/{version}/{fetchID}/{filename}. Each fetch writes its own +// object, so no fetch overwrites or deletes another's. +func FetchPath(ecosystem, name, version, fetchID, filename string) string { + return ecosystem + "/" + name + "/" + version + "/" + fetchID + "/" + filename +} + +// fetchIDBytes sizes a fetch id: 16 hex characters keeps paths short, which +// matters on Windows, and a collision between fetches of one artifact version +// is out of reach. +const fetchIDBytes = 8 + +// NewFetchID returns a random id for FetchPath. +func NewFetchID() string { + b := make([]byte, fetchIDBytes) + _, _ = rand.Read(b) + return hex.EncodeToString(b) +} + +// isFetchDir reports whether a directory name has the shape of a fetch id. +// Only one fetch ever writes into such a directory. +func isFetchDir(name string) bool { + if len(name) != 2*fetchIDBytes { + return false + } + for _, c := range name { + if (c < '0' || c > '9') && (c < 'a' || c > 'f') { + return false + } + } + return true +} diff --git a/internal/storage/storage_test.go b/internal/storage/storage_test.go index 65f5a23f..67745d0c 100644 --- a/internal/storage/storage_test.go +++ b/internal/storage/storage_test.go @@ -9,6 +9,19 @@ import ( "testing" ) +func TestIsFetchDir(t *testing.T) { + for range 10 { + if id := NewFetchID(); !isFetchDir(id) { + t.Errorf("isFetchDir(%q) = false for a fetch id", id) + } + } + for _, name := range []string{"1.0.0", "0123456789abcde", "0123456789abcdef0", "0123456789ABCDEF", "0123456789abcdeg", ""} { + if isFetchDir(name) { + t.Errorf("isFetchDir(%q) = true", name) + } + } +} + func TestArtifactPath(t *testing.T) { tests := []struct { ecosystem string From 436155e50cf474783a7acdcdc88e940f57c18f77 Mon Sep 17 00:00:00 2001 From: montehurd Date: Tue, 29 Sep 2026 10:24:22 -0700 Subject: [PATCH 2/2] Clear before deleting on eviction, and requeue failed reclaims Eviction deleted an object before checking its record still pointed at it. If a refetch had moved the record, that object was queued and a request that read the record earlier could still be opening it. Eviction now clears first, skips the delete if the record moved, and queues the object when the delete fails after a clear. Reclaim queues a failed delete again, which moves it behind the rest, so objects the backend keeps refusing cannot fill every batch. Also drop the clearArtifactCache passthrough and ClearCachedArtifact's unused context, note that only tests call ArtifactPath, and document that max_size does not count objects waiting to be reclaimed. --- docs/configuration.md | 2 + internal/database/queries.go | 29 ++++----- internal/handler/coalesce_semantics_test.go | 2 +- internal/handler/handler.go | 2 +- internal/handler/helm.go | 8 +-- internal/server/eviction.go | 19 +++--- internal/server/eviction_test.go | 69 +++++++++++++++++---- internal/server/reclaim.go | 9 ++- internal/server/reclaim_test.go | 15 ++++- internal/storage/storage.go | 3 +- 10 files changed, 113 insertions(+), 45 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index aea28f60..f09e0108 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -44,6 +44,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.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. + ### Amazon S3 ```yaml diff --git a/internal/database/queries.go b/internal/database/queries.go index 4417c171..12d076e2 100644 --- a/internal/database/queries.go +++ b/internal/database/queries.go @@ -462,21 +462,6 @@ func (db *DB) GetCachedArtifactCount() (int64, error) { // elsewhere is left alone, so a clear never orphans the object that fetch // committed. func (db *DB) ClearArtifactCache(versionPURL, filename, storagePath string) (bool, error) { - return db.clearArtifactCache(versionPURL, filename, storagePath) -} - -// DiscardArtifact clears the record as ClearArtifactCache does and queues -// storagePath for deletion, for callers that must not delete an object another -// request may be reading. -func (db *DB) DiscardArtifact(versionPURL, filename, storagePath string) error { - cleared, err := db.clearArtifactCache(versionPURL, filename, storagePath) - if err != nil || !cleared { - return err - } - return db.QueuePendingDelete(storagePath) -} - -func (db *DB) clearArtifactCache(versionPURL, filename, storagePath string) (bool, error) { query := db.Rebind(` UPDATE artifacts SET storage_path = NULL, content_hash = NULL, size = NULL, @@ -491,8 +476,20 @@ func (db *DB) clearArtifactCache(versionPURL, filename, storagePath string) (boo return n > 0, err } +// DiscardArtifact clears the record as ClearArtifactCache does and queues +// storagePath for deletion, for callers that must not delete an object another +// request may be reading. +func (db *DB) DiscardArtifact(versionPURL, filename, storagePath string) error { + cleared, err := db.ClearArtifactCache(versionPURL, filename, storagePath) + if err != nil || !cleared { + return err + } + return db.QueuePendingDelete(storagePath) +} + // QueuePendingDelete queues a storage path no record points at any more. -// Queueing a path again restarts its grace period. +// Queueing a path again, as reclaim does after a failed delete, restarts its +// grace period and moves it behind the rest. func (db *DB) QueuePendingDelete(path string) error { query := db.Rebind(` INSERT INTO pending_deletes (path, queued_at) VALUES (?, ?) diff --git a/internal/handler/coalesce_semantics_test.go b/internal/handler/coalesce_semantics_test.go index c391d743..2ca52022 100644 --- a/internal/handler/coalesce_semantics_test.go +++ b/internal/handler/coalesce_semantics_test.go @@ -402,7 +402,7 @@ func TestCoalesce_KeyIsReleasedAfterFetch(t *testing.T) { // A fresh miss for the same key must start a new fetch, not rejoin the old // entry. Clearing the cache record forces the miss path again. - if err := proxy.ClearCachedArtifact(context.Background(), "npm", "pkg", "1.0.0", "pkg-1.0.0.tgz"); err != nil { + if err := proxy.ClearCachedArtifact("npm", "pkg", "1.0.0", "pkg-1.0.0.tgz"); err != nil { t.Fatalf("clear cached artifact: %v", err) } res, err := proxy.GetOrFetchArtifactFromURL(context.Background(), diff --git a/internal/handler/handler.go b/internal/handler/handler.go index 88e08bd9..d7b4a202 100644 --- a/internal/handler/handler.go +++ b/internal/handler/handler.go @@ -269,7 +269,7 @@ func (p *Proxy) GetCachedArtifact(ctx context.Context, ecosystem, name, version, // ClearCachedArtifact clears an artifact cache record after an external // integrity check fails, and queues its stored bytes for deletion. -func (p *Proxy) ClearCachedArtifact(_ context.Context, ecosystem, name, version, filename string) error { +func (p *Proxy) ClearCachedArtifact(ecosystem, name, version, filename string) error { if p.DB == nil || p.Storage == nil { return nil } diff --git a/internal/handler/helm.go b/internal/handler/helm.go index 7f207a99..8cacdb67 100644 --- a/internal/handler/helm.go +++ b/internal/handler/helm.go @@ -94,7 +94,7 @@ func (h *HelmHandler) handleChart(w http.ResponseWriter, r *http.Request) { return } if cached != nil { - h.serveChart(w, r, repository, digest, filename, cached) + h.serveChart(w, repository, digest, filename, cached) return } @@ -121,15 +121,15 @@ func (h *HelmHandler) handleChart(w http.ResponseWriter, r *http.Request) { h.proxy.serveArtifactError(w, err, "failed to fetch chart") return } - h.serveChart(w, r, repository, digest, filename, result) + h.serveChart(w, repository, digest, filename, result) } -func (h *HelmHandler) serveChart(w http.ResponseWriter, r *http.Request, repository, digest, filename string, result *CacheResult) { +func (h *HelmHandler) serveChart(w http.ResponseWriter, repository, digest, filename string, result *CacheResult) { if !strings.EqualFold(result.Artifact.Digest.Encoded(), digest) { if result.Reader != nil { _ = result.Reader.Close() } - if clearErr := h.proxy.ClearCachedArtifact(r.Context(), helmMetadataEcosystem, repository, digest, filename); clearErr != nil { + if clearErr := h.proxy.ClearCachedArtifact(helmMetadataEcosystem, repository, digest, filename); clearErr != nil { h.proxy.Logger.Warn("failed to clear Helm chart with invalid digest", "error", clearErr) } http.Error(w, "chart digest verification failed", http.StatusBadGateway) diff --git a/internal/server/eviction.go b/internal/server/eviction.go index c4639f7a..f9a706ad 100644 --- a/internal/server/eviction.go +++ b/internal/server/eviction.go @@ -116,24 +116,27 @@ func evictBatch(ctx context.Context, db *database.DB, store storage.Storage, log continue } - if err := store.Delete(ctx, art.StoragePath.String); err != nil { - logger.Warn("eviction: failed to delete from storage", - "path", art.StoragePath.String, "error", err) - continue - } + path := art.StoragePath.String - recordCleared, err := db.ClearArtifactCache(art.VersionPURL, art.Filename, art.StoragePath.String) + // Clear before deleting: a record a newer fetch moved has left this + // path queued, and a request may still open it within the grace period. + recordCleared, err := db.ClearArtifactCache(art.VersionPURL, art.Filename, path) if err != nil { logger.Warn("eviction: failed to clear artifact record", "version_purl", art.VersionPURL, "filename", art.Filename, "error", err) continue } if !recordCleared { - // The record no longer points here, so this delete freed nothing - // the recorded size counts. continue } + if err := store.Delete(ctx, path); err != nil { + logger.Warn("eviction: failed to delete from storage, queueing it", "path", path, "error", err) + if err := db.QueuePendingDelete(path); err != nil { + logger.Warn("eviction: failed to queue object for deletion", "path", path, "error", err) + } + } + if art.Size.Valid { freed += art.Size.Int64 } diff --git a/internal/server/eviction_test.go b/internal/server/eviction_test.go index 58cd1dce..15355864 100644 --- a/internal/server/eviction_test.go +++ b/internal/server/eviction_test.go @@ -7,6 +7,7 @@ import ( "io" "log/slog" "path/filepath" + "slices" "strings" "sync/atomic" "testing" @@ -326,10 +327,10 @@ func runEvictionWithDeadline(t *testing.T, ctx context.Context, db *database.DB, } } -// TestEvictLRU_EndsPassWhenNothingCanBeEvicted is the loop that would otherwise -// never end. Records that fail to delete stay eligible, so the same batch comes -// back forever while the recorded size never drops. -func TestEvictLRU_EndsPassWhenNothingCanBeEvicted(t *testing.T) { +// TestEvictLRU_QueuesObjectsStorageRefusesToDelete evicts against a backend +// that refuses every delete. The records are cleared, so the pass makes +// progress, and the objects wait in the queue for reclaim to retry. +func TestEvictLRU_QueuesObjectsStorageRefusesToDelete(t *testing.T) { db, store := setupEvictionTest(t) ctx := context.Background() @@ -340,17 +341,55 @@ func TestEvictLRU_EndsPassWhenNothingCanBeEvicted(t *testing.T) { undeletable := &undeletableStorage{Storage: store} runEvictionWithDeadline(t, ctx, db, undeletable, 100) - // One attempt per record, then the pass ends. Both records survive for the - // next pass to retry. if got := undeletable.deletes.Load(); got != 2 { - t.Errorf("delete attempts = %d, want 2: one per record in the single batch", got) + t.Errorf("delete attempts = %d, want 2", got) } + count, err := db.GetCachedArtifactCount() + if err != nil { + t.Fatalf("failed to get count: %v", err) + } + if count != 0 { + t.Errorf("cached artifacts = %d, want 0", count) + } + want := []string{ + storage.ArtifactPath("npm", "", "new-pkg", "1.0.0", "new-pkg-1.0.0.tgz"), + storage.ArtifactPath("npm", "", "old-pkg", "1.0.0", "old-pkg-1.0.0.tgz"), + } + got := queuedPaths(t, db) + slices.Sort(got) + if !slices.Equal(got, want) { + t.Errorf("queue = %v, want %v", got, want) + } +} + +// TestEvictLRU_EndsPassWhenNothingCanBeCleared is the loop that would otherwise +// never end. A record that fails to clear stays eligible, so the same batch +// comes back forever while the recorded size never drops. +func TestEvictLRU_EndsPassWhenNothingCanBeCleared(t *testing.T) { + db, store := setupEvictionTest(t) + ctx := context.Background() + + now := time.Now() + seedArtifact(t, ctx, db, store, "old-pkg", 500, now.Add(-2*time.Hour)) + seedArtifact(t, ctx, db, store, "new-pkg", 500, now) + if _, err := db.Exec(`CREATE TRIGGER refuse_clear BEFORE UPDATE OF storage_path ON artifacts + BEGIN SELECT RAISE(FAIL, 'clear refused'); END`); err != nil { + t.Fatalf("creating trigger: %v", err) + } + + runEvictionWithDeadline(t, ctx, db, store, 0) + count, err := db.GetCachedArtifactCount() if err != nil { t.Fatalf("failed to get count: %v", err) } if count != 2 { - t.Errorf("cached artifacts = %d, want 2: a failed delete must not clear the record", count) + t.Errorf("cached artifacts = %d, want 2", count) + } + for _, name := range []string{"old-pkg", "new-pkg"} { + if ok, _ := store.Exists(ctx, storage.ArtifactPath("npm", "", name, "1.0.0", name+"-1.0.0.tgz")); !ok { + t.Errorf("%s deleted although its record was not cleared", name) + } } } @@ -405,9 +444,10 @@ func clearRecordedSize(t *testing.T, db *database.DB, versionPURL, filename stri } // TestEvictBatch_SkipsRecordThatMoved evicts from a row read before a newer -// fetch moved the record. Only the old object goes, and nothing is counted as -// freed, since the size in use is unchanged and counting it would end the pass -// while the cache is still over its limit. +// fetch moved the record. Neither object goes: the old one is queued, and a +// request that read the record before it moved may still open it. Nothing is +// counted as freed either, or the pass would end with the cache still over its +// limit. func TestEvictBatch_SkipsRecordThatMoved(t *testing.T) { db, store := setupEvictionTest(t) ctx := context.Background() @@ -439,4 +479,11 @@ func TestEvictBatch_SkipsRecordThatMoved(t *testing.T) { if ok, _ := store.Exists(ctx, newer); !ok { t.Error("the newer fetch's object was deleted") } + old := stale[0].StoragePath.String + if ok, _ := store.Exists(ctx, old); !ok { + t.Error("the old object was deleted during its grace period") + } + if got := queuedPaths(t, db); !slices.Equal(got, []string{old}) { + t.Errorf("queue = %v, want the old path kept for reclaim", got) + } } diff --git a/internal/server/reclaim.go b/internal/server/reclaim.go index deaeaa8a..72a4e8f7 100644 --- a/internal/server/reclaim.go +++ b/internal/server/reclaim.go @@ -36,7 +36,8 @@ func (s *Server) startReclaimLoop(ctx context.Context) { } // reclaimStorage deletes up to one batch of objects queued before cutoff. A -// delete that fails stays queued for the next pass to retry. +// 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) { paths, err := db.GetDuePendingDeletes(cutoff, reclaimBatch) if err != nil { @@ -49,7 +50,13 @@ func reclaimStorage(ctx context.Context, db *database.DB, store storage.Storage, return } if err := store.Delete(ctx, path); err != nil { + if ctx.Err() != nil { + return + } 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) + } continue } if err := db.RemovePendingDelete(path); err != nil { diff --git a/internal/server/reclaim_test.go b/internal/server/reclaim_test.go index cda8a316..679dcdfd 100644 --- a/internal/server/reclaim_test.go +++ b/internal/server/reclaim_test.go @@ -61,14 +61,22 @@ func TestReclaimStorageWaitsForGracePeriod(t *testing.T) { } } -func TestReclaimStorageKeepsFailedDeletesQueued(t *testing.T) { +// TestReclaimStorageRequeuesFailedDeletes checks a refused delete goes behind +// the rest of the queue, so objects the backend keeps refusing cannot hold every +// slot in a batch. +func TestReclaimStorageRequeuesFailedDeletes(t *testing.T) { db, store := setupEvictionTest(t) logger := slog.New(slog.NewTextHandler(io.Discard, nil)) const path = "npm/old/1.0.0/a/old-1.0.0.tgz" storeQueued(t, db, store, path) + queued := time.Now().UTC().Add(-2 * time.Hour) + if _, err := db.Exec(db.Rebind(`UPDATE pending_deletes SET queued_at = ? WHERE path = ?`), queued, path); err != nil { + t.Fatalf("backdating queue entry: %v", err) + } + cutoff := time.Now().Add(-time.Hour) undeletable := &undeletableStorage{Storage: store} - reclaimStorage(context.Background(), db, undeletable, logger, time.Now().Add(time.Hour)) + reclaimStorage(context.Background(), db, undeletable, logger, cutoff) if got := undeletable.deletes.Load(); got != 1 { t.Errorf("delete attempts = %d, want 1", got) @@ -76,6 +84,9 @@ func TestReclaimStorageKeepsFailedDeletesQueued(t *testing.T) { if got := queuedPaths(t, db); !slices.Equal(got, []string{path}) { t.Errorf("queue = %v, want the failed path kept for retry", got) } + if due, err := db.GetDuePendingDeletes(cutoff, 100); err != nil || len(due) != 0 { + t.Errorf("due after a failed delete = %v (err %v), want it moved behind the cutoff", due, err) + } } func TestReclaimStorageStopsWhenContextCanceled(t *testing.T) { diff --git a/internal/storage/storage.go b/internal/storage/storage.go index cb807f07..0f64ed77 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -75,7 +75,8 @@ type Storage interface { } // ArtifactPath builds the storage path artifacts were cached under before each -// fetch got its own; records from then still point at such paths. +// fetch got its own; records from then still point at such paths. Only tests +// call it, to build such records. // Format: {ecosystem}/{namespace}/{name}/{version}/{filename} // For packages without namespace: {ecosystem}/{name}/{version}/{filename} func ArtifactPath(ecosystem, namespace, name, version, filename string) string {