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.
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
Vendored
+25
@@ -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(`\`, `\\`, `*`, `\*`, `?`, `\?`, `[`, `\[`, `]`, `\]`)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user