From 7f772bccbbcc31130bdd94988b60aff6b051d958 Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Thu, 24 Sep 2026 22:46:37 +0800 Subject: [PATCH] =?UTF-8?q?feat(registry):=20=E5=BC=80=E5=90=AF=20manifest?= =?UTF-8?q?=20=E5=88=A0=E9=99=A4=EF=BC=8CGC=20sidecar=20=E5=9C=A8=20gate?= =?UTF-8?q?=20=E5=8F=AA=E8=AF=BB=E7=AA=97=E5=8F=A3=E5=86=85=E5=9B=9E?= =?UTF-8?q?=E6=94=B6=20blob?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/felis/registrygate.go | 57 +++++++- deploy/bootstrap.sh | 19 ++- internal/imagepush/push.go | 61 ++++++-- internal/imagepush/push_test.go | 78 +++++++++++ internal/platform/workloads.go | 114 ++++++++++++++- internal/platform/workloads_test.go | 49 ++++++- internal/registrygate/gate.go | 46 +++++- internal/registrygate/maint.go | 196 ++++++++++++++++++++++++++ internal/registrygate/maint_test.go | 208 ++++++++++++++++++++++++++++ 9 files changed, 806 insertions(+), 22 deletions(-) create mode 100644 internal/registrygate/maint.go create mode 100644 internal/registrygate/maint_test.go diff --git a/cmd/felis/registrygate.go b/cmd/felis/registrygate.go index 6a2a1e2..0404821 100644 --- a/cmd/felis/registrygate.go +++ b/cmd/felis/registrygate.go @@ -7,6 +7,7 @@ import ( "fmt" "io" "log/slog" + "net" "net/http" "net/url" "os" @@ -26,19 +27,30 @@ import ( // for an authenticated principal allowed to write that repository. See // internal/registrygate for the policy. // -// Tokens are files under --auth-dir, one per principal (platform, build), mounted -// from the registry-auth Secret. A missing file disables that principal: writes -// fail closed while every pull keeps working, which is the right way round for a -// registry the running workloads depend on. +// Tokens are files under --auth-dir, one per principal (platform, build, prune), +// mounted from the registry-auth Secret. A missing file disables that principal: +// writes fail closed while every pull keeps working, which is the right way round +// for a registry the running workloads depend on. +// +// --maint-listen is the GC sidecar's read-only handshake (registrygate.MaintHandler). +// It has no authentication, so it must name a loopback address; --maint-dir keeps +// an open window across a gate restart. func cmdRegistryGate(args []string, _, stderr io.Writer) int { fs := flag.NewFlagSet("registry-gate", flag.ContinueOnError) fs.SetOutput(stderr) listen := fs.String("listen", ":5000", "address the gate serves the registry API on") upstream := fs.String("upstream", "http://127.0.0.1:5001", "the loopback registry the gate forwards to") authDir := fs.String("auth-dir", "/etc/felis-registry-auth", "directory holding one token file per principal") + 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") if err := fs.Parse(args); err != nil { return 2 } + if *maintListen != "" && !loopbackAddr(*maintListen) { + fmt.Fprintf(stderr, "felis registry-gate: --maint-listen %q must be a loopback address: the handshake has no authentication\n", *maintListen) + return 2 + } target, err := url.Parse(*upstream) if err != nil || target.Scheme == "" || target.Host == "" { fmt.Fprintf(stderr, "felis registry-gate: bad --upstream %q\n", *upstream) @@ -46,7 +58,7 @@ func cmdRegistryGate(args []string, _, stderr io.Writer) int { } log := slog.New(slog.NewTextHandler(stderr, nil)) tokens := map[string]string{} - for _, p := range []string{registrygate.PrincipalPlatform, registrygate.PrincipalBuild} { + for _, p := range registrygate.Principals { b, err := os.ReadFile(filepath.Join(*authDir, p)) tok := strings.TrimSpace(string(b)) if err != nil || tok == "" { @@ -56,11 +68,28 @@ func cmdRegistryGate(args []string, _, stderr io.Writer) int { tokens[p] = tok } + gate := registrygate.New(target, tokens, log) + gate.SetQuiet(*quiet) + if *maintDir != "" { + if err := gate.SetMaintenanceState(registrygate.MaintStatePath(*maintDir)); err != nil { + // A corrupt file must not keep the registry from serving pulls. + log.Warn("ignoring the saved read-only window", "err", err) + } + } srv := &http.Server{ Addr: *listen, - Handler: registrygate.New(target, tokens, log), + Handler: gate, ReadHeaderTimeout: 10 * time.Second, } + var maint *http.Server + if *maintListen != "" { + maint = &http.Server{Addr: *maintListen, Handler: gate.MaintHandler(), ReadHeaderTimeout: 10 * time.Second} + go func() { + if err := maint.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { + log.Error("maintenance listener stopped; garbage collection cannot get a read-only window", "err", err) + } + }() + } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() go func() { @@ -68,6 +97,9 @@ func cmdRegistryGate(args []string, _, stderr io.Writer) int { shutdown, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() _ = srv.Shutdown(shutdown) + if maint != nil { + _ = maint.Shutdown(shutdown) + } }() log.Info("registry gate listening", "addr", *listen, "upstream", target.String(), "principals", len(tokens)) if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { @@ -77,6 +109,19 @@ func cmdRegistryGate(args []string, _, stderr io.Writer) int { return 0 } +// loopbackAddr reports whether a host:port listen address binds loopback only. +func loopbackAddr(addr string) bool { + host, _, err := net.SplitHostPort(addr) + if err != nil { + return false + } + if host == "localhost" { + return true + } + ip := net.ParseIP(host) + return ip != nil && ip.IsLoopback() +} + // cmdPushImage is the build Job's publish step. It runs after Kaniko built the // image into a tarball (--no-push) and Trivy passed that tarball, and it is the // only container of the build pod that holds the registry credential — the one diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 0cb47c7..9d732e7 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -2869,10 +2869,27 @@ push_image_to_registry() { esac log "mirroring ${ref} into the internal registry" docker tag "$ref" "$push_ref" || die "could not tag ${ref} as ${push_ref} — is docker healthy?" - docker --config "$REGISTRY_DOCKER_CONFIG" push "$push_ref" || die "could not mirror ${ref} into the internal registry — check the registry Deployment/pod (both the registry and registry-gate containers) and its PVC" + local attempt=1 + until docker --config "$REGISTRY_DOCKER_CONFIG" push "$push_ref"; do + # The gate answers writes 503 while the registry-gc sidecar sweeps; wait + # that out, and fail at once on anything else. + if [ "$attempt" -ge 40 ] || ! registry_read_only; then + die "could not mirror ${ref} into the internal registry — check the registry Deployment/pod (the registry, registry-gate and registry-gc containers) and its PVC" + fi + warn "the registry is read-only for garbage collection; retrying the push of ${push_ref} in 30s (${attempt}/40)" + attempt=$((attempt + 1)) + sleep 30 + done docker rmi "$push_ref" >/dev/null 2>&1 || true } +# registry_read_only asks the gate, over its pod-loopback maintenance listener, +# whether a garbage-collection window is open. +registry_read_only() { + kubectl -n "$CONTROL_NS" exec deploy/registry -c registry-gc -- \ + wget -q -O /dev/null "http://127.0.0.1:$((${REGISTRY_URL##*:} + 2))/readonly" >/dev/null 2>&1 +} + # registry_docker_login logs a throwaway docker config into the registry gate as # the platform principal: writes are refused anonymously, and this identity is # the only one allowed under felis/. The config lives in a 0700 temp dir that the diff --git a/internal/imagepush/push.go b/internal/imagepush/push.go index da16dec..49fdf5d 100644 --- a/internal/imagepush/push.go +++ b/internal/imagepush/push.go @@ -25,6 +25,7 @@ import ( "net/http" "net/url" "os" + "strconv" "strings" "time" ) @@ -57,6 +58,13 @@ type Pusher struct { Log io.Writer // Attempts bounds retries of one blob upload or the manifest PUT. Zero means 3. Attempts int + // MaxWait bounds the total time a push waits out 503 answers that carry + // Retry-After — the gate's read-only window while garbage collection runs. + // Those waits do not use up Attempts. Zero means 20 minutes. + MaxWait time.Duration + + // after is time.After, replaced in tests. + after func(time.Duration) <-chan time.Time } // Ref is a parsed host/repository:tag reference. @@ -296,33 +304,62 @@ func (p *Pusher) do(ctx context.Context, method, target string, body io.Reader, } // retry runs fn up to Attempts times, backing off between tries. A refusal the -// registry will repeat (401/403/4xx other than 408/429) is returned at once. +// registry will repeat (401/403/4xx other than 408/429) is returned at once. A 503 +// with Retry-After is waited out without using up an attempt, for as long as +// MaxWait allows. func (p *Pusher) retry(ctx context.Context, fn func() error) error { attempts := p.Attempts if attempts <= 0 { attempts = 3 } - var err error - for i := 0; i < attempts; i++ { + maxWait := p.MaxWait + if maxWait <= 0 { + maxWait = 20 * time.Minute + } + var ( + err error + waited time.Duration + ) + for i := 0; i < attempts; { if err = fn(); err == nil { return nil } var se *StatusError + if errors.As(err, &se) && se.Code == http.StatusServiceUnavailable && se.RetryAfter > 0 && waited+se.RetryAfter <= maxWait { + p.logf("registry unavailable, waiting %s: %s", se.RetryAfter, se.Body) + if err := p.sleep(ctx, se.RetryAfter); err != nil { + return err + } + waited += se.RetryAfter + continue + } if errors.As(err, &se) && se.Code >= 400 && se.Code < 500 && se.Code != http.StatusRequestTimeout && se.Code != http.StatusTooManyRequests { return err } - if i+1 < attempts { + i++ + if i < attempts { p.logf("retrying after: %v", err) - select { - case <-ctx.Done(): - return ctx.Err() - case <-time.After(time.Duration(i+1) * time.Second): + if err := p.sleep(ctx, time.Duration(i)*time.Second); err != nil { + return err } } } return err } +func (p *Pusher) sleep(ctx context.Context, d time.Duration) error { + after := p.after + if after == nil { + after = time.After + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-after(d): + return nil + } +} + func (p *Pusher) logf(format string, args ...any) { if p.Log != nil { fmt.Fprintf(p.Log, format+"\n", args...) @@ -334,6 +371,8 @@ type StatusError struct { Op string Code int Body string + // RetryAfter is the registry's Retry-After, when it sent one in seconds. + RetryAfter time.Duration } func (e *StatusError) Error() string { @@ -342,7 +381,11 @@ func (e *StatusError) Error() string { func statusError(op string, resp *http.Response) error { b, _ := io.ReadAll(io.LimitReader(resp.Body, 2048)) - return &StatusError{Op: op, Code: resp.StatusCode, Body: strings.TrimSpace(string(b))} + se := &StatusError{Op: op, Code: resp.StatusCode, Body: strings.TrimSpace(string(b))} + if sec, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil && sec > 0 { + se.RetryAfter = time.Duration(sec) * time.Second + } + return se } // openEntry returns a reader positioned at the named tar entry. The caller closes diff --git a/internal/imagepush/push_test.go b/internal/imagepush/push_test.go index 69ab804..0cd1a96 100644 --- a/internal/imagepush/push_test.go +++ b/internal/imagepush/push_test.go @@ -19,6 +19,7 @@ import ( "strings" "sync" "testing" + "time" "felis.lolicon.best/internal/registrygate" ) @@ -191,6 +192,83 @@ func TestPushThroughTheGate(t *testing.T) { } } +// TestPushWaitsOutTheGCWindow: while the gate holds the registry read-only for +// garbage collection, a push waits on Retry-After (without spending its attempts) +// and completes once the window closes; past MaxWait it gives up with the 503. +func TestPushWaitsOutTheGCWindow(t *testing.T) { + reg := newFakeRegistry() + upstream := httptest.NewServer(reg) + t.Cleanup(upstream.Close) + u, _ := url.Parse(upstream.URL) + g := registrygate.New(u, map[string]string{registrygate.PrincipalBuild: "build-secret"}, nil) + g.SetQuiet(0) + gate := httptest.NewServer(g) + t.Cleanup(gate.Close) + maint := httptest.NewServer(g.MaintHandler()) + t.Cleanup(maint.Close) + host := strings.TrimPrefix(gate.URL, "http://") + setReadOnly := func(on bool) { + t.Helper() + path := "/readwrite" + if on { + path = "/readonly?lease=600" + } + resp, err := http.Post(maint.URL+path, "text/plain", nil) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("POST %s = %d", path, resp.StatusCode) + } + } + + ref := host + "/user-uploads/sub-gc:latest" + tarPath := writeTarball(t, ref, []byte("layer")) + setReadOnly(true) + var waits []time.Duration + p := &Pusher{Scheme: "http", Username: registrygate.PrincipalBuild, Password: "build-secret", Attempts: 1} + p.after = func(d time.Duration) <-chan time.Time { + waits = append(waits, d) + if len(waits) == 3 { + setReadOnly(false) // the sweep finished + } + ch := make(chan time.Time, 1) + ch <- time.Time{} + return ch + } + if _, err := p.Push(context.Background(), tarPath, ref); err != nil { + t.Fatalf("push across the window: %v", err) + } + if len(waits) != 3 || waits[0] != 30*time.Second { + t.Fatalf("waits = %v, want three Retry-After (30s) waits", waits) + } + if reg.manifests["user-uploads/sub-gc:latest"] == nil { + t.Fatal("no manifest after the window closed") + } + + // A window that outlasts MaxWait fails the push with the 503. + setReadOnly(true) + ref2 := host + "/user-uploads/sub-gc2:latest" + tar2 := writeTarball(t, ref2, []byte("other layer")) + waits = nil + p.MaxWait = time.Minute + p.after = func(d time.Duration) <-chan time.Time { + waits = append(waits, d) + ch := make(chan time.Time, 1) + ch <- time.Time{} + return ch + } + _, err := p.Push(context.Background(), tar2, ref2) + var se *StatusError + if !errors.As(err, &se) || se.Code != http.StatusServiceUnavailable { + t.Fatalf("push past MaxWait = %v, want the 503", err) + } + if len(waits) != 2 { + t.Fatalf("waits = %v, want 2 (60s of MaxWait at 30s each)", waits) + } +} + func TestPushIntoAReservedRepoIsRefusedWithoutRetry(t *testing.T) { reg, host := startStack(t) ref := host + "/felis/felis:v0.1.0" diff --git a/internal/platform/workloads.go b/internal/platform/workloads.go index e486bf5..9753785 100644 --- a/internal/platform/workloads.go +++ b/internal/platform/workloads.go @@ -4,6 +4,7 @@ import ( "crypto/sha256" "encoding/hex" "fmt" + "time" "felis.lolicon.best/internal/naming" appsv1 "k8s.io/api/apps/v1" @@ -91,6 +92,14 @@ const ( registryGateName = "registry-gate" registryAuthVolume = "registry-auth" registryAuthMountPath = "/etc/felis-registry-auth" + // registryGCName is the garbage-collection sidecar, and registryMaint* the + // emptyDir where the gate keeps an open read-only window across its own restart + // (registrygate.SetMaintenanceState). registryGCInterval is how often a sweep + // runs; the pruner in felis-api deletes manifests between sweeps. + registryGCName = "registry-gc" + registryMaintVolume = "maint" + registryMaintMountPath = "/run/felis-maint" + registryGCInterval = 24 * time.Hour configVolume = "config" tmpVolume = "tmp" @@ -804,13 +813,15 @@ func controlPlaneDeployment(p Params, sa string, container corev1.Container, vol // target real. The registry never calls the K8s API, so its token auto-mount is // disabled (matching the weak build/restore SA hygiene). // -// The pod has two containers. registry:2 itself has no auth and listens on the +// The pod has three containers. registry:2 itself has no auth and listens on the // pod's loopback only (registryUpstreamPort), so nothing outside the pod can reach // it directly. The gate sidecar (felis registry-gate, internal/registrygate) owns // the registry port: reads pass anonymously, writes need the platform or build // credential from the registry-auth Secret, and the build credential cannot touch // the platform's own repositories. Before the gate any pod that could reach the -// registry could overwrite felis/felis. +// registry could overwrite felis/felis. The registry-gc sidecar reclaims the +// blobs of deleted manifests inside a read-only window the gate grants +// (registryGCScript). // // The gate's port also carries a loopback hostPort (registryLoopbackHost): it is // the node-side pull path. The node's containerd cannot dial the Service VIP, so @@ -830,6 +841,15 @@ func registryDeployment(p Params) *appsv1.Deployment { Image: p.RegistryImage, Env: []corev1.EnvVar{ {Name: "REGISTRY_HTTP_ADDR", Value: fmt.Sprintf("127.0.0.1:%d", upstreamPort)}, + // Manifest DELETE is how felis-api's pruner releases an image; the gate + // lets only the prune and platform principals send it. + {Name: "REGISTRY_STORAGE_DELETE_ENABLED", Value: "true"}, + // The in-memory blob descriptor cache outlives a garbage-collect run: a + // blob the sweep deleted would still answer HEAD, a push would skip + // uploading it, and the manifest pushed after it would name a blob that + // is gone. Any value other than inmemory/redis turns the cache off + // (registry 2.8 logs "unknown cache type ... caching disabled"). + {Name: "REGISTRY_STORAGE_CACHE_BLOBDESCRIPTOR", Value: "none"}, }, VolumeMounts: []corev1.VolumeMount{ {Name: registryVolume, MountPath: registryDataPath}, @@ -847,6 +867,8 @@ func registryDeployment(p Params) *appsv1.Deployment { fmt.Sprintf("--listen=:%d", p.RegistryPort), fmt.Sprintf("--upstream=http://127.0.0.1:%d", upstreamPort), "--auth-dir=" + registryAuthMountPath, + fmt.Sprintf("--maint-listen=127.0.0.1:%d", registryMaintPort(p)), + "--maint-dir=" + registryMaintMountPath, }, Ports: []corev1.ContainerPort{ { @@ -858,6 +880,7 @@ func registryDeployment(p Params) *appsv1.Deployment { }, VolumeMounts: []corev1.VolumeMount{ {Name: registryAuthVolume, MountPath: registryAuthMountPath, ReadOnly: true}, + {Name: registryMaintVolume, MountPath: registryMaintMountPath}, }, // /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 @@ -898,7 +921,7 @@ func registryDeployment(p Params) *appsv1.Deployment { AutomountServiceAccountToken: boolPtr(false), PriorityClassName: controlPlanePriorityName, SecurityContext: hardenedPodSecurityContext(), - Containers: []corev1.Container{registry, gate}, + Containers: []corev1.Container{registry, gate, registryGCContainer(p)}, Volumes: []corev1.Volume{ { Name: registryVolume, @@ -907,6 +930,7 @@ func registryDeployment(p Params) *appsv1.Deployment { }, }, {Name: tmpVolume, VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{}}}, + {Name: registryMaintVolume, VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{}}}, { Name: registryAuthVolume, VolumeSource: corev1.VolumeSource{Secret: &corev1.SecretVolumeSource{ @@ -929,6 +953,90 @@ func registryDeployment(p Params) *appsv1.Deployment { // the next port after the public one. func registryUpstreamPort(p Params) int32 { return p.RegistryPort + 1 } +// registryMaintPort is the gate's loopback-only maintenance listener, the one +// after the upstream port. +func registryMaintPort(p Params) int32 { return p.RegistryPort + 2 } + +// registryGCScript is the registry-gc sidecar's loop. Once per interval (the last +// run is stamped on the data volume, so a pod restart does not reset the clock) it +// asks the gate for a read-only window, waiting up to 30 minutes for pushes to go +// quiet, runs registry garbage-collect, and hands the window back. The window is +// a lease, so a sidecar killed mid-sweep leaves the registry writable again within +// the hour. +// +// No --delete-untagged: a server's spec pins its image by digest, and the tag it +// was created from moves with every rebuild, so an untagged manifest may be the +// exact build a sleeping server boots. Manifests go only when felis-api's pruner +// has found no whitelist entry, server or recent build naming them and deleted +// them; this sweep then frees the blobs nothing references any more. +const registryGCScript = `set -u +# sh as PID 1 ignores SIGTERM unless it traps it, and runs the trap only once the +# foreground child exits: sleep in the background and wait, so a pod delete does +# not sit out the grace period. +trap 'exit 0' TERM +nap() { sleep "$1" & wait $!; } +maint="http://127.0.0.1:${FELIS_GC_MAINT_PORT}" +stamp=/var/lib/registry/.felis-last-gc +while :; do + now=$(date +%s) + last=$(cat "$stamp" 2>/dev/null || echo 0) + case "$last" in ''|*[!0-9]*) last=0 ;; esac + if [ $((now - last)) -ge "$FELIS_GC_INTERVAL_SECONDS" ]; then + tries=0 + until wget -q -O /dev/null --post-data '' "$maint/readonly?lease=3600"; do + tries=$((tries + 1)) + [ "$tries" -ge 180 ] && break + nap 10 + done + if [ "$tries" -lt 180 ]; then + echo "felis-gc: registry is read-only; collecting" + if registry garbage-collect /etc/docker/registry/config.yml > /tmp/gc.log 2>&1; then + date +%s > "$stamp" + echo "felis-gc: done ($(grep -c 'blob eligible for deletion' /tmp/gc.log) blob(s) deleted)" + else + echo "felis-gc: garbage-collect failed:" >&2 + tail -n 20 /tmp/gc.log >&2 + fi + wget -q -O /dev/null --post-data '' "$maint/readwrite" \ + || echo "felis-gc: could not hand the read-only window back; it lapses with its lease" >&2 + else + echo "felis-gc: pushes never went quiet for 30 minutes; trying again later" >&2 + fi + fi + nap 600 +done +` + +// registryGCContainer renders the garbage-collection sidecar. It runs the +// registry image (garbage-collect is a subcommand of the registry binary) against +// the same data volume, under the same non-root identity and read-only root. +func registryGCContainer(p Params) corev1.Container { + return corev1.Container{ + Name: registryGCName, + Image: p.RegistryImage, + Command: []string{"/bin/sh", "-c", registryGCScript}, + Env: []corev1.EnvVar{ + {Name: "FELIS_GC_MAINT_PORT", Value: fmt.Sprint(registryMaintPort(p))}, + {Name: "FELIS_GC_INTERVAL_SECONDS", Value: fmt.Sprint(int64(registryGCInterval / time.Second))}, + }, + VolumeMounts: []corev1.VolumeMount{ + {Name: registryVolume, MountPath: registryDataPath}, + {Name: tmpVolume, MountPath: "/tmp"}, + }, + Resources: corev1.ResourceRequirements{ + Requests: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("10m"), + corev1.ResourceMemory: resource.MustParse("16Mi"), + }, + Limits: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("500m"), + corev1.ResourceMemory: resource.MustParse("512Mi"), + }, + }, + SecurityContext: hardenedContainerSecurityContext(), + } +} + // registryGateResources sizes the gate sidecar: a streaming reverse proxy that // holds no layer in memory. func registryGateResources() corev1.ResourceRequirements { diff --git a/internal/platform/workloads_test.go b/internal/platform/workloads_test.go index 7e521c8..23cb719 100644 --- a/internal/platform/workloads_test.go +++ b/internal/platform/workloads_test.go @@ -2,6 +2,7 @@ package platform import ( "fmt" + "strings" "testing" appsv1 "k8s.io/api/apps/v1" @@ -462,8 +463,8 @@ func TestRegistry_DeploymentServicePVC(t *testing.T) { pvc := registryPVC(p) ps, c := namedContainer(t, dep, registryName) - if len(ps.Containers) != 2 { - t.Fatalf("registry pod containers = %d, want registry + gate", len(ps.Containers)) + if len(ps.Containers) != 3 { + t.Fatalf("registry pod containers = %d, want registry + gate + gc", len(ps.Containers)) } if c.Image != defaultRegistryImage { t.Errorf("registry image = %q, want default %q", c.Image, defaultRegistryImage) @@ -476,6 +477,14 @@ func TestRegistry_DeploymentServicePVC(t *testing.T) { if len(c.Ports) != 0 { t.Errorf("registry container ports = %+v, want none (the gate owns the port)", c.Ports) } + // Deletion on, and the blob descriptor cache off: a cached descriptor for a + // blob the GC removed would let the next push skip uploading it. + if v := envValue(c.Env, "REGISTRY_STORAGE_DELETE_ENABLED"); v != "true" { + t.Errorf("REGISTRY_STORAGE_DELETE_ENABLED = %q, want true", v) + } + if v := envValue(c.Env, "REGISTRY_STORAGE_CACHE_BLOBDESCRIPTOR"); v == "inmemory" || v == "redis" || v == "" { + t.Errorf("REGISTRY_STORAGE_CACHE_BLOBDESCRIPTOR = %q, want the cache disabled", v) + } // The registry's limits are deliberately NOT the control-plane template's: audit // #46 caught the registry OOM-killed mid-upload at 256Mi on a real 475MB-layer push. if mem := c.Resources.Limits[corev1.ResourceMemory]; mem.Value() < 2*1024*1024*1024 { @@ -495,11 +504,47 @@ func TestRegistry_DeploymentServicePVC(t *testing.T) { fmt.Sprintf("--listen=:%d", p.RegistryPort), fmt.Sprintf("--upstream=http://127.0.0.1:%d", p.RegistryPort+1), "--auth-dir=" + registryAuthMountPath, + // The GC handshake has no authentication: loopback only. + fmt.Sprintf("--maint-listen=127.0.0.1:%d", p.RegistryPort+2), + "--maint-dir=" + registryMaintMountPath, } { if !contains(gate.Args, want) { t.Errorf("gate args = %v, want %s", gate.Args, want) } } + + // The GC sidecar: the registry image on the same data volume, hardened like + // the rest, pointed at the gate's maintenance port, never --delete-untagged + // (digest-pinned servers may boot an untagged manifest). + _, gc := namedContainer(t, dep, registryGCName) + if gc.Image != p.RegistryImage { + t.Errorf("gc image = %q, want the registry image %q", gc.Image, p.RegistryImage) + } + if v := envValue(gc.Env, "FELIS_GC_MAINT_PORT"); v != fmt.Sprint(p.RegistryPort+2) { + t.Errorf("FELIS_GC_MAINT_PORT = %q, want %d", v, p.RegistryPort+2) + } + script := strings.Join(gc.Command, " ") + if !strings.Contains(script, "garbage-collect") || strings.Contains(script, "delete-untagged") { + t.Errorf("gc command must run garbage-collect without --delete-untagged: %s", script) + } + if !strings.Contains(script, "/readonly?lease=") || !strings.Contains(script, "/readwrite") { + t.Errorf("gc command must take and hand back the gate's read-only window: %s", script) + } + gcData := false + for _, m := range gc.VolumeMounts { + if m.Name == registryAuthVolume { + t.Error("the gc sidecar must not mount the write tokens") + } + if m.Name == registryVolume && m.MountPath == registryDataPath { + gcData = true + } + } + if !gcData { + t.Errorf("gc sidecar must mount the registry data at %s, mounts=%v", registryDataPath, gc.VolumeMounts) + } + if gc.SecurityContext == nil || gc.SecurityContext.ReadOnlyRootFilesystem == nil || !*gc.SecurityContext.ReadOnlyRootFilesystem { + t.Error("gc sidecar must run with a read-only root filesystem") + } // The node-side pull path: exactly one container port, mirrored by a LOOPBACK // hostPort. Node containerd cannot dial the Service VIP, so its registries.yaml // mirror rewrites the Service name onto 127.0.0.1:; nothing else may be diff --git a/internal/registrygate/gate.go b/internal/registrygate/gate.go index f5066af..72c8447 100644 --- a/internal/registrygate/gate.go +++ b/internal/registrygate/gate.go @@ -11,7 +11,13 @@ // and none of them should carry a credential; // - the "platform" principal (the installer) may write anything; // - the "build" principal (the push step of a build Job) may write any repository -// outside the platform-reserved ones (felis/…, mirror/…), and may not delete. +// outside the platform-reserved ones (felis/…, mirror/…), and may not delete; +// - the "prune" principal (felis-api's registry pruner) may only delete a +// manifest by digest, the one write that frees space. +// +// The gate also holds the registry read-only while the GC sidecar runs registry +// garbage-collect (maint.go): a blob pushed during the sweep could be deleted +// under a manifest that is about to reference it. // // Before this gate any pod that could reach the registry — a game server running a // tenant's plugin, or a Dockerfile RUN step inside Kaniko — could overwrite @@ -27,17 +33,24 @@ import ( "net/http" "net/http/httputil" "net/url" + "regexp" "strings" "time" ) +var digestRE = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`) + // Principal names. They are the basic-auth usernames and the file names under the // gate's auth directory (cmd/felis registry-gate --auth-dir). const ( PrincipalPlatform = "platform" PrincipalBuild = "build" + PrincipalPrune = "prune" ) +// Principals lists every principal the gate knows, for loading their tokens. +var Principals = []string{PrincipalPlatform, PrincipalBuild, PrincipalPrune} + // ReservedRepoRoots are the first path components the build principal may never // write: felis/ holds the control-plane and game images the platform runs, mirror/ // holds the Trivy DB mirrors the scan gate trusts. A build that could overwrite @@ -59,6 +72,7 @@ type Gate struct { proxy *httputil.ReverseProxy health *http.Client + maint maintenance } // New builds a Gate for upstream. @@ -76,6 +90,8 @@ func New(upstream *url.URL, tokens map[string]string, log *slog.Logger) *Gate { rp.FlushInterval = -1 g.proxy = rp g.health = &http.Client{Timeout: 3 * time.Second} + g.maint.now = time.Now + g.maint.quiet = DefaultQuiet return g } @@ -135,6 +151,15 @@ func (g *Gate) ServeHTTP(w http.ResponseWriter, r *http.Request) { writeError(w, http.StatusForbidden, "DENIED", reason) return } + if !g.maint.beginWrite() { + // Retry-After is what imagepush and docker push back off on; a GC sweep + // over a few GiB takes well under a minute. + w.Header().Set("Retry-After", "30") + writeError(w, http.StatusServiceUnavailable, "UNAVAILABLE", + "the registry is read-only while garbage collection runs; retry shortly") + return + } + defer g.maint.endWrite() g.proxy.ServeHTTP(w, r) } @@ -181,11 +206,30 @@ func Authorize(principal, method, path string) string { } } return "" + case PrincipalPrune: + // Deleting a manifest only unlinks it; the blobs go at the next GC. The + // pruner never needs anything else, so a leaked prune token can neither + // plant an image nor delete a blob a live manifest still references. + if method != http.MethodDelete || !isManifestDigestPath(path) { + return "the prune principal may only delete a manifest by digest" + } + return "" default: return "unknown principal" } } +// isManifestDigestPath reports whether path is /v2//manifests/sha256:. +func isManifestDigestPath(path string) bool { + rest, ok := strings.CutPrefix(path, "/v2/") + if !ok { + return false + } + seg := strings.Split(rest, "/") + n := len(seg) + return n >= 3 && seg[n-2] == "manifests" && digestRE.MatchString(seg[n-1]) +} + // RepoFromPath extracts the repository name from a registry API v2 path, or "" // when the path addresses no repository (/v2/, /v2/_catalog). The shapes are the // distribution API's: /manifests/, /blobs/, diff --git a/internal/registrygate/maint.go b/internal/registrygate/maint.go new file mode 100644 index 0000000..a14db28 --- /dev/null +++ b/internal/registrygate/maint.go @@ -0,0 +1,196 @@ +package registrygate + +import ( + "fmt" + "net/http" + "os" + "path/filepath" + "strconv" + "strings" + "sync" + "time" +) + +// registry garbage-collect is only safe on a registry nobody writes to: it marks +// the blobs every manifest references, then deletes the rest, and a layer pushed +// between the two phases is deleted under the manifest that arrives next. The GC +// sidecar therefore asks the gate for a read-only window over a loopback-only +// maintenance listener (MaintHandler), and the gate grants it only once writes +// have been quiet for a while, so a push that is between two of its requests is +// not cut in half. +// +// While the window is open every write answers 503 with Retry-After; reads keep +// working, so running servers and kubelet re-pulls never notice. The window is a +// lease: a GC sidecar that dies mid-sweep cannot leave the registry read-only for +// longer than the lease it asked for. + +// DefaultQuiet is how long writes must have been idle before a read-only window is +// granted. A push issues its requests back to back; two minutes of silence means +// no push is mid-way. +const DefaultQuiet = 2 * time.Minute + +// maxLease bounds a read-only window. A sweep over a few GiB takes seconds; an +// hour covers a large registry on a slow disk. +const maxLease = time.Hour + +type maintenance struct { + mu sync.Mutex + inflight int + lastWrite time.Time + until time.Time + quiet time.Duration + stateFile string + now func() time.Time +} + +// beginWrite admits a write unless a read-only window is open. +func (m *maintenance) beginWrite() bool { + m.mu.Lock() + defer m.mu.Unlock() + now := m.now() + if now.Before(m.until) { + return false + } + m.inflight++ + m.lastWrite = now + return true +} + +func (m *maintenance) endWrite() { + m.mu.Lock() + defer m.mu.Unlock() + m.inflight-- + m.lastWrite = m.now() +} + +// acquire opens (or extends) a read-only window for lease. It refuses while a +// write is in flight or the last one finished less than quiet ago. +func (m *maintenance) acquire(lease time.Duration) error { + m.mu.Lock() + defer m.mu.Unlock() + now := m.now() + if !now.Before(m.until) { + if m.inflight > 0 { + return fmt.Errorf("%d write(s) in flight", m.inflight) + } + if idle := now.Sub(m.lastWrite); idle < m.quiet { + return fmt.Errorf("last write %s ago, waiting for %s of quiet", idle.Round(time.Second), m.quiet) + } + } + m.until = now.Add(lease) + m.persist() + return nil +} + +func (m *maintenance) release() { + m.mu.Lock() + defer m.mu.Unlock() + m.until = time.Time{} + m.persist() +} + +func (m *maintenance) readOnly() (bool, time.Time) { + m.mu.Lock() + defer m.mu.Unlock() + return m.now().Before(m.until), m.until +} + +// persist records the window's end in stateFile, so a gate container restarted +// mid-sweep comes back read-only instead of admitting writes into a running GC. +// Called with mu held. +func (m *maintenance) persist() { + if m.stateFile == "" { + return + } + if m.until.IsZero() { + _ = os.Remove(m.stateFile) + return + } + tmp := m.stateFile + ".tmp" + if err := os.WriteFile(tmp, []byte(strconv.FormatInt(m.until.Unix(), 10)), 0o600); err == nil { + _ = os.Rename(tmp, m.stateFile) + } +} + +// SetMaintenanceState makes the gate keep its read-only window in path (on a +// volume that outlives the container) and resumes a window a previous run of the +// gate left open. +func (g *Gate) SetMaintenanceState(path string) error { + g.maint.mu.Lock() + defer g.maint.mu.Unlock() + g.maint.stateFile = path + b, err := os.ReadFile(path) + if os.IsNotExist(err) { + return nil + } + if err != nil { + return err + } + sec, err := strconv.ParseInt(strings.TrimSpace(string(b)), 10, 64) + if err != nil { + return fmt.Errorf("maintenance state %s: %w", path, err) + } + if until := time.Unix(sec, 0); g.maint.now().Before(until) { + g.maint.until = until + } + return nil +} + +// SetQuiet overrides DefaultQuiet (tests and drills shorten it). +func (g *Gate) SetQuiet(d time.Duration) { + g.maint.mu.Lock() + g.maint.quiet = d + g.maint.mu.Unlock() +} + +// MaintHandler serves the GC sidecar's side of the handshake. It carries no +// authentication, so it must only ever listen on the pod's loopback: +// +// POST /readonly?lease= 200 once the window is open, 409 while writes are not quiet +// POST /readwrite 200, the window is closed +// GET /readonly 200 while read-only, 409 otherwise +// +// POST-only verbs keep the sidecar's busybox wget (which cannot send DELETE) able +// to drive it. +func (g *Gate) MaintHandler() http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("POST /readonly", func(w http.ResponseWriter, r *http.Request) { + lease := maxLease + if s := r.URL.Query().Get("lease"); s != "" { + sec, err := strconv.Atoi(s) + if err != nil || sec <= 0 { + http.Error(w, "lease must be a positive number of seconds", http.StatusBadRequest) + return + } + lease = min(time.Duration(sec)*time.Second, maxLease) + } + if err := g.maint.acquire(lease); err != nil { + http.Error(w, "busy: "+err.Error(), http.StatusConflict) + return + } + _, until := g.maint.readOnly() + if g.Log != nil { + g.Log.Info("registry read-only for garbage collection", "until", until.UTC().Format(time.RFC3339)) + } + fmt.Fprintf(w, "read-only until %s\n", until.UTC().Format(time.RFC3339)) + }) + mux.HandleFunc("POST /readwrite", func(w http.ResponseWriter, r *http.Request) { + g.maint.release() + if g.Log != nil { + g.Log.Info("registry writable again") + } + fmt.Fprintln(w, "writable") + }) + mux.HandleFunc("GET /readonly", func(w http.ResponseWriter, r *http.Request) { + if ro, until := g.maint.readOnly(); ro { + fmt.Fprintf(w, "read-only until %s\n", until.UTC().Format(time.RFC3339)) + return + } + http.Error(w, "writable", http.StatusConflict) + }) + return mux +} + +// MaintStatePath is where cmd/felis keeps the window inside the maintenance +// volume. +func MaintStatePath(dir string) string { return filepath.Join(dir, "readonly-until") } diff --git a/internal/registrygate/maint_test.go b/internal/registrygate/maint_test.go new file mode 100644 index 0000000..fe66845 --- /dev/null +++ b/internal/registrygate/maint_test.go @@ -0,0 +1,208 @@ +package registrygate + +import ( + "net/http" + "net/http/httptest" + "net/url" + "os" + "strings" + "sync" + "testing" + "time" +) + +const testDigest = "sha256:0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef" + +type clock struct { + mu sync.Mutex + t time.Time +} + +func (c *clock) now() time.Time { + c.mu.Lock() + defer c.mu.Unlock() + return c.t +} + +func (c *clock) advance(d time.Duration) { + c.mu.Lock() + c.t = c.t.Add(d) + c.mu.Unlock() +} + +func newMaintGate(t *testing.T) (*Gate, *httptest.Server, *httptest.Server, *upstreamLog, *clock) { + t.Helper() + up := &upstreamLog{} + upSrv := httptest.NewServer(up.handler()) + t.Cleanup(upSrv.Close) + target, _ := url.Parse(upSrv.URL) + g := New(target, map[string]string{ + PrincipalPlatform: "plat-secret", PrincipalBuild: "build-secret", PrincipalPrune: "prune-secret", + }, nil) + clk := &clock{t: time.Date(2026, 9, 24, 3, 0, 0, 0, time.UTC)} + g.maint.now = clk.now + gs := httptest.NewServer(g) + t.Cleanup(gs.Close) + ms := httptest.NewServer(g.MaintHandler()) + t.Cleanup(ms.Close) + return g, gs, ms, up, clk +} + +func post(t *testing.T, srv *httptest.Server, path string) int { + t.Helper() + resp, err := http.Post(srv.URL+path, "text/plain", nil) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + return resp.StatusCode +} + +func TestPrunePrincipalMayOnlyDeleteManifestsByDigest(t *testing.T) { + up2 := &upstreamLog{} + upSrv := httptest.NewServer(up2.handler()) + t.Cleanup(upSrv.Close) + target, _ := url.Parse(upSrv.URL) + g := New(target, map[string]string{PrincipalPrune: "prune-secret"}, nil) + srv := httptest.NewServer(g) + t.Cleanup(srv.Close) + + for _, p := range []string{ + "/v2/user-uploads/s1/manifests/" + testDigest, + "/v2/felis/felis/manifests/" + testDigest, + } { + if resp := do(t, srv, http.MethodDelete, p, PrincipalPrune, "prune-secret"); resp.StatusCode != http.StatusOK { + t.Errorf("prune DELETE %s = %d, want it proxied", p, resp.StatusCode) + } + } + for _, c := range []struct{ method, path string }{ + {http.MethodDelete, "/v2/user-uploads/s1/manifests/latest"}, + {http.MethodDelete, "/v2/user-uploads/s1/blobs/" + testDigest}, + {http.MethodDelete, "/v2/user-uploads/s1/manifests/sha256:short"}, + {http.MethodPut, "/v2/user-uploads/s1/manifests/" + testDigest}, + {http.MethodPost, "/v2/user-uploads/s1/blobs/uploads/"}, + {http.MethodPatch, "/v2/user-uploads/s1/blobs/uploads/abc"}, + } { + if resp := do(t, srv, c.method, c.path, PrincipalPrune, "prune-secret"); resp.StatusCode != http.StatusForbidden { + t.Errorf("prune %s %s = %d, want 403", c.method, c.path, resp.StatusCode) + } + } + if n := up2.count(); n != 2 { + t.Fatalf("upstream saw %d requests, want the 2 allowed deletes: %v", n, up2.seen) + } +} + +func TestReadOnlyWindowWaitsForQuietThenRefusesWrites(t *testing.T) { + _, gs, ms, up, clk := newMaintGate(t) + put := "/v2/user-uploads/s1/manifests/latest" + + if resp := do(t, gs, http.MethodPut, put, PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusCreated { + t.Fatalf("write before any window = %d, want 201", resp.StatusCode) + } + // A write just finished: a push may be between two of its requests. + if code := post(t, ms, "/readonly?lease=600"); code != http.StatusConflict { + t.Fatalf("readonly right after a write = %d, want 409", code) + } + clk.advance(DefaultQuiet) + if code := post(t, ms, "/readonly?lease=600"); code != http.StatusOK { + t.Fatalf("readonly after %s of quiet = %d, want 200", DefaultQuiet, code) + } + + before := up.count() + resp := do(t, gs, http.MethodPut, put, PrincipalPlatform, "plat-secret") + if resp.StatusCode != http.StatusServiceUnavailable || resp.Header.Get("Retry-After") == "" { + t.Fatalf("write during the window = %d Retry-After=%q, want 503 with Retry-After", resp.StatusCode, resp.Header.Get("Retry-After")) + } + if resp := do(t, gs, http.MethodGet, "/v2/felis/felis/manifests/b1", "", ""); resp.StatusCode != http.StatusOK { + t.Fatalf("read during the window = %d, want 200", resp.StatusCode) + } + if up.count() != before+1 { + t.Fatalf("the refused write reached the registry: %v", up.seen) + } + resp, err := http.Get(ms.URL + "/readonly") + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("GET /readonly during the window = %d, want 200", resp.StatusCode) + } + + if code := post(t, ms, "/readwrite"); code != http.StatusOK { + t.Fatalf("readwrite = %d", code) + } + if resp := do(t, gs, http.MethodPut, put, PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusCreated { + t.Fatalf("write after the window = %d, want 201", resp.StatusCode) + } +} + +func TestReadOnlyWindowIsALease(t *testing.T) { + _, gs, ms, _, clk := newMaintGate(t) + if code := post(t, ms, "/readonly?lease=60"); code != http.StatusOK { + t.Fatalf("readonly on an idle gate = %d, want 200", code) + } + put := "/v2/user-uploads/s1/manifests/latest" + if resp := do(t, gs, http.MethodPut, put, PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("write inside the lease = %d, want 503", resp.StatusCode) + } + // A GC sidecar that died mid-sweep never releases; the lease does. + clk.advance(61 * time.Second) + if resp := do(t, gs, http.MethodPut, put, PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusCreated { + t.Fatalf("write after the lease expired = %d, want 201", resp.StatusCode) + } + for _, bad := range []string{"0", "-5", "soon"} { + if code := post(t, ms, "/readonly?lease="+bad); code != http.StatusBadRequest { + t.Errorf("lease=%s = %d, want 400", bad, code) + } + } +} + +func TestReadOnlyWindowSurvivesAGateRestart(t *testing.T) { + dir := t.TempDir() + state := MaintStatePath(dir) + g, _, ms, _, clk := newMaintGate(t) + if err := g.SetMaintenanceState(state); err != nil { + t.Fatal(err) + } + if code := post(t, ms, "/readonly?lease=600"); code != http.StatusOK { + t.Fatalf("readonly = %d", code) + } + if _, err := os.Stat(state); err != nil { + t.Fatalf("window not persisted: %v", err) + } + + // The restarted gate reads the window back and keeps refusing writes. + g2, gs2, _, _, _ := newMaintGate(t) + g2.maint.now = clk.now + if err := g2.SetMaintenanceState(state); err != nil { + t.Fatal(err) + } + if resp := do(t, gs2, http.MethodPut, "/v2/user-uploads/s1/manifests/latest", PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("write on the restarted gate = %d, want 503", resp.StatusCode) + } + + if code := post(t, ms, "/readwrite"); code != http.StatusOK { + t.Fatalf("readwrite = %d", code) + } + if _, err := os.Stat(state); !os.IsNotExist(err) { + t.Fatalf("state file left after release: %v", err) + } + + // An expired window in the file is ignored. + if err := os.WriteFile(state, []byte("1"), 0o600); err != nil { + t.Fatal(err) + } + g3, gs3, _, _, _ := newMaintGate(t) + if err := g3.SetMaintenanceState(state); err != nil { + t.Fatal(err) + } + if resp := do(t, gs3, http.MethodPut, "/v2/user-uploads/s1/manifests/latest", PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusCreated { + t.Fatalf("write with an expired window on file = %d, want 201", resp.StatusCode) + } + if err := os.WriteFile(state, []byte("not-a-number"), 0o600); err != nil { + t.Fatal(err) + } + if err := g3.SetMaintenanceState(state); err == nil || !strings.Contains(err.Error(), "maintenance state") { + t.Fatalf("a corrupt state file = %v, want an error naming it", err) + } +}