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") + } +}