diff --git a/internal/api/v2/objects.go b/internal/api/v2/objects.go
index 4505c3e..7d66451 100644
--- a/internal/api/v2/objects.go
+++ b/internal/api/v2/objects.go
@@ -1,6 +1,8 @@
package v2
import (
+ "context"
+ "errors"
"fmt"
"net/http"
"strconv"
@@ -8,8 +10,14 @@ import (
"github.com/go-chi/chi/v5"
"git.unkin.net/unkin/artifactapi/internal/database"
+ "git.unkin.net/unkin/artifactapi/internal/proxy"
)
+// Evictor drops a remote path from every cache layer.
+type Evictor interface {
+ Evict(ctx context.Context, remoteName, path string) error
+}
+
type ObjectsHandler struct {
db *database.DB
}
@@ -18,10 +26,13 @@ func NewObjectsHandler(db *database.DB) *ObjectsHandler {
return &ObjectsHandler{db: db}
}
-func (h *ObjectsHandler) Routes() chi.Router {
+// Routes lists and evicts objects for remote repos; evictor serves the DELETE.
+func (h *ObjectsHandler) Routes(evictor Evictor) chi.Router {
r := chi.NewRouter()
r.Get("/", h.list)
- r.Delete("/*", h.evict)
+ r.Delete("/*", func(w http.ResponseWriter, r *http.Request) {
+ evict(w, r, evictor)
+ })
return r
}
@@ -82,11 +93,16 @@ func (h *ObjectsHandler) evictLocal(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusNoContent)
}
-func (h *ObjectsHandler) evict(w http.ResponseWriter, r *http.Request) {
+func evict(w http.ResponseWriter, r *http.Request, evictor Evictor) {
remoteName := chi.URLParam(r, "name")
path := chi.URLParam(r, "*")
- if err := h.db.DeleteArtifact(r.Context(), remoteName, path); err != nil {
+ if err := evictor.Evict(r.Context(), remoteName, path); err != nil {
+ var proxyErr *proxy.ProxyError
+ if errors.As(err, &proxyErr) {
+ http.Error(w, proxyErr.Message, proxyErr.Status)
+ return
+ }
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..75f38c2
--- /dev/null
+++ b/internal/api/v2/objects_evict_test.go
@@ -0,0 +1,59 @@
+package v2
+
+import (
+ "context"
+ "errors"
+ "fmt"
+ "net/http"
+ "net/http/httptest"
+ "testing"
+
+ "github.com/go-chi/chi/v5"
+
+ "git.unkin.net/unkin/artifactapi/internal/proxy"
+)
+
+type fakeEvictor struct {
+ remote, path string
+ err error
+}
+
+func (f *fakeEvictor) Evict(_ context.Context, remote, path string) error {
+ f.remote, f.path = remote, path
+ return f.err
+}
+
+func deleteObject(ev Evictor, path string) int {
+ router := chi.NewRouter()
+ router.Route("/remotes/{name}/objects", func(r chi.Router) {
+ r.Delete("/*", NewObjectsHandler(nil).Routes(ev).ServeHTTP)
+ })
+ w := httptest.NewRecorder()
+ router.ServeHTTP(w, httptest.NewRequest("DELETE", "/remotes/epel/objects/"+path, nil))
+ return w.Code
+}
+
+func TestRemoteEvictDelegatesToEvictor(t *testing.T) {
+ for _, path := range []string{"8/Everything/x86_64/repodata/repomd.xml", "8/Everything/x86_64/repodata/*"} {
+ ev := &fakeEvictor{}
+ if code := deleteObject(ev, path); code != 204 || ev.remote != "epel" || ev.path != path {
+ t.Errorf("DELETE %s: code=%d evicted=%q/%q", path, code, ev.remote, ev.path)
+ }
+ }
+}
+
+func TestRemoteEvictMapsErrorStatus(t *testing.T) {
+ for name, tc := range map[string]struct {
+ err error
+ want int
+ }{
+ "bad wildcard": {&proxy.ProxyError{Status: http.StatusBadRequest, Message: "wildcard evict must be
/*"}, http.StatusBadRequest},
+ "unknown remote": {&proxy.ProxyError{Status: http.StatusNotFound, Message: "remote not found"}, http.StatusNotFound},
+ "lock busy": {fmt.Errorf("wrapped: %w", &proxy.ProxyError{Status: http.StatusServiceUnavailable, Message: "retry"}), http.StatusServiceUnavailable},
+ "other error": {errors.New("db down"), http.StatusInternalServerError},
+ } {
+ if code := deleteObject(&fakeEvictor{err: tc.err}, "x"); code != tc.want {
+ t.Errorf("%s: code = %d, want %d", name, code, tc.want)
+ }
+ }
+}
diff --git a/internal/cache/cache_test.go b/internal/cache/cache_test.go
index 5edbf41..4093a60 100644
--- a/internal/cache/cache_test.go
+++ b/internal/cache/cache_test.go
@@ -131,3 +131,72 @@ func TestFlushRemote(t *testing.T) {
t.Error("expected keys flushed")
}
}
+
+func setPathKeys(t *testing.T, remote string, paths ...string) {
+ t.Helper()
+ ctx := context.Background()
+ for _, p := range paths {
+ if err := testRedis.SetTTL(ctx, remote, p, time.Minute); err != nil {
+ t.Fatal(err)
+ }
+ if err := testRedis.SetETag(ctx, remote, p, `"e"`, time.Minute); err != nil {
+ t.Fatal(err)
+ }
+ }
+}
+
+func pathKeysExist(t *testing.T, remote, path string) (ttl, etag bool) {
+ t.Helper()
+ ctx := context.Background()
+ ttl, _ = testRedis.CheckTTL(ctx, remote, path)
+ e, _ := testRedis.GetETag(ctx, remote, path)
+ return ttl, e != ""
+}
+
+func TestForgetPath(t *testing.T) {
+ requireRedis(t)
+ const meta = `repo/a*b?[c]\d.xml`
+ setPathKeys(t, "fp", meta, "repo/aXb.xml")
+ setPathKeys(t, "fp-other", meta)
+
+ if err := testRedis.ForgetPath(context.Background(), "fp", meta); err != nil {
+ t.Fatal(err)
+ }
+ if ttl, etag := pathKeysExist(t, "fp", meta); ttl || etag {
+ t.Errorf("forgotten path keys remain: ttl=%v etag=%v", ttl, etag)
+ }
+ for _, k := range [][2]string{{"fp", "repo/aXb.xml"}, {"fp-other", meta}} {
+ if ttl, etag := pathKeysExist(t, k[0], k[1]); !ttl || !etag {
+ t.Errorf("%s:%s lost keys: ttl=%v etag=%v", k[0], k[1], ttl, etag)
+ }
+ }
+}
+
+func TestForgetPrefix(t *testing.T) {
+ requireRedis(t)
+ const prefix = `r*[1]?\/`
+ under := []string{prefix + "repomd.xml", prefix + "sub/x.rpm"}
+ // Each would match the prefix if its glob metacharacters were left unescaped.
+ globMatches := []string{`rX11/repomd.xml`, `r[1]?\/x`, `r*1Z/x`}
+ setPathKeys(t, "fx", append(under, globMatches...)...)
+ setPathKeys(t, "fx-other", under...)
+
+ if err := testRedis.ForgetPrefix(context.Background(), "fx", prefix); err != nil {
+ t.Fatal(err)
+ }
+ for _, p := range under {
+ if ttl, etag := pathKeysExist(t, "fx", p); ttl || etag {
+ t.Errorf("%s keys remain: ttl=%v etag=%v", p, ttl, etag)
+ }
+ }
+ for _, p := range globMatches {
+ if ttl, etag := pathKeysExist(t, "fx", p); !ttl || !etag {
+ t.Errorf("%s outside prefix lost keys: ttl=%v etag=%v", p, ttl, etag)
+ }
+ }
+ for _, p := range under {
+ if ttl, etag := pathKeysExist(t, "fx-other", p); !ttl || !etag {
+ t.Errorf("other remote %s lost keys: ttl=%v etag=%v", p, ttl, etag)
+ }
+ }
+}
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..45a1b87 100644
--- a/internal/proxy/engine.go
+++ b/internal/proxy/engine.go
@@ -16,6 +16,8 @@ import (
"sync/atomic"
"time"
+ "github.com/jackc/pgx/v5"
+
"git.unkin.net/unkin/artifactapi/internal/cache"
"git.unkin.net/unkin/artifactapi/internal/database"
"git.unkin.net/unkin/artifactapi/internal/provider"
@@ -47,6 +49,8 @@ type Engine struct {
// mirror strategy to prefer the mirror currently handling the fewest
// requests. Per-replica and approximate, which is fine.
inflight sync.Map
+ // evictLockWait bounds how long Evict waits on a held fetch lock.
+ evictLockWait time.Duration
}
func NewEngine(db *database.DB, c *cache.Redis, s *storage.S3) *Engine {
@@ -57,6 +61,8 @@ func NewEngine(db *database.DB, c *cache.Redis, s *storage.S3) *Engine {
cas: storage.NewCAS(s),
circuit: NewCircuitBreaker(c),
accessLog: make(chan database.AccessLogEntry, accessLogBufferSize),
+
+ evictLockWait: fetchLockTTL,
}
go e.runAccessLogWriter()
return e
@@ -208,6 +214,71 @@ 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 directory.
+func (e *Engine) Evict(ctx context.Context, remoteName, path string) error {
+ prefix, wildcard := strings.CutSuffix(path, "*")
+ if wildcard && !strings.HasSuffix(prefix, "/") {
+ return &ProxyError{Status: http.StatusBadRequest, Message: "wildcard evict must be /*"}
+ }
+ if _, err := e.db.GetRemote(ctx, remoteName); errors.Is(err, pgx.ErrNoRows) {
+ return &ProxyError{Status: http.StatusNotFound, Message: fmt.Sprintf("remote %q not found", remoteName)}
+ } else if err != nil {
+ return fmt.Errorf("get remote: %w", err)
+ }
+ if !wildcard {
+ if err := e.waitForLock(ctx, remoteName, path); err != nil {
+ return err
+ }
+ defer func() { _ = e.cache.ReleaseLock(context.WithoutCancel(ctx), remoteName, path) }()
+ 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)
+ }
+ // ponytail: no lock for wildcards; a Fetch already in flight under the
+ // prefix can re-cache its path after the evict. Per-path locks over a
+ // directory would close it if that ever matters.
+ 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)
+}
+
+// waitForLock takes the per-path fetch lock so an in-flight Fetch cannot
+// re-set TTL/ETag keys after an evict. It fails with a 503 when the lock
+// cannot be taken within evictLockWait or Redis errors.
+func (e *Engine) waitForLock(ctx context.Context, remoteName, path string) error {
+ deadline := time.Now().Add(e.evictLockWait)
+ for {
+ ok, err := e.cache.AcquireLock(ctx, remoteName, path, fetchLockTTL)
+ if ok {
+ return nil
+ }
+ if ctx.Err() != nil {
+ return ctx.Err()
+ }
+ if err != nil {
+ return &ProxyError{Status: http.StatusServiceUnavailable, Message: fmt.Sprintf("fetch lock: %v", err)}
+ }
+ if time.Now().After(deadline) {
+ return &ProxyError{Status: http.StatusServiceUnavailable, Message: "fetch in progress, retry evict"}
+ }
+ select {
+ case <-ctx.Done():
+ return ctx.Err()
+ case <-time.After(50 * time.Millisecond):
+ }
+ }
+}
+
// 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..15b8def
--- /dev/null
+++ b/internal/proxy/evict_test.go
@@ -0,0 +1,271 @@
+package proxy
+
+import (
+ "context"
+ "errors"
+ "net/http"
+ "net/http/httptest"
+ "sync/atomic"
+ "testing"
+ "time"
+
+ _ "git.unkin.net/unkin/artifactapi/internal/provider/rpm"
+ "git.unkin.net/unkin/artifactapi/internal/storage"
+ "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. It answers 304 to a matching
+// If-None-Match and counts conditional requests.
+func changingUpstream(t *testing.T) (*httptest.Server, *atomic.Value, *atomic.Int32) {
+ t.Helper()
+ var rev atomic.Value
+ var conditional atomic.Int32
+ rev.Store("rev1")
+ srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
+ v := rev.Load().(string)
+ etag := `"` + v + `"`
+ w.Header().Set("ETag", etag)
+ if inm := r.Header.Get("If-None-Match"); inm != "" {
+ conditional.Add(1)
+ if inm == etag {
+ w.WriteHeader(http.StatusNotModified)
+ return
+ }
+ }
+ _, _ = w.Write([]byte(v + ":" + r.URL.Path))
+ }))
+ t.Cleanup(srv.Close)
+ return srv, &rev, &conditional
+}
+
+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})
+}
+
+// cached reports which cache layers hold path: artifact row, index object,
+// Redis TTL key, Redis ETag key.
+type cached struct{ row, index, ttl, etag bool }
+
+// Mutable indexes live in the S3 index; immutable blobs get an artifact row.
+var (
+ indexCached = cached{index: true, ttl: true, etag: true}
+ blobCached = cached{row: true, ttl: true, etag: true}
+)
+
+func layers(t *testing.T, remote, path string) cached {
+ t.Helper()
+ ctx := context.Background()
+ var c cached
+ _, err := testDB.GetArtifact(ctx, remote, path)
+ c.row = err == nil
+ c.index, err = testEngine.store.Exists(ctx, storage.IndexKey(remote, path))
+ if err != nil {
+ t.Fatalf("stat index %s: %v", path, err)
+ }
+ c.ttl, _ = testCache.CheckTTL(ctx, remote, path)
+ etag, _ := testCache.GetETag(ctx, remote, path)
+ c.etag = etag != ""
+ return c
+}
+
+func TestEvictMutableIndexRefetches(t *testing.T) {
+ requireStack(t)
+ srv, rev, conditional := 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)
+ }
+ if c := layers(t, r.Name, path); c != indexCached {
+ t.Fatalf("before evict = %+v, want %+v", c, indexCached)
+ }
+ 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 c := layers(t, r.Name, path); c != (cached{}) {
+ t.Fatalf("after evict = %+v, want every layer gone", c)
+ }
+ conditional.Store(0)
+ if got := fetchBody(t, r, path); got != "rev2:/"+path {
+ t.Fatalf("after evict = %q, want rev2", got)
+ }
+ if n := conditional.Load(); n != 0 {
+ t.Errorf("post-evict fetch revalidated with the evicted ETag (%d conditional requests)", n)
+ }
+}
+
+func TestEvictImmutableBlobDropsRow(t *testing.T) {
+ requireStack(t)
+ srv, _, _ := changingUpstream(t)
+ r := rpmRemote(t, "evict-blob", srv.URL)
+ const path = "8/Everything/x86_64/Packages/a/a-1.0-1.el8.x86_64.rpm"
+ fetchBody(t, r, path)
+ if c := layers(t, r.Name, path); c != blobCached {
+ t.Fatalf("before evict = %+v, want %+v", c, blobCached)
+ }
+ if err := testEngine.Evict(context.Background(), r.Name, path); err != nil {
+ t.Fatalf("evict: %v", err)
+ }
+ if c := layers(t, r.Name, path); c != (cached{}) {
+ t.Fatalf("after evict = %+v, want every layer gone", c)
+ }
+}
+
+func TestEvictRejectsNonDirectoryWildcard(t *testing.T) {
+ requireStack(t)
+ srv, _, _ := changingUpstream(t)
+ r := rpmRemote(t, "evict-bare", srv.URL)
+ const path = "8/Everything/x86_64/repodata/repomd.xml"
+ fetchBody(t, r, path)
+
+ for _, bad := range []string{"*", "8*", "8/Every*"} {
+ var pe *ProxyError
+ if err := testEngine.Evict(context.Background(), r.Name, bad); !errors.As(err, &pe) || pe.Status != http.StatusBadRequest {
+ t.Errorf("evict %s = %v, want 400", bad, err)
+ }
+ }
+ if c := layers(t, r.Name, path); c != (indexCached) {
+ t.Errorf("rejected wildcard evicted layers: %+v", c)
+ }
+}
+
+func TestEvictUnknownRemote(t *testing.T) {
+ requireStack(t)
+ var pe *ProxyError
+ if err := testEngine.Evict(context.Background(), "evict-no-such-remote", "a/b"); !errors.As(err, &pe) || pe.Status != http.StatusNotFound {
+ t.Fatalf("evict unknown remote = %v, want 404", err)
+ }
+}
+
+func holdLock(t *testing.T, remote, path string) {
+ t.Helper()
+ ctx := context.Background()
+ if ok, err := testCache.AcquireLock(ctx, remote, path, time.Minute); !ok || err != nil {
+ t.Fatalf("acquire: %v %v", ok, err)
+ }
+ t.Cleanup(func() { _ = testCache.ReleaseLock(ctx, remote, path) })
+}
+
+func TestEvictWaitsForFetchLock(t *testing.T) {
+ requireStack(t)
+ srv, _, _ := changingUpstream(t)
+ r := rpmRemote(t, "evict-lock", srv.URL)
+ ctx := context.Background()
+ const path = "8/Everything/x86_64/repodata/repomd.xml"
+ holdLock(t, r.Name, path)
+ done := make(chan error, 1)
+ go func() { done <- testEngine.Evict(ctx, r.Name, path) }()
+ select {
+ case err := <-done:
+ t.Fatalf("evict returned while fetch lock held: %v", err)
+ case <-time.After(200 * time.Millisecond):
+ }
+ _ = testCache.ReleaseLock(ctx, r.Name, path)
+ if err := <-done; err != nil {
+ t.Fatalf("evict: %v", err)
+ }
+ ok, err := testCache.AcquireLock(ctx, r.Name, path, time.Second)
+ if !ok || err != nil {
+ t.Fatalf("lock still held after evict: %v %v", ok, err)
+ }
+}
+
+func TestEvictCancelledWhileWaitingDeletesNothing(t *testing.T) {
+ requireStack(t)
+ srv, _, _ := changingUpstream(t)
+ r := rpmRemote(t, "evict-cancel", srv.URL)
+ const path = "8/Everything/x86_64/repodata/repomd.xml"
+ fetchBody(t, r, path)
+ holdLock(t, r.Name, path)
+
+ ctx, cancel := context.WithCancel(context.Background())
+ done := make(chan error, 1)
+ go func() { done <- testEngine.Evict(ctx, r.Name, path) }()
+ time.Sleep(100 * time.Millisecond)
+ cancel()
+ select {
+ case err := <-done:
+ if err == nil {
+ t.Fatal("cancelled evict returned nil")
+ }
+ case <-time.After(2 * time.Second):
+ t.Fatal("cancelled evict did not return")
+ }
+ if c := layers(t, r.Name, path); c != (indexCached) {
+ t.Errorf("cancelled evict deleted layers: %+v", c)
+ }
+}
+
+func TestEvictLockTimeoutIs503(t *testing.T) {
+ requireStack(t)
+ srv, _, _ := changingUpstream(t)
+ r := rpmRemote(t, "evict-timeout", srv.URL)
+ const path = "8/Everything/x86_64/repodata/repomd.xml"
+ fetchBody(t, r, path)
+ holdLock(t, r.Name, path)
+ prev := testEngine.evictLockWait
+ testEngine.evictLockWait = 100 * time.Millisecond
+ t.Cleanup(func() { testEngine.evictLockWait = prev })
+
+ var pe *ProxyError
+ if err := testEngine.Evict(context.Background(), r.Name, path); !errors.As(err, &pe) || pe.Status != http.StatusServiceUnavailable {
+ t.Fatalf("evict = %v, want 503", err)
+ }
+ if c := layers(t, r.Name, path); c != (indexCached) {
+ t.Errorf("timed-out evict deleted layers: %+v", c)
+ }
+}
+
+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)
+ }
+ if c := layers(t, r.Name, rpm); c != blobCached {
+ t.Fatalf("%s before evict = %+v, want %+v", rpm, c, blobCached)
+ }
+ rev.Store("rev2")
+
+ if err := testEngine.Evict(context.Background(), r.Name, "8/Everything/x86_64/*"); err != nil {
+ t.Fatalf("evict: %v", err)
+ }
+ for _, p := range []string{repomd, rpm} {
+ if c := layers(t, r.Name, p); c != (cached{}) {
+ t.Errorf("%s after evict = %+v, want every layer gone", p, c)
+ }
+ }
+ if c := layers(t, r.Name, other); c != (indexCached) {
+ t.Errorf("%s outside prefix = %+v, want %+v", other, c, indexCached)
+ }
+ 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..19e3677 100644
--- a/internal/server/server.go
+++ b/internal/server/server.go
@@ -196,8 +196,8 @@ func (s *Server) routes() chi.Router {
r.Route("/remotes/{name}/objects", func(r chi.Router) {
objHandler := v2.NewObjectsHandler(s.db)
- r.Get("/", objHandler.Routes().ServeHTTP)
- r.Delete("/*", objHandler.Routes().ServeHTTP)
+ r.Get("/", objHandler.Routes(s.engine).ServeHTTP)
+ r.Delete("/*", objHandler.Routes(s.engine).ServeHTTP)
})
r.Route("/locals/{name}/objects", func(r chi.Router) {
diff --git a/internal/storage/s3.go b/internal/storage/s3.go
index 97e7da2..817db94 100644
--- a/internal/storage/s3.go
+++ b/internal/storage/s3.go
@@ -99,6 +99,30 @@ func (s *S3) Stat(ctx context.Context, key string) (*minio.ObjectInfo, error) {
return &info, nil
}
+// DeletePrefix removes every object whose key starts with prefix. It lists
+// first so a failed list returns its error instead of deleting nothing.
+func (s *S3) DeletePrefix(ctx context.Context, prefix string) error {
+ var objs []minio.ObjectInfo
+ for obj := range s.client.ListObjects(ctx, s.bucket, minio.ListObjectsOptions{Prefix: prefix, Recursive: true}) {
+ if obj.Err != nil {
+ return obj.Err
+ }
+ objs = append(objs, obj)
+ }
+ ch := make(chan minio.ObjectInfo, len(objs))
+ for _, obj := range objs {
+ ch <- obj
+ }
+ close(ch)
+ var err error
+ for res := range s.client.RemoveObjects(ctx, s.bucket, ch, minio.RemoveObjectsOptions{}) {
+ if res.Err != nil && err == nil {
+ err = res.Err
+ }
+ }
+ return err
+}
+
// 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) {
diff --git a/internal/storage/storage_test.go b/internal/storage/storage_test.go
index 0379900..a5987e0 100644
--- a/internal/storage/storage_test.go
+++ b/internal/storage/storage_test.go
@@ -4,11 +4,15 @@ import (
"bytes"
"context"
"io"
+ "net/http"
+ "net/http/httptest"
"os"
"strings"
"testing"
"time"
+ "github.com/minio/minio-go/v7"
+
"git.unkin.net/unkin/artifactapi/internal/testsupport"
)
@@ -158,3 +162,26 @@ func TestCASStore(t *testing.T) {
t.Errorf("stored content mismatch: %q", got)
}
}
+
+// The fake endpoint fails every list but accepts every delete, as an S3 that
+// tolerates deleting an empty key would.
+func TestDeletePrefixReturnsListError(t *testing.T) {
+ srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
+ if r.Method == http.MethodPost {
+ _, _ = io.WriteString(w, ``)
+ return
+ }
+ w.WriteHeader(http.StatusInternalServerError)
+ _, _ = io.WriteString(w, `InternalErrorlist failed`)
+ }))
+ defer srv.Close()
+ client, err := minio.New(strings.TrimPrefix(srv.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
+ if err != nil {
+ t.Fatal(err)
+ }
+ s := &S3{client: client, bucket: "bucket"}
+ err = s.DeletePrefix(context.Background(), "indexes/r/")
+ if err == nil {
+ t.Fatal("DeletePrefix on a failing list returned nil")
+ }
+}