feat(registry): api 定期删除无引用 manifest,gate 提供 manifest 索引

This commit is contained in:
Lemon-miaow committed 2026-09-24 22:58:07 +08:00
1 parent 05e8c64e47
commit 151c9d2e30
14 files changed
+1061 -8

No files matched your search

+80
View File
@@ -5,6 +5,7 @@ import (
"flag"
"fmt"
"io"
"log/slog"
"net/http"
"os"
"regexp"
@@ -26,6 +27,7 @@ import (
"felis.lolicon.best/internal/passkey"
"felis.lolicon.best/internal/platform"
"felis.lolicon.best/internal/reaper"
"felis.lolicon.best/internal/registryprune"
"felis.lolicon.best/internal/restore"
"felis.lolicon.best/internal/store"
"felis.lolicon.best/internal/submit"
@@ -399,6 +401,10 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int {
// reconciles it, but this loop converges builds nobody is polling.
go reconcileBuilds(ctx, builder, stderr)
if pruner := registryPruner(cfg, builder.Store, a.Cluster, stderr); pruner != nil {
go pruner.Loop(ctx, registryPruneInterval)
}
select {
case <-ctx.Done():
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
@@ -613,6 +619,80 @@ func reconcileBuilds(ctx context.Context, b *build.Builder, stderr io.Writer) {
}
}
// registryPruneInterval spaces the registry pruner's runs. The registry-gc
// sidecar sweeps once a day, so pruning more often only changes which sweep frees
// a layer.
const registryPruneInterval = 6 * time.Hour
// registryPruner deletes the registry manifests nothing references
// (internal/registryprune); the registry-gc sidecar frees their layers on its next
// sweep. It acts as the gate's prune principal, whose token the api Deployment
// injects from felis-registry-auth. Without the token the registry only grows,
// which is said once here.
func registryPruner(cfg *config.Config, store imageRefStore, servers serverLister, stderr io.Writer) *registryprune.Pruner {
if cfg.Registry.URL == "" {
return nil
}
token := os.Getenv(platform.RegistryPruneTokenEnv)
if token == "" {
fmt.Fprintf(stderr, "felis api: registry pruner disabled (%s unset) — images nothing uses are never deleted from the registry\n", platform.RegistryPruneTokenEnv)
return nil
}
static := []string{
os.Getenv("FELIS_IMAGE"),
cfg.Registry.KanikoImage, cfg.Registry.TrivyImage,
cfg.Registry.TrivyDBRepository, cfg.Registry.TrivyJavaDBRepository,
}
return &registryprune.Pruner{
Registry: &registryprune.Client{Endpoint: "http://" + cfg.Registry.URL, Token: token},
Host: cfg.Registry.URL,
Refs: func(ctx context.Context) ([]string, error) {
return inUseImageRefs(ctx, store, servers, static)
},
Log: slog.New(slog.NewTextHandler(stderr, nil)),
}
}
type imageRefStore interface {
ListImages(ctx context.Context) ([]build.Image, error)
ListUnfinishedBuilds(ctx context.Context) ([]build.Build, error)
}
type serverLister interface {
ListServers(ctx context.Context) ([]api.ServerInfo, error)
}
// inUseImageRefs lists every image reference the platform still depends on: the
// whitelist (disabled rows too, an admin may enable them again), every server's
// spec, builds still running, and the images the control plane and the build
// Jobs run. Any source failing fails the whole list, so the pruner never decides
// on a partial view.
func inUseImageRefs(ctx context.Context, store imageRefStore, servers serverLister, static []string) ([]string, error) {
refs := append([]string(nil), static...)
images, err := store.ListImages(ctx)
if err != nil {
return nil, fmt.Errorf("image whitelist: %w", err)
}
for _, img := range images {
refs = append(refs, img.ImageRef)
}
srvs, err := servers.ListServers(ctx)
if err != nil {
return nil, fmt.Errorf("servers: %w", err)
}
for _, s := range srvs {
refs = append(refs, s.Image)
}
builds, err := store.ListUnfinishedBuilds(ctx)
if err != nil {
return nil, fmt.Errorf("running builds: %w", err)
}
for _, b := range builds {
refs = append(refs, b.ImageRef)
}
return refs, nil
}
// mailLimit turns smtp.max_per_hour into the API's install-wide mail bucket:
// the hourly cap as the refill rate, with a quarter of it (at least 5) allowed
// at once so a burst of real sign-ins is not queued behind the average.
+49
View File
@@ -1,9 +1,14 @@
package main
import (
"context"
"errors"
"fmt"
"net/http"
"testing"
"felis.lolicon.best/internal/api"
"felis.lolicon.best/internal/build"
"felis.lolicon.best/internal/config"
)
@@ -91,3 +96,47 @@ func TestNewAPIServerSetsHardenedTimeouts(t *testing.T) {
t.Errorf("ReadTimeout = %v, want 0 (unset) so a slow SSE attach is not capped", srv.ReadTimeout)
}
}
type fakeRefStore struct {
images []build.Image
builds []build.Build
err error
}
func (f fakeRefStore) ListImages(context.Context) ([]build.Image, error) { return f.images, f.err }
func (f fakeRefStore) ListUnfinishedBuilds(context.Context) ([]build.Build, error) {
return f.builds, nil
}
type fakeServers []api.ServerInfo
func (f fakeServers) ListServers(context.Context) ([]api.ServerInfo, error) { return f, nil }
// The registry pruner deletes whatever this list does not name, so every source of
// a reference has to be in it, and a failing source must fail the list.
func TestInUseImageRefsCoversEverySource(t *testing.T) {
const reg = "registry.felis.svc:5000/"
store := fakeRefStore{
images: []build.Image{{ImageRef: reg + "modpacks/pack:*"}, {ImageRef: reg + "felis/paper:demo"}},
builds: []build.Build{{ImageRef: reg + "user-uploads/sub-9:latest"}},
}
servers := fakeServers{{Name: "s1", Image: reg + "felis/paper:demo@sha256:" + fmt.Sprintf("%064d", 1)}}
got, err := inUseImageRefs(context.Background(), store, servers, []string{reg + "felis/felis:b60"})
if err != nil {
t.Fatal(err)
}
want := []string{
reg + "felis/felis:b60",
reg + "modpacks/pack:*", reg + "felis/paper:demo",
reg + "felis/paper:demo@sha256:" + fmt.Sprintf("%064d", 1),
reg + "user-uploads/sub-9:latest",
}
if fmt.Sprint(got) != fmt.Sprint(want) {
t.Fatalf("refs = %v\nwant %v", got, want)
}
store.err = errors.New("db down")
if _, err := inUseImageRefs(context.Background(), store, servers, nil); err == nil {
t.Fatal("a failing whitelist read produced a reference list")
}
}
+2
View File
@@ -44,6 +44,7 @@ func cmdRegistryGate(args []string, _, stderr io.Writer) int {
maintListen := fs.String("maint-listen", "", "loopback address for the GC sidecar's read-only handshake (empty disables it)")
maintDir := fs.String("maint-dir", "", "directory that keeps an open read-only window across a gate restart")
quiet := fs.Duration("maint-quiet", registrygate.DefaultQuiet, "how long writes must be idle before a read-only window is granted")
dataDir := fs.String("data-dir", "", "the registry's storage root, mounted read-only, for the manifest index (empty disables it)")
if err := fs.Parse(args); err != nil {
return 2
}
@@ -70,6 +71,7 @@ func cmdRegistryGate(args []string, _, stderr io.Writer) int {
gate := registrygate.New(target, tokens, log)
gate.SetQuiet(*quiet)
gate.DataDir = *dataDir
if *maintDir != "" {
if err := gate.SetMaintenanceState(registrygate.MaintStatePath(*maintDir)); err != nil {
// A corrupt file must not keep the registry from serving pulls.
+6 -1
View File
@@ -411,10 +411,12 @@ apply_registry_secrets() {
remember_temp "$dir"
printf '%s' "$REGISTRY_PLATFORM_TOKEN" > "${dir}/platform"
printf '%s' "$REGISTRY_BUILD_TOKEN" > "${dir}/build"
printf '%s' "$REGISTRY_PRUNE_TOKEN" > "${dir}/prune"
printf '%s' build > "${dir}/username"
kube -n "$CONTROL_NS" create secret generic felis-registry-auth \
--from-file=platform="${dir}/platform" \
--from-file=build="${dir}/build" \
--from-file=prune="${dir}/prune" \
--dry-run=client -o yaml | kube apply -f -
kube -n "$BUILD_NS" create secret generic felis-registry-push \
--from-file=username="${dir}/username" \
@@ -2210,9 +2212,11 @@ load_or_make_secrets() {
# Registry write credentials, one per principal the registry gate knows
# (internal/registrygate): platform pushes the installer's own images and the
# Trivy DB mirrors, build is what a build Job's push container presents and may
# never write under felis/ or mirror/. Reads stay anonymous.
# never write under felis/ or mirror/, prune is felis-api deleting manifests
# nothing references (internal/registryprune). Reads stay anonymous.
REGISTRY_PLATFORM_TOKEN="${REGISTRY_PLATFORM_TOKEN:-$(openssl rand -hex 32)}"
REGISTRY_BUILD_TOKEN="${REGISTRY_BUILD_TOKEN:-$(openssl rand -hex 32)}"
REGISTRY_PRUNE_TOKEN="${REGISTRY_PRUNE_TOKEN:-$(openssl rand -hex 32)}"
(
umask 077
cat > "$SECRETS_ENV" <<EOF
@@ -2222,6 +2226,7 @@ SESSION_SECRET=${SESSION_SECRET}
FORWARDING_SECRET=${FORWARDING_SECRET}
REGISTRY_PLATFORM_TOKEN=${REGISTRY_PLATFORM_TOKEN}
REGISTRY_BUILD_TOKEN=${REGISTRY_BUILD_TOKEN}
REGISTRY_PRUNE_TOKEN=${REGISTRY_PRUNE_TOKEN}
EOF
)
chmod 0600 "$SECRETS_ENV"
+2 -2
View File
@@ -72,8 +72,8 @@ const (
)
// Registry write credentials (internal/registrygate). The registry namespace holds
// RegistryAuthSecretName with one key per principal (platform, build), mounted into
// the gate sidecar. The build namespace holds RegistryPushSecretName with the build
// RegistryAuthSecretName with one key per principal (platform, build, prune),
// mounted into the gate sidecar; felis-api reads the prune key into env. The build namespace holds RegistryPushSecretName with the build
// principal's username/password, read only by a build Job's push container. Both
// are provisioned out-of-band by deploy/bootstrap.sh.
const (
+7 -5
View File
@@ -236,11 +236,13 @@ func loginToInternalAPI(p Params) *networkingv1.NetworkPolicy {
// RegistryIngressPolicy fences the registry pod: only build pods and felis-api
// reach its port. Build pods push what they build; felis-api reads a manifest
// digest to pin a new server's image (internal/imagepin). Everything else that
// uses the registry runs on the node — containerd's pulls and the installer's
// pushes both arrive through the loopback hostPort — and Kubernetes never blocks
// resident-node traffic. Write authorization is the gate's job (felis-api holds no
// registry credential); this policy keeps every other pod from even trying.
// digest to pin a new server's image (internal/imagepin) and deletes manifests
// nothing references (internal/registryprune). Everything else that uses the
// registry runs on the node — containerd's pulls and the installer's pushes both
// arrive through the loopback hostPort — and Kubernetes never blocks
// resident-node traffic. Write authorization is the gate's job (felis-api holds
// only the prune principal, which can delete a manifest by digest and nothing
// else); this policy keeps every other pod from even trying.
func RegistryIngressPolicy(p Params) *networkingv1.NetworkPolicy {
p = p.withDefaults()
tcp := corev1.ProtocolTCP
+16
View File
@@ -151,6 +151,16 @@ const (
// [smtp] password_ref defaults to this name.
SMTPPasswordEnv = "FELIS_SMTP_PASSWORD"
// RegistryPruneTokenEnv carries the registry gate's prune principal token into
// felis-api, whose pruner deletes the manifests nothing references
// (internal/registryprune). It comes from the prune key of
// naming.RegistryAuthSecretName, optionally: that Secret lives in the registry
// namespace, which is the control namespace on every install the bootstrap
// makes, and an install that splits them or predates the key runs without the
// pruner.
RegistryPruneTokenEnv = "FELIS_REGISTRY_PRUNE_TOKEN"
registryPruneTokenKey = "prune"
// worldsMountPath is where the reaper CronJob mounts the worlds-root (read-only).
// It is the default of `felis reaper --worlds-root`; the resolver then reads each
// world at <worldsMountPath>/<pvc>. Single-sourced with cmd/felis/reaper.go.
@@ -339,6 +349,9 @@ func APIDeployment(p Params) *appsv1.Deployment {
corev1.EnvVar{Name: SMTPPasswordEnv, ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{
LocalObjectReference: corev1.LocalObjectReference{Name: SMTPSecretName}, Key: SMTPSecretPasswordKey, Optional: optional,
}}},
corev1.EnvVar{Name: RegistryPruneTokenEnv, ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{
LocalObjectReference: corev1.LocalObjectReference{Name: naming.RegistryAuthSecretName}, Key: registryPruneTokenKey, Optional: optional,
}}},
)
container := corev1.Container{
@@ -869,6 +882,8 @@ func registryDeployment(p Params) *appsv1.Deployment {
"--auth-dir=" + registryAuthMountPath,
fmt.Sprintf("--maint-listen=127.0.0.1:%d", registryMaintPort(p)),
"--maint-dir=" + registryMaintMountPath,
// The manifest index felis-api's pruner reads (registrygate/index.go).
"--data-dir=" + registryDataPath,
},
Ports: []corev1.ContainerPort{
{
@@ -881,6 +896,7 @@ func registryDeployment(p Params) *appsv1.Deployment {
VolumeMounts: []corev1.VolumeMount{
{Name: registryAuthVolume, MountPath: registryAuthMountPath, ReadOnly: true},
{Name: registryMaintVolume, MountPath: registryMaintMountPath},
{Name: registryVolume, MountPath: registryDataPath, ReadOnly: true},
},
// /healthz answers 200 only while registry:2 answers GET /v2/ on loopback,
// so a registry whose storage broke shows up as an unready pod instead of a
+10
View File
@@ -11,6 +11,8 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/util/intstr"
"felis.lolicon.best/internal/naming"
)
// podSpec returns the single container and the pod template of a Deployment,
@@ -289,6 +291,7 @@ func TestAPIDeployment_UploadsStorage(t *testing.T) {
{UploadsS3AccessKeyEnv, UploadsS3SecretName, UploadsS3SecretAccessKey},
{UploadsS3SecretKeyEnv, UploadsS3SecretName, UploadsS3SecretSecretKey},
{SMTPPasswordEnv, SMTPSecretName, SMTPSecretPasswordKey},
{RegistryPruneTokenEnv, naming.RegistryAuthSecretName, "prune"},
} {
e := envVar(c.Env, ev.name)
if e == nil || e.ValueFrom == nil || e.ValueFrom.SecretKeyRef == nil {
@@ -507,6 +510,7 @@ func TestRegistry_DeploymentServicePVC(t *testing.T) {
// The GC handshake has no authentication: loopback only.
fmt.Sprintf("--maint-listen=127.0.0.1:%d", p.RegistryPort+2),
"--maint-dir=" + registryMaintMountPath,
"--data-dir=" + registryDataPath,
} {
if !contains(gate.Args, want) {
t.Errorf("gate args = %v, want %s", gate.Args, want)
@@ -575,6 +579,12 @@ func TestRegistry_DeploymentServicePVC(t *testing.T) {
t.Error("registry:2 must not mount the write tokens")
}
}
// The gate reads the manifest index off the data volume, and never writes it.
for _, m := range gate.VolumeMounts {
if m.Name == registryVolume && !m.ReadOnly {
t.Error("the gate must mount the registry data read-only")
}
}
if !mounted {
t.Errorf("gate must mount the tokens read-only, mounts=%v", gate.VolumeMounts)
}
+7
View File
@@ -69,6 +69,9 @@ type Gate struct {
Upstream *url.URL
// Log receives one line per refused write. Nil discards.
Log *slog.Logger
// DataDir is the registry's storage root, mounted read-only, for the manifest
// index (index.go). Empty turns the index off.
DataDir string
proxy *httputil.ReverseProxy
health *http.Client
@@ -105,6 +108,10 @@ func (g *Gate) ServeHTTP(w http.ResponseWriter, r *http.Request) {
g.serveHealth(w, r)
return
}
if strings.HasPrefix(r.URL.Path, IndexPathPrefix) {
g.serveIndex(w, r)
return
}
if !strings.HasPrefix(r.URL.Path, "/v2/") && r.URL.Path != "/v2" {
writeError(w, http.StatusNotFound, "NAME_UNKNOWN", "not a registry API path")
return
+124
View File
@@ -0,0 +1,124 @@
package registrygate
import (
"encoding/json"
"errors"
"io/fs"
"net/http"
"os"
"path/filepath"
"regexp"
"sort"
"strings"
"time"
)
// IndexPathPrefix serves the manifest index of one repository:
//
// GET /felis/manifests/<repo> {"revisions":[{"digest":…,"pushed":…}],"tags":{"<tag>":"<digest>"}}
//
// The registry API can list a repository's tags but not its manifests, so a
// manifest a tag moved off (every rebuild of :demo or :latest leaves one) is
// invisible to it, while it still holds every layer it names. felis-api's pruner
// needs the full set to decide which to delete; the gate reads it off the data
// volume it mounts read-only, from the filesystem driver's layout that registry
// 2.x and 3.x share. "pushed" is the modification time of the revision link, which
// the registry rewrites on every push of that manifest.
//
// Anonymous like every read: a digest list says no more than tags/list and
// _catalog already do.
const IndexPathPrefix = "/felis/manifests/"
// Index is the manifest index of one repository.
type Index struct {
Revisions []Revision `json:"revisions"`
Tags map[string]string `json:"tags"`
}
// Revision is one manifest stored in a repository.
type Revision struct {
Digest string `json:"digest"`
Pushed time.Time `json:"pushed"`
}
// repoNameRE is the distribution reference grammar for a repository path. Every
// component starts and ends alphanumeric, so no match can hold "." or "..".
var repoNameRE = regexp.MustCompile(`^[a-z0-9]+(?:(?:[._]|__|[-]*)[a-z0-9]+)*(?:/[a-z0-9]+(?:(?:[._]|__|[-]*)[a-z0-9]+)*)*$`)
var hexDigestRE = regexp.MustCompile(`^[0-9a-f]{64}$`)
func (g *Gate) serveIndex(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeError(w, http.StatusMethodNotAllowed, "UNSUPPORTED", "method not allowed")
return
}
if g.DataDir == "" {
writeError(w, http.StatusNotFound, "UNSUPPORTED", "this gate does not mount the registry data")
return
}
repo := strings.TrimPrefix(r.URL.Path, IndexPathPrefix)
if r.URL.RawPath != "" || !repoNameRE.MatchString(repo) {
writeError(w, http.StatusBadRequest, "NAME_INVALID", "invalid repository name")
return
}
idx, err := ReadIndex(g.DataDir, repo)
if err != nil {
writeError(w, http.StatusInternalServerError, "UNKNOWN", err.Error())
return
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(idx)
}
// ReadIndex reads repo's manifests and tags from a registry filesystem root. A
// repository that does not exist has an empty index.
func ReadIndex(root, repo string) (*Index, error) {
if !repoNameRE.MatchString(repo) {
return nil, errors.New("invalid repository name")
}
base := filepath.Join(root, "docker", "registry", "v2", "repositories", filepath.FromSlash(repo), "_manifests")
idx := &Index{Revisions: []Revision{}, Tags: map[string]string{}}
revDir := filepath.Join(base, "revisions", "sha256")
entries, err := os.ReadDir(revDir)
if err != nil && !errors.Is(err, fs.ErrNotExist) {
return nil, err
}
for _, e := range entries {
if !e.IsDir() || !hexDigestRE.MatchString(e.Name()) {
continue
}
// A deleted manifest keeps its directory and loses its link.
st, err := os.Stat(filepath.Join(revDir, e.Name(), "link"))
if errors.Is(err, fs.ErrNotExist) {
continue
}
if err != nil {
return nil, err
}
idx.Revisions = append(idx.Revisions, Revision{Digest: "sha256:" + e.Name(), Pushed: st.ModTime().UTC()})
}
sort.Slice(idx.Revisions, func(i, j int) bool { return idx.Revisions[i].Digest < idx.Revisions[j].Digest })
tagDir := filepath.Join(base, "tags")
tags, err := os.ReadDir(tagDir)
if err != nil && !errors.Is(err, fs.ErrNotExist) {
return nil, err
}
for _, e := range tags {
if !e.IsDir() {
continue
}
b, err := os.ReadFile(filepath.Join(tagDir, e.Name(), "current", "link"))
if errors.Is(err, fs.ErrNotExist) {
continue
}
if err != nil {
return nil, err
}
if d := strings.TrimSpace(string(b)); digestRE.MatchString(d) {
idx.Tags[e.Name()] = d
}
}
return idx, nil
}
+70
View File
@@ -1,6 +1,7 @@
package registrygate
import (
"encoding/json"
"net/http"
"net/http/httptest"
"net/url"
@@ -206,3 +207,72 @@ func TestReadOnlyWindowSurvivesAGateRestart(t *testing.T) {
t.Fatalf("a corrupt state file = %v, want an error naming it", err)
}
}
func TestManifestIndexListsUntaggedRevisions(t *testing.T) {
root := t.TempDir()
base := root + "/docker/registry/v2/repositories/felis/paper/_manifests"
hexA := strings.Repeat("a", 64)
hexB := strings.Repeat("b", 64)
hexGone := strings.Repeat("c", 64)
for _, h := range []string{hexA, hexB} {
if err := os.MkdirAll(base+"/revisions/sha256/"+h, 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(base+"/revisions/sha256/"+h+"/link", []byte("sha256:"+h), 0o644); err != nil {
t.Fatal(err)
}
}
// A deleted manifest leaves its directory without a link.
if err := os.MkdirAll(base+"/revisions/sha256/"+hexGone, 0o755); err != nil {
t.Fatal(err)
}
if err := os.MkdirAll(base+"/tags/demo/current", 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(base+"/tags/demo/current/link", []byte("sha256:"+hexB), 0o644); err != nil {
t.Fatal(err)
}
old := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
if err := os.Chtimes(base+"/revisions/sha256/"+hexA+"/link", old, old); err != nil {
t.Fatal(err)
}
up := httptest.NewServer(http.NotFoundHandler())
t.Cleanup(up.Close)
target, _ := url.Parse(up.URL)
g := New(target, nil, nil)
g.DataDir = root
srv := httptest.NewServer(g)
t.Cleanup(srv.Close)
resp, err := http.Get(srv.URL + IndexPathPrefix + "felis/paper")
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
var idx Index
if err := json.NewDecoder(resp.Body).Decode(&idx); err != nil {
t.Fatal(err)
}
if len(idx.Revisions) != 2 || idx.Revisions[0].Digest != "sha256:"+hexA || !idx.Revisions[0].Pushed.Equal(old) {
t.Fatalf("revisions = %+v, want a (pushed %s) and b, without the deleted c", idx.Revisions, old)
}
if idx.Tags["demo"] != "sha256:"+hexB || len(idx.Tags) != 1 {
t.Fatalf("tags = %v, want demo -> b", idx.Tags)
}
// A repository that does not exist is empty; a traversal is refused.
if idx, err := ReadIndex(root, "user-uploads/none"); err != nil || len(idx.Revisions) != 0 || len(idx.Tags) != 0 {
t.Fatalf("missing repo = %+v, %v, want empty", idx, err)
}
for _, bad := range []string{"felis/../felis/paper", "felis/./paper", "Felis/paper", "felis//paper", "felis/paper/"} {
resp, err := http.Get(srv.URL + IndexPathPrefix + bad)
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != http.StatusBadRequest {
t.Errorf("index of %q = %d, want 400", bad, resp.StatusCode)
}
}
}
+115
View File
@@ -0,0 +1,115 @@
package registryprune
import (
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"strings"
"time"
"felis.lolicon.best/internal/registrygate"
)
// Client talks to the registry through its gate: the catalog and the manifest
// index anonymously, deletes as the prune principal.
type Client struct {
// Endpoint is the gate's base URL, e.g. http://registry.felis.svc:5000.
Endpoint string
// Token is the prune principal's secret.
Token string
// HTTP makes the requests; nil uses a client with a 30s timeout.
HTTP *http.Client
}
func (c *Client) http() *http.Client {
if c.HTTP != nil {
return c.HTTP
}
return &http.Client{Timeout: 30 * time.Second}
}
// Repositories pages through /v2/_catalog.
func (c *Client) Repositories(ctx context.Context) ([]string, error) {
var out []string
next := c.Endpoint + "/v2/_catalog?n=1000"
for next != "" {
var page struct {
Repositories []string `json:"repositories"`
}
resp, err := c.get(ctx, next, &page)
if err != nil {
return nil, err
}
out = append(out, page.Repositories...)
next = ""
// Link: </v2/_catalog?last=x&n=1000>; rel="next"
if link := resp.Header.Get("Link"); strings.Contains(link, `rel="next"`) {
start, end := strings.Index(link, "<"), strings.Index(link, ">")
if start < 0 || end <= start {
return nil, fmt.Errorf("malformed catalog Link header %q", link)
}
u, err := url.Parse(c.Endpoint)
if err != nil {
return nil, err
}
ref, err := url.Parse(link[start+1 : end])
if err != nil {
return nil, err
}
next = u.ResolveReference(ref).String()
}
}
return out, nil
}
// Index reads the gate's manifest index of repo.
func (c *Client) Index(ctx context.Context, repo string) (*registrygate.Index, error) {
var idx registrygate.Index
if _, err := c.get(ctx, c.Endpoint+registrygate.IndexPathPrefix+repo, &idx); err != nil {
return nil, err
}
return &idx, nil
}
// Delete deletes one manifest. A manifest already gone counts as deleted.
func (c *Client) Delete(ctx context.Context, repo, digest string) error {
req, err := http.NewRequestWithContext(ctx, http.MethodDelete, c.Endpoint+"/v2/"+repo+"/manifests/"+digest, nil)
if err != nil {
return err
}
req.SetBasicAuth(registrygate.PrincipalPrune, c.Token)
resp, err := c.http().Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
switch resp.StatusCode {
case http.StatusAccepted, http.StatusOK, http.StatusNotFound:
return nil
}
b, _ := io.ReadAll(io.LimitReader(resp.Body, 1024))
return fmt.Errorf("registry answered %s: %s", resp.Status, strings.TrimSpace(string(b)))
}
func (c *Client) get(ctx context.Context, u string, into any) (*http.Response, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
return nil, err
}
resp, err := c.http().Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
b, _ := io.ReadAll(io.LimitReader(resp.Body, 1024))
return nil, fmt.Errorf("GET %s: registry answered %s: %s", u, resp.Status, strings.TrimSpace(string(b)))
}
if err := json.NewDecoder(io.LimitReader(resp.Body, 64<<20)).Decode(into); err != nil {
return nil, fmt.Errorf("GET %s: %w", u, err)
}
return resp, nil
}
+274
View File
@@ -0,0 +1,274 @@
// Package registryprune deletes the manifests in the platform registry that
// nothing uses any more, so the registry-gc sidecar's next sweep can free their
// layers. Without it the registry only grows: every rebuild of a tag leaves the
// previous manifest behind, a whitelist entry an admin removes keeps its image,
// and each installer run adds a platform image.
//
// A manifest is kept when any of these hold:
//
// - an image reference the platform still depends on names it: a whitelist
// entry (by tag, by wildcard tag, or by digest), a server's spec (servers are
// pinned by digest, so a sleeping server keeps the exact build it was created
// with even after its tag moved on), or the images the control plane and the
// build Jobs run;
// - it was pushed less than Grace ago, which covers a build that pushed but has
// not been admitted to the whitelist yet;
// - it lives under a platform-reserved root (felis/, mirror/), is tagged, and is
// among the KeepTagged most recently pushed tagged manifests of its
// repository. The installer owns those tags; the newest few stay for a
// rollback, the older ones go.
//
// Everything else is deleted through the registry gate as the prune principal,
// which may do nothing but delete a manifest by digest.
package registryprune
import (
"context"
"errors"
"fmt"
"log/slog"
"sort"
"strings"
"time"
"felis.lolicon.best/internal/registrygate"
)
// Defaults for Pruner's zero fields.
const (
DefaultGrace = 24 * time.Hour
DefaultKeepTagged = 5
)
// Registry is the pruner's view of the registry (Client in production).
type Registry interface {
Repositories(ctx context.Context) ([]string, error)
Index(ctx context.Context, repo string) (*registrygate.Index, error)
Delete(ctx context.Context, repo, digest string) error
}
// Pruner decides and deletes.
type Pruner struct {
Registry Registry
// Host is the registry host[:port] image refs spell (registry.felis.svc:5000).
// Refs under any other host are not this registry's and are ignored.
Host string
// Refs lists every image reference still in use. An error aborts the run: a
// partial view could delete an image something depends on.
Refs func(ctx context.Context) ([]string, error)
// Grace keeps manifests pushed this recently. Zero means DefaultGrace.
Grace time.Duration
// KeepTagged is how many tagged manifests per reserved repository survive by
// recency alone. Zero means DefaultKeepTagged.
KeepTagged int
// Now is the clock; nil means time.Now.
Now func() time.Time
// Log receives one line per deletion and a summary. Nil discards.
Log *slog.Logger
}
// Target is one manifest the pruner deletes.
type Target struct {
Repo, Digest string
}
func (t Target) String() string { return t.Repo + "@" + t.Digest }
// Report summarizes one run.
type Report struct {
Repositories int
Manifests int
Deleted []Target
// Failed counts deletions the registry refused; they are retried next run.
Failed int
}
// Run lists the registry, decides, and deletes.
func (p *Pruner) Run(ctx context.Context) (Report, error) {
var rep Report
if p.Registry == nil || p.Refs == nil || p.Host == "" {
return rep, errors.New("registryprune: Registry, Refs and Host are required")
}
refs, err := p.Refs(ctx)
if err != nil {
return rep, fmt.Errorf("registryprune: list image references: %w", err)
}
repos, err := p.Registry.Repositories(ctx)
if err != nil {
return rep, fmt.Errorf("registryprune: list repositories: %w", err)
}
indexes := make(map[string]*registrygate.Index, len(repos))
for _, repo := range repos {
idx, err := p.Registry.Index(ctx, repo)
if err != nil {
return rep, fmt.Errorf("registryprune: index %s: %w", repo, err)
}
indexes[repo] = idx
rep.Manifests += len(idx.Revisions)
}
rep.Repositories = len(repos)
now := time.Now
if p.Now != nil {
now = p.Now
}
for _, t := range Plan(indexes, refs, p.Host, now(), p.grace(), p.keepTagged()) {
if err := p.Registry.Delete(ctx, t.Repo, t.Digest); err != nil {
rep.Failed++
p.log().Warn("registry prune: delete failed", "manifest", t.String(), "err", err)
continue
}
rep.Deleted = append(rep.Deleted, t)
p.log().Info("registry prune: deleted manifest", "manifest", t.String())
}
return rep, nil
}
// Loop runs Run every interval until ctx ends, starting after one interval's
// tenth so a restart loop does not hammer the registry.
func (p *Pruner) Loop(ctx context.Context, every time.Duration) {
t := time.NewTimer(every / 10)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
}
rep, err := p.Run(ctx)
if err != nil {
p.log().Error("registry prune failed", "err", err)
} else {
p.log().Info("registry prune finished", "repositories", rep.Repositories,
"manifests", rep.Manifests, "deleted", len(rep.Deleted), "failed", rep.Failed)
}
t.Reset(every)
}
}
func (p *Pruner) grace() time.Duration {
if p.Grace > 0 {
return p.Grace
}
return DefaultGrace
}
func (p *Pruner) keepTagged() int {
if p.KeepTagged > 0 {
return p.KeepTagged
}
return DefaultKeepTagged
}
func (p *Pruner) log() *slog.Logger {
if p.Log != nil {
return p.Log
}
return slog.New(slog.DiscardHandler)
}
// Plan decides which manifests to delete. It is pure so every keep rule is
// tested without a registry.
func Plan(indexes map[string]*registrygate.Index, refs []string, host string, now time.Time, grace time.Duration, keepTagged int) []Target {
keep := map[Target]bool{}
for _, ref := range refs {
repo, tag, digest, ok := parseRef(ref, host)
if !ok {
continue
}
idx := indexes[repo]
switch {
case digest != "":
keep[Target{repo, digest}] = true
case idx == nil:
case tag == "*":
for _, d := range idx.Tags {
keep[Target{repo, d}] = true
}
default:
if d, ok := idx.Tags[tag]; ok {
keep[Target{repo, d}] = true
}
}
}
repos := make([]string, 0, len(indexes))
for repo := range indexes {
repos = append(repos, repo)
}
sort.Strings(repos)
var out []Target
for _, repo := range repos {
idx := indexes[repo]
pushed := make(map[string]time.Time, len(idx.Revisions))
for _, r := range idx.Revisions {
pushed[r.Digest] = r.Pushed
}
if reserved(repo) {
for _, d := range newestTagged(idx, pushed, keepTagged) {
keep[Target{repo, d}] = true
}
}
for _, r := range idx.Revisions {
t := Target{repo, r.Digest}
if keep[t] || now.Sub(r.Pushed) < grace {
continue
}
out = append(out, t)
}
}
return out
}
// newestTagged returns the n most recently pushed distinct digests any tag in idx
// names.
func newestTagged(idx *registrygate.Index, pushed map[string]time.Time, n int) []string {
seen := map[string]bool{}
var tagged []string
for _, d := range idx.Tags {
if !seen[d] {
seen[d] = true
tagged = append(tagged, d)
}
}
sort.Slice(tagged, func(i, j int) bool {
pi, pj := pushed[tagged[i]], pushed[tagged[j]]
if !pi.Equal(pj) {
return pi.After(pj)
}
return tagged[i] < tagged[j]
})
if len(tagged) > n {
tagged = tagged[:n]
}
return tagged
}
func reserved(repo string) bool {
root, _, _ := strings.Cut(repo, "/")
for _, r := range registrygate.ReservedRepoRoots {
if root == r {
return true
}
}
return false
}
// parseRef splits host/repo[:tag][@digest] for refs under host. A ref with
// neither tag nor digest means :latest, as it does for every image client.
func parseRef(ref, host string) (repo, tag, digest string, ok bool) {
rest, ok := strings.CutPrefix(strings.TrimSpace(ref), host+"/")
if !ok || rest == "" {
return "", "", "", false
}
name, digest, _ := strings.Cut(rest, "@")
repo, tag = name, ""
if colon := strings.LastIndex(name, ":"); colon > strings.LastIndex(name, "/") {
repo, tag = name[:colon], name[colon+1:]
}
if tag == "" && digest == "" {
tag = "latest"
}
return repo, tag, digest, repo != ""
}
+299
View File
@@ -0,0 +1,299 @@
package registryprune
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"net/url"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"testing"
"time"
"felis.lolicon.best/internal/registrygate"
)
const host = "registry.felis.svc:5000"
var now = time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC)
func dg(c string) string { return "sha256:" + strings.Repeat(c, 64) }
func rev(c string, age time.Duration) registrygate.Revision {
return registrygate.Revision{Digest: dg(c), Pushed: now.Add(-age)}
}
const day = 24 * time.Hour
func targets(ts []Target) []string {
out := make([]string, 0, len(ts))
for _, t := range ts {
out = append(out, t.Repo+"@"+t.Digest[7:8])
}
sort.Strings(out)
return out
}
func TestPlanKeepsWhatIsReferencedOrRecent(t *testing.T) {
indexes := map[string]*registrygate.Index{
// Rebuilt twice: a is untagged and pinned by a server, b untagged and
// unused, c is :latest and whitelisted.
"user-uploads/sub-1": {
Revisions: []registrygate.Revision{rev("a", 30*day), rev("b", 20*day), rev("c", 10*day)},
Tags: map[string]string{"latest": dg("c")},
},
// Whitelist entry removed, no server: goes.
"user-uploads/sub-2": {
Revisions: []registrygate.Revision{rev("d", 5*day)},
Tags: map[string]string{"latest": dg("d")},
},
// Pushed an hour ago, not admitted yet: the grace keeps it.
"user-uploads/sub-3": {
Revisions: []registrygate.Revision{rev("e", time.Hour)},
Tags: map[string]string{"latest": dg("e")},
},
// A wildcard whitelist entry keeps every tag, not an untagged leftover.
"modpacks/pack": {
Revisions: []registrygate.Revision{rev("f", 9*day), rev("1", 8*day), rev("2", 7*day)},
Tags: map[string]string{"1.0": dg("f"), "1.1": dg("1")},
},
}
refs := []string{
host + "/user-uploads/sub-1:latest",
host + "/user-uploads/sub-1:latest@" + dg("a"),
host + "/modpacks/pack:*",
"docker.io/library/nginx:latest",
"",
}
got := targets(Plan(indexes, refs, host, now, day, 5))
want := []string{"modpacks/pack@2", "user-uploads/sub-1@b", "user-uploads/sub-2@d"}
if fmt.Sprint(got) != fmt.Sprint(want) {
t.Fatalf("plan = %v, want %v", got, want)
}
}
func TestPlanKeepsTheNewestTaggedPlatformImages(t *testing.T) {
indexes := map[string]*registrygate.Index{
// Seven installer runs: b1..b7 tagged, plus untagged leftovers x (pinned by
// a server) and y.
"felis/felis": {
Revisions: []registrygate.Revision{
rev("1", 70*day), rev("2", 60*day), rev("3", 50*day), rev("4", 40*day),
rev("5", 30*day), rev("6", 20*day), rev("7", 10*day),
rev("x", 80*day), rev("y", 90*day),
},
Tags: map[string]string{
"b1": dg("1"), "b2": dg("2"), "b3": dg("3"), "b4": dg("4"),
"b5": dg("5"), "b6": dg("6"), "b7": dg("7"),
},
},
// The scanner DB mirror: one tag, kept however old.
"mirror/trivy-db": {
Revisions: []registrygate.Revision{rev("8", 200*day), rev("9", 300*day)},
Tags: map[string]string{"2": dg("8")},
},
}
refs := []string{
// The running control plane is older than the newest three.
host + "/felis/felis:b2",
host + "/felis/felis:b5@" + dg("x"),
}
got := targets(Plan(indexes, refs, host, now, day, 3))
want := []string{"felis/felis@1", "felis/felis@3", "felis/felis@4", "felis/felis@y", "mirror/trivy-db@9"}
if fmt.Sprint(got) != fmt.Sprint(want) {
t.Fatalf("plan = %v, want %v", got, want)
}
}
func TestParseRef(t *testing.T) {
for _, c := range []struct {
ref, repo, tag, digest string
ok bool
}{
{host + "/felis/paper:demo", "felis/paper", "demo", "", true},
{host + "/felis/paper:demo@" + dg("a"), "felis/paper", "demo", dg("a"), true},
{host + "/felis/paper@" + dg("a"), "felis/paper", "", dg("a"), true},
{host + "/felis/paper", "felis/paper", "latest", "", true},
{host + "/pack:*", "pack", "*", "", true},
{"other:5000/felis/paper:demo", "", "", "", false},
{host + "/", "", "", "", false},
} {
repo, tag, digest, ok := parseRef(c.ref, host)
if repo != c.repo || tag != c.tag || digest != c.digest || ok != c.ok {
t.Errorf("parseRef(%q) = %q %q %q %v, want %q %q %q %v", c.ref, repo, tag, digest, ok, c.repo, c.tag, c.digest, c.ok)
}
}
}
// fakeRegistry serves _catalog (paged by one), the manifest index and DELETE.
type fakeRegistry struct {
mu sync.Mutex
indexes map[string]*registrygate.Index
deletes []string
auth []string
}
func (f *fakeRegistry) ServeHTTP(w http.ResponseWriter, r *http.Request) {
f.mu.Lock()
defer f.mu.Unlock()
switch {
case r.URL.Path == "/v2/_catalog":
var repos []string
for repo := range f.indexes {
repos = append(repos, repo)
}
sort.Strings(repos)
last := r.URL.Query().Get("last")
var page []string
for _, repo := range repos {
if repo > last {
page = append(page, repo)
break
}
}
if len(page) == 1 && page[0] != repos[len(repos)-1] {
w.Header().Set("Link", `</v2/_catalog?last=`+url.QueryEscape(page[0])+`&n=1>; rel="next"`)
}
_ = json.NewEncoder(w).Encode(map[string][]string{"repositories": page})
case strings.HasPrefix(r.URL.Path, registrygate.IndexPathPrefix):
idx := f.indexes[strings.TrimPrefix(r.URL.Path, registrygate.IndexPathPrefix)]
_ = json.NewEncoder(w).Encode(idx)
case r.Method == http.MethodDelete:
user, pass, _ := r.BasicAuth()
f.auth = append(f.auth, user+":"+pass)
f.deletes = append(f.deletes, r.URL.Path)
w.WriteHeader(http.StatusAccepted)
default:
http.NotFound(w, r)
}
}
func TestRunDeletesThroughTheGateAsThePrunePrincipal(t *testing.T) {
reg := &fakeRegistry{indexes: map[string]*registrygate.Index{
"e2e/probe": {Revisions: []registrygate.Revision{rev("a", 3*day)}, Tags: map[string]string{"latest": dg("a")}},
"user-uploads/sub-1": {Revisions: []registrygate.Revision{rev("b", 3*day)}, Tags: map[string]string{"latest": dg("b")}},
"felis/paper": {Revisions: []registrygate.Revision{rev("c", 3*day)}, Tags: map[string]string{"demo": dg("c")}},
}}
srv := httptest.NewServer(reg)
t.Cleanup(srv.Close)
p := &Pruner{
Registry: &Client{Endpoint: srv.URL, Token: "prune-secret"},
Host: host,
Refs: func(context.Context) ([]string, error) {
return []string{host + "/user-uploads/sub-1:latest"}, nil
},
Now: func() time.Time { return now },
}
rep, err := p.Run(context.Background())
if err != nil {
t.Fatal(err)
}
if rep.Repositories != 3 || rep.Manifests != 3 || len(rep.Deleted) != 1 || rep.Failed != 0 {
t.Fatalf("report = %+v, want 3 repos, 3 manifests, 1 deleted", rep)
}
if want := "/v2/e2e/probe/manifests/" + dg("a"); len(reg.deletes) != 1 || reg.deletes[0] != want {
t.Fatalf("deletes = %v, want [%s]", reg.deletes, want)
}
if reg.auth[0] != registrygate.PrincipalPrune+":prune-secret" {
t.Fatalf("delete sent as %q, want the prune principal", reg.auth[0])
}
}
func TestRunAbortsOnAPartialView(t *testing.T) {
reg := &fakeRegistry{indexes: map[string]*registrygate.Index{
"e2e/probe": {Revisions: []registrygate.Revision{rev("a", 3*day)}, Tags: map[string]string{"latest": dg("a")}},
}}
srv := httptest.NewServer(reg)
t.Cleanup(srv.Close)
p := &Pruner{
Registry: &Client{Endpoint: srv.URL, Token: "t"},
Host: host,
Refs: func(context.Context) ([]string, error) {
return nil, errors.New("cluster unreachable")
},
Now: func() time.Time { return now },
}
if _, err := p.Run(context.Background()); err == nil {
t.Fatal("run with no view of the servers succeeded")
}
if len(reg.deletes) != 0 {
t.Fatalf("deleted %v without knowing what is in use", reg.deletes)
}
}
// Through the real gate: the index comes off the data directory, the delete is
// authorized as the prune principal and reaches the registry without the
// credential.
func TestRunThroughTheGate(t *testing.T) {
root := t.TempDir()
repoDir := filepath.Join(root, "docker/registry/v2/repositories/e2e/probe/_manifests")
hex := strings.Repeat("a", 64)
for _, dir := range []string{"revisions/sha256/" + hex, "tags/latest/current"} {
if err := os.MkdirAll(filepath.Join(repoDir, dir), 0o755); err != nil {
t.Fatal(err)
}
}
for _, f := range []string{"revisions/sha256/" + hex + "/link", "tags/latest/current/link"} {
if err := os.WriteFile(filepath.Join(repoDir, f), []byte("sha256:"+hex), 0o644); err != nil {
t.Fatal(err)
}
}
var mu sync.Mutex
var upstream []string
up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
mu.Lock()
upstream = append(upstream, r.Method+" "+r.URL.Path+" auth="+r.Header.Get("Authorization"))
mu.Unlock()
switch {
case r.URL.Path == "/v2/_catalog":
_, _ = w.Write([]byte(`{"repositories":["e2e/probe"]}`))
case r.Method == http.MethodDelete:
w.WriteHeader(http.StatusAccepted)
default:
http.NotFound(w, r)
}
}))
t.Cleanup(up.Close)
target, _ := url.Parse(up.URL)
g := registrygate.New(target, map[string]string{
registrygate.PrincipalPrune: "prune-secret",
registrygate.PrincipalBuild: "build-secret",
}, nil)
g.DataDir = root
gate := httptest.NewServer(g)
t.Cleanup(gate.Close)
refs := func(context.Context) ([]string, error) { return nil, nil }
later := func() time.Time { return time.Now().Add(48 * time.Hour) }
// The build principal's token is not good for a delete.
wrong := &Pruner{Registry: &Client{Endpoint: gate.URL, Token: "build-secret"}, Host: host, Refs: refs, Now: later}
if rep, err := wrong.Run(context.Background()); err != nil || rep.Failed != 1 || len(rep.Deleted) != 0 {
t.Fatalf("run with the wrong token = %+v, %v; want one refused delete", rep, err)
}
p := &Pruner{Registry: &Client{Endpoint: gate.URL, Token: "prune-secret"}, Host: host, Refs: refs, Now: later}
rep, err := p.Run(context.Background())
if err != nil || len(rep.Deleted) != 1 || rep.Failed != 0 {
t.Fatalf("run = %+v, %v; want one deletion", rep, err)
}
mu.Lock()
defer mu.Unlock()
want := "DELETE /v2/e2e/probe/manifests/sha256:" + hex + " auth="
var deletes []string
for _, u := range upstream {
if strings.HasPrefix(u, "DELETE") {
deletes = append(deletes, u)
}
}
if len(deletes) != 1 || deletes[0] != want {
t.Fatalf("upstream deletes = %q, want [%q]", deletes, want)
}
}