fix(images): 服务器镜像在创建时固定到仓库 digest,更换镜像需确认备份,安装器重建前先固定旧服并推送不可变版本标签

This commit is contained in:
Lemon-miaow committed 2026-09-24 16:57:38 +08:00
1 parent 857b6da886
commit 215bfd78d7
29 files changed
+1149 -29

No files matched your search

+5
View File
@@ -33,6 +33,11 @@ type API struct {
// still exercised even before the subsystem is wired in.
Builder ImageBuilder
// Images pins a whitelisted image ref to the digest it names when a server is
// created or its image is changed (internal/imagepin), so a later push over
// the same tag never reaches an existing world. Nil stores refs as given.
Images ImagePinner
// Console is the synchronous RCON write channel (spec §8 写=RCON). It is
// wired in production (cmd/felis); a nil Console makes the command route report
// 503 rather than panic, so the ownership boundary is still exercised in tests.
+13
View File
@@ -64,6 +64,19 @@ func (a *API) audit(r *http.Request, action, target string) {
a.auditEntry(r, e)
}
// auditImageChange records a confirmed image change as server.patch with the
// image it replaced and the one it set, so the audit log alone can say which
// build a world ran before it was moved.
func (a *API) auditImageChange(r *http.Request, server, from, to string) {
p := principalFromContext(r.Context())
e := AuditEntry{Actor: auditActor(p), Action: "server.patch", ServerName: server}
if p != nil {
e.ActorUserID = p.UserID
}
e.Payload = auditPayload(map[string]any{"image_from": from, "image_to": to})
a.auditEntry(r, e)
}
// auditAccount records an action a pre-session door took for the account it
// resolved (u nil: none was). The username is the actor: the door has not yet
// proven anything about the address.
+1 -1
View File
@@ -75,7 +75,7 @@ func TestPatchServerAutostartPolicy(t *testing.T) {
func TestPatchServerImageReAdmitted(t *testing.T) {
api, _, cl, _ := newPatchAPI()
w := patchSurvival(api, `{"image":"`+admittedImage+`"}`)
w := patchSurvival(api, `{"image":"`+admittedImage+`","confirmImageChange":true}`)
if w.Code != http.StatusOK {
t.Fatalf("code = %d, want 200 (%s)", w.Code, w.Body.String())
}
+50 -9
View File
@@ -375,6 +375,13 @@ func (a *API) handleCreateServer(w http.ResponseWriter, r *http.Request) {
"image %q is not on the whitelist", body.Image))
return
}
// The spec keeps the digest the tag names now, not the tag: the world is
// created on this build and stays on it until an admin changes the image.
image, err := a.pinImage(r.Context(), body.Image)
if err != nil {
writeError(w, r, err)
return
}
// Quota is intentionally NOT enforced here. §15 creates an UNOWNED server
// (owner_id NULL); the per-user quota is charged at claim time (spec §9.3 /
@@ -432,7 +439,7 @@ func (a *API) handleCreateServer(w http.ResponseWriter, r *http.Request) {
Name: body.Name,
Subdomain: body.Subdomain,
DisplayName: body.DisplayName,
Image: body.Image,
Image: image,
JavaMemory: javaMemory,
StorageSize: storage,
AutostartPolicy: policy,
@@ -613,11 +620,16 @@ func parsePositiveQuantity(s, field string) (resource.Quantity, error) {
// the dual-write routing identity (name is the immutable object key; subdomain
// would desync the Postgres alias) nor for the world PVC size (see below).
type patchServerRequest struct {
DisplayName *string `json:"displayName,omitempty"`
AutostartPolicy *string `json:"autostartPolicy,omitempty"`
Image *string `json:"image,omitempty"`
Memory *string `json:"memory,omitempty"`
Resources *resourceRequest `json:"resources,omitempty"`
DisplayName *string `json:"displayName,omitempty"`
AutostartPolicy *string `json:"autostartPolicy,omitempty"`
Image *string `json:"image,omitempty"`
// ConfirmImageChange acknowledges that a new image opens the world with
// whatever Minecraft version it carries. Chunks a newer version has upgraded
// cannot be read by the older one again, so without it an image change that
// would actually move the server is refused (image_change_unconfirmed).
ConfirmImageChange bool `json:"confirmImageChange,omitempty"`
Memory *string `json:"memory,omitempty"`
Resources *resourceRequest `json:"resources,omitempty"`
// IdleStopSeconds sets idle auto-stop: 0 turns it off, otherwise the server
// stops after that many seconds with nobody online (60 to 86400).
IdleStopSeconds *int32 `json:"idleStopSeconds,omitempty"`
@@ -672,6 +684,8 @@ func (a *API) handlePatchServer(w http.ResponseWriter, r *http.Request) {
// so the response and audit name the real mutation.
var patch ServerSpecPatch
var changed []string
// imageFrom is the image a confirmed image change replaced, for the audit row.
var imageFrom string
if body.DisplayName != nil {
patch.DisplayName = body.DisplayName
@@ -730,8 +744,31 @@ func (a *API) handlePatchServer(w http.ResponseWriter, r *http.Request) {
"image %q is not on the whitelist", *body.Image))
return
}
patch.Image = body.Image
changed = append(changed, "image")
image, err := a.pinImage(r.Context(), *body.Image)
if err != nil {
writeError(w, r, err)
return
}
info, err := a.Cluster.GetServer(r.Context(), name)
if err != nil {
a.writeLookupError(w, r, err)
return
}
// Re-picking the tag a server was created from resolves to that tag's
// newest build, which is as much a version move as picking another image.
// Only a pin that lands on exactly the current image is no change at all.
if image != info.Image {
if !body.ConfirmImageChange {
writeError(w, r, newError(http.StatusConflict, "image_change_unconfirmed",
"changing the image from %q to %q opens this world with the new image's Minecraft version, "+
"and chunks it upgrades cannot be opened by the old one again; back the world up first, "+
"then resend with confirmImageChange", info.Image, image))
return
}
patch.Image = &image
changed = append(changed, "image")
imageFrom = info.Image
}
}
// Memory and the resource overrides move together: resolveResources derives the
@@ -814,7 +851,11 @@ func (a *API) handlePatchServer(w http.ResponseWriter, r *http.Request) {
}
}
a.audit(r, "server.patch", name)
if patch.Image != nil {
a.auditImageChange(r, name, imageFrom, *patch.Image)
} else {
a.audit(r, "server.patch", name)
}
writeJSON(w, http.StatusOK, map[string]any{
"name": name,
"patched": changed,
+27
View File
@@ -6,6 +6,7 @@ import (
"net/http"
"felis.lolicon.best/internal/build"
"felis.lolicon.best/internal/imagepin"
"k8s.io/apimachinery/pkg/util/validation"
)
@@ -252,3 +253,29 @@ func writeBuildError(w http.ResponseWriter, r *http.Request, err error) {
writeError(w, r, err)
}
}
// ImagePinner resolves an image ref to the immutable form a server's spec keeps
// (imagepin.Resolver). A ref it does not manage comes back unchanged.
type ImagePinner interface {
Pin(ctx context.Context, ref string) (string, error)
}
// pinImage pins an admitted ref for a server spec. A tag the registry does not
// hold is the caller's to fix (build or push it first); any other failure is the
// registry being unreachable, and the server is not created or changed without a
// pin, since an unpinned ref is exactly what lets a later push move its world.
func (a *API) pinImage(ctx context.Context, ref string) (string, error) {
if a.Images == nil {
return ref, nil
}
pinned, err := a.Images.Pin(ctx, ref)
switch {
case errors.Is(err, imagepin.ErrNotFound):
return "", newError(http.StatusBadRequest, "image_not_in_registry",
"image %q is whitelisted but the registry does not hold it; build or push it first", ref)
case err != nil:
return "", newError(http.StatusServiceUnavailable, "registry_unavailable",
"could not resolve image %q to a digest: %v", ref, err)
}
return pinned, nil
}
+162
View File
@@ -0,0 +1,162 @@
package api
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"testing"
"felis.lolicon.best/internal/imagepin"
)
const pinnedDigest = "sha256:1111111111111111111111111111111111111111111111111111111111111111"
// fakePinner pins every unpinned ref to pinnedDigest, or fails with err. A ref
// that already names a digest comes back as is, like imagepin.Resolver.
type fakePinner struct {
err error
seen []string
}
func (f *fakePinner) Pin(_ context.Context, ref string) (string, error) {
f.seen = append(f.seen, ref)
if f.err != nil {
return "", f.err
}
if imagepin.Pinned(ref) {
return ref, nil
}
return ref + "@" + pinnedDigest, nil
}
// TestCreateServerStoresPinnedImage: the spec a server is created with names the
// digest its tag resolved to, so a later push over the tag cannot move it.
func TestCreateServerStoresPinnedImage(t *testing.T) {
api, _, cl, _ := newCreateAPI()
pin := &fakePinner{}
api.Images = pin
w := do(api.ExternalHandler(), "POST", "/api/v1/servers", validCreateBody, nil)
if w.Code != http.StatusCreated {
t.Fatalf("code = %d, want 201 (%s)", w.Code, w.Body.String())
}
if got, want := cl.created["survival"].Image, admittedImage+"@"+pinnedDigest; got != want {
t.Errorf("created image = %q, want %q", got, want)
}
if len(pin.seen) != 1 || pin.seen[0] != admittedImage {
t.Errorf("pinned refs = %v, want the admitted ref once", pin.seen)
}
}
// TestCreateServerPinErrors: a tag the registry does not hold is the caller's
// problem (400); a registry that cannot answer is the platform's (503). Neither
// creates anything.
func TestCreateServerPinErrors(t *testing.T) {
for _, tc := range []struct {
err error
code int
wantCode string
}{
{fmt.Errorf("resolve: %w", imagepin.ErrNotFound), http.StatusBadRequest, "image_not_in_registry"},
{errors.New("dial tcp: connection refused"), http.StatusServiceUnavailable, "registry_unavailable"},
} {
api, _, cl, _ := newCreateAPI()
api.Images = &fakePinner{err: tc.err}
w := do(api.ExternalHandler(), "POST", "/api/v1/servers", validCreateBody, nil)
if w.Code != tc.code || decodeErr(t, w) != tc.wantCode {
t.Errorf("%v: got %d %s, want %d %s", tc.err, w.Code, w.Body.String(), tc.code, tc.wantCode)
}
if _, ok := cl.created["survival"]; ok {
t.Errorf("%v: server created despite the pin failure", tc.err)
}
}
}
// TestPatchServerImageChangeUnconfirmed: moving a world to another build is
// refused until the caller acknowledges the chunk upgrade cannot be undone.
func TestPatchServerImageChangeUnconfirmed(t *testing.T) {
api, repo, cl, _ := newPatchAPI()
cl.byName["survival"].Image = "registry.felis.svc:5000/mc:0@" + pinnedDigest
api.Images = &fakePinner{}
w := patchSurvival(api, `{"image":"`+admittedImage+`"}`)
if w.Code != http.StatusConflict || decodeErr(t, w) != "image_change_unconfirmed" {
t.Fatalf("got %d %s, want 409 image_change_unconfirmed", w.Code, w.Body.String())
}
if _, ok := cl.patched["survival"]; ok {
t.Error("spec patched without confirmation")
}
if len(repo.audits) != 0 {
t.Errorf("refused change audited: %+v", repo.audits)
}
}
// TestPatchServerImageConfirmedAudited: a confirmed change stores the pinned ref
// and the audit row names both builds.
func TestPatchServerImageConfirmedAudited(t *testing.T) {
api, repo, cl, _ := newPatchAPI()
from := "registry.felis.svc:5000/mc:0@" + pinnedDigest
cl.byName["survival"].Image = from
api.Images = &fakePinner{}
w := patchSurvival(api, `{"image":"`+admittedImage+`","confirmImageChange":true}`)
if w.Code != http.StatusOK {
t.Fatalf("code = %d, want 200 (%s)", w.Code, w.Body.String())
}
to := admittedImage + "@" + pinnedDigest
if p := cl.patched["survival"]; p.Image == nil || *p.Image != to {
t.Fatalf("patched image = %v, want %q", p.Image, to)
}
if len(repo.audits) != 1 || repo.audits[0].Action != "server.patch" {
t.Fatalf("audits = %+v, want one server.patch", repo.audits)
}
var payload map[string]string
if err := json.Unmarshal(repo.audits[0].Payload, &payload); err != nil {
t.Fatalf("audit payload %q: %v", repo.audits[0].Payload, err)
}
if payload["image_from"] != from || payload["image_to"] != to {
t.Errorf("audit payload = %v, want image_from %q image_to %q", payload, from, to)
}
}
// TestPatchServerImageSamePinIsNoChange: re-picking the tag a server runs, while
// the tag still names the same build, changes nothing and needs no confirmation.
func TestPatchServerImageSamePinIsNoChange(t *testing.T) {
api, repo, cl, _ := newPatchAPI()
cl.byName["survival"].Image = admittedImage + "@" + pinnedDigest
api.Images = &fakePinner{}
w := patchSurvival(api, `{"image":"`+admittedImage+`"}`)
if w.Code != http.StatusOK {
t.Fatalf("code = %d, want 200 (%s)", w.Code, w.Body.String())
}
if p := cl.patched["survival"]; p.Image != nil {
t.Errorf("patched image = %q, want no image change", *p.Image)
}
if len(repo.audits) != 1 || repo.audits[0].Payload != nil {
t.Errorf("audits = %+v, want a plain server.patch", repo.audits)
}
}
// TestPatchServerImageRestoresPinnedBuild: the exact build an earlier change
// replaced is admitted by its tag and set as is, so a world restored from a
// backup can go back to the build that wrote it.
func TestPatchServerImageRestoresPinnedBuild(t *testing.T) {
api, _, cl, fb := newPatchAPI()
cl.byName["survival"].Image = admittedImage + "@" + pinnedDigest
pin := &fakePinner{}
api.Images = pin
old := admittedImage + "@sha256:" + fmt.Sprintf("%064d", 0)
fb.admitted[old] = true
w := patchSurvival(api, `{"image":"`+old+`","confirmImageChange":true}`)
if w.Code != http.StatusOK {
t.Fatalf("code = %d, want 200 (%s)", w.Code, w.Body.String())
}
if p := cl.patched["survival"]; p.Image == nil || *p.Image != old {
t.Errorf("patched image = %v, want %q", p.Image, old)
}
}
+5 -1
View File
@@ -543,8 +543,12 @@ func (b *Builder) RemoveImage(ctx context.Context, imageRef string) error {
// image). A disabled row never admits. A wildcard whitelist entry
// ("registry/foo:*") admits any concrete tag on that repo (imageMatches); the
// caller always passes a concrete ref, never a wildcard. An empty ref is never
// admitted.
// admitted. A ref pinned to a digest (name:tag@sha256:…, what a server's spec
// carries once created) is admitted by its name:tag: the digest only fixes which
// build of that tag runs, so an admin can set a server back to the exact build an
// earlier image change replaced.
func (b *Builder) ImageAdmitted(ctx context.Context, imageRef string) (bool, error) {
imageRef, _, _ = strings.Cut(imageRef, "@")
if strings.TrimSpace(imageRef) == "" {
return false, nil
}
+7
View File
@@ -3,6 +3,7 @@ package build
import (
"context"
"errors"
"strings"
"testing"
"time"
@@ -548,6 +549,12 @@ func TestImageAdmitted(t *testing.T) {
{"registry.felis.svc:5000/unknown:1", false}, // not on the list
{"", false}, // empty ref
{" ", false}, // blank ref
// A pinned ref is admitted by its name:tag, wildcard or exact.
{"registry.felis.svc:5000/exact:1@sha256:" + strings.Repeat("a", 64), true},
{"registry.felis.svc:5000/wild:99@sha256:" + strings.Repeat("b", 64), true},
{"registry.felis.svc:5000/exact:2@sha256:" + strings.Repeat("a", 64), false},
{"registry.felis.svc:5000/off:1@sha256:" + strings.Repeat("a", 64), false},
{"@sha256:" + strings.Repeat("a", 64), false},
}
for _, c := range cases {
got, err := b.ImageAdmitted(context.Background(), c.ref)
+167
View File
@@ -0,0 +1,167 @@
// Package imagepin fixes a server's image to the exact build it was created
// with. The platform's own game images are published under mutable tags
// (registry.felis.svc:5000/felis/paper:demo): every installer run rebuilds them
// against the newest Paper/Limbo release and pushes over the same tag. A server
// whose spec.image names that tag would boot whatever the tag points at on its
// next wake, so re-running the installer would silently move a sleeping world to
// a newer Minecraft version. Chunk upgrades are one-way, so that move can never
// be taken back.
//
// Pin resolves such a tag to the manifest digest it names right now and appends
// it (name:tag@sha256:…). Kubernetes pulls a reference carrying a digest by the
// digest alone, so the tag stays only as a readable label of where the build
// came from. A pinned server changes image only when an admin changes
// spec.image, which the API makes an explicit, confirmed step.
//
// Only refs in the platform registry are resolved. That registry is reachable
// anonymously for reads from inside the cluster; a public registry would need a
// token exchange per vendor and egress felis-api does not otherwise have, and an
// external image is one an admin whitelisted by an exact tag of their choosing.
package imagepin
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"net/http"
"regexp"
"strings"
"time"
)
// ErrNotFound means the registry answered and does not hold the tag: the image
// was whitelisted but never pushed, or has been deleted since.
var ErrNotFound = errors.New("imagepin: tag not found in the registry")
// manifestAccept lists every manifest shape the registry may hold for a tag. A
// registry asked without an Accept it can satisfy answers with a converted
// schema-1 manifest, whose digest is not the one kubelet would pull.
var manifestAccept = strings.Join([]string{
"application/vnd.oci.image.index.v1+json",
"application/vnd.docker.distribution.manifest.list.v2+json",
"application/vnd.oci.image.manifest.v1+json",
"application/vnd.docker.distribution.manifest.v2+json",
}, ", ")
var digestRE = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`)
// maxManifestBytes bounds the body read when the registry sends no digest
// header; a manifest or index is a few KiB.
const maxManifestBytes = 4 << 20
// Pinned reports whether ref already names a digest.
func Pinned(ref string) bool { return strings.Contains(ref, "@") }
// Resolver pins refs that live in one registry.
type Resolver struct {
// Registry is the host[:port] the refs spell, the [registry] url
// (registry.felis.svc:5000). Refs under any other host are left as they are.
Registry string
// Endpoint is the host[:port] to dial for it. Empty means Registry, which is
// right inside the cluster; a command on the node reaches the same registry
// through its loopback hostPort instead.
Endpoint string
// Client makes the request. Nil uses a client with a 10s timeout.
Client *http.Client
}
// Covers reports whether ref lives in the resolver's registry.
func (r Resolver) Covers(ref string) bool {
return r.Registry != "" && strings.HasPrefix(ref, r.Registry+"/")
}
// Pin returns ref with the digest its tag names now appended. A ref that already
// carries a digest, or lives outside the registry, comes back unchanged.
func (r Resolver) Pin(ctx context.Context, ref string) (string, error) {
if Pinned(ref) || !r.Covers(ref) {
return ref, nil
}
repo, tag := splitTag(strings.TrimPrefix(ref, r.Registry+"/"))
if repo == "" {
return "", fmt.Errorf("imagepin: %q names no repository", ref)
}
digest, err := r.digest(ctx, repo, tag)
if err != nil {
return "", fmt.Errorf("imagepin: resolve %s: %w", ref, err)
}
if !strings.Contains(ref[strings.LastIndex(ref, "/")+1:], ":") {
ref += ":" + tag // spell the implied tag out, so the label reads as what was pinned
}
return ref + "@" + digest, nil
}
// splitTag splits "felis/paper:demo" into ("felis/paper", "demo"); a path with no
// tag means "latest", as it does for every image client.
func splitTag(path string) (repo, tag string) {
slash := strings.LastIndex(path, "/")
if colon := strings.LastIndex(path, ":"); colon > slash {
return path[:colon], path[colon+1:]
}
return path, "latest"
}
func (r Resolver) digest(ctx context.Context, repo, tag string) (string, error) {
endpoint := r.Endpoint
if endpoint == "" {
endpoint = r.Registry
}
// Plain HTTP: the platform registry is an in-cluster Service and the node's
// loopback hostPort, and containerd's mirror for it is configured the same way.
url := "http://" + endpoint + "/v2/" + repo + "/manifests/" + tag
client := r.Client
if client == nil {
client = &http.Client{Timeout: 10 * time.Second}
}
// HEAD first: the registry answers it with Docker-Content-Digest and no body.
// GET covers a registry that leaves the header off, by hashing the manifest.
for _, method := range []string{http.MethodHead, http.MethodGet} {
req, err := http.NewRequestWithContext(ctx, method, url, nil)
if err != nil {
return "", err
}
req.Header.Set("Accept", manifestAccept)
resp, err := client.Do(req)
if err != nil {
return "", err
}
d, err := readDigest(resp, method == http.MethodGet)
resp.Body.Close()
if err != nil || d != "" {
return d, err
}
}
return "", errors.New("registry sent neither a digest header nor a manifest")
}
// readDigest takes the digest from a manifest response, hashing the body when the
// header is missing and hash is set. An empty digest with a nil error means "try
// the next method".
func readDigest(resp *http.Response, hash bool) (string, error) {
switch {
case resp.StatusCode == http.StatusNotFound:
return "", ErrNotFound
case resp.StatusCode != http.StatusOK:
return "", fmt.Errorf("registry answered %s", resp.Status)
}
if d := resp.Header.Get("Docker-Content-Digest"); d != "" {
if !digestRE.MatchString(d) {
return "", fmt.Errorf("registry sent a malformed digest %q", d)
}
return d, nil
}
if !hash {
return "", nil
}
body, err := io.ReadAll(io.LimitReader(resp.Body, maxManifestBytes+1))
if err != nil {
return "", err
}
if len(body) > maxManifestBytes {
return "", errors.New("manifest is implausibly large")
}
sum := sha256.Sum256(body)
return "sha256:" + hex.EncodeToString(sum[:]), nil
}
+136
View File
@@ -0,0 +1,136 @@
package imagepin
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
const testDigest = "sha256:d2fcc09d2caa108678c540c99703db96d63038a5fc9e366402d7ef1712ec4d95"
// fakeRegistry serves /v2/felis/paper/manifests/demo and 404s everything else.
func fakeRegistry(t *testing.T, header bool) (*httptest.Server, *[]string) {
t.Helper()
var seen []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
seen = append(seen, r.Method+" "+r.URL.Path)
if !strings.Contains(r.Header.Get("Accept"), "application/vnd.oci.image.index.v1+json") {
t.Errorf("request without an OCI index Accept: %q", r.Header.Get("Accept"))
}
if r.URL.Path != "/v2/felis/paper/manifests/demo" {
http.NotFound(w, r)
return
}
if header {
w.Header().Set("Docker-Content-Digest", testDigest)
}
if r.Method == http.MethodGet {
w.Write([]byte(`{"schemaVersion":2}`))
}
}))
t.Cleanup(srv.Close)
return srv, &seen
}
func resolverFor(srv *httptest.Server) Resolver {
return Resolver{
Registry: "registry.felis.svc:5000",
Endpoint: strings.TrimPrefix(srv.URL, "http://"),
Client: srv.Client(),
}
}
func TestPinResolvesPlatformTag(t *testing.T) {
srv, seen := fakeRegistry(t, true)
got, err := resolverFor(srv).Pin(context.Background(), "registry.felis.svc:5000/felis/paper:demo")
if err != nil {
t.Fatalf("Pin: %v", err)
}
if want := "registry.felis.svc:5000/felis/paper:demo@" + testDigest; got != want {
t.Errorf("Pin = %q, want %q", got, want)
}
if len(*seen) != 1 || (*seen)[0] != "HEAD /v2/felis/paper/manifests/demo" {
t.Errorf("requests = %v, want a single HEAD", *seen)
}
}
func TestPinHashesManifestWithoutDigestHeader(t *testing.T) {
srv, seen := fakeRegistry(t, false)
got, err := resolverFor(srv).Pin(context.Background(), "registry.felis.svc:5000/felis/paper:demo")
if err != nil {
t.Fatalf("Pin: %v", err)
}
sum := sha256.Sum256([]byte(`{"schemaVersion":2}`))
if want := "registry.felis.svc:5000/felis/paper:demo@sha256:" + hex.EncodeToString(sum[:]); got != want {
t.Errorf("Pin = %q, want %q", got, want)
}
if len(*seen) != 2 {
t.Errorf("requests = %v, want HEAD then GET", *seen)
}
}
func TestPinLeavesOtherRefsAlone(t *testing.T) {
srv, seen := fakeRegistry(t, true)
r := resolverFor(srv)
for _, ref := range []string{
"registry.felis.svc:5000/felis/paper:demo@" + testDigest, // already pinned
"docker.io/itzg/minecraft-server:java21", // external
"registry.felis.svc:50000/felis/paper:demo", // a different port is a different registry
} {
got, err := r.Pin(context.Background(), ref)
if err != nil || got != ref {
t.Errorf("Pin(%q) = %q, %v; want it unchanged", ref, got, err)
}
}
if len(*seen) != 0 {
t.Errorf("requests = %v, want none", *seen)
}
}
func TestPinImpliedLatest(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/v2/felis/paper/manifests/latest" {
http.NotFound(w, r)
return
}
w.Header().Set("Docker-Content-Digest", testDigest)
}))
defer srv.Close()
got, err := resolverFor(srv).Pin(context.Background(), "registry.felis.svc:5000/felis/paper")
if err != nil {
t.Fatalf("Pin: %v", err)
}
if want := "registry.felis.svc:5000/felis/paper:latest@" + testDigest; got != want {
t.Errorf("Pin = %q, want %q", got, want)
}
}
func TestPinErrors(t *testing.T) {
srv, _ := fakeRegistry(t, true)
_, err := resolverFor(srv).Pin(context.Background(), "registry.felis.svc:5000/felis/paper:gone")
if !errors.Is(err, ErrNotFound) {
t.Errorf("missing tag: err = %v, want ErrNotFound", err)
}
bad := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Docker-Content-Digest", "sha256:nothex")
}))
defer bad.Close()
if _, err := resolverFor(bad).Pin(context.Background(), "registry.felis.svc:5000/felis/paper:demo"); err == nil {
t.Error("malformed digest accepted")
}
down := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusServiceUnavailable)
}))
defer down.Close()
_, err = resolverFor(down).Pin(context.Background(), "registry.felis.svc:5000/felis/paper:demo")
if err == nil || errors.Is(err, ErrNotFound) {
t.Errorf("503: err = %v, want a non-NotFound error", err)
}
}
+22 -9
View File
@@ -234,11 +234,13 @@ func loginToInternalAPI(p Params) *networkingv1.NetworkPolicy {
}
}
// RegistryIngressPolicy fences the registry pod: only build pods reach its port.
// Everything else that uses the registry runs on the node — containerd's pulls and
// the installer's pushes both arrive through the loopback hostPort — and Kubernetes
// never blocks resident-node traffic. Write authorization is the gate's job; this
// policy keeps every other pod from even trying.
// RegistryIngressPolicy fences the registry pod: only build pods and felis-api
// reach its port. Build pods push what they build; felis-api reads a manifest
// digest to pin a new server's image (internal/imagepin). Everything else that
// uses the registry runs on the node — containerd's pulls and the installer's
// pushes both arrive through the loopback hostPort — and Kubernetes never blocks
// resident-node traffic. Write authorization is the gate's job (felis-api holds no
// registry credential); this policy keeps every other pod from even trying.
func RegistryIngressPolicy(p Params) *networkingv1.NetworkPolicy {
p = p.withDefaults()
tcp := corev1.ProtocolTCP
@@ -246,11 +248,22 @@ func RegistryIngressPolicy(p Params) *networkingv1.NetworkPolicy {
np := netpol("felis-registry-ingress", p.RegistryNamespace,
metav1.LabelSelector{MatchLabels: registryLabels()},
[]networkingv1.NetworkPolicyIngressRule{{
From: []networkingv1.NetworkPolicyPeer{{
NamespaceSelector: &metav1.LabelSelector{
MatchLabels: map[string]string{"kubernetes.io/metadata.name": p.BuildNamespace},
From: []networkingv1.NetworkPolicyPeer{
{
NamespaceSelector: &metav1.LabelSelector{
MatchLabels: map[string]string{"kubernetes.io/metadata.name": p.BuildNamespace},
},
},
}},
{
NamespaceSelector: &metav1.LabelSelector{
MatchLabels: map[string]string{"kubernetes.io/metadata.name": p.ControlNamespace},
},
PodSelector: &metav1.LabelSelector{MatchLabels: map[string]string{
LabelPartOf: controlPlanePartOf,
LabelComponent: ComponentAPI,
}},
},
},
Ports: []networkingv1.NetworkPolicyPort{{Protocol: &tcp, Port: &port}},
}},
)
+16 -3
View File
@@ -307,7 +307,7 @@ func TestLoginToInternalAPI_SelectsOnlyTheSystemLoginPod(t *testing.T) {
// TestRegistryIngress_BuildNamespaceOnly pins who may dial the registry pod: build
// pods, on the registry port. Game servers and the control plane never pull
// through the Service — containerd pulls over the node's loopback hostPort.
func TestRegistryIngress_BuildNamespaceOnly(t *testing.T) {
func TestRegistryIngress_BuildNamespaceAndAPI(t *testing.T) {
p := testParams().withDefaults()
np := RegistryIngressPolicy(p)
if np.Namespace != p.RegistryNamespace {
@@ -319,14 +319,27 @@ func TestRegistryIngress_BuildNamespaceOnly(t *testing.T) {
if mapSelectorMatches(np.Spec.PodSelector.MatchLabels, APIDeployment(p).Spec.Template.Labels) {
t.Error("registry ingress must not also fence the api pod")
}
if len(np.Spec.Ingress) != 1 || len(np.Spec.Ingress[0].From) != 1 {
t.Fatalf("registry ingress shape = %+v, want one rule, one peer", np.Spec.Ingress)
if len(np.Spec.Ingress) != 1 || len(np.Spec.Ingress[0].From) != 2 {
t.Fatalf("registry ingress shape = %+v, want one rule, two peers", np.Spec.Ingress)
}
peer := np.Spec.Ingress[0].From[0]
if peer.PodSelector != nil || peer.IPBlock != nil || peer.NamespaceSelector == nil ||
peer.NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != p.BuildNamespace {
t.Errorf("registry ingress peer = %+v, want the whole %s namespace", peer, p.BuildNamespace)
}
// The second peer is felis-api alone: it selects the api pod and not the
// operator's, both of which live in the control namespace.
api := np.Spec.Ingress[0].From[1]
if api.IPBlock != nil || api.NamespaceSelector == nil || api.PodSelector == nil ||
api.NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != p.ControlNamespace {
t.Fatalf("registry ingress api peer = %+v, want pods in %s", api, p.ControlNamespace)
}
if !mapSelectorMatches(api.PodSelector.MatchLabels, APIDeployment(p).Spec.Template.Labels) {
t.Errorf("registry ingress api peer %v does not select the api pod", api.PodSelector)
}
if mapSelectorMatches(api.PodSelector.MatchLabels, OperatorDeployment(p).Spec.Template.Labels) {
t.Errorf("registry ingress api peer %v also selects the operator pod", api.PodSelector)
}
assertSinglePort(t, np.Spec.Ingress[0].Ports, int(p.RegistryPort))
}