Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| bda762b10e | |||
| a71d126239 | |||
| 6a08539a78 | |||
| 0b159dad90 | |||
| 7a4f4054cc |
@@ -3,14 +3,28 @@ when:
|
|||||||
|
|
||||||
steps:
|
steps:
|
||||||
- name: pre-commit
|
- name: pre-commit
|
||||||
image: git.unkin.net/unkin/almalinux9-gobuilder:20260606
|
# gobuilder trusts the internal CA, which the S3 build cache endpoint needs.
|
||||||
|
# go-cache-plugin is baked into the image; S3 errors degrade to cache misses.
|
||||||
|
image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
|
||||||
commands:
|
commands:
|
||||||
- uvx pre-commit run --all-files
|
- uvx pre-commit run --all-files
|
||||||
environment:
|
environment:
|
||||||
# golib lives on Gitea; skip the public proxy/sum db.
|
# golib lives on Gitea; skip the public proxy/sum db.
|
||||||
GOPRIVATE: git.unkin.net
|
GOPRIVATE: git.unkin.net
|
||||||
|
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
|
||||||
|
GOCACHE_S3_BUCKET: gocache
|
||||||
|
# Explicit region skips a GetBucketLocation probe RGW handles poorly.
|
||||||
|
GOCACHE_S3_REGION: us-east-1
|
||||||
|
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
|
||||||
|
GOCACHE_S3_PATH_STYLE: "true"
|
||||||
|
GOCACHE_KEY_PREFIX: ci-artifactapi
|
||||||
|
AWS_ACCESS_KEY_ID:
|
||||||
|
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
|
||||||
|
AWS_SECRET_ACCESS_KEY:
|
||||||
|
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
|
||||||
backend_options:
|
backend_options:
|
||||||
kubernetes:
|
kubernetes:
|
||||||
|
serviceAccountName: default
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: 512Mi
|
memory: 512Mi
|
||||||
|
|||||||
+24
-1
@@ -3,9 +3,32 @@ when:
|
|||||||
|
|
||||||
steps:
|
steps:
|
||||||
- name: test
|
- name: test
|
||||||
image: golang:1.25
|
# gobuilder trusts the internal CA, which the S3 build cache endpoint needs.
|
||||||
|
# go-cache-plugin is baked into the image; S3 errors degrade to cache misses.
|
||||||
|
image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
|
||||||
commands:
|
commands:
|
||||||
- go test -race -count=1 ./pkg/... ./internal/...
|
- go test -race -count=1 ./pkg/... ./internal/...
|
||||||
environment:
|
environment:
|
||||||
# golib lives on Gitea; skip the public proxy/sum db.
|
# golib lives on Gitea; skip the public proxy/sum db.
|
||||||
GOPRIVATE: git.unkin.net
|
GOPRIVATE: git.unkin.net
|
||||||
|
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
|
||||||
|
GOCACHE_S3_BUCKET: gocache
|
||||||
|
# Explicit region skips a GetBucketLocation probe RGW handles poorly.
|
||||||
|
GOCACHE_S3_REGION: us-east-1
|
||||||
|
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
|
||||||
|
GOCACHE_S3_PATH_STYLE: "true"
|
||||||
|
GOCACHE_KEY_PREFIX: ci-artifactapi
|
||||||
|
AWS_ACCESS_KEY_ID:
|
||||||
|
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
|
||||||
|
AWS_SECRET_ACCESS_KEY:
|
||||||
|
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
|
||||||
|
backend_options:
|
||||||
|
kubernetes:
|
||||||
|
serviceAccountName: default
|
||||||
|
resources:
|
||||||
|
requests:
|
||||||
|
memory: 1Gi
|
||||||
|
cpu: 1
|
||||||
|
limits:
|
||||||
|
memory: 4Gi
|
||||||
|
cpu: 2
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
package v2
|
package v2
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
"strconv"
|
"strconv"
|
||||||
@@ -8,8 +10,14 @@ import (
|
|||||||
"github.com/go-chi/chi/v5"
|
"github.com/go-chi/chi/v5"
|
||||||
|
|
||||||
"git.unkin.net/unkin/artifactapi/internal/database"
|
"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 {
|
type ObjectsHandler struct {
|
||||||
db *database.DB
|
db *database.DB
|
||||||
}
|
}
|
||||||
@@ -18,10 +26,13 @@ func NewObjectsHandler(db *database.DB) *ObjectsHandler {
|
|||||||
return &ObjectsHandler{db: db}
|
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 := chi.NewRouter()
|
||||||
r.Get("/", h.list)
|
r.Get("/", h.list)
|
||||||
r.Delete("/*", h.evict)
|
r.Delete("/*", func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
evict(w, r, evictor)
|
||||||
|
})
|
||||||
return r
|
return r
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -51,7 +62,7 @@ func (h *ObjectsHandler) list(w http.ResponseWriter, r *http.Request) {
|
|||||||
remoteName := chi.URLParam(r, "name")
|
remoteName := chi.URLParam(r, "name")
|
||||||
limit, offset := pageBounds(r)
|
limit, offset := pageBounds(r)
|
||||||
|
|
||||||
artifacts, err := h.db.ListArtifacts(r.Context(), remoteName, limit, offset)
|
artifacts, err := h.db.ListArtifacts(r.Context(), remoteName, r.URL.Query().Get("prefix"), limit, offset)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||||
return
|
return
|
||||||
@@ -82,11 +93,16 @@ func (h *ObjectsHandler) evictLocal(w http.ResponseWriter, r *http.Request) {
|
|||||||
w.WriteHeader(http.StatusNoContent)
|
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")
|
remoteName := chi.URLParam(r, "name")
|
||||||
path := chi.URLParam(r, "*")
|
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)
|
http.Error(w, fmt.Sprintf("evict failed: %v", err), http.StatusInternalServerError)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 <dir>/*"}, 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,68 @@
|
|||||||
|
package v2
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"net/http/httptest"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/go-chi/chi/v5"
|
||||||
|
|
||||||
|
"git.unkin.net/unkin/artifactapi/internal/database"
|
||||||
|
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestObjectsListPrefix verifies the remote objects listing passes ?prefix=
|
||||||
|
// through to the database filter.
|
||||||
|
func TestObjectsListPrefix(t *testing.T) {
|
||||||
|
if testDSN == "" {
|
||||||
|
t.Skip("Docker unavailable")
|
||||||
|
}
|
||||||
|
ctx := context.Background()
|
||||||
|
db, err := database.New(testDSN)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer db.Close()
|
||||||
|
|
||||||
|
const remote = "generic-objs-prefix"
|
||||||
|
if err := db.CreateRemote(ctx, &models.Remote{
|
||||||
|
Name: remote, PackageType: models.PackageGeneric, RepoType: models.RepoTypeRemote,
|
||||||
|
BaseURL: "https://example.com", MutableTTL: 3600,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
const hash = "sha256:bb22"
|
||||||
|
if err := db.UpsertBlob(ctx, hash, "blobs/bb/22", 10, "text/plain"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, p := range []string{"a/one.txt", "b/two.txt"} {
|
||||||
|
if err := db.UpsertArtifact(ctx, remote, p, hash, ""); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
router := chi.NewRouter()
|
||||||
|
router.Mount("/remotes/{name}/objects", NewObjectsHandler(db).Routes(&fakeEvictor{}))
|
||||||
|
|
||||||
|
list := func(query string) []models.Artifact {
|
||||||
|
t.Helper()
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
router.ServeHTTP(w, httptest.NewRequest("GET", "/remotes/"+remote+"/objects"+query, nil))
|
||||||
|
if w.Code != 200 {
|
||||||
|
t.Fatalf("list%s = %d, want 200", query, w.Code)
|
||||||
|
}
|
||||||
|
var got []models.Artifact
|
||||||
|
if err := json.Unmarshal(w.Body.Bytes(), &got); err != nil {
|
||||||
|
t.Fatalf("decode: %v", err)
|
||||||
|
}
|
||||||
|
return got
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := list(""); len(got) != 2 {
|
||||||
|
t.Fatalf("unfiltered listing returned %d objects, want 2", len(got))
|
||||||
|
}
|
||||||
|
if got := list("?prefix=b/"); len(got) != 1 || got[0].Path != "b/two.txt" {
|
||||||
|
t.Fatalf("prefix=b/ listing = %+v, want only b/two.txt", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
Vendored
+69
@@ -131,3 +131,72 @@ func TestFlushRemote(t *testing.T) {
|
|||||||
t.Error("expected keys flushed")
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Vendored
+25
@@ -3,6 +3,7 @@ package cache
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/redis/go-redis/v9"
|
"github.com/redis/go-redis/v9"
|
||||||
@@ -115,3 +116,27 @@ func (r *Redis) FlushRemote(ctx context.Context, remote string) error {
|
|||||||
}
|
}
|
||||||
return iter.Err()
|
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(`\`, `\\`, `*`, `\*`, `?`, `\?`, `[`, `\[`, `]`, `\]`)
|
||||||
|
|||||||
@@ -65,7 +65,8 @@ func (db *DB) TouchArtifactAccess(ctx context.Context, remoteName, path string)
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (db *DB) ListArtifacts(ctx context.Context, remoteName string, limit, offset int) ([]models.Artifact, error) {
|
// ListArtifacts pages a remote's artifacts whose path starts with prefix ("" lists all).
|
||||||
|
func (db *DB) ListArtifacts(ctx context.Context, remoteName, prefix string, limit, offset int) ([]models.Artifact, error) {
|
||||||
rows, err := db.Pool.Query(ctx, `
|
rows, err := db.Pool.Query(ctx, `
|
||||||
SELECT a.id, a.remote_name, a.path, a.content_hash, a.upstream_etag,
|
SELECT a.id, a.remote_name, a.path, a.content_hash, a.upstream_etag,
|
||||||
a.upstream_last_modified, a.first_seen_at, a.last_fetched_at,
|
a.upstream_last_modified, a.first_seen_at, a.last_fetched_at,
|
||||||
@@ -73,10 +74,10 @@ func (db *DB) ListArtifacts(ctx context.Context, remoteName string, limit, offse
|
|||||||
b.size_bytes, b.content_type
|
b.size_bytes, b.content_type
|
||||||
FROM artifacts a
|
FROM artifacts a
|
||||||
JOIN blobs b ON a.content_hash = b.content_hash
|
JOIN blobs b ON a.content_hash = b.content_hash
|
||||||
WHERE a.remote_name = $1
|
WHERE a.remote_name = $1 AND left(a.path, length($2)) = $2
|
||||||
ORDER BY a.path
|
ORDER BY a.path
|
||||||
LIMIT $2 OFFSET $3
|
LIMIT $3 OFFSET $4
|
||||||
`, remoteName, limit, offset)
|
`, remoteName, prefix, limit, offset)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -103,6 +104,13 @@ func (db *DB) DeleteArtifact(ctx context.Context, remoteName, path string) error
|
|||||||
return err
|
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 {
|
func (db *DB) InsertAccessLog(ctx context.Context, remoteName, path string, cacheHit bool, sizeBytes int64, upstreamMS int, clientIP string) error {
|
||||||
_, err := db.Pool.Exec(ctx, `
|
_, err := db.Pool.Exec(ctx, `
|
||||||
INSERT INTO access_log (remote_name, path, cache_hit, size_bytes, upstream_ms, client_ip)
|
INSERT INTO access_log (remote_name, path, cache_hit, size_bytes, upstream_ms, client_ip)
|
||||||
|
|||||||
@@ -168,10 +168,17 @@ func TestArtifactsAndBlobs(t *testing.T) {
|
|||||||
if err := testDB.TouchArtifactAccess(ctx(), "r-art", "path/a.txt"); err != nil {
|
if err := testDB.TouchArtifactAccess(ctx(), "r-art", "path/a.txt"); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
arts, err := testDB.ListArtifacts(ctx(), "r-art", 10, 0)
|
if err := testDB.UpsertArtifact(ctx(), "r-art", "other/b.txt", hash, ""); err != nil {
|
||||||
if err != nil || len(arts) != 1 {
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
arts, err := testDB.ListArtifacts(ctx(), "r-art", "", 10, 0)
|
||||||
|
if err != nil || len(arts) != 2 {
|
||||||
t.Fatalf("list artifacts: %v %v", len(arts), err)
|
t.Fatalf("list artifacts: %v %v", len(arts), err)
|
||||||
}
|
}
|
||||||
|
arts, err = testDB.ListArtifacts(ctx(), "r-art", "path/", 10, 0)
|
||||||
|
if err != nil || len(arts) != 1 || arts[0].Path != "path/a.txt" {
|
||||||
|
t.Fatalf("list artifacts with prefix: %+v %v", arts, err)
|
||||||
|
}
|
||||||
if err := testDB.InsertAccessLog(ctx(), "r-art", "path/a.txt", true, 10, 5, "1.2.3.4"); err != nil {
|
if err := testDB.InsertAccessLog(ctx(), "r-art", "path/a.txt", true, 10, 5, "1.2.3.4"); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -188,6 +195,45 @@ func TestArtifactsAndBlobs(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestListArtifactsPrefix(t *testing.T) {
|
||||||
|
requireDB(t)
|
||||||
|
seedRemote(t, "r-prefix")
|
||||||
|
seedBlob(t, "prefixhash")
|
||||||
|
for _, p := range []string{"a%b/y", "a1b/y", "a_b/x", "aXb/x", "pkg/1", "pkg/2", "pkg/3", "pkgx/4"} {
|
||||||
|
if err := testDB.UpsertArtifact(ctx(), "r-prefix", p, "sha256:prefixhash", ""); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
paths := func(prefix string, limit, offset int) []string {
|
||||||
|
t.Helper()
|
||||||
|
arts, err := testDB.ListArtifacts(ctx(), "r-prefix", prefix, limit, offset)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list %q: %v", prefix, err)
|
||||||
|
}
|
||||||
|
out := make([]string, len(arts))
|
||||||
|
for i, a := range arts {
|
||||||
|
out[i] = a.Path
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// LIKE wildcards in the prefix must match literally.
|
||||||
|
if got := paths("a_b/", 10, 0); len(got) != 1 || got[0] != "a_b/x" {
|
||||||
|
t.Fatalf("prefix a_b/ = %v, want [a_b/x]", got)
|
||||||
|
}
|
||||||
|
if got := paths("a%", 10, 0); len(got) != 1 || got[0] != "a%b/y" {
|
||||||
|
t.Fatalf("prefix a%% = %v, want [a%%b/y]", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
// limit/offset page the filtered set, not the whole remote.
|
||||||
|
if got := paths("pkg/", 2, 0); len(got) != 2 || got[0] != "pkg/1" || got[1] != "pkg/2" {
|
||||||
|
t.Fatalf("prefix pkg/ page 1 = %v, want [pkg/1 pkg/2]", got)
|
||||||
|
}
|
||||||
|
if got := paths("pkg/", 2, 2); len(got) != 1 || got[0] != "pkg/3" {
|
||||||
|
t.Fatalf("prefix pkg/ page 2 = %v, want [pkg/3]", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestOrphanAndColdCleanup(t *testing.T) {
|
func TestOrphanAndColdCleanup(t *testing.T) {
|
||||||
requireDB(t)
|
requireDB(t)
|
||||||
seedBlob(t, "orphanhash")
|
seedBlob(t, "orphanhash")
|
||||||
@@ -320,7 +366,7 @@ func TestDatabaseErrorPaths(t *testing.T) {
|
|||||||
if _, err := bad.ListVirtuals(ctx); err == nil {
|
if _, err := bad.ListVirtuals(ctx); err == nil {
|
||||||
t.Error("ListVirtuals should error")
|
t.Error("ListVirtuals should error")
|
||||||
}
|
}
|
||||||
if _, err := bad.ListArtifacts(ctx, "r", 10, 0); err == nil {
|
if _, err := bad.ListArtifacts(ctx, "r", "", 10, 0); err == nil {
|
||||||
t.Error("ListArtifacts should error")
|
t.Error("ListArtifacts should error")
|
||||||
}
|
}
|
||||||
if _, err := bad.ListLocalFiles(ctx, "r", 10, 0); err == nil {
|
if _, err := bad.ListLocalFiles(ctx, "r", 10, 0); err == nil {
|
||||||
|
|||||||
@@ -16,6 +16,8 @@ import (
|
|||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5"
|
||||||
|
|
||||||
"git.unkin.net/unkin/artifactapi/internal/cache"
|
"git.unkin.net/unkin/artifactapi/internal/cache"
|
||||||
"git.unkin.net/unkin/artifactapi/internal/database"
|
"git.unkin.net/unkin/artifactapi/internal/database"
|
||||||
"git.unkin.net/unkin/artifactapi/internal/provider"
|
"git.unkin.net/unkin/artifactapi/internal/provider"
|
||||||
@@ -47,6 +49,8 @@ type Engine struct {
|
|||||||
// mirror strategy to prefer the mirror currently handling the fewest
|
// mirror strategy to prefer the mirror currently handling the fewest
|
||||||
// requests. Per-replica and approximate, which is fine.
|
// requests. Per-replica and approximate, which is fine.
|
||||||
inflight sync.Map
|
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 {
|
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),
|
cas: storage.NewCAS(s),
|
||||||
circuit: NewCircuitBreaker(c),
|
circuit: NewCircuitBreaker(c),
|
||||||
accessLog: make(chan database.AccessLogEntry, accessLogBufferSize),
|
accessLog: make(chan database.AccessLogEntry, accessLogBufferSize),
|
||||||
|
|
||||||
|
evictLockWait: fetchLockTTL,
|
||||||
}
|
}
|
||||||
go e.runAccessLogWriter()
|
go e.runAccessLogWriter()
|
||||||
return e
|
return e
|
||||||
@@ -208,6 +214,71 @@ func (e *Engine) Fetch(ctx context.Context, remote models.Remote, path string, p
|
|||||||
return result, nil
|
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 <dir>/*"}
|
||||||
|
}
|
||||||
|
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.
|
// HeadResult carries artifact metadata for a HEAD request. There is no body.
|
||||||
type HeadResult struct {
|
type HeadResult struct {
|
||||||
ContentType string
|
ContentType string
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -196,8 +196,8 @@ func (s *Server) routes() chi.Router {
|
|||||||
|
|
||||||
r.Route("/remotes/{name}/objects", func(r chi.Router) {
|
r.Route("/remotes/{name}/objects", func(r chi.Router) {
|
||||||
objHandler := v2.NewObjectsHandler(s.db)
|
objHandler := v2.NewObjectsHandler(s.db)
|
||||||
r.Get("/", objHandler.Routes().ServeHTTP)
|
r.Get("/", objHandler.Routes(s.engine).ServeHTTP)
|
||||||
r.Delete("/*", objHandler.Routes().ServeHTTP)
|
r.Delete("/*", objHandler.Routes(s.engine).ServeHTTP)
|
||||||
})
|
})
|
||||||
|
|
||||||
r.Route("/locals/{name}/objects", func(r chi.Router) {
|
r.Route("/locals/{name}/objects", func(r chi.Router) {
|
||||||
|
|||||||
@@ -99,6 +99,30 @@ func (s *S3) Stat(ctx context.Context, key string) (*minio.ObjectInfo, error) {
|
|||||||
return &info, nil
|
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
|
// ListStaleObjects returns keys under prefix last modified before cutoff. Used
|
||||||
// by the GC to reap abandoned staging objects (e.g. cancelled docker pushes).
|
// 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) {
|
func (s *S3) ListStaleObjects(ctx context.Context, prefix string, cutoff time.Time) ([]string, error) {
|
||||||
|
|||||||
@@ -4,11 +4,15 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"io"
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
"os"
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/minio/minio-go/v7"
|
||||||
|
|
||||||
"git.unkin.net/unkin/artifactapi/internal/testsupport"
|
"git.unkin.net/unkin/artifactapi/internal/testsupport"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -158,3 +162,26 @@ func TestCASStore(t *testing.T) {
|
|||||||
t.Errorf("stored content mismatch: %q", got)
|
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, `<DeleteResult></DeleteResult>`)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
w.WriteHeader(http.StatusInternalServerError)
|
||||||
|
_, _ = io.WriteString(w, `<Error><Code>InternalError</Code><Message>list failed</Message></Error>`)
|
||||||
|
}))
|
||||||
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user