From 5802fd2e2c9c2bbb60040e71f7000e6982801fbd Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sun, 4 Oct 2026 15:45:30 +1100 Subject: [PATCH] Evict remote objects from every cache layer DELETE /objects only removed the artifacts row, so mutable indexes (S3 index object + Redis TTL/ETag keys) kept being served. Clear all layers and treat a trailing * as a prefix. --- internal/api/v2/local_evict_cleanup_test.go | 2 +- internal/api/v2/local_objects_test.go | 2 +- internal/api/v2/objects.go | 17 ++-- internal/api/v2/objects_evict_test.go | 32 ++++++++ internal/cache/redis.go | 25 ++++++ internal/database/artifacts.go | 7 ++ internal/proxy/engine.go | 23 ++++++ internal/proxy/evict_test.go | 89 +++++++++++++++++++++ internal/server/server.go | 4 +- internal/storage/s3.go | 11 +++ 10 files changed, 203 insertions(+), 9 deletions(-) create mode 100644 internal/api/v2/objects_evict_test.go create mode 100644 internal/proxy/evict_test.go diff --git a/internal/api/v2/local_evict_cleanup_test.go b/internal/api/v2/local_evict_cleanup_test.go index a6d493d..dc4d0f2 100644 --- a/internal/api/v2/local_evict_cleanup_test.go +++ b/internal/api/v2/local_evict_cleanup_test.go @@ -49,7 +49,7 @@ func TestLocalEvictCleansRPMMetadata(t *testing.T) { t.Fatal(err) } - h := NewObjectsHandler(db) + h := NewObjectsHandler(db, nil) router := chi.NewRouter() router.Route("/locals/{name}/objects", func(r chi.Router) { r.Delete("/*", h.LocalRoutes().ServeHTTP) diff --git a/internal/api/v2/local_objects_test.go b/internal/api/v2/local_objects_test.go index 9f0d029..e908f18 100644 --- a/internal/api/v2/local_objects_test.go +++ b/internal/api/v2/local_objects_test.go @@ -40,7 +40,7 @@ func TestLocalObjectsListing(t *testing.T) { t.Fatal(err) } - h := NewObjectsHandler(db) + h := NewObjectsHandler(db, nil) router := chi.NewRouter() router.Route("/locals/{name}/objects", func(r chi.Router) { r.Get("/", h.LocalRoutes().ServeHTTP) diff --git a/internal/api/v2/objects.go b/internal/api/v2/objects.go index 4505c3e..361a628 100644 --- a/internal/api/v2/objects.go +++ b/internal/api/v2/objects.go @@ -1,6 +1,7 @@ package v2 import ( + "context" "fmt" "net/http" "strconv" @@ -10,12 +11,18 @@ import ( "git.unkin.net/unkin/artifactapi/internal/database" ) -type ObjectsHandler struct { - db *database.DB +// Evictor drops a remote path from every cache layer. +type Evictor interface { + Evict(ctx context.Context, remoteName, path string) error } -func NewObjectsHandler(db *database.DB) *ObjectsHandler { - return &ObjectsHandler{db: db} +type ObjectsHandler struct { + db *database.DB + evictor Evictor +} + +func NewObjectsHandler(db *database.DB, evictor Evictor) *ObjectsHandler { + return &ObjectsHandler{db: db, evictor: evictor} } func (h *ObjectsHandler) Routes() chi.Router { @@ -86,7 +93,7 @@ func (h *ObjectsHandler) evict(w http.ResponseWriter, r *http.Request) { remoteName := chi.URLParam(r, "name") path := chi.URLParam(r, "*") - if err := h.db.DeleteArtifact(r.Context(), remoteName, path); err != nil { + if err := h.evictor.Evict(r.Context(), remoteName, path); err != nil { http.Error(w, fmt.Sprintf("evict failed: %v", err), http.StatusInternalServerError) return } diff --git a/internal/api/v2/objects_evict_test.go b/internal/api/v2/objects_evict_test.go new file mode 100644 index 0000000..b06c9f4 --- /dev/null +++ b/internal/api/v2/objects_evict_test.go @@ -0,0 +1,32 @@ +package v2 + +import ( + "context" + "net/http/httptest" + "testing" + + "github.com/go-chi/chi/v5" +) + +type fakeEvictor struct{ remote, path string } + +func (f *fakeEvictor) Evict(_ context.Context, remote, path string) error { + f.remote, f.path = remote, path + return nil +} + +func TestRemoteEvictDelegatesToEvictor(t *testing.T) { + for _, path := range []string{"8/Everything/x86_64/repodata/repomd.xml", "8/Everything/x86_64/repodata/*"} { + ev := &fakeEvictor{} + h := NewObjectsHandler(nil, ev) + router := chi.NewRouter() + router.Route("/remotes/{name}/objects", func(r chi.Router) { + r.Delete("/*", h.Routes().ServeHTTP) + }) + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("DELETE", "/remotes/epel/objects/"+path, nil)) + if w.Code != 204 || ev.remote != "epel" || ev.path != path { + t.Errorf("DELETE %s: code=%d evicted=%q/%q", path, w.Code, ev.remote, ev.path) + } + } +} diff --git a/internal/cache/redis.go b/internal/cache/redis.go index 529d88e..b33e0cc 100644 --- a/internal/cache/redis.go +++ b/internal/cache/redis.go @@ -3,6 +3,7 @@ package cache import ( "context" "fmt" + "strings" "time" "github.com/redis/go-redis/v9" @@ -115,3 +116,27 @@ func (r *Redis) FlushRemote(ctx context.Context, remote string) error { } return iter.Err() } + +// ForgetPath drops the freshness and ETag keys of one cached path. +func (r *Redis) ForgetPath(ctx context.Context, remote, path string) error { + return r.client.Del(ctx, fmt.Sprintf("ttl:%s:%s", remote, path), fmt.Sprintf("etag:%s:%s", remote, path)).Err() +} + +// ForgetPrefix drops the freshness and ETag keys of every cached path under prefix. +func (r *Redis) ForgetPrefix(ctx context.Context, remote, prefix string) error { + glob := globEscaper.Replace(remote + ":" + prefix) + for _, kind := range []string{"ttl:", "etag:"} { + iter := r.client.Scan(ctx, 0, kind+glob+"*", 100).Iterator() + for iter.Next(ctx) { + if err := r.client.Del(ctx, iter.Val()).Err(); err != nil { + return err + } + } + if err := iter.Err(); err != nil { + return err + } + } + return nil +} + +var globEscaper = strings.NewReplacer(`\`, `\\`, `*`, `\*`, `?`, `\?`, `[`, `\[`, `]`, `\]`) diff --git a/internal/database/artifacts.go b/internal/database/artifacts.go index ab1046f..da24e46 100644 --- a/internal/database/artifacts.go +++ b/internal/database/artifacts.go @@ -103,6 +103,13 @@ func (db *DB) DeleteArtifact(ctx context.Context, remoteName, path string) error return err } +// DeleteArtifactsByPrefix removes every artifact row of a remote whose path +// starts with prefix. +func (db *DB) DeleteArtifactsByPrefix(ctx context.Context, remoteName, prefix string) error { + _, err := db.Pool.Exec(ctx, `DELETE FROM artifacts WHERE remote_name = $1 AND left(path, length($2)) = $2`, remoteName, prefix) + return err +} + func (db *DB) InsertAccessLog(ctx context.Context, remoteName, path string, cacheHit bool, sizeBytes int64, upstreamMS int, clientIP string) error { _, err := db.Pool.Exec(ctx, ` INSERT INTO access_log (remote_name, path, cache_hit, size_bytes, upstream_ms, client_ip) diff --git a/internal/proxy/engine.go b/internal/proxy/engine.go index 5040ead..c59af9b 100644 --- a/internal/proxy/engine.go +++ b/internal/proxy/engine.go @@ -208,6 +208,29 @@ func (e *Engine) Fetch(ctx context.Context, remote models.Remote, path string, p return result, nil } +// Evict drops path from every cache layer (artifact row, index object, Redis +// freshness and ETag keys) so the next request refetches from upstream. A +// trailing "*" evicts every path under that prefix. +func (e *Engine) Evict(ctx context.Context, remoteName, path string) error { + prefix, wildcard := strings.CutSuffix(path, "*") + if !wildcard { + if err := e.db.DeleteArtifact(ctx, remoteName, path); err != nil { + return fmt.Errorf("delete artifact: %w", err) + } + if err := e.store.Delete(ctx, storage.IndexKey(remoteName, path)); err != nil { + return fmt.Errorf("delete index: %w", err) + } + return e.cache.ForgetPath(ctx, remoteName, path) + } + if err := e.db.DeleteArtifactsByPrefix(ctx, remoteName, prefix); err != nil { + return fmt.Errorf("delete artifacts: %w", err) + } + if err := e.store.DeletePrefix(ctx, storage.IndexKey(remoteName, prefix)); err != nil { + return fmt.Errorf("delete indexes: %w", err) + } + return e.cache.ForgetPrefix(ctx, remoteName, prefix) +} + // HeadResult carries artifact metadata for a HEAD request. There is no body. type HeadResult struct { ContentType string diff --git a/internal/proxy/evict_test.go b/internal/proxy/evict_test.go new file mode 100644 index 0000000..fc5f5dd --- /dev/null +++ b/internal/proxy/evict_test.go @@ -0,0 +1,89 @@ +package proxy + +import ( + "context" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + + _ "git.unkin.net/unkin/artifactapi/internal/provider/rpm" + "git.unkin.net/unkin/artifactapi/pkg/models" +) + +// changingUpstream serves every path with the current revision, as a mirror +// does after a sync replaces its repodata. +func changingUpstream(t *testing.T) (*httptest.Server, *atomic.Value) { + t.Helper() + var rev atomic.Value + rev.Store("rev1") + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + v := rev.Load().(string) + w.Header().Set("ETag", `"`+v+`"`) + _, _ = w.Write([]byte(v + ":" + r.URL.Path)) + })) + t.Cleanup(srv.Close) + return srv, &rev +} + +func fetchBody(t *testing.T, r models.Remote, path string) string { + t.Helper() + res, err := testEngine.Fetch(context.Background(), r, path, prov(t, models.PackageRPM)) + if err != nil { + t.Fatalf("fetch %s: %v", path, err) + } + return readAll(t, res) +} + +func rpmRemote(t *testing.T, name, baseURL string) models.Remote { + return seed(t, models.Remote{Name: name, PackageType: models.PackageRPM, RepoType: models.RepoTypeRemote, BaseURL: baseURL, MutableTTL: 7200, CheckMutable: true}) +} + +func TestEvictMutableIndexRefetches(t *testing.T) { + requireStack(t) + srv, rev := changingUpstream(t) + r := rpmRemote(t, "evict-idx", srv.URL) + const path = "8/Everything/x86_64/repodata/repomd.xml" + + if got := fetchBody(t, r, path); got != "rev1:/"+path { + t.Fatalf("initial fetch = %q", got) + } + rev.Store("rev2") + if got := fetchBody(t, r, path); got != "rev1:/"+path { + t.Fatalf("within TTL = %q, want cached rev1", got) + } + if err := testEngine.Evict(context.Background(), r.Name, path); err != nil { + t.Fatalf("evict: %v", err) + } + if got := fetchBody(t, r, path); got != "rev2:/"+path { + t.Fatalf("after evict = %q, want rev2", got) + } +} + +func TestEvictWildcardClearsPrefixOnly(t *testing.T) { + requireStack(t) + srv, rev := changingUpstream(t) + r := rpmRemote(t, "evict-wild", srv.URL) + const ( + repomd = "8/Everything/x86_64/repodata/repomd.xml" + rpm = "8/Everything/x86_64/Packages/a/a-1.0-1.el8.x86_64.rpm" + other = "9/Everything/x86_64/repodata/repomd.xml" + ) + for _, p := range []string{repomd, rpm, other} { + fetchBody(t, r, p) + } + rev.Store("rev2") + + if err := testEngine.Evict(context.Background(), r.Name, "8/Everything/x86_64/*"); err != nil { + t.Fatalf("evict: %v", err) + } + if got := fetchBody(t, r, repomd); got != "rev2:/"+repomd { + t.Errorf("index under prefix = %q, want rev2", got) + } + if got := fetchBody(t, r, rpm); got != "rev2:/"+rpm { + t.Errorf("artifact under prefix = %q, want rev2", got) + } + if got := fetchBody(t, r, other); got != "rev1:/"+other { + t.Errorf("path outside prefix = %q, want cached rev1", got) + } +} diff --git a/internal/server/server.go b/internal/server/server.go index 18cf55d..85622c5 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -195,13 +195,13 @@ func (s *Server) routes() chi.Router { r.Mount("/probe", probeHandler.Routes()) r.Route("/remotes/{name}/objects", func(r chi.Router) { - objHandler := v2.NewObjectsHandler(s.db) + objHandler := v2.NewObjectsHandler(s.db, s.engine) r.Get("/", objHandler.Routes().ServeHTTP) r.Delete("/*", objHandler.Routes().ServeHTTP) }) r.Route("/locals/{name}/objects", func(r chi.Router) { - objHandler := v2.NewObjectsHandler(s.db) + objHandler := v2.NewObjectsHandler(s.db, s.engine) r.Get("/", objHandler.LocalRoutes().ServeHTTP) r.Delete("/*", objHandler.LocalRoutes().ServeHTTP) }) diff --git a/internal/storage/s3.go b/internal/storage/s3.go index 97e7da2..1aa2ecb 100644 --- a/internal/storage/s3.go +++ b/internal/storage/s3.go @@ -99,6 +99,17 @@ func (s *S3) Stat(ctx context.Context, key string) (*minio.ObjectInfo, error) { return &info, nil } +// DeletePrefix removes every object whose key starts with prefix. +func (s *S3) DeletePrefix(ctx context.Context, prefix string) error { + objects := s.client.ListObjects(ctx, s.bucket, minio.ListObjectsOptions{Prefix: prefix, Recursive: true}) + for res := range s.client.RemoveObjects(ctx, s.bucket, objects, minio.RemoveObjectsOptions{}) { + if res.Err != nil { + return res.Err + } + } + return nil +} + // ListStaleObjects returns keys under prefix last modified before cutoff. Used // by the GC to reap abandoned staging objects (e.g. cancelled docker pushes). func (s *S3) ListStaleObjects(ctx context.Context, prefix string, cutoff time.Time) ([]string, error) {