diff --git a/cmd/felis/api.go b/cmd/felis/api.go index ce54cf9..bd8acde 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -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 ®istryprune.Pruner{ + Registry: ®istryprune.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. diff --git a/cmd/felis/api_test.go b/cmd/felis/api_test.go index ca6e752..f18c0b4 100644 --- a/cmd/felis/api_test.go +++ b/cmd/felis/api_test.go @@ -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") + } +} diff --git a/cmd/felis/registrygate.go b/cmd/felis/registrygate.go index 0404821..d8089f4 100644 --- a/cmd/felis/registrygate.go +++ b/cmd/felis/registrygate.go @@ -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. diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 9d732e7..6653583 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -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" </. 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 diff --git a/internal/platform/workloads_test.go b/internal/platform/workloads_test.go index 23cb719..4308b47 100644 --- a/internal/platform/workloads_test.go +++ b/internal/platform/workloads_test.go @@ -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) } diff --git a/internal/registrygate/gate.go b/internal/registrygate/gate.go index 72c8447..d0ea51b 100644 --- a/internal/registrygate/gate.go +++ b/internal/registrygate/gate.go @@ -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 diff --git a/internal/registrygate/index.go b/internal/registrygate/index.go new file mode 100644 index 0000000..d53de79 --- /dev/null +++ b/internal/registrygate/index.go @@ -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/ {"revisions":[{"digest":…,"pushed":…}],"tags":{"":""}} +// +// 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 +} diff --git a/internal/registrygate/maint_test.go b/internal/registrygate/maint_test.go index fe66845..83cfa9b 100644 --- a/internal/registrygate/maint_test.go +++ b/internal/registrygate/maint_test.go @@ -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) + } + } +} diff --git a/internal/registryprune/client.go b/internal/registryprune/client.go new file mode 100644 index 0000000..fbf735c --- /dev/null +++ b/internal/registryprune/client.go @@ -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: ; 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 +} diff --git a/internal/registryprune/prune.go b/internal/registryprune/prune.go new file mode 100644 index 0000000..02cc73a --- /dev/null +++ b/internal/registryprune/prune.go @@ -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 != "" +} diff --git a/internal/registryprune/prune_test.go b/internal/registryprune/prune_test.go new file mode 100644 index 0000000..69409cf --- /dev/null +++ b/internal/registryprune/prune_test.go @@ -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", `; 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) + } +}