fix(build): 构建 pod 等出口策略生效再运行,加 seccomp、可选 user namespace 与磁盘上限,上下文解包限总字节与条目数
This commit is contained in:
17 files changed
+948
-29
No files matched your search
+32
-1
@@ -9,6 +9,7 @@ import (
|
||||
"os"
|
||||
"regexp"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/api"
|
||||
@@ -152,11 +153,13 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int {
|
||||
// The fetch initContainer runs THIS image's fetch-context entrypoint, so the
|
||||
// build config carries the api's own image (the platform sets FELIS_IMAGE).
|
||||
buildCfg.FelisImage = os.Getenv("FELIS_IMAGE")
|
||||
buildJobs := build.NewK8sJobs(cl, buildCfg)
|
||||
builder := &build.Builder{
|
||||
Store: build.NewPGStore(drv.DB()),
|
||||
Jobs: build.NewK8sJobs(cl, buildCfg),
|
||||
Jobs: buildJobs,
|
||||
Config: buildCfg,
|
||||
}
|
||||
go probeBuildUserNamespaces(ctx, buildJobs, buildCfg, stderr)
|
||||
|
||||
// User-modpack approval lane (user-directed extension over §16; see
|
||||
// internal/submit). An ordinary user may only SUBMIT a
|
||||
@@ -457,6 +460,11 @@ func buildConfig(cfg *config.Config) build.Config {
|
||||
TrivyImage: cfg.Registry.TrivyImage,
|
||||
CPULimit: cfg.Registry.BuildCPULimit,
|
||||
MemLimit: cfg.Registry.BuildMemLimit,
|
||||
DiskLimit: cfg.Registry.BuildDiskLimit,
|
||||
// "auto" follows the startup probe (see probeBuildUserNamespaces).
|
||||
UserNamespaces: cfg.Registry.BuildUserNamespaces,
|
||||
UserNamespacesProbe: new(atomic.Bool),
|
||||
RuntimeClass: cfg.Registry.BuildRuntimeClass,
|
||||
// Empty keeps Trivy's own default; an install with builds points this at
|
||||
// the internal DB mirror (see config.RegistryConfig.TrivyDBRepository).
|
||||
TrivyDBRepository: cfg.Registry.TrivyDBRepository,
|
||||
@@ -464,6 +472,29 @@ func buildConfig(cfg *config.Config) build.Config {
|
||||
}
|
||||
}
|
||||
|
||||
// probeBuildUserNamespaces settles build_user_namespaces = "auto": one probe
|
||||
// pod with hostUsers: false tells whether this node's kernel and runtime can run
|
||||
// build pods in a user namespace. Builds submitted before it answers run without.
|
||||
func probeBuildUserNamespaces(ctx context.Context, jobs *build.K8sJobs, cfg build.Config, stderr io.Writer) {
|
||||
if mode := cfg.UserNamespaces; mode != "" && mode != build.UserNamespacesAuto {
|
||||
return
|
||||
}
|
||||
if cfg.FelisImage == "" {
|
||||
fmt.Fprintln(stderr, "felis api: FELIS_IMAGE unset — build pods run without a user namespace")
|
||||
return
|
||||
}
|
||||
ok, err := jobs.ProbeUserNamespaces(ctx, cfg.FelisImage)
|
||||
cfg.UserNamespacesProbe.Store(ok)
|
||||
switch {
|
||||
case ok:
|
||||
fmt.Fprintln(stderr, "felis api: build pods run in a user namespace (hostUsers: false)")
|
||||
case err != nil:
|
||||
fmt.Fprintf(stderr, "felis api: build pods run without a user namespace: the probe failed: %v\n", err)
|
||||
default:
|
||||
fmt.Fprintln(stderr, "felis api: build pods run without a user namespace: this node cannot start a pod with hostUsers: false")
|
||||
}
|
||||
}
|
||||
|
||||
// internalAPIBaseURL resolves the platform's internal-face base URL: the address
|
||||
// the platform rendered into this pod (felis API base URL env), or — for a
|
||||
// hand-rolled deployment that set none — the platform default control namespace,
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Vars so tests can shrink them. A dial that neither connects nor is refused
|
||||
// within egressDialTimeout counts as blocked: a policy that drops packets looks
|
||||
// exactly like that.
|
||||
var (
|
||||
egressDialTimeout = 500 * time.Millisecond
|
||||
egressPollInterval = 200 * time.Millisecond
|
||||
)
|
||||
|
||||
// cmdEgressGate is the first initContainer of every build pod. The pod's
|
||||
// NetworkPolicy is programmed asynchronously after the pod starts (live on k3s:
|
||||
// a build-labelled pod reached the internet and the Kubernetes API for its first
|
||||
// ~0.7 s), so the gate dials a destination the policy denies until it stops
|
||||
// answering, and only then lets the pod's next container, eventually the
|
||||
// untrusted Dockerfile, start.
|
||||
//
|
||||
// The default probe is the Kubernetes API Service, which the kubelet names in
|
||||
// every pod's environment and the build policy never admits. A probe that still
|
||||
// answers after --wait means the policy is not enforced at all (a CNI without
|
||||
// NetworkPolicy support, or k3s run with --disable-network-policy), and the
|
||||
// build fails closed.
|
||||
func cmdEgressGate(args []string, stdout, stderr io.Writer) int {
|
||||
fs := flag.NewFlagSet("egress-gate", flag.ContinueOnError)
|
||||
fs.SetOutput(stderr)
|
||||
probe := fs.String("probe", "", "host:port the build NetworkPolicy denies (default: the Kubernetes API Service from KUBERNETES_SERVICE_HOST/PORT)")
|
||||
wait := fs.Duration("wait", 2*time.Minute, "how long the probe may keep answering before the build is refused")
|
||||
if err := fs.Parse(args); err != nil {
|
||||
return 2
|
||||
}
|
||||
if *probe == "" {
|
||||
host, port := os.Getenv("KUBERNETES_SERVICE_HOST"), os.Getenv("KUBERNETES_SERVICE_PORT")
|
||||
if host == "" || port == "" {
|
||||
fmt.Fprintln(stderr, "felis egress-gate: no --probe and no KUBERNETES_SERVICE_HOST/PORT to default to")
|
||||
return 2
|
||||
}
|
||||
*probe = net.JoinHostPort(host, port)
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
for {
|
||||
conn, err := net.DialTimeout("tcp", *probe, egressDialTimeout)
|
||||
if err != nil {
|
||||
fmt.Fprintf(stdout, "felis egress-gate: %s is unreachable after %s (%v); the egress lock is in effect\n",
|
||||
*probe, time.Since(start).Round(time.Millisecond), err)
|
||||
return 0
|
||||
}
|
||||
_ = conn.Close()
|
||||
if time.Since(start) >= *wait {
|
||||
fmt.Fprintf(stderr, "felis egress-gate: %s still answers after %s: the build namespace's NetworkPolicy is not enforced "+
|
||||
"(a CNI without NetworkPolicy support, or k3s started with --disable-network-policy); refusing to run the build\n",
|
||||
*probe, *wait)
|
||||
return 1
|
||||
}
|
||||
time.Sleep(egressPollInterval)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"net"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func shrinkEgressGate(t *testing.T) {
|
||||
t.Helper()
|
||||
dial, poll := egressDialTimeout, egressPollInterval
|
||||
egressDialTimeout, egressPollInterval = 200*time.Millisecond, 10*time.Millisecond
|
||||
t.Cleanup(func() { egressDialTimeout, egressPollInterval = dial, poll })
|
||||
}
|
||||
|
||||
// The gate holds while the probe answers and lets the pod go on once the policy
|
||||
// lands, which the test plays by closing the listener.
|
||||
func TestEgressGateWaitsForTheLock(t *testing.T) {
|
||||
shrinkEgressGate(t)
|
||||
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
accepted := make(chan struct{}, 100)
|
||||
go func() {
|
||||
for {
|
||||
c, err := ln.Accept()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_ = c.Close()
|
||||
accepted <- struct{}{}
|
||||
}
|
||||
}()
|
||||
go func() {
|
||||
for i := 0; i < 3; i++ {
|
||||
<-accepted
|
||||
}
|
||||
_ = ln.Close()
|
||||
}()
|
||||
var out, errb bytes.Buffer
|
||||
if code := cmdEgressGate([]string{"--probe", ln.Addr().String(), "--wait", "10s"}, &out, &errb); code != 0 {
|
||||
t.Fatalf("exit %d: %s", code, errb.String())
|
||||
}
|
||||
if !strings.Contains(out.String(), "egress lock is in effect") {
|
||||
t.Errorf("stdout = %q", out.String())
|
||||
}
|
||||
}
|
||||
|
||||
// A probe that keeps answering means no policy is enforced: the build must not run.
|
||||
func TestEgressGateRefusesAnOpenNetwork(t *testing.T) {
|
||||
shrinkEgressGate(t)
|
||||
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer ln.Close()
|
||||
go func() {
|
||||
for {
|
||||
c, err := ln.Accept()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_ = c.Close()
|
||||
}
|
||||
}()
|
||||
var out, errb bytes.Buffer
|
||||
if code := cmdEgressGate([]string{"--probe", ln.Addr().String(), "--wait", "100ms"}, &out, &errb); code != 1 {
|
||||
t.Fatalf("exit %d, want 1", code)
|
||||
}
|
||||
if !strings.Contains(errb.String(), "not enforced") {
|
||||
t.Errorf("stderr = %q", errb.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestEgressGateDefaultsToTheKubernetesService(t *testing.T) {
|
||||
shrinkEgressGate(t)
|
||||
t.Setenv("KUBERNETES_SERVICE_HOST", "")
|
||||
t.Setenv("KUBERNETES_SERVICE_PORT", "")
|
||||
var out, errb bytes.Buffer
|
||||
if code := cmdEgressGate(nil, &out, &errb); code != 2 {
|
||||
t.Fatalf("exit %d without a probe, want 2", code)
|
||||
}
|
||||
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
host, port, _ := net.SplitHostPort(ln.Addr().String())
|
||||
_ = ln.Close() // closed: the lock reads as in effect at once
|
||||
t.Setenv("KUBERNETES_SERVICE_HOST", host)
|
||||
t.Setenv("KUBERNETES_SERVICE_PORT", port)
|
||||
out.Reset()
|
||||
if code := cmdEgressGate(nil, &out, &errb); code != 0 || !strings.Contains(out.String(), ln.Addr().String()) {
|
||||
t.Fatalf("exit %d, stdout %q", code, out.String())
|
||||
}
|
||||
}
|
||||
@@ -136,13 +136,27 @@ func fetchContextWithRetry(ctx context.Context, client *http.Client, url, token
|
||||
}
|
||||
}
|
||||
|
||||
// maxContextBytes / maxContextEntries bound what one context may expand to. The
|
||||
// compressed upload is capped at 1 GiB, but gzip turns that into hundreds of GiB
|
||||
// or millions of empty files, and the emptyDir's 4 GiB sizeLimit is only
|
||||
// enforced by the kubelet's periodic sweep, after the disk has filled. The byte
|
||||
// cap matches that sizeLimit; the entry cap is far above any real modpack (a
|
||||
// large one is a few thousand files) and far below an inode exhaustion.
|
||||
//
|
||||
// Vars, not consts, so tests can shrink them.
|
||||
var (
|
||||
maxContextBytes int64 = 4 << 30
|
||||
maxContextEntries = 200_000
|
||||
)
|
||||
|
||||
// extractTarGz streams a gzip'd tarball into root, creating directories as
|
||||
// needed. Every entry is vetted BEFORE anything is written: a path that is
|
||||
// absolute or escapes root (via ".."), a link (symlink or hardlink), or any
|
||||
// special file kind aborts the whole extraction. Refusing rather than skipping is
|
||||
// deliberate — a context that needs one of those constructs is not a context this
|
||||
// transport carries, and silently dropping entries would build from a corpus the
|
||||
// submitter did not upload.
|
||||
// submitter did not upload. The whole extraction is also bounded by
|
||||
// maxContextBytes and maxContextEntries.
|
||||
func extractTarGz(r io.Reader, root string) error {
|
||||
if err := os.MkdirAll(root, 0o755); err != nil {
|
||||
return fmt.Errorf("create context dir: %w", err)
|
||||
@@ -153,6 +167,8 @@ func extractTarGz(r io.Reader, root string) error {
|
||||
}
|
||||
defer zr.Close()
|
||||
tr := tar.NewReader(zr)
|
||||
var written int64
|
||||
entries := 0
|
||||
for {
|
||||
hdr, err := tr.Next()
|
||||
if errors.Is(err, io.EOF) {
|
||||
@@ -161,6 +177,9 @@ func extractTarGz(r io.Reader, root string) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("read context tarball: %w", err)
|
||||
}
|
||||
if entries++; entries > maxContextEntries {
|
||||
return fmt.Errorf("the build context has more than %d entries", maxContextEntries)
|
||||
}
|
||||
name := filepath.Clean(hdr.Name)
|
||||
if name == "." {
|
||||
continue
|
||||
@@ -188,10 +207,16 @@ func extractTarGz(r io.Reader, root string) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("create %q: %w", name, err)
|
||||
}
|
||||
if _, err := io.Copy(f, tr); err != nil {
|
||||
n, err := io.Copy(f, io.LimitReader(tr, maxContextBytes-written+1))
|
||||
written += n
|
||||
if err != nil {
|
||||
_ = f.Close()
|
||||
return fmt.Errorf("write %q: %w", name, err)
|
||||
}
|
||||
if written > maxContextBytes {
|
||||
_ = f.Close()
|
||||
return fmt.Errorf("the build context expands past %d bytes", maxContextBytes)
|
||||
}
|
||||
if err := f.Close(); err != nil {
|
||||
return fmt.Errorf("close %q: %w", name, err)
|
||||
}
|
||||
|
||||
@@ -121,6 +121,30 @@ func TestExtractTarGzRefusesEscapes(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A context that expands past the byte or entry cap is refused, however small
|
||||
// it was compressed: gzip bombs and inode floods stop at the cap.
|
||||
func TestExtractTarGzCapsExpansion(t *testing.T) {
|
||||
bytesCap, entriesCap := maxContextBytes, maxContextEntries
|
||||
t.Cleanup(func() { maxContextBytes, maxContextEntries = bytesCap, entriesCap })
|
||||
maxContextBytes, maxContextEntries = 1000, 5
|
||||
|
||||
fits := tgzBody(t, tarEntry{name: "a", body: strings.Repeat("x", 600)}, tarEntry{name: "b", body: strings.Repeat("y", 400)})
|
||||
if err := extractTarGz(bytes.NewReader(fits), t.TempDir()); err != nil {
|
||||
t.Fatalf("a context exactly at the byte cap: %v", err)
|
||||
}
|
||||
big := tgzBody(t, tarEntry{name: "a", body: strings.Repeat("x", 600)}, tarEntry{name: "b", body: strings.Repeat("y", 401)})
|
||||
if err := extractTarGz(bytes.NewReader(big), t.TempDir()); err == nil || !strings.Contains(err.Error(), "expands past") {
|
||||
t.Fatalf("one byte over the cap: err = %v", err)
|
||||
}
|
||||
var many []tarEntry
|
||||
for i := 0; i < 6; i++ {
|
||||
many = append(many, tarEntry{name: "d" + string(rune('0'+i)) + "/", typ: tar.TypeDir})
|
||||
}
|
||||
if err := extractTarGz(bytes.NewReader(tgzBody(t, many...)), t.TempDir()); err == nil || !strings.Contains(err.Error(), "entries") {
|
||||
t.Fatalf("six entries over a cap of five: err = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// The command end to end: it dials the URL with the bearer token from the
|
||||
// environment, and refuses to run without it (the internal face would 401
|
||||
// anyway; failing at parse time is the honest earlier error).
|
||||
|
||||
@@ -21,6 +21,7 @@ Commands:
|
||||
restore Extract a world archive into a world volume (internal Job entrypoint)
|
||||
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)
|
||||
egress-gate Hold a build pod until its egress NetworkPolicy is enforced (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)
|
||||
@@ -56,6 +57,7 @@ var commands = map[string]func(args []string, stdout, stderr io.Writer) int{
|
||||
"restore": cmdRestore,
|
||||
"backup": cmdBackup,
|
||||
"files": cmdFiles,
|
||||
"egress-gate": cmdEgressGate,
|
||||
"fetch-context": cmdFetchContext,
|
||||
"push-image": cmdPushImage,
|
||||
"registry-gate": cmdRegistryGate,
|
||||
|
||||
Reference in new issue
Block a user