Unverified Commit 7f772bcc authored by Lemon-miaow's avatar Lemon-miaow
Browse files

feat(registry): 开启 manifest 删除,GC sidecar 在 gate 只读窗口内回收 blob

parent c49336ba
Loading
Loading
Loading
Loading
+51 −6
Changes for cmd/felis/registrygate.go: 51 added lines, 6 removed lines.
Original line number Diff line number Diff line
@@ -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
+18 −1
Changes for deploy/bootstrap.sh: 18 added lines, 1 removed line.
Original line number Diff line number Diff line
@@ -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
+52 −9
Changes for internal/imagepush/push.go: 52 added lines, 9 removed lines.
Original line number Diff line number Diff line
@@ -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
+78 −0
Changes for internal/imagepush/push_test.go: 78 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -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"
+111 −3
Changes for internal/platform/workloads.go: 111 added lines, 3 removed lines.
Original line number Diff line number Diff line
@@ -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 {
Loading