feat(submit): local + S3 backends for modpack upload contexts, installer-selectable

This commit is contained in:
Lemon-miaow committed 2026-07-02 23:38:46 +08:00
1 parent 191640c3f5
commit 598f3d31f4
22 files changed
+2177 -34

No files matched your search

+115
View File
@@ -0,0 +1,115 @@
package submit
import (
"context"
"fmt"
"io"
"os"
"path/filepath"
"regexp"
)
// contextBlobName is the fixed object name of a submission's build context under
// its id-namespaced prefix. It is the single source of truth for both the
// derived context ref (deriveContextRef) and the on-disk write target
// (LocalContextStore), so the blob always lands exactly where Kaniko's
// --context points (build/jobspec.go).
const contextBlobName = "context.tar.gz"
// idRE re-validates a submission id at the storage boundary. The Manager only
// ever passes ids it loaded from the Store (already the crypto-hex ids newID
// mints), but LocalContextStore is a standalone component that treats the id as
// untrusted path input: a lowercase-alphanumeric-with-dashes id can contain no
// path separator and no "..", so it can never escape Base. This is the same
// defense-in-depth stance as internal/backup's zip-slip guard.
var idRE = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,127}$`)
// LocalContextStore is the filesystem-backed build-context blob store: it writes
// each submission's uploaded modpack to {Base}/{id}/context.tar.gz on a mounted
// PVC. It is one of two implemented Blobs backends — cmd/felis selects it when
// user_uploads_context is a local path and S3ContextStore when it is an s3:// base;
// a base that is neither (or an s3:// base with no credentials configured) leaves
// Manager.Blobs nil so the upload endpoint returns 503 rather than pretending to
// accept a file it cannot persist.
//
// Base MUST equal the Manager's ContextStore so the blob lands exactly where
// deriveContextRef points Kaniko's --context; cmd/felis wires both from the one
// config field (registry.user_uploads_context).
//
// INTEGRATION-ONLY seam (out of scope of the upload transport): persisting the
// blob is end-to-end only once the same uploads PVC is mounted into the Kaniko
// build Pod and Kaniko is told to read a local context (build/jobspec.go passes
// the ref straight into --context). The transport here makes the file durable at
// the derived location; wiring that path into the sandboxed build Job is a
// separate deployment integration, exactly like the restore executor's PVC mount.
type LocalContextStore struct {
// Base is the directory (uploads PVC mount) submission contexts are written
// under. Each submission gets its own {Base}/{id}/ subdirectory.
Base string
}
// dir returns the per-submission directory, rejecting an id that could escape
// Base. Every path the store touches is rooted here.
func (s *LocalContextStore) dir(id string) (string, error) {
if !idRE.MatchString(id) {
return "", fmt.Errorf("submit: invalid submission id %q", id)
}
return filepath.Join(s.Base, id), nil
}
// Put writes r to {Base}/{id}/context.tar.gz atomically: it streams into a temp
// file in the same directory and renames it over any previous upload only on a
// fully successful copy. So a failed, truncated, or oversize upload never
// replaces a good context and never leaves a half-written blob for Kaniko to
// read; a re-upload while the submission is still pending simply supersedes the
// previous one. It returns the number of bytes stored.
func (s *LocalContextStore) Put(_ context.Context, id string, r io.Reader) (int64, error) {
dir, err := s.dir(id)
if err != nil {
return 0, err
}
if err := os.MkdirAll(dir, 0o750); err != nil {
return 0, fmt.Errorf("submit: mkdir context dir: %w", err)
}
tmp, err := os.CreateTemp(dir, contextBlobName+".*.tmp")
if err != nil {
return 0, fmt.Errorf("submit: create temp context: %w", err)
}
tmpName := tmp.Name()
n, err := io.Copy(tmp, r)
if err != nil {
tmp.Close()
os.Remove(tmpName)
return 0, fmt.Errorf("submit: write context blob: %w", err)
}
if err := tmp.Close(); err != nil {
os.Remove(tmpName)
return 0, fmt.Errorf("submit: close context blob: %w", err)
}
if err := os.Rename(tmpName, filepath.Join(dir, contextBlobName)); err != nil {
os.Remove(tmpName)
return 0, fmt.Errorf("submit: commit context blob: %w", err)
}
return n, nil
}
// Exists reports whether a context blob has been stored for id. Approve consults
// it so a submission whose context was never uploaded is refused BEFORE the CAS,
// instead of being approved into a build Kaniko cannot pull.
func (s *LocalContextStore) Exists(_ context.Context, id string) (bool, error) {
dir, err := s.dir(id)
if err != nil {
return false, err
}
switch _, err := os.Stat(filepath.Join(dir, contextBlobName)); {
case err == nil:
return true, nil
case os.IsNotExist(err):
return false, nil
default:
return false, fmt.Errorf("submit: stat context blob: %w", err)
}
}
// Compile-time proof that the filesystem store satisfies the Blobs transport.
var _ Blobs = (*LocalContextStore)(nil)
+94
View File
@@ -0,0 +1,94 @@
package submit
import (
"context"
"os"
"path/filepath"
"strings"
"testing"
)
func TestLocalContextStorePutAndExists(t *testing.T) {
base := t.TempDir()
s := &LocalContextStore{Base: base}
ctx := context.Background()
if ok, err := s.Exists(ctx, "sub-abc"); err != nil || ok {
t.Fatalf("Exists before Put = (%v, %v), want (false, nil)", ok, err)
}
payload := "\x1f\x8b\x08\x00the modpack context"
n, err := s.Put(ctx, "sub-abc", strings.NewReader(payload))
if err != nil {
t.Fatalf("Put: %v", err)
}
if n != int64(len(payload)) {
t.Fatalf("Put returned %d bytes, want %d", n, len(payload))
}
// The blob lands at exactly {base}/{id}/context.tar.gz — where deriveContextRef
// points Kaniko's --context.
dest := filepath.Join(base, "sub-abc", contextBlobName)
got, err := os.ReadFile(dest)
if err != nil {
t.Fatalf("read stored blob: %v", err)
}
if string(got) != payload {
t.Fatalf("stored %q, want %q", got, payload)
}
if ok, err := s.Exists(ctx, "sub-abc"); err != nil || !ok {
t.Fatalf("Exists after Put = (%v, %v), want (true, nil)", ok, err)
}
// The write is atomic: no leftover temp files beside the committed blob.
entries, err := os.ReadDir(filepath.Join(base, "sub-abc"))
if err != nil {
t.Fatalf("read dir: %v", err)
}
if len(entries) != 1 || entries[0].Name() != contextBlobName {
var names []string
for _, e := range entries {
names = append(names, e.Name())
}
t.Fatalf("dir entries = %v, want only %q (no temp files)", names, contextBlobName)
}
}
func TestLocalContextStorePutOverwrites(t *testing.T) {
base := t.TempDir()
s := &LocalContextStore{Base: base}
ctx := context.Background()
if _, err := s.Put(ctx, "sub-1", strings.NewReader("\x1f\x8bfirst")); err != nil {
t.Fatalf("first Put: %v", err)
}
if _, err := s.Put(ctx, "sub-1", strings.NewReader("\x1f\x8bsecond upload")); err != nil {
t.Fatalf("second Put: %v", err)
}
got, err := os.ReadFile(filepath.Join(base, "sub-1", contextBlobName))
if err != nil {
t.Fatalf("read: %v", err)
}
if string(got) != "\x1f\x8bsecond upload" {
t.Fatalf("stored %q, want the second upload (a re-upload supersedes)", got)
}
}
func TestLocalContextStoreRejectsUnsafeID(t *testing.T) {
base := t.TempDir()
s := &LocalContextStore{Base: base}
ctx := context.Background()
for _, id := range []string{"../evil", "sub/../../etc", "SUB-UPPER", "has space", "", "a/b"} {
if _, err := s.Put(ctx, id, strings.NewReader("\x1f\x8bx")); err == nil {
t.Errorf("Put(%q) succeeded, want rejection", id)
}
if _, err := s.Exists(ctx, id); err == nil {
t.Errorf("Exists(%q) succeeded, want rejection", id)
}
}
// Nothing escaped the base directory.
if _, err := os.Stat(filepath.Join(filepath.Dir(base), "evil")); !os.IsNotExist(err) {
t.Fatal("an unsafe id wrote outside Base")
}
}
+217
View File
@@ -0,0 +1,217 @@
package submit
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"path"
"strings"
"github.com/minio/minio-go/v7"
"github.com/minio/minio-go/v7/pkg/credentials"
)
// s3Client is the minimal object-store surface S3ContextStore needs. *minio.Client
// satisfies it, and a fake satisfies it in tests — so the store's key derivation
// and not-found handling are unit-verifiable without a live bucket.
type s3Client interface {
PutObject(ctx context.Context, bucket, object string, reader io.Reader, size int64, opts minio.PutObjectOptions) (minio.UploadInfo, error)
StatObject(ctx context.Context, bucket, object string, opts minio.StatObjectOptions) (minio.ObjectInfo, error)
}
// S3ContextStore is the object-store-backed build-context blob store: it writes
// each submission's uploaded modpack to {prefix}/{id}/context.tar.gz inside an S3
// bucket. It is the second implemented Blobs backend (alongside LocalContextStore),
// selected by cmd/felis when user_uploads_context is an s3:// base.
//
// The bucket + key prefix are parsed from that same base (parseS3Base), so an
// object written here lands at exactly s3://{bucket}/{prefix}/{id}/context.tar.gz —
// the ref deriveContextRef records and Kaniko's native s3:// --context reads.
// Credentials are static V4 keys resolved by cmd/felis from the environment (the
// setup wizard injects them into felis-api from the felis-uploads-s3 Secret); they
// never touch felis.toml.
//
// Kaniko reading the S3 context at build time needs its own credentials + egress
// on the sandboxed build Job — a separate deployment integration, exactly like the
// LocalContextStore PVC mount. This transport only makes the upload durable at the
// derived location.
type S3ContextStore struct {
client s3Client
bucket string
prefix string // key prefix within the bucket; may be empty
}
// S3StoreConfig is the resolved input for NewS3ContextStore. Base is the s3://
// user_uploads_context (bucket + optional prefix are parsed from it, so the write
// path matches deriveContextRef); Endpoint may carry an http:// or https:// scheme
// (a bare host defaults to TLS); the keys come from the environment.
type S3StoreConfig struct {
Base string
Endpoint string
Region string
AccessKey string
SecretKey string
}
// NewS3ContextStore builds a store backed by a real minio client. It fails fast
// when the base is malformed or the credentials are missing, so cmd/felis leaves
// Manager.Blobs nil (upload endpoint → 503) rather than wiring a store that cannot
// authenticate.
func NewS3ContextStore(cfg S3StoreConfig) (*S3ContextStore, error) {
bucket, prefix, err := parseS3Base(cfg.Base)
if err != nil {
return nil, err
}
if cfg.AccessKey == "" || cfg.SecretKey == "" {
return nil, errors.New("submit: s3 store requires credentials")
}
host, secure, err := splitS3Endpoint(cfg.Endpoint)
if err != nil {
return nil, fmt.Errorf("submit: s3 endpoint: %w", err)
}
client, err := minio.New(host, &minio.Options{
Creds: credentials.NewStaticV4(cfg.AccessKey, cfg.SecretKey, ""),
Secure: secure,
Region: cfg.Region,
})
if err != nil {
return nil, fmt.Errorf("submit: s3 client: %w", err)
}
return &S3ContextStore{client: client, bucket: bucket, prefix: prefix}, nil
}
// CheckS3Access verifies the S3 coordinates before they are committed to config:
// it builds a client from the entered endpoint/credentials and probes the bucket.
// It is the install-time preflight that turns a mistyped key, wrong endpoint, or
// missing bucket into an immediate, legible error at the keyboard instead of a 503
// at the first real upload.
func CheckS3Access(ctx context.Context, cfg S3StoreConfig) error {
store, err := NewS3ContextStore(cfg)
if err != nil {
return err
}
return checkBucketAccess(ctx, store.client, store.bucket)
}
// checkBucketAccess probes the bucket with the SAME object-level HEAD the upload
// path uses (StatObject on a key that will not exist), not a bucket-level
// HeadBucket. This matters: a least-privilege key scoped to object Put/Get may lack
// s3:ListBucket, so a HeadBucket would falsely reject a key that uploads fine. A
// NoSuchKey/absent result means the endpoint is reachable and the credentials are
// accepted — exactly the runtime dependency Exists() relies on. Split from
// CheckS3Access so the error mapping is unit-testable against a fake.
func checkBucketAccess(ctx context.Context, client s3Client, bucket string) error {
const probe = "felis-access-probe/does-not-exist"
if _, err := client.StatObject(ctx, bucket, probe, minio.StatObjectOptions{}); err != nil {
resp := minio.ToErrorResponse(err)
switch resp.Code {
case "NoSuchKey", "NotFound":
return nil // reachable + authorized; the probe object is simply absent
case "NoSuchBucket":
return fmt.Errorf("submit: bucket %q not found", bucket)
case "AccessDenied", "SignatureDoesNotMatch", "InvalidAccessKeyId":
return fmt.Errorf("submit: s3 credentials rejected: %w", err)
default:
// A bare 404 with no bucket-specific code = object absent in a live bucket.
if resp.StatusCode == http.StatusNotFound {
return nil
}
return fmt.Errorf("submit: cannot reach s3 (endpoint unreachable or credentials rejected): %w", err)
}
}
return nil // the probe object improbably exists — access clearly works
}
// keyFor derives the object key for a submission, re-validating the id at the
// storage boundary (the same defense-in-depth as LocalContextStore: a validated id
// carries no path separator, so it cannot alter the key layout).
func (s *S3ContextStore) keyFor(id string) (string, error) {
if !idRE.MatchString(id) {
return "", fmt.Errorf("submit: invalid submission id %q", id)
}
return path.Join(s.prefix, id, contextBlobName), nil
}
// Put streams r to the derived object key. Size is unknown (the Manager hands us a
// size-capped reader), so it is uploaded with size -1 (multipart). PutObject is
// atomic from a reader's perspective — a partial upload never becomes a readable
// object — so a failed or oversize upload never replaces a good context. It
// returns the number of bytes stored.
func (s *S3ContextStore) Put(ctx context.Context, id string, r io.Reader) (int64, error) {
key, err := s.keyFor(id)
if err != nil {
return 0, err
}
info, err := s.client.PutObject(ctx, s.bucket, key, r, -1, minio.PutObjectOptions{ContentType: "application/gzip"})
if err != nil {
return 0, fmt.Errorf("submit: put context blob: %w", err)
}
return info.Size, nil
}
// Exists reports whether a context blob has been stored for id. Approve consults
// it so a submission whose context was never uploaded is refused BEFORE the CAS.
func (s *S3ContextStore) Exists(ctx context.Context, id string) (bool, error) {
key, err := s.keyFor(id)
if err != nil {
return false, err
}
if _, err := s.client.StatObject(ctx, s.bucket, key, minio.StatObjectOptions{}); err != nil {
if isS3NotFound(err) {
return false, nil
}
return false, fmt.Errorf("submit: stat context blob: %w", err)
}
return true, nil
}
// isS3NotFound recognizes the "object is absent" outcome across S3
// implementations: a GET-shaped NoSuchKey code or a bare 404 from the HEAD that
// StatObject issues.
func isS3NotFound(err error) bool {
resp := minio.ToErrorResponse(err)
return resp.Code == "NoSuchKey" || resp.StatusCode == http.StatusNotFound
}
// splitS3Endpoint separates a configured endpoint into the host[:port] minio.New
// wants and a TLS flag. A bare host defaults to TLS (the safe default); an
// explicit http:// opts out for a plaintext dev store.
func splitS3Endpoint(ep string) (host string, secure bool, err error) {
ep = strings.TrimSpace(ep)
switch {
case ep == "":
return "", false, errors.New("empty endpoint")
case strings.HasPrefix(ep, "https://"):
return strings.Trim(strings.TrimPrefix(ep, "https://"), "/"), true, nil
case strings.HasPrefix(ep, "http://"):
return strings.Trim(strings.TrimPrefix(ep, "http://"), "/"), false, nil
default:
return strings.Trim(ep, "/"), true, nil
}
}
// parseS3Base splits an s3://bucket[/prefix] base into its bucket and key prefix.
// It is the single source of truth for how a user_uploads_context s3:// base maps
// onto object storage, kept beside the store so the write path and deriveContextRef
// can never disagree about where the blob lands.
func parseS3Base(base string) (bucket, prefix string, err error) {
rest := base
if i := strings.Index(strings.ToLower(rest), "://"); i >= 0 {
rest = rest[i+3:]
}
rest = strings.Trim(rest, "/")
if rest == "" {
return "", "", fmt.Errorf("submit: s3 base %q has no bucket", base)
}
parts := strings.SplitN(rest, "/", 2)
bucket = parts[0]
if len(parts) == 2 {
prefix = strings.Trim(parts[1], "/")
}
return bucket, prefix, nil
}
// Compile-time proof that the object store satisfies the Blobs transport.
var _ Blobs = (*S3ContextStore)(nil)
+224
View File
@@ -0,0 +1,224 @@
package submit
import (
"context"
"errors"
"io"
"net/http"
"strings"
"testing"
"github.com/minio/minio-go/v7"
)
// fakeS3 is an in-memory s3Client: it records puts and returns a NoSuchKey/404
// error for a missing stat, so S3ContextStore's key derivation and not-found
// handling are exercised without a live bucket.
type fakeS3 struct {
objects map[string][]byte
putErr error
statErr error // when set, StatObject returns it (e.g. auth rejected / bucket missing)
}
func (f *fakeS3) PutObject(_ context.Context, bucket, object string, r io.Reader, _ int64, _ minio.PutObjectOptions) (minio.UploadInfo, error) {
if f.putErr != nil {
return minio.UploadInfo{}, f.putErr
}
data, err := io.ReadAll(r)
if err != nil {
return minio.UploadInfo{}, err
}
if f.objects == nil {
f.objects = map[string][]byte{}
}
f.objects[bucket+"/"+object] = data
return minio.UploadInfo{Bucket: bucket, Key: object, Size: int64(len(data))}, nil
}
func (f *fakeS3) StatObject(_ context.Context, bucket, object string, _ minio.StatObjectOptions) (minio.ObjectInfo, error) {
if f.statErr != nil {
return minio.ObjectInfo{}, f.statErr
}
if data, ok := f.objects[bucket+"/"+object]; ok {
return minio.ObjectInfo{Key: object, Size: int64(len(data))}, nil
}
return minio.ObjectInfo{}, minio.ErrorResponse{Code: "NoSuchKey", StatusCode: http.StatusNotFound}
}
func TestCheckBucketAccess(t *testing.T) {
ctx := context.Background()
// A reachable bucket where the probe object is simply absent (NoSuchKey) — the
// object-scoped key case that a bucket-level HeadBucket would wrongly reject.
if err := checkBucketAccess(ctx, &fakeS3{}, "b"); err != nil {
t.Fatalf("reachable bucket, probe absent = %v, want nil", err)
}
// A bare 404 (object absent, no bucket-specific code) is also success.
if err := checkBucketAccess(ctx, &fakeS3{statErr: minio.ErrorResponse{StatusCode: http.StatusNotFound}}, "b"); err != nil {
t.Fatalf("bare 404 = %v, want nil (object absent in a live bucket)", err)
}
// A missing bucket is a real error.
if err := checkBucketAccess(ctx, &fakeS3{statErr: minio.ErrorResponse{Code: "NoSuchBucket", StatusCode: http.StatusNotFound}}, "b"); err == nil {
t.Fatal("missing bucket = nil, want error")
}
// Rejected credentials are a real error.
if err := checkBucketAccess(ctx, &fakeS3{statErr: minio.ErrorResponse{Code: "AccessDenied", StatusCode: http.StatusForbidden}}, "b"); err == nil {
t.Fatal("rejected credentials = nil, want error")
}
// An unreachable endpoint (non-HTTP error) is a real error.
if err := checkBucketAccess(ctx, &fakeS3{statErr: errors.New("dial tcp: connection refused")}, "b"); err == nil {
t.Fatal("unreachable endpoint = nil, want error")
}
}
func TestCheckS3AccessValidatesConfigFirst(t *testing.T) {
// A malformed base fails at construction, before any network probe is attempted.
if err := CheckS3Access(context.Background(), S3StoreConfig{Base: "s3://", Endpoint: "x", AccessKey: "a", SecretKey: "b"}); err == nil {
t.Fatal("CheckS3Access with no bucket succeeded, want error")
}
}
func TestS3ContextStorePutAndExists(t *testing.T) {
fake := &fakeS3{}
s := &S3ContextStore{client: fake, bucket: "felis-uploads", prefix: "builds"}
ctx := context.Background()
if ok, err := s.Exists(ctx, "sub-abc"); err != nil || ok {
t.Fatalf("Exists before Put = (%v, %v), want (false, nil)", ok, err)
}
payload := "\x1f\x8b\x08\x00the modpack context"
n, err := s.Put(ctx, "sub-abc", strings.NewReader(payload))
if err != nil {
t.Fatalf("Put: %v", err)
}
if n != int64(len(payload)) {
t.Fatalf("Put returned %d bytes, want %d", n, len(payload))
}
// The blob lands at exactly {prefix}/{id}/context.tar.gz — the key half of the
// s3://bucket/prefix/id/context.tar.gz ref deriveContextRef records.
wantKey := "felis-uploads/builds/sub-abc/" + contextBlobName
if got := string(fake.objects[wantKey]); got != payload {
t.Fatalf("object at %q = %q, want %q", wantKey, got, payload)
}
if ok, err := s.Exists(ctx, "sub-abc"); err != nil || !ok {
t.Fatalf("Exists after Put = (%v, %v), want (true, nil)", ok, err)
}
}
func TestS3ContextStoreEmptyPrefix(t *testing.T) {
fake := &fakeS3{}
s := &S3ContextStore{client: fake, bucket: "b", prefix: ""}
if _, err := s.Put(context.Background(), "sub-1", strings.NewReader("\x1f\x8bx")); err != nil {
t.Fatalf("Put: %v", err)
}
// No prefix ⇒ the key is just {id}/context.tar.gz (no leading slash).
if _, ok := fake.objects["b/sub-1/"+contextBlobName]; !ok {
t.Fatalf("object not at expected key; got keys %v", s3KeysOf(fake.objects))
}
}
func TestS3ContextStoreRejectsUnsafeID(t *testing.T) {
fake := &fakeS3{}
s := &S3ContextStore{client: fake, bucket: "b", prefix: "p"}
ctx := context.Background()
for _, id := range []string{"../evil", "sub/../../etc", "SUB-UPPER", "has space", "", "a/b"} {
if _, err := s.Put(ctx, id, strings.NewReader("\x1f\x8bx")); err == nil {
t.Errorf("Put(%q) succeeded, want rejection", id)
}
if _, err := s.Exists(ctx, id); err == nil {
t.Errorf("Exists(%q) succeeded, want rejection", id)
}
}
if len(fake.objects) != 0 {
t.Fatalf("an unsafe id wrote an object: %v", s3KeysOf(fake.objects))
}
}
func TestParseS3Base(t *testing.T) {
cases := []struct {
base string
bucket, prefix string
wantErr bool
}{
{"s3://felis-user-uploads", "felis-user-uploads", "", false},
{"s3://bucket/builds", "bucket", "builds", false},
{"s3://bucket/a/b/c", "bucket", "a/b/c", false},
{"S3://Bucket/", "Bucket", "", false},
{"s3://bucket/pre/", "bucket", "pre", false},
{"s3://", "", "", true},
{"s3:///onlyslash", "", "", false}, // trims to "onlyslash" bucket
}
for _, c := range cases {
bucket, prefix, err := parseS3Base(c.base)
if (err != nil) != c.wantErr {
t.Errorf("parseS3Base(%q) err = %v, wantErr %v", c.base, err, c.wantErr)
continue
}
if err != nil {
continue
}
if c.base == "s3:///onlyslash" {
if bucket != "onlyslash" {
t.Errorf("parseS3Base(%q) bucket = %q, want onlyslash", c.base, bucket)
}
continue
}
if bucket != c.bucket || prefix != c.prefix {
t.Errorf("parseS3Base(%q) = (%q, %q), want (%q, %q)", c.base, bucket, prefix, c.bucket, c.prefix)
}
}
}
func TestSplitS3Endpoint(t *testing.T) {
cases := []struct {
ep string
host string
secure bool
wantErr bool
}{
{"https://s3.amazonaws.com", "s3.amazonaws.com", true, false},
{"http://minio:9000", "minio:9000", false, false},
{"minio.example.com:9000", "minio.example.com:9000", true, false},
{"https://s3.example.com/", "s3.example.com", true, false},
{"", "", false, true},
}
for _, c := range cases {
host, secure, err := splitS3Endpoint(c.ep)
if (err != nil) != c.wantErr {
t.Errorf("splitS3Endpoint(%q) err = %v, wantErr %v", c.ep, err, c.wantErr)
continue
}
if err != nil {
continue
}
if host != c.host || secure != c.secure {
t.Errorf("splitS3Endpoint(%q) = (%q, %v), want (%q, %v)", c.ep, host, secure, c.host, c.secure)
}
}
}
func TestNewS3ContextStoreValidation(t *testing.T) {
if _, err := NewS3ContextStore(S3StoreConfig{Base: "s3://", AccessKey: "a", SecretKey: "b", Endpoint: "x"}); err == nil {
t.Error("NewS3ContextStore with no bucket succeeded, want error")
}
if _, err := NewS3ContextStore(S3StoreConfig{Base: "s3://b", AccessKey: "", SecretKey: "", Endpoint: "x"}); err == nil {
t.Error("NewS3ContextStore with no credentials succeeded, want error")
}
if _, err := NewS3ContextStore(S3StoreConfig{Base: "s3://b", AccessKey: "a", SecretKey: "b", Endpoint: ""}); err == nil {
t.Error("NewS3ContextStore with no endpoint succeeded, want error")
}
if _, err := NewS3ContextStore(S3StoreConfig{Base: "s3://b/pre", AccessKey: "a", SecretKey: "b", Endpoint: "minio:9000"}); err != nil {
t.Errorf("NewS3ContextStore with valid config: %v", err)
}
}
func s3KeysOf(m map[string][]byte) []string {
var out []string
for k := range m {
out = append(out, k)
}
return out
}
+159 -7
View File
@@ -55,11 +55,13 @@
package submit
import (
"bufio"
"context"
"crypto/rand"
"encoding/hex"
"errors"
"fmt"
"io"
"regexp"
"strings"
"time"
@@ -76,13 +78,19 @@ const (
StatusRejected Status = "rejected"
)
// Sentinels. The API layer maps ErrInvalid→400, ErrNotFound→404 and
// ErrAlreadyReviewed→409; they are kept distinct from store/cluster failures so
// those surface as 500.
// Sentinels. The API layer maps ErrInvalid→400, ErrNotFound→404,
// ErrAlreadyReviewed→409 and ErrUploadsUnavailable→503; they are kept distinct
// from store/cluster failures so those surface as 500.
var (
ErrInvalid = errors.New("submit: invalid request")
ErrNotFound = errors.New("submit: submission not found")
ErrAlreadyReviewed = errors.New("submit: submission already reviewed")
// ErrUploadsUnavailable means this deployment configured a context store with
// no implemented upload transport (a nil Manager.Blobs — e.g. an object-store
// base with no client wired). UploadContext returns it so the endpoint reports
// an honest 503, never a 500, exactly as the restore executor does when its
// integration is not wired.
ErrUploadsUnavailable = errors.New("submit: context upload transport not configured")
)
// invalidf wraps ErrInvalid so every malformed-request case maps to one 400.
@@ -90,9 +98,19 @@ func invalidf(format string, a ...any) error {
return fmt.Errorf("%w: "+format, append([]any{ErrInvalid}, a...)...)
}
// errContextTooLarge trips when an upload exceeds the size cap. It wraps
// ErrInvalid so an oversize upload maps to a 400; the storage layer's %w wrapping
// preserves that chain through to the API error mapper.
var errContextTooLarge = fmt.Errorf("%w: build context exceeds the maximum allowed size", ErrInvalid)
const (
maxDisplayName = 200
maxRejectReason = 1000
// defaultMaxContextBytes caps an uploaded build-context blob. Modpack contexts
// (mods, configs, an occasional bundled world) are large, so the cap is
// generous; it bounds what one untrusted upload can write to the uploads PVC,
// not a tight quota. Override per-Manager via MaxContextBytes.
defaultMaxContextBytes = 1 << 30 // 1 GiB
)
// displayNameRE constrains the user-supplied label to a calm, single-line set:
@@ -156,6 +174,24 @@ type Builds interface {
Submit(ctx context.Context, req build.Request) (*build.Build, error)
}
// Blobs is the build-context blob transport the lane depends on to place a
// submitter's uploaded modpack at the platform-derived, id-namespaced location
// deriveContextRef points Kaniko at. It is the piece the package doc calls a
// "separate, deferred transport": creation only derives and records the ref, and
// the bytes behind it arrive through Put here. It is an interface so the Manager
// is unit-tested against an in-memory fake; the production implementation is the
// filesystem-backed LocalContextStore.
//
// Both methods key off the submission id, never a caller-supplied path, so the
// write target is as platform-pinned as the derived ref itself. Put stores (and
// atomically overwrites, while the submission is still pending) the blob; Exists
// reports whether one has been stored, so Approve can refuse to build a
// submission whose context was never uploaded.
type Blobs interface {
Put(ctx context.Context, id string, r io.Reader) (int64, error)
Exists(ctx context.Context, id string) (bool, error)
}
// Manager orchestrates the approval lane. It holds no mutable state; the clock
// and id generator are injectable for hermetic tests.
type Manager struct {
@@ -171,15 +207,30 @@ type Manager struct {
// field. The derived ref is {Registry}/user-uploads/{id}:latest.
Registry string
// ContextStore is the pinned Kaniko build-context base for user uploads, e.g.
// "s3://felis-user-uploads" (mirrors a configured object store). The derived
// context ref is {ContextStore}/{id}/context.tar.gz; the modpack blob is
// placed there by a separate upload transport (deferred — see package doc).
// "s3://felis-user-uploads" (an object store) or a local uploads PVC path. The
// derived context ref is {ContextStore}/{id}/context.tar.gz.
ContextStore string
// Blobs is the upload transport that persists the modpack behind the derived
// context ref. When nil (a store with no implemented transport, e.g. an
// object-store base with no client), UploadContext returns ErrUploadsUnavailable
// so the endpoint reports 503. Its backing MUST match ContextStore so the blob
// lands exactly where the derived ref points.
Blobs Blobs
// MaxContextBytes overrides the uploaded-context size cap; 0 uses
// defaultMaxContextBytes.
MaxContextBytes int64
Now func() time.Time
IDGen func() string
}
func (m *Manager) maxContextBytes() int64 {
if m.MaxContextBytes > 0 {
return m.MaxContextBytes
}
return defaultMaxContextBytes
}
func (m *Manager) now() time.Time {
if m.Now != nil {
return m.Now()
@@ -218,7 +269,7 @@ func (m *Manager) deriveImageRef(id string) string {
// selects nothing that reaches Kaniko's --context argument; only the blob behind
// this pinned, id-namespaced location (placed by the upload transport) varies.
func (m *Manager) deriveContextRef(id string) string {
return fmt.Sprintf("%s/%s/context.tar.gz", strings.TrimRight(m.ContextStore, "/"), id)
return fmt.Sprintf("%s/%s/%s", strings.TrimRight(m.ContextStore, "/"), id, contextBlobName)
}
// auditDockerfile is the audit-archive Dockerfile recorded on the build row. It
@@ -274,6 +325,91 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e
return s, nil
}
// UploadContext stores the caller's uploaded modpack as the build context for
// their OWN pending submission — the blob transport the package doc calls
// deferred. It places the bytes at exactly deriveContextRef(id), the platform-
// pinned, id-namespaced location Kaniko reads via --context, so the submitter
// selects nothing that reaches the executor beyond the modpack itself. The row
// is not mutated (there is no "uploaded" column): the blob store is the source of
// truth for presence, which Approve consults via Blobs.Exists.
//
// The gates mirror the lane's trust model:
// - only the submitter may upload; another user's id is invisible (404, not
// 403) so this endpoint cannot probe other users' submissions;
// - the context is mutable ONLY while pending_review — once approved the build
// has already consumed it, once rejected it is dead;
// - the body must be a gzip tarball (context.tar.gz) and is size-capped, so a
// wrong-format or oversize upload is rejected as a 400 without persisting.
//
// A re-upload while still pending atomically supersedes the previous blob, so a
// user can fix their pack before an admin reviews it.
func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r io.Reader) (*Submission, error) {
if strings.TrimSpace(submittedBy) == "" {
return nil, invalidf("submitter identity is required")
}
if m.Blobs == nil {
return nil, ErrUploadsUnavailable
}
sub, err := m.Store.GetSubmission(ctx, id)
if err != nil {
return nil, err
}
if sub.SubmittedBy != submittedBy {
// Not the owner: invisible, so the endpoint cannot confirm the id exists.
return nil, ErrNotFound
}
if sub.Status != StatusPendingReview {
return nil, ErrAlreadyReviewed
}
// Sniff the gzip magic before touching the store so a wrong-format upload fails
// fast, without persisting anything or reading the whole body.
br := bufio.NewReader(r)
if magic, err := br.Peek(2); err != nil || magic[0] != 0x1f || magic[1] != 0x8b {
return nil, invalidf("build context must be a gzip-compressed tarball (.tar.gz)")
}
// Cap the size: cappedReader trips errContextTooLarge on the first byte past
// the limit, so the store never persists an oversize blob (it removes its temp
// file on the copy error) and the failure surfaces as a 400, not a 500.
if _, err := m.Blobs.Put(ctx, id, &cappedReader{r: br, left: m.maxContextBytes()}); err != nil {
return nil, err
}
return sub, nil
}
// cappedReader passes through at most left bytes; the first byte beyond the limit
// trips errContextTooLarge. It reads one probe byte past the limit to tell an
// exactly-at-limit blob (accepted) from a larger one (rejected), so a stream of
// exactly the cap is never falsely rejected.
type cappedReader struct {
r io.Reader
left int64
}
func (c *cappedReader) Read(p []byte) (int, error) {
if c.left <= 0 {
// At the limit: peek one more byte. Any further data means too large; EOF
// means the blob was exactly the cap.
var probe [1]byte
n, err := c.r.Read(probe[:])
if n > 0 {
return 0, errContextTooLarge
}
if err == nil {
return 0, io.EOF
}
return 0, err
}
if int64(len(p)) > c.left {
p = p[:c.left]
}
n, err := c.r.Read(p)
c.left -= int64(n)
return n, err
}
// Approve is the admin gate. It atomically claims the pending_review -> approved
// transition (CAS) and ONLY the winner starts the build, so concurrent approvals
// can never double-build. The build runs through the SAME gated Builder.Submit as
@@ -314,6 +450,22 @@ func (m *Manager) Approve(ctx context.Context, id, reviewedBy string) (*Submissi
return nil, ErrAlreadyReviewed
}
// Refuse to approve a submission whose build context was never uploaded: the
// derived context ref would point Kaniko at nothing, failing the build after a
// committed CAS. This deterministic check runs BEFORE the CAS (like the
// build.Validate below), so a missing blob leaves the row pending, never
// stranded in approved. Skipped when no transport is wired (Blobs nil): the
// deferred/object-store case cannot be checked here and must not block approve.
if m.Blobs != nil {
ok, err := m.Blobs.Exists(ctx, id)
if err != nil {
return nil, err
}
if !ok {
return nil, invalidf("no build context has been uploaded for this submission")
}
}
imageRef := m.deriveImageRef(id)
req := build.Request{
ImageRef: imageRef,
+212
View File
@@ -3,6 +3,7 @@ package submit
import (
"context"
"errors"
"io"
"strings"
"testing"
"time"
@@ -10,6 +11,45 @@ import (
"felis.lolicon.best/internal/build"
)
// gzBody returns a minimal gzip-magic-prefixed blob standing in for a real
// context.tar.gz: UploadContext only sniffs the first two bytes, so the payload
// after the magic is opaque.
func gzBody(payload string) string { return "\x1f\x8b\x08\x00" + payload }
// fakeBlobs is an in-memory Blobs transport. It records what was stored so a test
// can assert the derived id was used, and can force Exists/Put outcomes.
type fakeBlobs struct {
stored map[string][]byte
putErr error
existsErr error
forceExists *bool // overrides the stored-map lookup for the approve-gate tests
}
func newFakeBlobs() *fakeBlobs { return &fakeBlobs{stored: map[string][]byte{}} }
func (f *fakeBlobs) Put(_ context.Context, id string, r io.Reader) (int64, error) {
if f.putErr != nil {
return 0, f.putErr
}
b, err := io.ReadAll(r)
if err != nil {
return 0, err // e.g. cappedReader tripping — persist nothing, mirror the real store
}
f.stored[id] = b
return int64(len(b)), nil
}
func (f *fakeBlobs) Exists(_ context.Context, id string) (bool, error) {
if f.existsErr != nil {
return false, f.existsErr
}
if f.forceExists != nil {
return *f.forceExists, nil
}
_, ok := f.stored[id]
return ok, nil
}
// testNow is the frozen clock for hermetic assertions.
var testNow = time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC)
@@ -452,3 +492,175 @@ func TestListBy(t *testing.T) {
t.Fatalf("empty submitter err = %v, want ErrInvalid", err)
}
}
func TestUploadContextStoresUnderDerivedID(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
payload := gzBody("the modpack bytes")
sub, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(payload))
if err != nil {
t.Fatalf("UploadContext: %v", err)
}
if sub.ID != seed.ID {
t.Fatalf("returned submission %q, want %q", sub.ID, seed.ID)
}
// The blob is stored under the submission id (the pinned, id-namespaced key),
// never a caller-supplied path.
got, ok := fb.stored[seed.ID]
if !ok {
t.Fatalf("nothing stored under id %q; stored keys: %v", seed.ID, keysOf(fb.stored))
}
if string(got) != payload {
t.Fatalf("stored %q, want the uploaded bytes", got)
}
}
func TestUploadContextRejectsNonGzip(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
_, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader("PK\x03\x04 a zip, not gzip"))
if !errors.Is(err, ErrInvalid) {
t.Fatalf("err = %v, want ErrInvalid", err)
}
if len(fb.stored) != 0 {
t.Fatal("a wrong-format upload must persist nothing")
}
}
func TestUploadContextOversizeRejectedAndNotPersisted(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
m.MaxContextBytes = 8 // tiny cap
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
// gzBody's 4-byte magic + payload well over 8 bytes total.
_, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(gzBody("this is far too large")))
if !errors.Is(err, ErrInvalid) {
t.Fatalf("err = %v, want ErrInvalid (too large)", err)
}
if len(fb.stored) != 0 {
t.Fatal("an oversize upload must persist nothing")
}
}
func TestUploadContextExactlyAtCapAccepted(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
body := gzBody("payload") // measure and cap at exactly this length
m.MaxContextBytes = int64(len(body))
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(body)); err != nil {
t.Fatalf("a blob of exactly the cap must be accepted, got %v", err)
}
if string(fb.stored[seed.ID]) != body {
t.Fatalf("stored %q, want the full body (no truncation at the cap)", fb.stored[seed.ID])
}
}
func TestUploadContextWrongOwnerIsNotFound(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
// A different user uploading to user-1's submission sees 404, not 403: the id
// is invisible so it cannot be probed.
_, err := m.UploadContext(context.Background(), seed.ID, "user-2", strings.NewReader(gzBody("x")))
if !errors.Is(err, ErrNotFound) {
t.Fatalf("err = %v, want ErrNotFound", err)
}
if len(fb.stored) != 0 {
t.Fatal("a non-owner upload must persist nothing")
}
}
func TestUploadContextNotPendingIsAlreadyReviewed(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.Reject(context.Background(), seed.ID, "admin@x", "nope"); err != nil {
t.Fatalf("Reject: %v", err)
}
_, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(gzBody("x")))
if !errors.Is(err, ErrAlreadyReviewed) {
t.Fatalf("err = %v, want ErrAlreadyReviewed (context is frozen once reviewed)", err)
}
}
func TestUploadContextUnknownSubmission(t *testing.T) {
m, _, _ := newManager()
m.Blobs = newFakeBlobs()
_, err := m.UploadContext(context.Background(), "sub-nope", "user-1", strings.NewReader(gzBody("x")))
if !errors.Is(err, ErrNotFound) {
t.Fatalf("err = %v, want ErrNotFound", err)
}
}
func TestUploadContextNoTransportUnavailable(t *testing.T) {
m, _, _ := newManager() // Blobs left nil
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
_, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(gzBody("x")))
if !errors.Is(err, ErrUploadsUnavailable) {
t.Fatalf("err = %v, want ErrUploadsUnavailable", err)
}
}
func TestApproveRefusesMissingContext(t *testing.T) {
// With a transport wired, approving a submission whose context was never
// uploaded fails BEFORE the CAS: the row stays pending and no build starts.
m, st, bl := newManager()
m.Blobs = newFakeBlobs() // empty → Exists=false
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
_, err := m.Approve(context.Background(), seed.ID, "admin@x")
if !errors.Is(err, ErrInvalid) {
t.Fatalf("err = %v, want ErrInvalid (no context uploaded)", err)
}
if got := st.subs[seed.ID]; got.Status != StatusPendingReview {
t.Fatalf("status = %q, want still pending_review (CAS not reached)", got.Status)
}
if bl.calls != 0 {
t.Fatalf("builds started = %d, want 0", bl.calls)
}
}
func TestApproveProceedsWithUploadedContext(t *testing.T) {
// The end-to-end user path: create -> upload -> admin approve -> exactly one
// build through the SAME gated Builder.
m, _, bl := newManager()
m.Blobs = newFakeBlobs()
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(gzBody("mods"))); err != nil {
t.Fatalf("UploadContext: %v", err)
}
sub, err := m.Approve(context.Background(), seed.ID, "admin@x")
if err != nil {
t.Fatalf("Approve: %v", err)
}
if sub.Status != StatusApproved {
t.Fatalf("status = %q, want approved", sub.Status)
}
if bl.calls != 1 {
t.Fatalf("builds started = %d, want 1", bl.calls)
}
}
// keysOf lists a map's keys for test diagnostics.
func keysOf(m map[string][]byte) []string {
out := make([]string, 0, len(m))
for k := range m {
out = append(out, k)
}
return out
}