From 3424852a390573184837c4c1a87cac03e6c4d7f3 Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Thu, 24 Sep 2026 14:25:17 +0800 Subject: [PATCH] =?UTF-8?q?feat(registry):=20=E5=86=99=E5=85=A5=E6=94=B9?= =?UTF-8?q?=E8=B5=B0=E9=89=B4=E6=9D=83=E7=BD=91=E5=85=B3=EF=BC=8C=E6=9E=84?= =?UTF-8?q?=E5=BB=BA=E5=85=88=E6=89=AB=E6=8F=8F=E5=86=8D=E6=8E=A8=E9=80=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/felis/registrygate.go | 117 ++++++++ cmd/felis/run.go | 4 + internal/api/logstream.go | 2 +- internal/build/build.go | 26 +- internal/build/build_test.go | 30 ++ internal/build/jobspec.go | 147 +++++++--- internal/build/jobspec_test.go | 161 ++++++++--- internal/build/k8sjobs.go | 4 +- internal/build/validate.go | 11 + internal/imagepush/push.go | 424 ++++++++++++++++++++++++++++ internal/imagepush/push_test.go | 261 +++++++++++++++++ internal/naming/naming.go | 12 + internal/platform/workloads.go | 95 +++++-- internal/platform/workloads_test.go | 135 ++++++--- internal/registrygate/gate.go | 272 ++++++++++++++++++ internal/registrygate/gate_test.go | 258 +++++++++++++++++ 16 files changed, 1814 insertions(+), 145 deletions(-) create mode 100644 cmd/felis/registrygate.go create mode 100644 internal/imagepush/push.go create mode 100644 internal/imagepush/push_test.go create mode 100644 internal/registrygate/gate.go create mode 100644 internal/registrygate/gate_test.go diff --git a/cmd/felis/registrygate.go b/cmd/felis/registrygate.go new file mode 100644 index 0000000..6a2a1e2 --- /dev/null +++ b/cmd/felis/registrygate.go @@ -0,0 +1,117 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "io" + "log/slog" + "net/http" + "net/url" + "os" + "os/signal" + "path/filepath" + "strings" + "syscall" + "time" + + "felis.lolicon.best/internal/imagepush" + "felis.lolicon.best/internal/registrygate" +) + +// cmdRegistryGate is the sidecar entrypoint in the registry pod: it owns the +// registry port (and the loopback hostPort containerd pulls through), lets reads +// through anonymously, and forwards writes to the loopback-only registry:2 only +// 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. +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") + if err := fs.Parse(args); err != nil { + 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) + return 2 + } + log := slog.New(slog.NewTextHandler(stderr, nil)) + tokens := map[string]string{} + for _, p := range []string{registrygate.PrincipalPlatform, registrygate.PrincipalBuild} { + b, err := os.ReadFile(filepath.Join(*authDir, p)) + tok := strings.TrimSpace(string(b)) + if err != nil || tok == "" { + log.Warn("registry principal disabled: no token", "principal", p, "dir", *authDir) + continue + } + tokens[p] = tok + } + + srv := &http.Server{ + Addr: *listen, + Handler: registrygate.New(target, tokens, log), + ReadHeaderTimeout: 10 * time.Second, + } + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + go func() { + <-ctx.Done() + shutdown, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + _ = srv.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) { + fmt.Fprintf(stderr, "felis registry-gate: %v\n", err) + return 1 + } + return 0 +} + +// 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 +// executing the untrusted Dockerfile never sees it. +func cmdPushImage(args []string, stdout, stderr io.Writer) int { + fs := flag.NewFlagSet("push-image", flag.ContinueOnError) + fs.SetOutput(stderr) + tarPath := fs.String("tar", "", "image tarball Kaniko wrote with --tar-path") + ref := fs.String("ref", "", "host/repository:tag to publish it as") + scheme := fs.String("scheme", "http", "registry scheme: http for the in-cluster registry, https otherwise") + if err := fs.Parse(args); err != nil { + return 2 + } + if *tarPath == "" || *ref == "" { + fmt.Fprintln(stderr, "felis push-image: --tar and --ref are required") + return 2 + } + if *scheme != "http" && *scheme != "https" { + fmt.Fprintf(stderr, "felis push-image: bad --scheme %q\n", *scheme) + return 2 + } + user := os.Getenv("FELIS_REGISTRY_USERNAME") + pass := os.Getenv("FELIS_REGISTRY_PASSWORD") + if user == "" || pass == "" { + fmt.Fprintln(stderr, "felis push-image: FELIS_REGISTRY_USERNAME/FELIS_REGISTRY_PASSWORD are empty — the registry refuses anonymous writes") + return 2 + } + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + p := &imagepush.Pusher{Scheme: *scheme, Username: user, Password: pass, Log: stderr} + digest, err := p.Push(ctx, *tarPath, *ref) + if err != nil { + fmt.Fprintf(stderr, "felis push-image: %v\n", err) + return 1 + } + fmt.Fprintln(stdout, digest) + return 0 +} diff --git a/cmd/felis/run.go b/cmd/felis/run.go index d0906e7..ed42a52 100644 --- a/cmd/felis/run.go +++ b/cmd/felis/run.go @@ -20,6 +20,8 @@ Commands: backup Archive a world into the backup store and record it (internal Job entrypoint) files List/read/write one file in a stopped server's world (internal Job entrypoint) fetch-context Fetch and extract a submission's build context (internal Job entrypoint) + push-image Push a scanned image tarball to the registry (internal Job entrypoint) + registry-gate Authorize registry writes in front of registry:2 (internal sidecar entrypoint) manifests Render the control-plane RBAC + NetworkPolicy install bundle as YAML apply Create a MinecraftServer CRD (direct K8s write; use -f server.json) setup Run host bootstrap + first-run setup console (TUI; requires root/sudo) @@ -50,6 +52,8 @@ var commands = map[string]func(args []string, stdout, stderr io.Writer) int{ "backup": cmdBackup, "files": cmdFiles, "fetch-context": cmdFetchContext, + "push-image": cmdPushImage, + "registry-gate": cmdRegistryGate, "manifests": cmdManifests, "apply": cmdApply, "setup": cmdSetup, diff --git a/internal/api/logstream.go b/internal/api/logstream.go index 2b595d8..df1aa7e 100644 --- a/internal/api/logstream.go +++ b/internal/api/logstream.go @@ -299,7 +299,7 @@ func (k *K8sLogStreamer) StreamLogs(ctx context.Context, name string) (io.ReadCl // the build-id label and follows the kaniko container's log — the build/push // output an admin watches live as a build runs. It deliberately does NOT reuse // K8sLogStreamer's PodRunning filter: a Pod running its kaniko *initContainer* is -// Phase=Pending (the main trivy container has not started), so a running filter +// Phase=Pending (the trivy scan and the push container have not started), so a running filter // would never match a live build. Trivy's CRITICAL-CVE verdict is the admission // gate, surfaced via the build status (handleGetBuild), not through this stream. // diff --git a/internal/build/build.go b/internal/build/build.go index dcde97c..cca9586 100644 --- a/internal/build/build.go +++ b/internal/build/build.go @@ -1,22 +1,24 @@ // Package build implements the image build subsystem (spec §16) — "the // platform's biggest security surface". A SysAdmin uploads a Dockerfile and a -// context tarball; felis-api starts an in-cluster Kaniko Job that builds and -// pushes to the internal registry, after which a Trivy scan gates admission to -// the image whitelist. +// context tarball; felis-api starts an in-cluster Job in which Kaniko builds the +// image into a tarball, Trivy scans that tarball, and only a clean image is +// pushed to the internal registry and admitted to the image whitelist. // // Trust model (spec §16, §22): we trust the SysAdmin at the *ingress* (only an // admin through Zero Trust may submit a build) but never trust the *Dockerfile // at runtime* — an arbitrary Dockerfile is build-time RCE whose victim is the // cluster, not the uploader. So the build Pod runs with a deliberately weak -// service account in an isolated namespace that can only push to the registry +// service account in an isolated namespace that can only reach the registry // and cannot touch the minecraft namespace, the felis database, or the K8s API // (spec §21). Those isolation guarantees live in the Job/NetworkPolicy specs // (jobspec.go) and are asserted by unit tests, since no cluster runs here. // // The Trivy gate is enforced as the build Pod's *exit code*: a kaniko -// initContainer builds and pushes, then a trivy container scans the pushed ref -// with `--exit-code 1 --severity CRITICAL`. Therefore "Job Succeeded" is -// equivalent to "pushed AND no CRITICAL CVE". felis-api observes the Job phase +// initContainer builds into a tarball (--no-push), a trivy initContainer scans it +// with `--exit-code 1 --severity CRITICAL`, and only then does the push container +// — the one holding the registry credential — publish it. Therefore "Job +// Succeeded" is equivalent to "no CRITICAL CVE AND pushed", and a rejected image +// never reaches the registry. felis-api observes the Job phase // and performs the database writes — the build Pod itself never has database // credentials (the weak-SA red line). On success the image is admitted to // image_whitelist with enabled=true (recording added_by); on failure the build @@ -70,11 +72,11 @@ const ( JobUnknown JobPhase = iota JobPending JobRunning - // JobSucceeded means kaniko pushed AND trivy found no CRITICAL CVE — the - // scan gate passed (spec §16). + // JobSucceeded means trivy found no CRITICAL CVE AND the image was pushed — + // the scan gate passed (spec §16). JobSucceeded - // JobFailed means kaniko failed OR trivy found a CRITICAL CVE — the build - // is rejected and nothing is admitted. + // JobFailed means kaniko failed, trivy found a CRITICAL CVE, or the push + // failed — the build is rejected and nothing is admitted. JobFailed ) @@ -390,7 +392,7 @@ func (b *Builder) Get(ctx context.Context, id string) (*Build, error) { // translation (spec §16). A terminal build is returned unchanged (idempotent). // // - JobSucceeded → status=succeeded AND the image is admitted to the whitelist -// with enabled=true (kaniko pushed and trivy found no CRITICAL CVE). +// with enabled=true (trivy found no CRITICAL CVE and the push landed). // - JobFailed / JobUnknown → status=failed, nothing admitted (a CRITICAL CVE // surfaces here as a failed Job, since trivy runs with --exit-code 1). // - JobPending / JobRunning → no change. diff --git a/internal/build/build_test.go b/internal/build/build_test.go index 1d7e1a0..ca3e53b 100644 --- a/internal/build/build_test.go +++ b/internal/build/build_test.go @@ -217,6 +217,36 @@ func TestSubmitRejectsExternalRegistryTarget(t *testing.T) { } } +// The platform's own images and the scanner's DB mirrors live under felis/ and +// mirror/; the registry gate refuses the build principal there, and Validate turns +// that into a 400 before a Job spends minutes building an image it cannot push. +func TestValidateRejectsReservedRepos(t *testing.T) { + cfg := Config{RegistryURL: "registry.felis.svc:5000"} + for _, ref := range []string{ + "registry.felis.svc:5000/felis/felis:v0.1.0", + "registry.felis.svc:5000/felis:latest", + "registry.felis.svc:5000/felis", + "registry.felis.svc:5000/mirror/trivy-db:2", + } { + req := goodRequest() + req.ImageRef = ref + if err := Validate(req, cfg); !errors.Is(err, ErrInvalid) { + t.Errorf("Validate(%q) = %v, want ErrInvalid", ref, err) + } + } + for _, ref := range []string{ + "registry.felis.svc:5000/user-uploads/s1:latest", + "registry.felis.svc:5000/felis-pack:1", + "registry.felis.svc:5000/builds/felis:1", + } { + req := goodRequest() + req.ImageRef = ref + if err := Validate(req, cfg); err != nil { + t.Errorf("Validate(%q) = %v, want accepted", ref, err) + } + } +} + func TestSubmitRejectsEmptyAndOversizeDockerfile(t *testing.T) { b, _, _ := newBuilder() req := goodRequest() diff --git a/internal/build/jobspec.go b/internal/build/jobspec.go index 9427319..06c7401 100644 --- a/internal/build/jobspec.go +++ b/internal/build/jobspec.go @@ -26,14 +26,16 @@ const ( ) // Container names within the build Pod. Kaniko is the initContainer that builds -// and pushes the image — its log IS the "build log" an admin watches (spec §16); -// Trivy is the main container whose CRITICAL-CVE verdict gates admission and is -// surfaced via the build status, not the log stream. Exported so the build-log -// streamer (internal/api.K8sBuildLogStreamer, spec §416 日志流复用 §8) follows the -// same container this Job defines — one source of truth for the name. +// the image into a tarball — its log IS the "build log" an admin watches (spec +// §16); Trivy is the next initContainer, whose CRITICAL-CVE verdict gates both the +// push and admission and is surfaced via the build status, not the log stream; +// Push is the main container that publishes the scanned tarball. Exported so the +// build-log streamer (internal/api.K8sBuildLogStreamer, spec §416 日志流复用 §8) +// follows the same container this Job defines — one source of truth for the name. const ( ContainerKaniko = "kaniko" ContainerTrivy = "trivy" + ContainerPush = "push" // ContainerFetch is the initContainer that pulls a submission's build context // from the felis-api internal face and extracts it into the shared emptyDir. // It exists only for an http(s) ContextRef (see BuildJob); a ref Kaniko can @@ -44,8 +46,19 @@ const ( // initContainer writes the extracted tree there, Kaniko reads it read-only. contextVolume = "context" contextMountPath = "/context" + + // imageVolume/imageTarPath carry the built image from Kaniko (--tar-path) to + // Trivy (--input) and then to the push container. + imageVolume = "image" + imageMountPath = "/image" + imageTarPath = imageMountPath + "/image.tar" ) +// imageSizeLimit bounds the built image tarball. A modpack image is typically a +// JRE, a server jar and a few hundred MiB of mods; 10 GiB leaves ample room while +// still stopping a runaway build from filling the node's disk. +var imageSizeLimit = resource.MustParse("10Gi") + // contextSizeLimit bounds the extracted (attacker-controlled) context tree so a // tarball bomb wedges the build pod instead of the node's disk. The compressed // upload is capped at 1 GiB by the submit lane; 4 GiB leaves expansion room. @@ -120,14 +133,25 @@ func buildLabels(p JobParams) map[string]string { // CRITICAL CVE fails the Pod and therefore the Job — the only retained // automatic admission gate (spec §16). // -// Sequencing: kaniko runs as an initContainer (build + push to the internal -// registry) and trivy as the main container (scan the pushed ref). The Pod -// succeeds only if kaniko pushed AND trivy found no CRITICAL CVE. +// Sequencing: kaniko builds with --no-push into a tarball, trivy scans that +// tarball, and only then does the push container publish it. So: +// +// - an image that fails the scan is never published — it used to be pushed to +// the final tag first and scanned after, overwriting whatever that tag held; +// - the registry credential lives in the push container alone. Kaniko executes +// the untrusted Dockerfile and holds no credential at all, and the registry +// refuses anonymous writes (internal/registrygate). +// +// The Pod succeeds only if kaniko built, trivy found no CRITICAL CVE, and the push +// landed. func BuildJob(p JobParams) (*batchv1.Job, error) { limits, err := resourceLimits(p.CPULimit, p.MemLimit) if err != nil { return nil, err } + if p.FelisImage == "" { + return nil, fmt.Errorf("build: FelisImage is required: the push container runs it") + } deadline := int64(p.Deadline / time.Second) if deadline <= 0 { deadline = int64(defaultDeadline / time.Second) @@ -164,12 +188,15 @@ func BuildJob(p JobParams) (*batchv1.Job, error) { // in place (s3://, or a path an installer pre-mounted) passes through untouched. contextPath := p.ContextRef initContainers := []corev1.Container{} - var kanikoMounts []corev1.VolumeMount - var podVolumes []corev1.Volume + imageMount := corev1.VolumeMount{Name: imageVolume, MountPath: imageMountPath} + kanikoMounts := []corev1.VolumeMount{imageMount} + podVolumes := []corev1.Volume{{ + Name: imageVolume, + VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{ + SizeLimit: quantityPtr(imageSizeLimit), + }}, + }} if isHTTPContextRef(p.ContextRef) { - if p.FelisImage == "" { - return nil, fmt.Errorf("build: context ref %q needs FelisImage for the fetch initContainer", p.ContextRef) - } contextPath = contextMountPath // The fetch container runs as root while Kaniko keeps the image default // (also root): Kaniko re-copies the Dockerfile out of the context and @@ -209,17 +236,17 @@ func BuildJob(p JobParams) (*batchv1.Job, error) { SecurityContext: fetchSec, } initContainers = append(initContainers, fetch) - kanikoMounts = []corev1.VolumeMount{{Name: contextVolume, MountPath: contextMountPath, ReadOnly: true}} - podVolumes = []corev1.Volume{{ + kanikoMounts = append(kanikoMounts, corev1.VolumeMount{Name: contextVolume, MountPath: contextMountPath, ReadOnly: true}) + podVolumes = append(podVolumes, corev1.Volume{ Name: contextVolume, VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{ // The extracted tree is attacker-controlled; bound it so a tarball // bomb wedges THIS pod (admitted failure) instead of filling the // node's disk. The compressed upload is capped at 1 GiB by the // submit lane, and 4 GiB leaves room for a typical expansion. - SizeLimit: sizeLimitPtr(), + SizeLimit: quantityPtr(contextSizeLimit), }}, - }} + }) } kaniko := corev1.Container{ @@ -228,16 +255,16 @@ func BuildJob(p JobParams) (*batchv1.Job, error) { Args: []string{ "--dockerfile=Dockerfile", "--context=" + contextPath, + // --destination only names the image inside the tarball; --no-push + // keeps Kaniko off the registry's write path entirely. "--destination=" + p.ImageRef, - // The internal registry is in-cluster only and may serve plain HTTP; - // it is never a public ingress (spec §17). Both directions need the - // insecure flags: --insecure/--skip-tls-verify cover the PUSH, while - // the pull side needs its own pair — a Dockerfile's `FROM - // registry.felis.svc:5000/...` otherwise fails with "server gave - // HTTP response to HTTPS client", breaking every build based on a + "--no-push", + "--tar-path=" + imageTarPath, + // The internal registry is in-cluster only and serves plain HTTP; it + // is never a public ingress (spec §17). A Dockerfile's `FROM + // registry.felis.svc:5000/...` fails with "server gave HTTP response + // to HTTPS client" without these, breaking every build based on a // platform image (the canonical modpack shape). - "--insecure", - "--skip-tls-verify", "--insecure-pull", "--skip-tls-verify-pull", }, @@ -249,6 +276,7 @@ func BuildJob(p JobParams) (*batchv1.Job, error) { trivyArgs := []string{ "image", + "--input", imageTarPath, "--exit-code", "1", "--severity", "CRITICAL", "--no-progress", @@ -265,14 +293,44 @@ func BuildJob(p JobParams) (*batchv1.Job, error) { if p.TrivyJavaDBRepository != "" { trivyArgs = append(trivyArgs, "--java-db-repository", p.TrivyJavaDBRepository) } - trivyArgs = append(trivyArgs, p.ImageRef) trivy := corev1.Container{ Name: ContainerTrivy, Image: p.TrivyImage, Args: trivyArgs, + VolumeMounts: []corev1.VolumeMount{{Name: imageVolume, MountPath: imageMountPath, ReadOnly: true}}, Resources: corev1.ResourceRequirements{Limits: limits, Requests: buildRequests(limits)}, SecurityContext: sec, } + initContainers = append(initContainers, trivy) + + // The publish step: the only container that holds the registry credential, + // read from a Secret the installer materializes in this namespace. It runs + // the felis binary (internal/imagepush), which only reads the tarball. + pushSec := sec.DeepCopy() + pushSec.ReadOnlyRootFilesystem = boolPtr(true) + secretEnv := func(name, key string) corev1.EnvVar { + return corev1.EnvVar{Name: name, ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{Name: naming.RegistryPushSecretName}, + Key: key, + }}} + } + push := corev1.Container{ + Name: ContainerPush, + Image: p.FelisImage, + Args: []string{ + "push-image", + "--tar=" + imageTarPath, + "--ref=" + p.ImageRef, + "--scheme=http", + }, + Env: []corev1.EnvVar{ + secretEnv("FELIS_REGISTRY_USERNAME", naming.RegistryPushUsernameKey), + secretEnv("FELIS_REGISTRY_PASSWORD", naming.RegistryPushPasswordKey), + }, + VolumeMounts: []corev1.VolumeMount{{Name: imageVolume, MountPath: imageMountPath, ReadOnly: true}}, + Resources: corev1.ResourceRequirements{Limits: limits, Requests: buildRequests(limits)}, + SecurityContext: pushSec, + } job := &batchv1.Job{ ObjectMeta: metav1.ObjectMeta{ @@ -293,7 +351,7 @@ func BuildJob(p JobParams) (*batchv1.Job, error) { ServiceAccountName: p.ServiceAccount, AutomountServiceAccountToken: boolPtr(false), InitContainers: initContainers, - Containers: []corev1.Container{trivy}, + Containers: []corev1.Container{push}, Volumes: podVolumes, }, }, @@ -309,6 +367,20 @@ func isHTTPContextRef(ref string) bool { return strings.HasPrefix(ref, "http://") || strings.HasPrefix(ref, "https://") } +// ClusterDNSPeer selects the cluster DNS pods (CoreDNS in kube-system, labelled +// k8s-app=kube-dns on k3s and upstream alike) — the only resolver a sandboxed pod +// needs. +func ClusterDNSPeer() networkingv1.NetworkPolicyPeer { + return networkingv1.NetworkPolicyPeer{ + NamespaceSelector: &metav1.LabelSelector{ + MatchLabels: map[string]string{"kubernetes.io/metadata.name": "kube-system"}, + }, + PodSelector: &metav1.LabelSelector{ + MatchLabels: map[string]string{"k8s-app": "kube-dns"}, + }, + } +} + // NetPolParams parameterises the build-namespace egress lock. type NetPolParams struct { Namespace string @@ -329,8 +401,8 @@ type NetPolParams struct { // BuildNetworkPolicy renders the default-deny egress policy for build Pods // (spec §16, §21: build ns egress 仅放 registry + 包源,默认拒外网). It selects // build Pods by the managed-by label, denies all ingress, and allows egress -// only to DNS, the internal registry, and any explicitly configured package -// mirrors. There is deliberately no allow-all egress rule. +// only to the cluster DNS pods, the internal registry, felis-api's internal +// face, and any explicitly configured package mirrors. There is deliberately no allow-all egress rule. func BuildNetworkPolicy(p NetPolParams) *networkingv1.NetworkPolicy { port := p.RegistryPort if port == 0 { @@ -351,9 +423,13 @@ func BuildNetworkPolicy(p NetPolParams) *networkingv1.NetworkPolicy { ctxPort := intstr.FromInt32(apiPort) egress := []networkingv1.NetworkPolicyEgressRule{ - // DNS resolution: port-restricted to 53, so this is not an open-internet - // hole — name resolution only. + // DNS resolution, to the cluster resolver only. Port 53 to ANY address + // would be an exfiltration channel out of an otherwise sealed sandbox + // (a Dockerfile RUN can speak DNS, or anything else, to a resolver it + // controls); the cluster DNS Service is DNATed to these pods before the + // policy is evaluated, so selecting them is what "resolve names" means. { + To: []networkingv1.NetworkPolicyPeer{ClusterDNSPeer()}, Ports: []networkingv1.NetworkPolicyPort{ {Protocol: &dnsUDP, Port: &dns53}, {Protocol: &dnsTCP, Port: &dns53}, @@ -487,11 +563,10 @@ func buildRequests(limits corev1.ResourceList) corev1.ResourceList { func boolPtr(b bool) *bool { return &b } func int32Ptr(i int32) *int32 { return &i } -// sizeLimitPtr returns a copy of contextSizeLimit for a VolumeSource (the API -// object only ever gets serialized, but a shared pointer across rendered Jobs -// invites accidental aliasing). -func sizeLimitPtr() *resource.Quantity { - q := contextSizeLimit +// quantityPtr returns a pointer to a copy of q for a VolumeSource (the API object +// only ever gets serialized, but a shared pointer across rendered Jobs invites +// accidental aliasing). +func quantityPtr(q resource.Quantity) *resource.Quantity { return &q } func int64Ptr(i int64) *int64 { return &i } diff --git a/internal/build/jobspec_test.go b/internal/build/jobspec_test.go index effae8d..4e31dd6 100644 --- a/internal/build/jobspec_test.go +++ b/internal/build/jobspec_test.go @@ -17,6 +17,7 @@ func sampleJobParams() JobParams { Namespace: defaultNamespace, ServiceAccount: defaultServiceAccount, RegistryURL: "registry.felis.svc:5000", + FelisImage: "felis:test", KanikoImage: defaultKanikoImage, TrivyImage: defaultTrivyImage, Deadline: 30 * time.Minute, @@ -159,55 +160,119 @@ func TestBuildJobRequestsAreASchedulableFloor(t *testing.T) { } } -// kaniko builds and pushes to the request's exact target; trivy gates admission -// with --exit-code 1 --severity CRITICAL on that same ref. -func TestBuildJobKanikoPushesAndTrivyGates(t *testing.T) { +// kaniko builds the request's exact target into a tarball and never pushes; +// trivy gates on that tarball with --exit-code 1 --severity CRITICAL; only then +// does the push container publish it. An image that fails the scan is therefore +// never in the registry, and the credential is never where the Dockerfile runs. +func TestBuildJobScansBeforePush(t *testing.T) { p := sampleJobParams() job, err := BuildJob(p) if err != nil { t.Fatalf("BuildJob: %v", err) } - if len(job.Spec.Template.Spec.InitContainers) != 1 { - t.Fatalf("expected exactly one (kaniko) initContainer") + inits := job.Spec.Template.Spec.InitContainers + if len(inits) != 2 || inits[0].Name != ContainerKaniko || inits[1].Name != ContainerTrivy { + t.Fatalf("initContainers = %v, want [kaniko trivy]", initNames(inits)) } - kaniko := job.Spec.Template.Spec.InitContainers[0] - if kaniko.Name != "kaniko" { - t.Errorf("init container = %q, want kaniko", kaniko.Name) + kaniko, trivy := inits[0], inits[1] + for _, want := range []string{"--destination=" + p.ImageRef, "--no-push", "--tar-path=" + imageTarPath} { + if !hasArg(kaniko.Args, want) { + t.Errorf("kaniko args = %v, want %s", kaniko.Args, want) + } } - if !hasArg(kaniko.Args, "--destination="+p.ImageRef) { - t.Errorf("kaniko must push to %q, args=%v", p.ImageRef, kaniko.Args) + for _, pushFlag := range []string{"--insecure", "--skip-tls-verify"} { + if hasArg(kaniko.Args, pushFlag) { + t.Errorf("kaniko args = %v still carry the push-side %s", kaniko.Args, pushFlag) + } } - // The pull direction needs its own flags: --insecure/--skip-tls-verify only - // cover the push, and without the pull pair a Dockerfile's `FROM` fails - // against the plain-HTTP registry ("server gave HTTP response to HTTPS - // client") — the live failure this guards. + // The pull direction needs its own flags: without the pull pair a + // Dockerfile's `FROM` fails against the plain-HTTP registry ("server gave + // HTTP response to HTTPS client") — the live failure this guards. for _, flag := range []string{"--insecure-pull", "--skip-tls-verify-pull"} { if !hasArg(kaniko.Args, flag) { t.Errorf("kaniko args = %v, want %s so base-image pulls use plain HTTP", kaniko.Args, flag) } } - if len(job.Spec.Template.Spec.Containers) != 1 { - t.Fatalf("expected exactly one (trivy) main container") + // The scan gate: a CRITICAL CVE must fail the Pod (and thus the Job) before + // the push container ever starts. + if !argPairPresent(trivy.Args, "--input", imageTarPath) { + t.Errorf("trivy must scan the built tarball, args=%v", trivy.Args) } - trivy := job.Spec.Template.Spec.Containers[0] - if trivy.Name != "trivy" { - t.Errorf("main container = %q, want trivy", trivy.Name) + if hasArg(trivy.Args, p.ImageRef) { + t.Errorf("trivy must not scan the registry ref (nothing is pushed yet), args=%v", trivy.Args) } - // The scan gate: a CRITICAL CVE must fail the Pod (and thus the Job). if !argPairPresent(trivy.Args, "--exit-code", "1") { t.Errorf("trivy must run with --exit-code 1, args=%v", trivy.Args) } if !argPairPresent(trivy.Args, "--severity", "CRITICAL") { t.Errorf("trivy must gate on --severity CRITICAL, args=%v", trivy.Args) } - if !hasArg(trivy.Args, p.ImageRef) { - t.Errorf("trivy must scan the pushed ref %q, args=%v", p.ImageRef, trivy.Args) - } // No DB repositories configured: Trivy keeps its own defaults. if hasArg(trivy.Args, "--db-repository") || hasArg(trivy.Args, "--java-db-repository") { t.Errorf("unset DB repositories must not render --db-repository/--java-db-repository, args=%v", trivy.Args) } + + if len(job.Spec.Template.Spec.Containers) != 1 { + t.Fatalf("expected exactly one (push) main container") + } + push := job.Spec.Template.Spec.Containers[0] + if push.Name != ContainerPush || push.Image != p.FelisImage { + t.Errorf("main container = %s (%s), want push running the platform image", push.Name, push.Image) + } + if !hasArg(push.Args, "push-image") || !hasArg(push.Args, "--ref="+p.ImageRef) || !hasArg(push.Args, "--tar="+imageTarPath) { + t.Errorf("push args = %v, want push-image --tar=%s --ref=%s", push.Args, imageTarPath, p.ImageRef) + } + creds := map[string]string{} + for _, e := range push.Env { + if e.Value != "" || e.ValueFrom == nil || e.ValueFrom.SecretKeyRef == nil { + t.Errorf("push env %s must come from a secretKeyRef, got %#v", e.Name, e) + continue + } + creds[e.Name] = e.ValueFrom.SecretKeyRef.Name + } + for _, name := range []string{"FELIS_REGISTRY_USERNAME", "FELIS_REGISTRY_PASSWORD"} { + if creds[name] != "felis-registry-push" { + t.Errorf("push env %s from secret %q, want felis-registry-push", name, creds[name]) + } + } + // Nothing else in the pod may hold the credential. + for _, c := range inits { + for _, e := range c.Env { + if e.ValueFrom != nil && e.ValueFrom.SecretKeyRef != nil && e.ValueFrom.SecretKeyRef.Name == "felis-registry-push" { + t.Errorf("container %s holds the registry credential", c.Name) + } + } + } + // The tarball is shared through a bounded emptyDir, read-only past kaniko. + var imgVol *corev1.Volume + for i := range job.Spec.Template.Spec.Volumes { + if job.Spec.Template.Spec.Volumes[i].Name == imageVolume { + imgVol = &job.Spec.Template.Spec.Volumes[i] + } + } + if imgVol == nil || imgVol.EmptyDir == nil || imgVol.EmptyDir.SizeLimit == nil { + t.Fatalf("image volume must be a size-limited emptyDir, got %#v", imgVol) + } + for _, c := range []corev1.Container{trivy, push} { + ro := false + for _, m := range c.VolumeMounts { + if m.Name == imageVolume && m.ReadOnly { + ro = true + } + } + if !ro { + t.Errorf("%s must mount the image tarball read-only, got %v", c.Name, c.VolumeMounts) + } + } +} + +func TestBuildJobNeedsFelisImage(t *testing.T) { + p := sampleJobParams() + p.FelisImage = "" + if _, err := BuildJob(p); err == nil { + t.Fatal("a build without FelisImage has no push container and must fail to render") + } } // Kaniko (and only kaniko) must carry back exactly the unpacking capability @@ -274,21 +339,23 @@ func TestBuildJobTrivyDBRepositoryOverride(t *testing.T) { if err != nil { t.Fatalf("BuildJob: %v", err) } - trivy := job.Spec.Template.Spec.Containers[0] + trivy := job.Spec.Template.Spec.InitContainers[1] + if trivy.Name != ContainerTrivy { + t.Fatalf("initContainers = %v, want trivy second", initNames(job.Spec.Template.Spec.InitContainers)) + } if !argPairPresent(trivy.Args, "--db-repository", p.TrivyDBRepository) { t.Errorf("trivy args = %v, want --db-repository %s", trivy.Args, p.TrivyDBRepository) } if !argPairPresent(trivy.Args, "--java-db-repository", p.TrivyJavaDBRepository) { t.Errorf("trivy args = %v, want --java-db-repository %s", trivy.Args, p.TrivyJavaDBRepository) } - // The scanned image ref must stay the last argument. - if last := trivy.Args[len(trivy.Args)-1]; last != p.ImageRef { - t.Errorf("image ref must remain the last argument, args=%v", trivy.Args) + if !argPairPresent(trivy.Args, "--input", imageTarPath) { + t.Errorf("trivy args = %v, want --input %s", trivy.Args, imageTarPath) } } // The build namespace egress lock must be default-deny: deny all ingress, and -// allow egress only to DNS + the internal registry — never an allow-all rule. +// allow egress only to cluster DNS + the internal registry — never an allow-all rule. func TestBuildNetworkPolicyIsDefaultDeny(t *testing.T) { np := BuildNetworkPolicy(NetPolParams{ Namespace: "felis-build", @@ -321,6 +388,19 @@ func TestBuildNetworkPolicyIsDefaultDeny(t *testing.T) { if !egressAllowsPort(np, 53) { t.Error("egress must allow DNS (port 53)") } + // ...but only to the cluster resolver: port 53 to any address is a way out + // of the sandbox for anything that speaks DNS (or anything at all) on 53. + for i, rule := range np.Spec.Egress { + for _, port := range rule.Ports { + if port.Port == nil || port.Port.IntVal != 53 { + continue + } + if len(rule.To) != 1 || rule.To[0].PodSelector == nil || rule.To[0].PodSelector.MatchLabels["k8s-app"] != "kube-dns" || + rule.To[0].NamespaceSelector == nil || rule.To[0].NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != "kube-system" { + t.Errorf("egress rule %d opens port 53 to %v, want only kube-system/k8s-app=kube-dns", i, rule.To) + } + } + } // The context fetch: build Pods stream submissions from the control // namespace's internal face (defaults: felis + 8081). if !egressAllowsNamespace(np, "felis") { @@ -343,8 +423,8 @@ func TestBuildJobFetchesHTTPContext(t *testing.T) { t.Fatalf("BuildJob: %v", err) } inits := job.Spec.Template.Spec.InitContainers - if len(inits) != 2 || inits[0].Name != ContainerFetch || inits[1].Name != ContainerKaniko { - t.Fatalf("initContainers = %v, want [%s %s]", initNames(inits), ContainerFetch, ContainerKaniko) + if len(inits) != 3 || inits[0].Name != ContainerFetch || inits[1].Name != ContainerKaniko || inits[2].Name != ContainerTrivy { + t.Fatalf("initContainers = %v, want [%s %s %s]", initNames(inits), ContainerFetch, ContainerKaniko, ContainerTrivy) } fetch, kaniko := inits[0], inits[1] if fetch.Image != p.FelisImage { @@ -405,17 +485,6 @@ func TestBuildJobFetchesHTTPContext(t *testing.T) { } } -// Without the platform image the fetch initContainer cannot run, so rendering an -// http(s) context must fail loudly at Job-creation time, not with an ImagePull -// error at 3am. -func TestBuildJobHTTPContextNeedsFelisImage(t *testing.T) { - p := sampleJobParams() - p.ContextRef = "https://example.invalid/sub-abc/context" - if _, err := BuildJob(p); err == nil { - t.Fatal("http(s) context without FelisImage must fail to render") - } -} - // A ref Kaniko reads natively (or an installer pre-mounted) must NOT grow the // fetch initContainer: the transport is for http(s) only. func TestBuildJobNativeContextNeedsNoFetch(t *testing.T) { @@ -425,11 +494,13 @@ func TestBuildJobNativeContextNeedsNoFetch(t *testing.T) { if err != nil { t.Fatalf("BuildJob: %v", err) } - if len(job.Spec.Template.Spec.InitContainers) != 1 || job.Spec.Template.Spec.InitContainers[0].Name != ContainerKaniko { - t.Errorf("a native ref must render just kaniko, got %v", initNames(job.Spec.Template.Spec.InitContainers)) + if inits := job.Spec.Template.Spec.InitContainers; len(inits) != 2 || inits[0].Name != ContainerKaniko { + t.Errorf("a native ref must render just kaniko + trivy, got %v", initNames(inits)) } - if len(job.Spec.Template.Spec.Volumes) != 0 { - t.Errorf("a native ref must render no context volume, got %v", job.Spec.Template.Spec.Volumes) + for _, v := range job.Spec.Template.Spec.Volumes { + if v.Name == contextVolume { + t.Errorf("a native ref must render no context volume, got %v", job.Spec.Template.Spec.Volumes) + } } } diff --git a/internal/build/k8sjobs.go b/internal/build/k8sjobs.go index 4aaa680..bf69052 100644 --- a/internal/build/k8sjobs.go +++ b/internal/build/k8sjobs.go @@ -61,8 +61,8 @@ func (k *K8sJobs) JobPhase(ctx context.Context, jobName string) (JobPhase, error case batchv1.JobComplete: return JobSucceeded, nil case batchv1.JobFailed: - // Covers a CRITICAL CVE (trivy --exit-code 1), a kaniko failure, and - // DeadlineExceeded — all are a rejected build. + // Covers a CRITICAL CVE (trivy --exit-code 1), a kaniko or push + // failure, and DeadlineExceeded — all are a rejected build. return JobFailed, nil } } diff --git a/internal/build/validate.go b/internal/build/validate.go index 686bf65..b96cfe2 100644 --- a/internal/build/validate.go +++ b/internal/build/validate.go @@ -4,6 +4,8 @@ import ( "fmt" "regexp" "strings" + + "felis.lolicon.best/internal/registrygate" ) // invalidf builds a validation error wrapping ErrInvalid so the API layer maps @@ -112,6 +114,15 @@ func validateRegistryTarget(ref, registryURL string) error { return invalidf("image reference %q must target the internal registry %q, not %q", ref, registryHost(registryURL), host) } + // The registry gate refuses the build principal these repositories anyway + // (they hold the platform's own images and the scanner's DB mirrors); refusing + // here turns a build that would fail at its last step into a 400 up front. + root, _, _ := strings.Cut(rest, "/") + for _, reserved := range registrygate.ReservedRepoRoots { + if root == reserved || strings.HasPrefix(root, reserved+":") { + return invalidf("image reference %q is in %s/, which is reserved for the platform's own images", ref, reserved) + } + } return nil } diff --git a/internal/imagepush/push.go b/internal/imagepush/push.go new file mode 100644 index 0000000..da16dec --- /dev/null +++ b/internal/imagepush/push.go @@ -0,0 +1,424 @@ +// Package imagepush uploads an image tarball — the docker-save layout Kaniko writes +// with --tar-path — to a registry over the distribution v2 API. +// +// It exists so a build Job can split "build" from "publish": Kaniko runs the +// untrusted Dockerfile with --no-push and no credential, Trivy scans the tarball, +// and only then does a separate container holding the registry credential push it. +// Before, Kaniko pushed straight to the final tag, so a scan failure left the +// unscanned image already published, and the push credential (had there been one) +// would have lived in the same container as the Dockerfile's RUN steps. +// +// The protocol subset is deliberately small: HEAD to skip blobs the repository +// already holds, POST + monolithic PUT to upload the rest, PUT for the manifest. +package imagepush + +import ( + "archive/tar" + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "os" + "strings" + "time" +) + +// Media types the pushed manifest uses. Gzip layers keep the Docker schema-2 shape +// every runtime pulls; any other layer compression switches the whole manifest to +// OCI, the only format that can describe it. +const ( + mediaDockerManifest = "application/vnd.docker.distribution.manifest.v2+json" + mediaDockerConfig = "application/vnd.docker.container.image.v1+json" + mediaDockerLayerGz = "application/vnd.docker.image.rootfs.diff.tar.gzip" + + mediaOCIManifest = "application/vnd.oci.image.manifest.v1+json" + mediaOCIConfig = "application/vnd.oci.image.config.v1+json" + mediaOCILayerGz = "application/vnd.oci.image.layer.v1.tar+gzip" + mediaOCILayerZstd = "application/vnd.oci.image.layer.v1.tar+zstd" + mediaOCILayerTar = "application/vnd.oci.image.layer.v1.tar" +) + +// Pusher uploads tarballs to one registry. +type Pusher struct { + // Client performs the requests; nil uses a client with sane timeouts. + Client *http.Client + // Scheme is "http" for the in-cluster registry (never a public ingress) or + // "https". + Scheme string + // Username/Password are sent as basic auth on every request. + Username, Password string + // Log receives progress lines. Nil discards. + Log io.Writer + // Attempts bounds retries of one blob upload or the manifest PUT. Zero means 3. + Attempts int +} + +// Ref is a parsed host/repository:tag reference. +type Ref struct { + Host, Repo, Tag string +} + +func (r Ref) String() string { return r.Host + "/" + r.Repo + ":" + r.Tag } + +// ParseRef splits host/repo:tag. The host must look like a registry host (contain +// a '.' or ':'), and a digest reference is refused: a push names a tag. +func ParseRef(ref string) (Ref, error) { + if strings.Contains(ref, "@") { + return Ref{}, fmt.Errorf("imagepush: %q is a digest reference; push needs a tag", ref) + } + host, rest, ok := strings.Cut(ref, "/") + if !ok || rest == "" || !strings.ContainsAny(host, ".:") { + return Ref{}, fmt.Errorf("imagepush: %q must be host/repository:tag", ref) + } + repo, tag := rest, "latest" + if i := strings.LastIndexByte(rest, ':'); i > strings.LastIndexByte(rest, '/') { + repo, tag = rest[:i], rest[i+1:] + } + if repo == "" || tag == "" { + return Ref{}, fmt.Errorf("imagepush: %q must be host/repository:tag", ref) + } + return Ref{Host: host, Repo: repo, Tag: tag}, nil +} + +// descriptor is one blob of the image as the manifest records it. +type descriptor struct { + MediaType string `json:"mediaType"` + Size int64 `json:"size"` + Digest string `json:"digest"` + + file string // entry name inside the tarball +} + +type manifest struct { + SchemaVersion int `json:"schemaVersion"` + MediaType string `json:"mediaType"` + Config descriptor `json:"config"` + Layers []descriptor `json:"layers"` +} + +// tarManifest is one entry of the tarball's manifest.json. +type tarManifest struct { + Config string + RepoTags []string + Layers []string +} + +// Push uploads the image in tarPath as ref and returns the manifest digest the +// registry recorded. +func (p *Pusher) Push(ctx context.Context, tarPath, ref string) (string, error) { + r, err := ParseRef(ref) + if err != nil { + return "", err + } + m, err := p.describe(tarPath) + if err != nil { + return "", err + } + blobs := append([]descriptor{m.Config}, m.Layers...) + for i, b := range blobs { + if err := p.retry(ctx, func() error { return p.pushBlob(ctx, r, tarPath, b) }); err != nil { + return "", fmt.Errorf("imagepush: blob %d/%d (%s): %w", i+1, len(blobs), b.Digest, err) + } + } + body, err := json.Marshal(m) + if err != nil { + return "", err + } + var digest string + err = p.retry(ctx, func() error { + d, err := p.putManifest(ctx, r, m.MediaType, body) + digest = d + return err + }) + if err != nil { + return "", fmt.Errorf("imagepush: manifest: %w", err) + } + p.logf("pushed %s@%s", r, digest) + return digest, nil +} + +// describe reads the tarball's manifest.json and digests every blob it names. +func (p *Pusher) describe(tarPath string) (*manifest, error) { + raw, err := readEntry(tarPath, "manifest.json", 1<<20) + if err != nil { + return nil, err + } + var entries []tarManifest + if err := json.Unmarshal(raw, &entries); err != nil { + return nil, fmt.Errorf("imagepush: manifest.json: %w", err) + } + if len(entries) != 1 { + return nil, fmt.Errorf("imagepush: manifest.json describes %d images, want exactly 1", len(entries)) + } + e := entries[0] + if e.Config == "" || len(e.Layers) == 0 { + return nil, errors.New("imagepush: manifest.json names no config or no layers") + } + + cfg, _, err := digestEntry(tarPath, e.Config) + if err != nil { + return nil, err + } + // Kaniko names the config after its own digest; a mismatch means the tarball + // is not what it claims to be. + if strings.HasPrefix(e.Config, "sha256:") && e.Config != cfg.Digest { + return nil, fmt.Errorf("imagepush: config %s hashes to %s", e.Config, cfg.Digest) + } + + m := &manifest{SchemaVersion: 2, MediaType: mediaDockerManifest} + oci := false + for _, name := range e.Layers { + d, magic, err := digestEntry(tarPath, name) + if err != nil { + return nil, err + } + switch { + case bytes.HasPrefix(magic, []byte{0x1f, 0x8b}): + d.MediaType = mediaDockerLayerGz + case bytes.HasPrefix(magic, []byte{0x28, 0xb5, 0x2f, 0xfd}): + d.MediaType, oci = mediaOCILayerZstd, true + default: + d.MediaType, oci = mediaOCILayerTar, true + } + m.Layers = append(m.Layers, d) + } + cfg.MediaType = mediaDockerConfig + if oci { + m.MediaType = mediaOCIManifest + cfg.MediaType = mediaOCIConfig + for i := range m.Layers { + if m.Layers[i].MediaType == mediaDockerLayerGz { + m.Layers[i].MediaType = mediaOCILayerGz + } + } + } + m.Config = cfg + return m, nil +} + +func (p *Pusher) pushBlob(ctx context.Context, r Ref, tarPath string, b descriptor) error { + head, err := p.do(ctx, http.MethodHead, p.url(r, "blobs/"+b.Digest), nil, 0, "") + if err != nil { + return err + } + head.Body.Close() + if head.StatusCode == http.StatusOK { + p.logf("exists %s", b.Digest) + return nil + } + + start, err := p.do(ctx, http.MethodPost, p.url(r, "blobs/uploads/"), nil, 0, "") + if err != nil { + return err + } + start.Body.Close() + if start.StatusCode != http.StatusAccepted { + return statusError("start upload", start) + } + loc, err := start.Location() + if err != nil { + return fmt.Errorf("start upload: %w", err) + } + q := loc.Query() + q.Set("digest", b.Digest) + loc.RawQuery = q.Encode() + + f, entry, err := openEntry(tarPath, b.file) + if err != nil { + return err + } + defer f.Close() + put, err := p.do(ctx, http.MethodPut, loc.String(), entry, b.Size, "application/octet-stream") + if err != nil { + return err + } + put.Body.Close() + if put.StatusCode != http.StatusCreated { + return statusError("upload", put) + } + p.logf("pushed %s (%d bytes)", b.Digest, b.Size) + return nil +} + +func (p *Pusher) putManifest(ctx context.Context, r Ref, mediaType string, body []byte) (string, error) { + resp, err := p.do(ctx, http.MethodPut, p.url(r, "manifests/"+r.Tag), bytes.NewReader(body), int64(len(body)), mediaType) + if err != nil { + return "", err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusCreated { + return "", statusError("put manifest", resp) + } + sum := sha256.Sum256(body) + want := "sha256:" + hex.EncodeToString(sum[:]) + if got := resp.Header.Get("Docker-Content-Digest"); got != "" && got != want { + return "", fmt.Errorf("registry recorded %s for a manifest that hashes to %s", got, want) + } + return want, nil +} + +func (p *Pusher) url(r Ref, tail string) string { + scheme := p.Scheme + if scheme == "" { + scheme = "https" + } + u := url.URL{Scheme: scheme, Host: r.Host, Path: "/v2/" + r.Repo + "/" + tail} + return u.String() +} + +func (p *Pusher) do(ctx context.Context, method, target string, body io.Reader, size int64, contentType string) (*http.Response, error) { + req, err := http.NewRequestWithContext(ctx, method, target, body) + if err != nil { + return nil, err + } + if body != nil { + req.ContentLength = size + } + if contentType != "" { + req.Header.Set("Content-Type", contentType) + } + if p.Username != "" { + req.SetBasicAuth(p.Username, p.Password) + } + c := p.Client + if c == nil { + // No overall timeout: a modpack layer can take minutes on a slow disk, and + // the Job's activeDeadlineSeconds is the real bound. + c = &http.Client{Transport: &http.Transport{ResponseHeaderTimeout: 2 * time.Minute}} + } + return c.Do(req) +} + +// 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. +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++ { + if err = fn(); err == nil { + return nil + } + var se *StatusError + 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 { + p.logf("retrying after: %v", err) + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(time.Duration(i+1) * time.Second): + } + } + } + return err +} + +func (p *Pusher) logf(format string, args ...any) { + if p.Log != nil { + fmt.Fprintf(p.Log, format+"\n", args...) + } +} + +// StatusError is a registry answer the push could not proceed past. +type StatusError struct { + Op string + Code int + Body string +} + +func (e *StatusError) Error() string { + return fmt.Sprintf("%s: registry answered %d: %s", e.Op, e.Code, e.Body) +} + +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))} +} + +// openEntry returns a reader positioned at the named tar entry. The caller closes +// the returned file. archive/tar seeks past the entries it skips, so reopening per +// blob costs no extra reads of the layers in between. +func openEntry(tarPath, name string) (*os.File, io.Reader, error) { + f, err := os.Open(tarPath) + if err != nil { + return nil, nil, err + } + tr := tar.NewReader(f) + for { + h, err := tr.Next() + if err == io.EOF { + f.Close() + return nil, nil, fmt.Errorf("imagepush: %s has no entry %q", tarPath, name) + } + if err != nil { + f.Close() + return nil, nil, fmt.Errorf("imagepush: reading %s: %w", tarPath, err) + } + if h.Name == name || strings.TrimPrefix(h.Name, "./") == name { + if h.Typeflag != tar.TypeReg { + f.Close() + return nil, nil, fmt.Errorf("imagepush: entry %q is not a regular file", name) + } + return f, tr, nil + } + } +} + +func readEntry(tarPath, name string, limit int64) ([]byte, error) { + f, r, err := openEntry(tarPath, name) + if err != nil { + return nil, err + } + defer f.Close() + b, err := io.ReadAll(io.LimitReader(r, limit+1)) + if err != nil { + return nil, err + } + if int64(len(b)) > limit { + return nil, fmt.Errorf("imagepush: %s exceeds %d bytes", name, limit) + } + return b, nil +} + +// digestEntry hashes one entry and returns its descriptor plus its first bytes, +// from which the caller infers the compression. +func digestEntry(tarPath, name string) (descriptor, []byte, error) { + f, r, err := openEntry(tarPath, name) + if err != nil { + return descriptor{}, nil, err + } + defer f.Close() + h := sha256.New() + var magic bytes.Buffer + n, err := io.Copy(io.MultiWriter(h, &prefixWriter{buf: &magic, max: 4}), r) + if err != nil { + return descriptor{}, nil, fmt.Errorf("imagepush: hashing %s: %w", name, err) + } + return descriptor{Size: n, Digest: "sha256:" + hex.EncodeToString(h.Sum(nil)), file: name}, magic.Bytes(), nil +} + +// prefixWriter keeps the first max bytes written to it. +type prefixWriter struct { + buf *bytes.Buffer + max int +} + +func (w *prefixWriter) Write(p []byte) (int, error) { + if room := w.max - w.buf.Len(); room > 0 { + if room > len(p) { + room = len(p) + } + w.buf.Write(p[:room]) + } + return len(p), nil +} diff --git a/internal/imagepush/push_test.go b/internal/imagepush/push_test.go new file mode 100644 index 0000000..69ab804 --- /dev/null +++ b/internal/imagepush/push_test.go @@ -0,0 +1,261 @@ +package imagepush + +import ( + "archive/tar" + "bytes" + "compress/gzip" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/http/httptest" + "net/url" + "os" + "path/filepath" + "strings" + "sync" + "testing" + + "felis.lolicon.best/internal/registrygate" +) + +// fakeRegistry is the distribution v2 subset the pusher speaks, backed by maps. It +// checks what a real registry checks: an upload's bytes hash to the digest the +// client claims, and a manifest references only blobs the repository holds. +type fakeRegistry struct { + mu sync.Mutex + blobs map[string][]byte // repo@digest -> bytes + manifests map[string][]byte // repo:tag -> manifest + types map[string]string // repo:tag -> content type + uploads int + puts int +} + +func newFakeRegistry() *fakeRegistry { + return &fakeRegistry{blobs: map[string][]byte{}, manifests: map[string][]byte{}, types: map[string]string{}} +} + +func (f *fakeRegistry) ServeHTTP(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + repo := registrygate.RepoFromPath(r.URL.Path) + p := r.URL.Path + switch { + case r.Method == http.MethodHead && strings.Contains(p, "/blobs/sha256:"): + d := p[strings.LastIndex(p, "/")+1:] + if _, ok := f.blobs[repo+"@"+d]; ok { + w.WriteHeader(http.StatusOK) + return + } + w.WriteHeader(http.StatusNotFound) + case r.Method == http.MethodPost && strings.HasSuffix(p, "/blobs/uploads/"): + f.uploads++ + w.Header().Set("Location", fmt.Sprintf("/v2/%s/blobs/uploads/u%d?_state=x", repo, f.uploads)) + w.WriteHeader(http.StatusAccepted) + case r.Method == http.MethodPut && strings.Contains(p, "/blobs/uploads/"): + if r.URL.Query().Get("_state") != "x" { + http.Error(w, "upload state lost", http.StatusBadRequest) + return + } + body, _ := io.ReadAll(r.Body) + sum := sha256.Sum256(body) + got := "sha256:" + hex.EncodeToString(sum[:]) + if want := r.URL.Query().Get("digest"); want != got { + http.Error(w, `{"errors":[{"code":"DIGEST_INVALID"}]}`, http.StatusBadRequest) + return + } + f.blobs[repo+"@"+got] = body + f.puts++ + w.WriteHeader(http.StatusCreated) + case r.Method == http.MethodPut && strings.Contains(p, "/manifests/"): + body, _ := io.ReadAll(r.Body) + var m manifest + if err := json.Unmarshal(body, &m); err != nil { + http.Error(w, "bad manifest", http.StatusBadRequest) + return + } + for _, d := range append([]descriptor{m.Config}, m.Layers...) { + b, ok := f.blobs[repo+"@"+d.Digest] + if !ok || int64(len(b)) != d.Size { + http.Error(w, `{"errors":[{"code":"MANIFEST_BLOB_UNKNOWN"}]}`, http.StatusBadRequest) + return + } + } + tag := p[strings.LastIndex(p, "/")+1:] + f.manifests[repo+":"+tag] = body + f.types[repo+":"+tag] = r.Header.Get("Content-Type") + sum := sha256.Sum256(body) + w.Header().Set("Docker-Content-Digest", "sha256:"+hex.EncodeToString(sum[:])) + w.WriteHeader(http.StatusCreated) + default: + http.Error(w, "unexpected "+r.Method+" "+p, http.StatusNotImplemented) + } +} + +// writeTarball writes the layout Kaniko's --tar-path produces (go-containerregistry +// tarball.Write): the config under its digest, gzip layers as .tar.gz, and a +// manifest.json tying them to the destination tag. +func writeTarball(t *testing.T, ref string, layers ...[]byte) string { + t.Helper() + cfg := []byte(`{"architecture":"arm64","os":"linux","rootfs":{"type":"layers","diff_ids":[]}}`) + cfgSum := sha256.Sum256(cfg) + cfgName := "sha256:" + hex.EncodeToString(cfgSum[:]) + + var buf bytes.Buffer + tw := tar.NewWriter(&buf) + add := func(name string, b []byte) { + if err := tw.WriteHeader(&tar.Header{Name: name, Mode: 0o644, Size: int64(len(b)), Typeflag: tar.TypeReg}); err != nil { + t.Fatal(err) + } + if _, err := tw.Write(b); err != nil { + t.Fatal(err) + } + } + add(cfgName, cfg) + var names []string + for _, l := range layers { + var gz bytes.Buffer + zw := gzip.NewWriter(&gz) + zw.Write(l) + zw.Close() + sum := sha256.Sum256(gz.Bytes()) + name := hex.EncodeToString(sum[:]) + ".tar.gz" + add(name, gz.Bytes()) + names = append(names, name) + } + mj, _ := json.Marshal([]tarManifest{{Config: cfgName, RepoTags: []string{ref}, Layers: names}}) + add("manifest.json", mj) + tw.Close() + + path := filepath.Join(t.TempDir(), "image.tar") + if err := os.WriteFile(path, buf.Bytes(), 0o644); err != nil { + t.Fatal(err) + } + return path +} + +func startStack(t *testing.T) (*fakeRegistry, string) { + t.Helper() + reg := newFakeRegistry() + upstream := httptest.NewServer(reg) + t.Cleanup(upstream.Close) + u, _ := url.Parse(upstream.URL) + gate := httptest.NewServer(registrygate.New(u, map[string]string{ + registrygate.PrincipalBuild: "build-secret", + registrygate.PrincipalPlatform: "plat-secret", + }, nil)) + t.Cleanup(gate.Close) + return reg, strings.TrimPrefix(gate.URL, "http://") +} + +func TestPushThroughTheGate(t *testing.T) { + reg, host := startStack(t) + ref := host + "/user-uploads/sub-1:latest" + tarPath := writeTarball(t, ref, []byte("layer one"), []byte("layer two")) + p := &Pusher{Scheme: "http", Username: registrygate.PrincipalBuild, Password: "build-secret", Attempts: 1} + + digest, err := p.Push(context.Background(), tarPath, ref) + if err != nil { + t.Fatalf("push: %v", err) + } + body := reg.manifests["user-uploads/sub-1:latest"] + if body == nil { + t.Fatal("no manifest recorded") + } + sum := sha256.Sum256(body) + if digest != "sha256:"+hex.EncodeToString(sum[:]) { + t.Fatalf("returned digest %s does not name the stored manifest", digest) + } + if ct := reg.types["user-uploads/sub-1:latest"]; ct != mediaDockerManifest { + t.Fatalf("manifest content type %q, want %q", ct, mediaDockerManifest) + } + var m manifest + json.Unmarshal(body, &m) + if len(m.Layers) != 2 || m.Layers[0].MediaType != mediaDockerLayerGz || m.Config.MediaType != mediaDockerConfig { + t.Fatalf("manifest shape: %+v", m) + } + if reg.puts != 3 { + t.Fatalf("uploaded %d blobs, want 3 (config + 2 layers)", reg.puts) + } + + // A second push of the same image re-sends only the manifest. + if _, err := p.Push(context.Background(), tarPath, ref); err != nil { + t.Fatalf("re-push: %v", err) + } + if reg.puts != 3 { + t.Fatalf("re-push uploaded %d blobs in total, want the 3 from the first push", reg.puts) + } +} + +func TestPushIntoAReservedRepoIsRefusedWithoutRetry(t *testing.T) { + reg, host := startStack(t) + ref := host + "/felis/felis:v0.1.0" + tarPath := writeTarball(t, ref, []byte("evil")) + p := &Pusher{Scheme: "http", Username: registrygate.PrincipalBuild, Password: "build-secret"} + + _, err := p.Push(context.Background(), tarPath, ref) + var se *StatusError + if !errors.As(err, &se) || se.Code != http.StatusForbidden { + t.Fatalf("push into felis/ = %v, want a 403 StatusError", err) + } + if reg.uploads != 0 || len(reg.manifests) != 0 { + t.Fatalf("a refused push reached the registry: uploads=%d manifests=%d", reg.uploads, len(reg.manifests)) + } +} + +func TestPushWithoutCredentialsIsRefused(t *testing.T) { + reg, host := startStack(t) + ref := host + "/user-uploads/sub-2:latest" + tarPath := writeTarball(t, ref, []byte("x")) + p := &Pusher{Scheme: "http"} + if _, err := p.Push(context.Background(), tarPath, ref); err == nil { + t.Fatal("anonymous push succeeded") + } + if len(reg.manifests) != 0 { + t.Fatal("anonymous push stored a manifest") + } +} + +func TestTamperedTarballIsRefused(t *testing.T) { + _, host := startStack(t) + ref := host + "/user-uploads/sub-3:latest" + dir := t.TempDir() + // A config whose name is not its digest. + var buf bytes.Buffer + tw := tar.NewWriter(&buf) + cfg := []byte(`{}`) + tw.WriteHeader(&tar.Header{Name: "sha256:" + strings.Repeat("0", 64), Mode: 0o644, Size: int64(len(cfg))}) + tw.Write(cfg) + mj, _ := json.Marshal([]tarManifest{{Config: "sha256:" + strings.Repeat("0", 64), Layers: []string{"l.tar.gz"}}}) + tw.WriteHeader(&tar.Header{Name: "manifest.json", Mode: 0o644, Size: int64(len(mj))}) + tw.Write(mj) + tw.Close() + path := filepath.Join(dir, "bad.tar") + os.WriteFile(path, buf.Bytes(), 0o644) + p := &Pusher{Scheme: "http", Username: registrygate.PrincipalBuild, Password: "build-secret"} + if _, err := p.Push(context.Background(), path, ref); err == nil || !strings.Contains(err.Error(), "hashes to") { + t.Fatalf("tampered config = %v, want a digest mismatch", err) + } +} + +func TestParseRef(t *testing.T) { + for in, want := range map[string]Ref{ + "registry.felis.svc:5000/user-uploads/s:latest": {"registry.felis.svc:5000", "user-uploads/s", "latest"}, + "127.0.0.1:5000/a/b/c:1.2": {"127.0.0.1:5000", "a/b/c", "1.2"}, + "registry.felis.svc:5000/untagged": {"registry.felis.svc:5000", "untagged", "latest"}, + } { + got, err := ParseRef(in) + if err != nil || got != want { + t.Errorf("ParseRef(%q) = %+v, %v; want %+v", in, got, err, want) + } + } + for _, bad := range []string{"", "busybox", "library/busybox:1", "r.io/x@sha256:00", "r.io/", "r.io/x:"} { + if _, err := ParseRef(bad); err == nil { + t.Errorf("ParseRef(%q) accepted", bad) + } + } +} diff --git a/internal/naming/naming.go b/internal/naming/naming.go index 86f14f1..9f91b66 100644 --- a/internal/naming/naming.go +++ b/internal/naming/naming.go @@ -71,6 +71,18 @@ const ( EnvAPIBaseURL = "FELIS_API_BASE_URL" ) +// Registry write credentials (internal/registrygate). The registry namespace holds +// RegistryAuthSecretName with one key per principal (platform, build), mounted into +// the gate sidecar. The build namespace holds RegistryPushSecretName with the build +// principal's username/password, read only by a build Job's push container. Both +// are provisioned out-of-band by deploy/bootstrap.sh. +const ( + RegistryAuthSecretName = "felis-registry-auth" + RegistryPushSecretName = "felis-registry-push" + RegistryPushUsernameKey = "username" + RegistryPushPasswordKey = "password" +) + // ForwardingSecretName / ForwardingSecretKey name the Velocity modern player-info // forwarding secret — the shared HMAC key the proxy signs each login handshake with // and every backend verifies. It is what makes a backend's idea of "who is this diff --git a/internal/platform/workloads.go b/internal/platform/workloads.go index fb5b903..5c5ad6c 100644 --- a/internal/platform/workloads.go +++ b/internal/platform/workloads.go @@ -84,6 +84,11 @@ const ( // pull lands here. Loopback-only is deliberate — the registry serves plain HTTP // and must never be reachable off the node. registryLoopbackHost = "127.0.0.1" + // registryGateName is the write-authorization sidecar in the registry pod, and + // registryAuth* mount its per-principal token files (naming.RegistryAuthSecretName). + registryGateName = "registry-gate" + registryAuthVolume = "registry-auth" + registryAuthMountPath = "/etc/felis-registry-auth" configVolume = "config" tmpVolume = "tmp" @@ -676,24 +681,51 @@ func controlPlaneDeployment(p Params, sa string, container corev1.Container, vol // Build Jobs push to registry..svc:, the destination the build // egress NetworkPolicy opens — so this Deployment+Service+PVC make that policy // target real. The registry never calls the K8s API, so its token auto-mount is -// disabled (matching the weak build/restore SA hygiene), and REGISTRY_HTTP_ADDR -// pins its listen port to the Service port instead of trusting the image default. +// disabled (matching the weak build/restore SA hygiene). // -// The container port also carries a loopback hostPort (registryLoopbackHost): it is +// The pod has two 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. +// +// 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 // deploy/bootstrap.sh writes a registries.yaml mirror rewriting // registry..svc: onto http://127.0.0.1:, and that request // arrives at this hostPort — which is what lets kubelet re-pull a garbage-collected -// platform image without an operator re-import. +// platform image without an operator re-import. The gate runs the felis image, so +// the installer pins that image in containerd: the registry must never depend on +// pulling its own gate from itself. func registryDeployment(p Params) *appsv1.Deployment { p = p.withDefaults() labels := registryLabels() + upstreamPort := registryUpstreamPort(p) - container := corev1.Container{ + registry := corev1.Container{ Name: registryName, Image: p.RegistryImage, Env: []corev1.EnvVar{ - {Name: "REGISTRY_HTTP_ADDR", Value: fmt.Sprintf(":%d", p.RegistryPort)}, + {Name: "REGISTRY_HTTP_ADDR", Value: fmt.Sprintf("127.0.0.1:%d", upstreamPort)}, + }, + VolumeMounts: []corev1.VolumeMount{ + {Name: registryVolume, MountPath: registryDataPath}, + {Name: tmpVolume, MountPath: "/tmp"}, + }, + Resources: registryResources(), + SecurityContext: hardenedContainerSecurityContext(), + } + + gate := corev1.Container{ + Name: registryGateName, + Image: p.FelisImage, + Command: []string{felisBinaryPath, "registry-gate"}, + Args: []string{ + fmt.Sprintf("--listen=:%d", p.RegistryPort), + fmt.Sprintf("--upstream=http://127.0.0.1:%d", upstreamPort), + "--auth-dir=" + registryAuthMountPath, }, Ports: []corev1.ContainerPort{ { @@ -704,32 +736,31 @@ func registryDeployment(p Params) *appsv1.Deployment { }, }, VolumeMounts: []corev1.VolumeMount{ - {Name: registryVolume, MountPath: registryDataPath}, - {Name: tmpVolume, MountPath: "/tmp"}, + {Name: registryAuthVolume, MountPath: registryAuthMountPath, ReadOnly: true}, }, - // Distribution serves GET /v2/ (200 = app + storage healthy) for any - // client, so both probes reuse it: without them a registry whose storage - // backend broke would stay "Running" and every build push would fail with - // nothing red in the Deployment status. + // /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 + // "Running" one every push fails against. Liveness checks the gate alone: a + // failing registry is not fixed by restarting the gate in front of it. ReadinessProbe: &corev1.Probe{ ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{ - Path: "/v2/", Port: intstr.FromString(registryName), + Path: "/healthz", Port: intstr.FromString(registryName), }}, InitialDelaySeconds: 5, PeriodSeconds: 10, - TimeoutSeconds: 3, + TimeoutSeconds: 5, FailureThreshold: 3, }, LivenessProbe: &corev1.Probe{ ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{ - Path: "/v2/", Port: intstr.FromString(registryName), + Path: "/livez", Port: intstr.FromString(registryName), }}, InitialDelaySeconds: 10, PeriodSeconds: 10, TimeoutSeconds: 3, FailureThreshold: 3, }, - Resources: registryResources(), + Resources: registryGateResources(), SecurityContext: hardenedContainerSecurityContext(), } @@ -746,7 +777,7 @@ func registryDeployment(p Params) *appsv1.Deployment { AutomountServiceAccountToken: boolPtr(false), PriorityClassName: controlPlanePriorityName, SecurityContext: hardenedPodSecurityContext(), - Containers: []corev1.Container{container}, + Containers: []corev1.Container{registry, gate}, Volumes: []corev1.Volume{ { Name: registryVolume, @@ -755,6 +786,17 @@ func registryDeployment(p Params) *appsv1.Deployment { }, }, {Name: tmpVolume, VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{}}}, + { + Name: registryAuthVolume, + VolumeSource: corev1.VolumeSource{Secret: &corev1.SecretVolumeSource{ + SecretName: naming.RegistryAuthSecretName, + // Optional: without the Secret the gate starts with no + // principals, so pulls keep working and every write is + // refused — never a registry that cannot start. + Optional: boolPtr(true), + DefaultMode: int32Ptr(0o440), + }}, + }, }, }, }, @@ -762,6 +804,25 @@ func registryDeployment(p Params) *appsv1.Deployment { } } +// registryUpstreamPort is the loopback port registry:2 listens on behind the gate: +// the next port after the public one. +func registryUpstreamPort(p Params) int32 { return p.RegistryPort + 1 } + +// registryGateResources sizes the gate sidecar: a streaming reverse proxy that +// holds no layer in memory. +func registryGateResources() corev1.ResourceRequirements { + return corev1.ResourceRequirements{ + Requests: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("20m"), + corev1.ResourceMemory: resource.MustParse("32Mi"), + }, + Limits: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("500m"), + corev1.ResourceMemory: resource.MustParse("128Mi"), + }, + } +} + // registryService renders the ClusterIP Service that gives the registry its // pinned DNS name registry..svc: — hardcoded across the build // subsystem and config. Its selector matches the registry pod labels; because diff --git a/internal/platform/workloads_test.go b/internal/platform/workloads_test.go index b3056ab..2a202ea 100644 --- a/internal/platform/workloads_test.go +++ b/internal/platform/workloads_test.go @@ -23,6 +23,20 @@ func podSpec(t *testing.T, d *appsv1.Deployment) (corev1.PodSpec, corev1.Contain return ps, ps.Containers[0] } +// namedContainer returns the pod template of a Deployment and its container +// called name, failing if there is none. +func namedContainer(t *testing.T, d *appsv1.Deployment, name string) (corev1.PodSpec, corev1.Container) { + t.Helper() + ps := d.Spec.Template.Spec + for _, c := range ps.Containers { + if c.Name == name { + return ps, c + } + } + t.Fatalf("%s: no container %q", d.Name, name) + return ps, corev1.Container{} +} + // rconPeerSelector returns the podSelector of the allow-rcon NetworkPolicy peer, // compiled into the same labels.Selector K8s evaluates at runtime. This is the // real gate: a pod reaches server RCON iff its labels Match this selector. @@ -121,7 +135,7 @@ func TestControlPlanePods_SatisfyRConPeer(t *testing.T) { func TestControlPlanePods_Hardened(t *testing.T) { p := testParams() for _, d := range []*appsv1.Deployment{APIDeployment(p), OperatorDeployment(p), registryDeployment(p)} { - ps, c := podSpec(t, d) + ps := d.Spec.Template.Spec if ps.SecurityContext == nil || ps.SecurityContext.RunAsNonRoot == nil || !*ps.SecurityContext.RunAsNonRoot { t.Errorf("%s: pod must set runAsNonRoot=true", d.Name) @@ -129,21 +143,23 @@ func TestControlPlanePods_Hardened(t *testing.T) { if ps.SecurityContext == nil || ps.SecurityContext.RunAsUser == nil || *ps.SecurityContext.RunAsUser != nonRootUID { t.Errorf("%s: pod runAsUser must be %d", d.Name, nonRootUID) } - sc := c.SecurityContext - if sc == nil { - t.Fatalf("%s: container has no SecurityContext", d.Name) - } - if sc.Privileged == nil || *sc.Privileged { - t.Errorf("%s: container must not be privileged", d.Name) - } - if sc.AllowPrivilegeEscalation == nil || *sc.AllowPrivilegeEscalation { - t.Errorf("%s: container must set allowPrivilegeEscalation=false", d.Name) - } - if sc.ReadOnlyRootFilesystem == nil || !*sc.ReadOnlyRootFilesystem { - t.Errorf("%s: container must set readOnlyRootFilesystem=true", d.Name) - } - if sc.Capabilities == nil || len(sc.Capabilities.Drop) == 0 || sc.Capabilities.Drop[0] != "ALL" { - t.Errorf("%s: container must drop ALL capabilities", d.Name) + for _, c := range ps.Containers { + sc := c.SecurityContext + if sc == nil { + t.Fatalf("%s/%s: container has no SecurityContext", d.Name, c.Name) + } + if sc.Privileged == nil || *sc.Privileged { + t.Errorf("%s/%s: container must not be privileged", d.Name, c.Name) + } + if sc.AllowPrivilegeEscalation == nil || *sc.AllowPrivilegeEscalation { + t.Errorf("%s/%s: container must set allowPrivilegeEscalation=false", d.Name, c.Name) + } + if sc.ReadOnlyRootFilesystem == nil || !*sc.ReadOnlyRootFilesystem { + t.Errorf("%s/%s: container must set readOnlyRootFilesystem=true", d.Name, c.Name) + } + if sc.Capabilities == nil || len(sc.Capabilities.Drop) == 0 || sc.Capabilities.Drop[0] != "ALL" { + t.Errorf("%s/%s: container must drop ALL capabilities", d.Name, c.Name) + } } } } @@ -416,30 +432,78 @@ func TestRegistry_DeploymentServicePVC(t *testing.T) { svc := registryService(p) pvc := registryPVC(p) - ps, c := podSpec(t, dep) + ps, c := namedContainer(t, dep, registryName) + if len(ps.Containers) != 2 { + t.Fatalf("registry pod containers = %d, want registry + gate", len(ps.Containers)) + } if c.Image != defaultRegistryImage { t.Errorf("registry image = %q, want default %q", c.Image, defaultRegistryImage) } - // REGISTRY_HTTP_ADDR pins the listen port to the Service port rather than - // trusting the image default. - if v := envValue(c.Env, "REGISTRY_HTTP_ADDR"); v != ":5000" { - t.Errorf("REGISTRY_HTTP_ADDR = %q, want :5000", v) + // registry:2 itself listens on loopback only and exposes nothing: every request + // from outside the pod passes the gate, which is what makes writes authorized. + if v := envValue(c.Env, "REGISTRY_HTTP_ADDR"); v != fmt.Sprintf("127.0.0.1:%d", p.RegistryPort+1) { + t.Errorf("REGISTRY_HTTP_ADDR = %q, want loopback 127.0.0.1:%d", v, p.RegistryPort+1) } - // 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 - // exposed (the registry serves plain HTTP). - if len(c.Ports) != 1 { - t.Fatalf("registry container ports = %+v, want exactly 1", c.Ports) - } - if p0 := c.Ports[0]; p0.ContainerPort != p.RegistryPort || p0.HostPort != p.RegistryPort || p0.HostIP != registryLoopbackHost { - t.Errorf("registry port = %+v, want container/host port %d bound to %s", p0, p.RegistryPort, registryLoopbackHost) + if len(c.Ports) != 0 { + t.Errorf("registry container ports = %+v, want none (the gate owns the port)", c.Ports) } // 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 { t.Errorf("registry memory limit = %s, want >= 2Gi (audit #46: 256Mi OOM-killed on a 475MB-layer push)", mem.String()) } + + // The gate: the platform image, forwarding to the loopback registry, reading + // the per-principal tokens from the optional Secret. + _, gate := namedContainer(t, dep, registryGateName) + if gate.Image != p.FelisImage { + t.Errorf("gate image = %q, want the platform image %q", gate.Image, p.FelisImage) + } + if len(gate.Command) != 2 || gate.Command[1] != "registry-gate" { + t.Errorf("gate command = %v, want felis registry-gate", gate.Command) + } + for _, want := range []string{ + fmt.Sprintf("--listen=:%d", p.RegistryPort), + fmt.Sprintf("--upstream=http://127.0.0.1:%d", p.RegistryPort+1), + "--auth-dir=" + registryAuthMountPath, + } { + if !contains(gate.Args, want) { + t.Errorf("gate args = %v, want %s", gate.Args, want) + } + } + // 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 + // exposed (the registry serves plain HTTP). + if len(gate.Ports) != 1 { + t.Fatalf("gate ports = %+v, want exactly 1", gate.Ports) + } + if p0 := gate.Ports[0]; p0.ContainerPort != p.RegistryPort || p0.HostPort != p.RegistryPort || p0.HostIP != registryLoopbackHost { + t.Errorf("gate port = %+v, want container/host port %d bound to %s", p0, p.RegistryPort, registryLoopbackHost) + } + auth := volumeByName(ps.Volumes, registryAuthVolume) + if auth == nil || auth.Secret == nil || auth.Secret.SecretName != "felis-registry-auth" { + t.Fatalf("gate token volume = %#v, want Secret felis-registry-auth", auth) + } + // Optional: a missing Secret must degrade to "reads only", never to a registry + // pod stuck in ContainerCreating that every game pull depends on. + if auth.Secret.Optional == nil || !*auth.Secret.Optional { + t.Error("the registry-auth Secret volume must be optional") + } + mounted := false + for _, m := range gate.VolumeMounts { + if m.Name == registryAuthVolume && m.ReadOnly { + mounted = true + } + } + for _, m := range c.VolumeMounts { + if m.Name == registryAuthVolume { + t.Error("registry:2 must not mount the write tokens") + } + } + if !mounted { + t.Errorf("gate must mount the tokens read-only, mounts=%v", gate.VolumeMounts) + } // Registry never calls the K8s API ⇒ no auto-mounted token. if ps.AutomountServiceAccountToken == nil || *ps.AutomountServiceAccountToken { t.Error("registry pod must set automountServiceAccountToken=false") @@ -492,7 +556,7 @@ func TestWorkloads_DeploymentsCarryProbes(t *testing.T) { }{ {APIDeployment(p), "/readyz", "/healthz", apiInternalPort}, {OperatorDeployment(p), "/readyz", "/healthz", operatorHealthPort}, - {registryDeployment(p), "/v2/", "/v2/", p.RegistryPort}, + {registryDeployment(p), "/healthz", "/livez", p.RegistryPort}, } // resolve maps a probe target (by number or container-port name) to the // declared container port it denotes. @@ -508,7 +572,14 @@ func TestWorkloads_DeploymentsCarryProbes(t *testing.T) { return 0, false } for _, tc := range cases { - _, c := podSpec(t, tc.dep) + var c corev1.Container + if tc.dep.Name == registryName { + // The gate owns the registry port; registry:2 behind it is probed + // through the gate's /healthz. + _, c = namedContainer(t, tc.dep, registryGateName) + } else { + _, c = podSpec(t, tc.dep) + } if c.ReadinessProbe == nil || c.ReadinessProbe.HTTPGet == nil { t.Fatalf("%s: readiness probe missing or not an HTTP GET", tc.dep.Name) } diff --git a/internal/registrygate/gate.go b/internal/registrygate/gate.go new file mode 100644 index 0000000..f5066af --- /dev/null +++ b/internal/registrygate/gate.go @@ -0,0 +1,272 @@ +// Package registrygate is the write-authorization front of the in-cluster image +// registry. registry:2 runs with no auth of its own and listens on the pod's +// loopback only; this gate owns the registry port and forwards to it. +// +// Why a gate rather than registry:2's own htpasswd auth: htpasswd is all-or-nothing +// (every principal may write every repository, and anonymous pulls stop working), +// while the platform needs two distinct writers and anonymous reads: +// +// - reads (GET/HEAD) stay anonymous, because the node's containerd pulls through +// the loopback hostPort, Kaniko pulls FROM images and Trivy pulls its DB mirror, +// 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. +// +// 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 +// felis/felis and take over the control plane on its next pull. +package registrygate + +import ( + "context" + "crypto/subtle" + "encoding/json" + "fmt" + "log/slog" + "net/http" + "net/http/httputil" + "net/url" + "strings" + "time" +) + +// 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" +) + +// 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 +// either would own the platform or blind its own scanner. +var ReservedRepoRoots = []string{"felis", "mirror"} + +// Realm is the basic-auth realm the gate challenges with. +const Realm = "felis-registry" + +// Gate authorizes registry requests and proxies the allowed ones upstream. +type Gate struct { + // Tokens maps a principal to its secret. A principal with an empty or missing + // token cannot authenticate: writes fail closed while reads keep working. + Tokens map[string]string + // Upstream is the loopback registry, e.g. http://127.0.0.1:5001. + Upstream *url.URL + // Log receives one line per refused write. Nil discards. + Log *slog.Logger + + proxy *httputil.ReverseProxy + health *http.Client +} + +// New builds a Gate for upstream. +func New(upstream *url.URL, tokens map[string]string, log *slog.Logger) *Gate { + g := &Gate{Tokens: tokens, Upstream: upstream, Log: log} + rp := httputil.NewSingleHostReverseProxy(upstream) + base := rp.Director + rp.Director = func(r *http.Request) { + base(r) + // The registry has no auth of its own; the credential stops here. + r.Header.Del("Authorization") + } + // Blob uploads and pulls are streamed; flush as bytes arrive so a large layer + // pull is not buffered in the gate. + rp.FlushInterval = -1 + g.proxy = rp + g.health = &http.Client{Timeout: 3 * time.Second} + return g +} + +// ServeHTTP implements http.Handler. +func (g *Gate) ServeHTTP(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/livez": + w.WriteHeader(http.StatusOK) + return + case "/healthz": + g.serveHealth(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 + } + // The gate and the registry must read the same path, or an authorization check + // on one repository could be spent on another. Refuse every shape that a later + // clean-up or decode could turn into a different path. + if !canonicalPath(r) { + writeError(w, http.StatusBadRequest, "NAME_INVALID", "non-canonical request path") + return + } + + principal, authErr := g.authenticate(r) + if authErr != nil { + challenge(w, "invalid credentials") + return + } + + switch r.Method { + case http.MethodGet, http.MethodHead: + // The API root is where Docker-compatible clients learn which auth scheme + // the registry wants; the daemon sends credentials on later writes only if + // this answer challenged it. Everything else stays anonymous for reads. + if isAPIRoot(r.URL.Path) && principal == "" { + challenge(w, "authentication required") + return + } + g.proxy.ServeHTTP(w, r) + return + case http.MethodPost, http.MethodPut, http.MethodPatch, http.MethodDelete: + default: + writeError(w, http.StatusMethodNotAllowed, "UNSUPPORTED", "method not allowed") + return + } + + if principal == "" { + challenge(w, "authentication required") + return + } + if reason := Authorize(principal, r.Method, r.URL.Path); reason != "" { + if g.Log != nil { + g.Log.Warn("registry write refused", "principal", principal, "method", r.Method, "path", r.URL.Path, "reason", reason) + } + writeError(w, http.StatusForbidden, "DENIED", reason) + return + } + g.proxy.ServeHTTP(w, r) +} + +// authenticate returns the principal the request's basic credentials name, "" for +// an anonymous request, and an error for credentials that do not verify. +func (g *Gate) authenticate(r *http.Request) (string, error) { + if r.Header.Get("Authorization") == "" { + return "", nil + } + user, pass, ok := r.BasicAuth() + if !ok { + return "", fmt.Errorf("malformed authorization") + } + want, known := g.Tokens[user] + if !known || want == "" { + // Still spend a comparison so an unknown user is not faster to refuse. + subtle.ConstantTimeCompare([]byte(pass), []byte("x")) + return "", fmt.Errorf("unknown principal") + } + if subtle.ConstantTimeCompare([]byte(pass), []byte(want)) != 1 { + return "", fmt.Errorf("bad secret") + } + return user, nil +} + +// Authorize decides whether an authenticated principal may send a write request +// for path. It returns "" to allow, or the refusal reason. +func Authorize(principal, method, path string) string { + switch principal { + case PrincipalPlatform: + return "" + case PrincipalBuild: + if method == http.MethodDelete { + return "the build principal may not delete" + } + repo := RepoFromPath(path) + if repo == "" { + return "the build principal may only write to a repository" + } + root, _, _ := strings.Cut(repo, "/") + for _, reserved := range ReservedRepoRoots { + if root == reserved { + return fmt.Sprintf("repository %s/ is reserved for the platform", reserved) + } + } + return "" + default: + return "unknown principal" + } +} + +// 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/, +// /blobs/uploads/[], /tags/list, /referrers/. +// Neither a reference, a digest nor an upload id contains '/', so the repository +// is everything before the fixed tail. +func RepoFromPath(path string) string { + rest, ok := strings.CutPrefix(path, "/v2/") + if !ok { + return "" + } + seg := strings.Split(rest, "/") + n := len(seg) + switch { + case n >= 4 && seg[n-3] == "blobs" && seg[n-2] == "uploads": + return strings.Join(seg[:n-3], "/") + case n >= 3 && (seg[n-2] == "manifests" || seg[n-2] == "blobs" || seg[n-2] == "tags" || seg[n-2] == "referrers"): + return strings.Join(seg[:n-2], "/") + } + return "" +} + +// canonicalPath rejects any request path the upstream could read differently from +// the gate: percent-escapes Go would decode (RawPath set), dot segments a router +// would clean, and empty segments other than the trailing slash of an upload POST. +func canonicalPath(r *http.Request) bool { + if r.URL.RawPath != "" && r.URL.RawPath != r.URL.Path { + return false + } + p := r.URL.Path + if strings.ContainsAny(p, "%\\") { + return false + } + seg := strings.Split(strings.TrimPrefix(p, "/"), "/") + for i, s := range seg { + if s == "." || s == ".." { + return false + } + if s == "" && i != len(seg)-1 { + return false + } + } + return true +} + +func isAPIRoot(p string) bool { return p == "/v2/" || p == "/v2" } + +// serveHealth answers 200 when the upstream registry answers its API root, the +// check registry:2's own probes used before the gate took its port. +func (g *Gate) serveHealth(w http.ResponseWriter, r *http.Request) { + ctx, cancel := context.WithTimeout(r.Context(), 3*time.Second) + defer cancel() + req, err := http.NewRequestWithContext(ctx, http.MethodGet, g.Upstream.JoinPath("/v2/").String(), nil) + if err != nil { + http.Error(w, err.Error(), http.StatusServiceUnavailable) + return + } + resp, err := g.health.Do(req) + if err != nil { + http.Error(w, "upstream: "+err.Error(), http.StatusServiceUnavailable) + return + } + resp.Body.Close() + if resp.StatusCode != http.StatusOK { + http.Error(w, fmt.Sprintf("upstream answered %d", resp.StatusCode), http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) +} + +func challenge(w http.ResponseWriter, msg string) { + w.Header().Set("WWW-Authenticate", `Basic realm="`+Realm+`"`) + writeError(w, http.StatusUnauthorized, "UNAUTHORIZED", msg) +} + +// writeError answers in the registry's own error envelope, which Docker and +// containerd both surface verbatim to the operator. +func writeError(w http.ResponseWriter, status int, code, msg string) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(map[string]any{ + "errors": []map[string]string{{"code": code, "message": msg}}, + }) +} diff --git a/internal/registrygate/gate_test.go b/internal/registrygate/gate_test.go new file mode 100644 index 0000000..5e9e887 --- /dev/null +++ b/internal/registrygate/gate_test.go @@ -0,0 +1,258 @@ +package registrygate + +import ( + "net/http" + "net/http/httptest" + "net/url" + "strings" + "sync" + "testing" +) + +type upstreamLog struct { + mu sync.Mutex + seen []string + auth []string +} + +func (u *upstreamLog) handler() http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + u.mu.Lock() + u.seen = append(u.seen, r.Method+" "+r.URL.RequestURI()) + u.auth = append(u.auth, r.Header.Get("Authorization")) + u.mu.Unlock() + switch r.Method { + case http.MethodPost: + w.Header().Set("Location", "/v2/x/blobs/uploads/abc") + w.WriteHeader(http.StatusAccepted) + case http.MethodPut: + w.WriteHeader(http.StatusCreated) + default: + w.WriteHeader(http.StatusOK) + } + }) +} + +func (u *upstreamLog) count() int { + u.mu.Lock() + defer u.mu.Unlock() + return len(u.seen) +} + +func newGate(t *testing.T) (*httptest.Server, *upstreamLog) { + 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"}, nil) + gs := httptest.NewServer(g) + t.Cleanup(gs.Close) + return gs, up +} + +func do(t *testing.T, srv *httptest.Server, method, path, user, pass string) *http.Response { + t.Helper() + req, err := http.NewRequest(method, srv.URL+path, strings.NewReader("")) + if err != nil { + t.Fatal(err) + } + if user != "" { + req.SetBasicAuth(user, pass) + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + return resp +} + +func TestAnonymousReadsPassButTheAPIRootChallenges(t *testing.T) { + gs, up := newGate(t) + for _, p := range []string{ + "/v2/felis/felis/manifests/v0.1.0", + "/v2/felis/felis/blobs/sha256:abc", + "/v2/mirror/trivy-db/manifests/2", + "/v2/_catalog", + "/v2/user-uploads/s1/tags/list", + } { + for _, m := range []string{http.MethodGet, http.MethodHead} { + if resp := do(t, gs, m, p, "", ""); resp.StatusCode != http.StatusOK { + t.Errorf("%s %s anonymous = %d, want 200", m, p, resp.StatusCode) + } + } + } + // Docker's daemon only sends credentials on a push if the API root challenged + // it, so the root must say 401 + Basic to anonymous callers. + resp := do(t, gs, http.MethodGet, "/v2/", "", "") + if resp.StatusCode != http.StatusUnauthorized || !strings.HasPrefix(resp.Header.Get("WWW-Authenticate"), "Basic ") { + t.Fatalf("anonymous GET /v2/ = %d %q, want 401 Basic challenge", resp.StatusCode, resp.Header.Get("WWW-Authenticate")) + } + if resp := do(t, gs, http.MethodGet, "/v2/", PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusOK { + t.Fatalf("authenticated GET /v2/ = %d, want 200", resp.StatusCode) + } + if resp := do(t, gs, http.MethodGet, "/v2/", PrincipalBuild, "wrong"); resp.StatusCode != http.StatusUnauthorized { + t.Fatalf("GET /v2/ with a bad secret = %d, want 401", resp.StatusCode) + } + _ = up +} + +func TestAnonymousWritesNeverReachTheRegistry(t *testing.T) { + gs, up := newGate(t) + for _, c := range []struct{ method, path string }{ + {http.MethodPost, "/v2/felis/felis/blobs/uploads/"}, + {http.MethodPut, "/v2/felis/felis/manifests/v0.1.0"}, + {http.MethodPatch, "/v2/user-uploads/s1/blobs/uploads/abc"}, + {http.MethodDelete, "/v2/user-uploads/s1/manifests/sha256:abc"}, + } { + resp := do(t, gs, c.method, c.path, "", "") + if resp.StatusCode != http.StatusUnauthorized { + t.Errorf("anonymous %s %s = %d, want 401", c.method, c.path, resp.StatusCode) + } + if resp := do(t, gs, c.method, c.path, PrincipalPlatform, "not-it"); resp.StatusCode != http.StatusUnauthorized { + t.Errorf("bad-secret %s %s = %d, want 401", c.method, c.path, resp.StatusCode) + } + if resp := do(t, gs, c.method, c.path, "intruder", "plat-secret"); resp.StatusCode != http.StatusUnauthorized { + t.Errorf("unknown-user %s %s = %d, want 401", c.method, c.path, resp.StatusCode) + } + } + if n := up.count(); n != 0 { + t.Fatalf("%d refused writes reached the registry: %v", n, up.seen) + } +} + +func TestBuildPrincipalIsFencedOffPlatformRepos(t *testing.T) { + gs, up := newGate(t) + for _, p := range []string{ + "/v2/felis/felis/manifests/v0.1.0", + "/v2/felis/limbo/blobs/uploads/", + "/v2/felis/manifests/latest", + "/v2/mirror/trivy-db/manifests/2", + "/v2/mirror/trivy-java-db/blobs/uploads/abc", + } { + for _, m := range []string{http.MethodPost, http.MethodPut, http.MethodPatch} { + if resp := do(t, gs, m, p, PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusForbidden { + t.Errorf("build %s %s = %d, want 403", m, p, resp.StatusCode) + } + } + } + if resp := do(t, gs, http.MethodDelete, "/v2/user-uploads/s1/manifests/sha256:abc", PrincipalBuild, "build-secret"); resp.StatusCode != http.StatusForbidden { + t.Errorf("build DELETE = %d, want 403", resp.StatusCode) + } + if n := up.count(); n != 0 { + t.Fatalf("%d refused writes reached the registry: %v", n, up.seen) + } + + for _, c := range []struct { + method, path string + want int + }{ + {http.MethodPost, "/v2/user-uploads/s1/blobs/uploads/", http.StatusAccepted}, + {http.MethodPatch, "/v2/user-uploads/s1/blobs/uploads/abc", http.StatusOK}, + {http.MethodPut, "/v2/user-uploads/s1/blobs/uploads/abc?digest=sha256:0", http.StatusCreated}, + {http.MethodPut, "/v2/user-uploads/s1/manifests/latest", http.StatusCreated}, + {http.MethodPut, "/v2/modpacks/pack/manifests/1.0", http.StatusCreated}, + // A repository merely containing "felis" deeper down is not reserved. + {http.MethodPut, "/v2/builds/felis/manifests/1", http.StatusCreated}, + } { + if resp := do(t, gs, c.method, c.path, PrincipalBuild, "build-secret"); resp.StatusCode != c.want { + t.Errorf("build %s %s = %d, want %d", c.method, c.path, resp.StatusCode, c.want) + } + } + for i, a := range up.auth { + if a != "" { + t.Errorf("request %d (%s) reached the registry with an Authorization header", i, up.seen[i]) + } + } +} + +func TestPlatformPrincipalMayWriteAndDeleteAnything(t *testing.T) { + gs, _ := newGate(t) + for _, c := range []struct { + method, path string + want int + }{ + {http.MethodPut, "/v2/felis/felis/manifests/v0.2.0", http.StatusCreated}, + {http.MethodPost, "/v2/mirror/trivy-db/blobs/uploads/", http.StatusAccepted}, + {http.MethodDelete, "/v2/user-uploads/s1/manifests/sha256:abc", http.StatusOK}, + } { + if resp := do(t, gs, c.method, c.path, PrincipalPlatform, "plat-secret"); resp.StatusCode != c.want { + t.Errorf("platform %s %s = %d, want %d", c.method, c.path, resp.StatusCode, c.want) + } + } +} + +func TestPathTricksAreRefusedBeforeAuthorization(t *testing.T) { + gs, up := newGate(t) + for _, p := range []string{ + "/v2/user-uploads/../felis/felis/manifests/v1", + "/v2/user-uploads/./x/manifests/v1", + "/v2/user-uploads%2F..%2Ffelis/felis/manifests/v1", + "/v2/user-uploads//felis/manifests/v1", + } { + req, _ := http.NewRequest(http.MethodPut, gs.URL, nil) + req.URL.Opaque = p // send the path exactly as written, no client-side cleaning + req.SetBasicAuth(PrincipalBuild, "build-secret") + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + if resp.StatusCode != http.StatusBadRequest { + t.Errorf("PUT %s = %d, want 400", p, resp.StatusCode) + } + } + if n := up.count(); n != 0 { + t.Fatalf("%d crafted paths reached the registry: %v", n, up.seen) + } +} + +func TestMissingTokenFailsClosed(t *testing.T) { + up := &upstreamLog{} + upSrv := httptest.NewServer(up.handler()) + defer upSrv.Close() + target, _ := url.Parse(upSrv.URL) + gs := httptest.NewServer(New(target, map[string]string{PrincipalBuild: ""}, nil)) + defer gs.Close() + if resp := do(t, gs, http.MethodPut, "/v2/user-uploads/s1/manifests/latest", PrincipalBuild, ""); resp.StatusCode != http.StatusUnauthorized { + t.Fatalf("empty configured token accepted an empty password: %d", resp.StatusCode) + } + if resp := do(t, gs, http.MethodGet, "/v2/user-uploads/s1/manifests/latest", "", ""); resp.StatusCode != http.StatusOK { + t.Fatalf("reads must keep working without tokens: %d", resp.StatusCode) + } +} + +func TestHealth(t *testing.T) { + gs, _ := newGate(t) + if resp := do(t, gs, http.MethodGet, "/healthz", "", ""); resp.StatusCode != http.StatusOK { + t.Fatalf("healthz = %d", resp.StatusCode) + } + dead, _ := url.Parse("http://127.0.0.1:1") + ds := httptest.NewServer(New(dead, nil, nil)) + defer ds.Close() + if resp := do(t, ds, http.MethodGet, "/healthz", "", ""); resp.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("healthz with a dead upstream = %d, want 503", resp.StatusCode) + } + if resp := do(t, ds, http.MethodGet, "/livez", "", ""); resp.StatusCode != http.StatusOK { + t.Fatalf("livez = %d", resp.StatusCode) + } +} + +func TestRepoFromPath(t *testing.T) { + for path, want := range map[string]string{ + "/v2/": "", + "/v2/_catalog": "", + "/v2/a/manifests/latest": "a", + "/v2/a/b/c/manifests/sha256:0": "a/b/c", + "/v2/a/blobs/sha256:0": "a", + "/v2/a/b/blobs/uploads/": "a/b", + "/v2/a/b/blobs/uploads/uuid-1": "a/b", + "/v2/a/tags/list": "a", + "/v2/user-uploads/blobs/manifests/one": "user-uploads/blobs", + } { + if got := RepoFromPath(path); got != want { + t.Errorf("RepoFromPath(%q) = %q, want %q", path, got, want) + } + } +}