fix(submit): bound the untrusted upload lane — per-user caps + throttles (#75)

A logged-in user could file submissions without bound and stream a 1 GiB
context per submission. The only limits were the single-blob size cap and the
5 GiB uploads PVC (platform/workloads.go); nothing counted a user's rows or
bytes, so one account could fill the volume and every other user's upload
would start failing.

- Create: per-user pending_review cap (default 5) — the review queue cannot
  be parked full of one account's rows. Check-then-insert, documented soft.
- UploadContext: per-user stored-context budget (default 2 GiB) charged
  against the blob store's REAL sizes (new Blobs.Size on local/S3 stores), so
  the sum cannot drift from the volume; the write is capped at the remaining
  budget, so the excess is refused before it is persisted, and a re-upload is
  charged only for its new bytes.
- API: per-user create/upload throttles (30s/15s, cmd/felis-wired) on a
  dedicated cooldown keyspace, reserve→release so a failed attempt never
  burns the window and a burst collapses to one winner; ErrQuotaExceeded →
  403 submission_quota_exceeded (distinct from the 400 an oversize blob
  gets), 429 submission_cooldown for the throttles.
- Panel: zh/en copy for both codes; openapi documents 403/429 on the two
  user routes; pgint covers the pending-queue count.

Unit tests: submit package (cap, budget boundary/exact-fit/replacement,
oversize-vs-quota split) and api handlers (quota 403 both paths, throttle
429 + recovery + failure-release). go vet/go test/gofmt clean; panel
vitest 118 + typecheck green.
This commit is contained in:
Lemon-miaow committed 2026-09-24 10:14:42 +08:00
1 parent 3ed8bd7be9
commit ad4d256d8f
15 files changed
+577 -15

No files matched your search

+26
View File
@@ -123,6 +123,15 @@ type API struct {
// on the wake lever). Zero disables throttling.
WakeCooldown time.Duration
// SubmitCreateCooldown / SubmitUploadCooldown throttle the user-modpack
// submission lane per user: create bounds how quickly review-queue rows can
// appear, upload bounds how often a user may stream a (up to 1 GiB) build
// context. The keys are separate, so the lane's normal shape — create, then
// upload — is never blocked by its own throttle. Zero disables each lever
// (the same idiom as WakeCooldown); cmd/felis wires positive values.
SubmitCreateCooldown time.Duration
SubmitUploadCooldown time.Duration
// MaxRunningServers caps how many servers may be desired-Running cluster-wide
// (spec §9.1: the concurrency-上限 lever hanging on the same wake chokepoint as
// cooldown and autostartPolicy). Zero — the default — disables it: §9.2 wires
@@ -158,6 +167,9 @@ type API struct {
otpCooldownOnce sync.Once
otpCooldown *cooldownLimiter
submitCooldownOnce sync.Once
submitCooldown *cooldownLimiter
streamCapOnce sync.Once
streamCap *streamLimiter
}
@@ -204,6 +216,20 @@ func (a *API) otpLimiter() *cooldownLimiter {
return a.otpCooldown
}
// submitLimiter lazily builds a SEPARATE cooldown limiter for the user-modpack
// submission lane, so its throttles never share state with the wake or OTP
// keyspaces. One limiter backs both levers with prefixed keys (see the
// submissionCreateKey/UploadKey constants), so create and upload never contend
// with each other. Like the other cooldowns it is process-local; with multiple
// api replicas the effective spacing is per-replica, the same accepted
// KNOWN-LIMITATION the OTP resend throttle carries.
func (a *API) submitLimiter() *cooldownLimiter {
a.submitCooldownOnce.Do(func() {
a.submitCooldown = &cooldownLimiter{now: a.now, last: map[string]time.Time{}}
})
return a.submitCooldown
}
// streamGate lazily builds the per-principal SSE stream cap bound to
// MaxStreamsPerPrincipal. A zero cap yields a disabled limiter that admits every
// stream, so a deployment (or test) that leaves it unset pays nothing.
+59 -3
View File
@@ -59,6 +59,15 @@ type createSubmissionRequest struct {
DisplayName string `json:"display_name"`
}
// The submission lane's two cooldown keys, prefixed into the shared submit
// limiter's per-user keys. Create and upload are separate levers on purpose:
// creating a submission and then immediately uploading its context is the lane's
// normal shape, so one must never consume the other's window.
const (
submissionCreateKey = "create:"
submissionUploadKey = "upload:"
)
// rejectSubmissionRequest is the POST /submissions/{id}/reject body. A reason is
// required (the submit layer rejects an empty one with 400).
type rejectSubmissionRequest struct {
@@ -74,6 +83,26 @@ func (a *API) handleCreateSubmission(w http.ResponseWriter, r *http.Request) {
return
}
p := principalFromContext(r.Context())
// Reserve the per-user create cooldown BEFORE the store write. Unlike the
// wake lever's allowed→record (whose real gate is the running cap and whose
// effect is idempotent), a create is a non-idempotent row insertion with no
// other bound on its rate, so a burst of truly concurrent creates must yield
// exactly one winner per window. The deferred rollback frees the window
// whenever the create fails — a 400 typo, a spent quota, a store error — so
// only a row that was actually recorded consumes it.
lim := a.submitLimiter()
reservedAt, ok := lim.reserve(submissionCreateKey+p.UserID, a.SubmitCreateCooldown)
if !ok {
writeError(w, r, newError(http.StatusTooManyRequests, "submission_cooldown",
"a submission was created recently; wait a moment before creating another"))
return
}
committed := false
defer func() {
if !committed {
lim.release(submissionCreateKey+p.UserID, reservedAt)
}
}()
var body createSubmissionRequest
if err := decodeJSON(w, r, &body); err != nil {
writeError(w, r, err)
@@ -87,6 +116,7 @@ func (a *API) handleCreateSubmission(w http.ResponseWriter, r *http.Request) {
writeSubmitError(w, r, err)
return
}
committed = true
a.audit(r, p.Email, "submission.create", sub.ID)
writeJSON(w, http.StatusCreated, sub)
}
@@ -107,12 +137,34 @@ func (a *API) handleUploadSubmissionContext(w http.ResponseWriter, r *http.Reque
return
}
p := principalFromContext(r.Context())
// Reserve the per-user upload cooldown BEFORE streaming. The body is the
// expensive part (up to the 1 GiB blob cap), so without a reservation the
// throttle would never bound the resource it exists for: a caller could
// repeatedly start long uploads and abort them. Reserving also collapses the
// lane's parallel overshoot — a burst of concurrent uploads from one user
// yields exactly one admitted stream per replica. The rollback keeps a failed
// upload (aborted transfer, wrong format, spent quota) from burning the
// window, so a legit retry after a genuine failure is not punished.
lim := a.submitLimiter()
reservedAt, ok := lim.reserve(submissionUploadKey+p.UserID, a.SubmitUploadCooldown)
if !ok {
writeError(w, r, newError(http.StatusTooManyRequests, "submission_cooldown",
"an upload was accepted recently; wait a moment before uploading again"))
return
}
committed := false
defer func() {
if !committed {
lim.release(submissionUploadKey+p.UserID, reservedAt)
}
}()
id := r.PathValue("id")
sub, err := a.Submissions.UploadContext(r.Context(), id, p.UserID, r.Body)
if err != nil {
writeSubmitError(w, r, err)
return
}
committed = true
a.audit(r, p.Email, "submission.upload", sub.ID)
writeJSON(w, http.StatusOK, sub)
}
@@ -250,9 +302,10 @@ var errSubmissionsUnavailable = newError(http.StatusServiceUnavailable, "submiss
// writeSubmitError maps submit-package errors onto HTTP status codes. Only the
// business sentinels are client-facing: a validation failure is 400, a missing
// submission is 404, an already-reviewed submission is 409, and an unconfigured
// upload transport is 503 (the store this deployment set has no implemented
// transport — an honest "not available here", not a client error). Everything
// submission is 404, an already-reviewed submission is 409, a spent per-user
// allowance is 403 (the same status the server-resource quota answers with), and
// an unconfigured upload transport is 503 (the store this deployment set has no
// implemented transport — an honest "not available here", not a client error). Everything
// else — including a build.ErrInvalid raised by the pre-CAS build.Validate (a
// platform registry/context MISCONFIGURATION, never client input, since every
// build input is platform-derived) and a post-CAS Submit hand-off failure — is a
@@ -267,6 +320,9 @@ func writeSubmitError(w http.ResponseWriter, r *http.Request, err error) {
case errors.Is(err, submit.ErrAlreadyReviewed):
writeError(w, r, newError(http.StatusConflict, "already_reviewed",
"submission has already been reviewed"))
case errors.Is(err, submit.ErrQuotaExceeded):
writeError(w, r, newError(http.StatusForbidden, "submission_quota_exceeded",
"submission quota reached"))
case errors.Is(err, submit.ErrBlobNotFound):
writeError(w, r, newError(http.StatusNotFound, "not_found", "no context uploaded for this submission"))
case errors.Is(err, submit.ErrUploadsUnavailable):
+110
View File
@@ -9,6 +9,7 @@ import (
"net/http"
"strings"
"testing"
"time"
"felis.lolicon.best/internal/build"
"felis.lolicon.best/internal/submit"
@@ -567,3 +568,112 @@ func TestAdminSubmissionContextRoute(t *testing.T) {
}
})
}
// A spent per-user allowance is 403 submission_quota_exceeded on both the create
// and the upload path — distinctly NOT the 400 a malformed request gets, and not
// the 429 the cooldown answers with.
func TestSubmissionQuotaIs403(t *testing.T) {
t.Run("create", func(t *testing.T) {
fs := &fakeSubmissions{createErr: fmt.Errorf("%w: 5 submissions are already awaiting review", submit.ErrQuotaExceeded)}
w := do(appSubAPI(fs).ExternalHandler(), "POST", "/api/v1/me/submissions", `{"display_name":"Pack"}`, nil)
if w.Code != http.StatusForbidden {
t.Fatalf("code = %d, want 403 (%s)", w.Code, w.Body.String())
}
if got := decodeErr(t, w); got != "submission_quota_exceeded" {
t.Errorf("error code = %q, want submission_quota_exceeded", got)
}
})
t.Run("upload", func(t *testing.T) {
fs := &fakeSubmissions{uploadErr: fmt.Errorf("%w: exceeds your remaining storage allowance", submit.ErrQuotaExceeded)}
w := do(appSubAPI(fs).ExternalHandler(), "POST", "/api/v1/me/submissions/sub-9/context", "\x1f\x8bdata", nil)
if w.Code != http.StatusForbidden {
t.Fatalf("code = %d, want 403 (%s)", w.Code, w.Body.String())
}
if got := decodeErr(t, w); got != "submission_quota_exceeded" {
t.Errorf("error code = %q, want submission_quota_exceeded", got)
}
})
}
// The per-user create cooldown bounds review-queue growth: a second create in
// the same window is 429 submission_cooldown and never reaches the service; the
// window recovers afterwards.
func TestCreateSubmissionRateLimited(t *testing.T) {
fs := &fakeSubmissions{}
api := appSubAPI(fs)
clock := time.Unix(1_700_000_000, 0)
api.Now = func() time.Time { return clock }
api.SubmitCreateCooldown = time.Minute
eh := api.ExternalHandler()
if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"First"}`, nil); w.Code != http.StatusCreated {
t.Fatalf("first create: code = %d, want 201 (%s)", w.Code, w.Body.String())
}
w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"Second"}`, nil)
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "submission_cooldown" {
t.Fatalf("immediate second create: code = %d body %s, want 429 submission_cooldown", w.Code, w.Body.String())
}
// The gate sits before the body handling: even a malformed request is
// refused while the window is closed, so it cannot be used to probe.
if w := do(eh, "POST", "/api/v1/me/submissions", `{`, nil); w.Code != http.StatusTooManyRequests {
t.Fatalf("malformed create during cooldown: code = %d, want 429", w.Code)
}
clock = clock.Add(time.Minute + time.Second)
if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"Third"}`, nil); w.Code != http.StatusCreated {
t.Fatalf("post-cooldown create: code = %d, want 201 (%s)", w.Code, w.Body.String())
}
}
// A failed create frees the window: only a row that was actually recorded burns
// the cooldown, so a validation typo is not punished with a wait.
func TestCreateSubmissionFailureDoesNotBurnCooldown(t *testing.T) {
fs := &fakeSubmissions{createErr: fmt.Errorf("%w: display name is required", submit.ErrInvalid)}
api := appSubAPI(fs)
api.Now = func() time.Time { return time.Unix(1_700_000_000, 0) }
api.SubmitCreateCooldown = time.Minute
eh := api.ExternalHandler()
if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":""}`, nil); w.Code != http.StatusBadRequest {
t.Fatalf("failed create: code = %d, want 400", w.Code)
}
fs.createErr = nil
if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"Fixed"}`, nil); w.Code != http.StatusCreated {
t.Fatalf("retry at the same instant: code = %d, want 201 (%s)", w.Code, w.Body.String())
}
}
// The per-user upload cooldown bounds context streaming: a second upload in the
// same window is 429 submission_cooldown, and a FAILED upload frees the window
// for an immediate retry.
func TestUploadSubmissionContextRateLimited(t *testing.T) {
fs := &fakeSubmissions{}
api := appSubAPI(fs)
clock := time.Unix(1_700_000_000, 0)
api.Now = func() time.Time { return clock }
api.SubmitUploadCooldown = time.Minute
eh := api.ExternalHandler()
body := "\x1f\x8b\x08\x00 the modpack bytes"
if w := do(eh, "POST", "/api/v1/me/submissions/sub-9/context", body, nil); w.Code != http.StatusOK {
t.Fatalf("first upload: code = %d, want 200 (%s)", w.Code, w.Body.String())
}
w := do(eh, "POST", "/api/v1/me/submissions/sub-9/context", body, nil)
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "submission_cooldown" {
t.Fatalf("immediate second upload: code = %d body %s, want 429 submission_cooldown", w.Code, w.Body.String())
}
// A failed upload releases its reservation, so the user is not punished for
// a genuine failure (aborted transfer, spent quota) with a cooldown wait.
fs2 := &fakeSubmissions{uploadErr: submit.ErrUploadsUnavailable}
api2 := appSubAPI(fs2)
api2.Now = func() time.Time { return clock }
api2.SubmitUploadCooldown = time.Minute
eh2 := api2.ExternalHandler()
if w := do(eh2, "POST", "/api/v1/me/submissions/sub-9/context", body, nil); w.Code != http.StatusServiceUnavailable {
t.Fatalf("failed upload: code = %d, want 503", w.Code)
}
fs2.uploadErr = nil
if w := do(eh2, "POST", "/api/v1/me/submissions/sub-9/context", body, nil); w.Code != http.StatusOK {
t.Fatalf("retry at the same instant after failure: code = %d, want 200 (%s)", w.Code, w.Body.String())
}
}
+12
View File
@@ -980,6 +980,14 @@ func TestSubmitStoreContract(t *testing.T) {
if err := s.CreateSubmission(ctx, sub); err != nil {
t.Fatalf("CreateSubmission: %v", err)
}
// The pending-queue count is the read behind the per-user pending cap: a
// fresh user starts at zero and a pending row counts.
if n, err := s.CountPendingSubmissionsBy(ctx, u.ID); err != nil || n != 1 {
t.Fatalf("CountPendingSubmissionsBy after create = (%d, %v), want (1, nil)", n, err)
}
if n, err := s.CountPendingSubmissionsBy(ctx, u.ID+"-nobody"); err != nil || n != 0 {
t.Fatalf("CountPendingSubmissionsBy for an unknown user = (%d, %v), want (0, nil)", n, err)
}
got, err := s.GetSubmission(ctx, id)
if err != nil {
t.Fatalf("GetSubmission: %v", err)
@@ -1012,6 +1020,10 @@ func TestSubmitStoreContract(t *testing.T) {
if ok, err := s.RejectSubmission(ctx, id, "[email protected]", "no", now); err != nil || ok {
t.Fatalf("reject after approve = (%v, %v), want (false, nil)", ok, err)
}
// A reviewed row leaves the pending queue.
if n, err := s.CountPendingSubmissionsBy(ctx, u.ID); err != nil || n != 0 {
t.Fatalf("CountPendingSubmissionsBy after approve = (%d, %v), want (0, nil)", n, err)
}
got, _ = s.GetSubmission(ctx, id)
if got.Status != submit.StatusApproved || got.ImageRef == "" || got.ReviewedBy != "[email protected]" || got.ReviewedAt == nil {
t.Fatalf("approved row = %+v", got)
+20
View File
@@ -129,5 +129,25 @@ func (s *LocalContextStore) Open(_ context.Context, id string) (io.ReadCloser, e
return f, nil
}
// Size reports the stored blob's size — the accounting read behind the per-user
// storage budget. A missing blob is (0, false, nil): absence is not an error
// here, it is simply no bytes to count (the same distinction Exists draws for
// Approve).
func (s *LocalContextStore) Size(_ context.Context, id string) (int64, bool, error) {
dir, err := s.dir(id)
if err != nil {
return 0, false, err
}
fi, err := os.Stat(filepath.Join(dir, contextBlobName))
switch {
case err == nil:
return fi.Size(), true, nil
case os.IsNotExist(err):
return 0, false, nil
default:
return 0, false, fmt.Errorf("submit: stat context blob: %w", err)
}
}
// Compile-time proof that the filesystem store satisfies the Blobs transport.
var _ Blobs = (*LocalContextStore)(nil)
+8
View File
@@ -33,6 +33,14 @@ func (s *PGStore) CreateSubmission(ctx context.Context, sub *Submission) error {
return err
}
func (s *PGStore) CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error) {
const q = `SELECT count(*) FROM image_submissions
WHERE submitted_by = $1 AND status = 'pending_review'`
var n int
err := s.db.QueryRowContext(ctx, q, submittedBy).Scan(&n)
return n, err
}
func (s *PGStore) GetSubmission(ctx context.Context, id string) (*Submission, error) {
const q = `SELECT ` + submissionColumns + ` FROM image_submissions WHERE id = $1`
return scanSubmission(s.db.QueryRowContext(ctx, q, id))
+18
View File
@@ -187,6 +187,24 @@ func (s *S3ContextStore) Exists(ctx context.Context, id string) (bool, error) {
return true, nil
}
// Size reports the stored blob's size — the accounting read behind the per-user
// storage budget. A missing object is (0, false, nil), the same distinction
// Exists draws.
func (s *S3ContextStore) Size(ctx context.Context, id string) (int64, bool, error) {
key, err := s.keyFor(id)
if err != nil {
return 0, false, err
}
info, err := s.client.StatObject(ctx, s.bucket, key, minio.StatObjectOptions{})
if err != nil {
if isS3NotFound(err) {
return 0, false, nil
}
return 0, false, fmt.Errorf("submit: stat context blob: %w", err)
}
return info.Size, true, nil
}
// Open returns the stored context blob for id — the read side of the transport the
// build Pod's fetch initContainer uses. minio's GetObject returns only once the
// server answered with an object (it surfaces NoSuchKey up front), so a missing
+134 -10
View File
@@ -86,6 +86,12 @@ var (
ErrInvalid = errors.New("submit: invalid request")
ErrNotFound = errors.New("submit: submission not found")
ErrAlreadyReviewed = errors.New("submit: submission already reviewed")
// ErrQuotaExceeded reports that the caller's upload allowance is spent —
// either too many of their submissions are already awaiting review, or their
// stored contexts already fill the per-user byte budget. The request is not
// malformed; the allowance is exhausted. The API maps it to 403, matching the
// server-resource quota's status (spec §7).
ErrQuotaExceeded = errors.New("submit: quota exceeded")
// 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
@@ -108,14 +114,37 @@ func invalidf(format string, a ...any) error {
// preserves that chain through to the API error mapper.
var errContextTooLarge = fmt.Errorf("%w: build context exceeds the maximum allowed size", ErrInvalid)
// errStorageQuota trips when an upload would push the user past their per-user
// storage budget. It wraps ErrQuotaExceeded so the API answers 403, distinctly
// from errContextTooLarge's 400: the blob is fine, the allowance is spent.
var errStorageQuota = fmt.Errorf("%w: the upload exceeds your remaining storage allowance", ErrQuotaExceeded)
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.
// one blob at a time. The per-user budget below bounds the SUM across a user's
// uploads; this cap is what keeps any single write bounded. Override per-Manager
// via MaxContextBytes.
defaultMaxContextBytes = 1 << 30 // 1 GiB
// defaultMaxPendingPerUser caps how many pending_review submissions one user
// may hold at once. Untrusted ingress has no natural bound — a logged-in
// player could otherwise file rows all day — and every pending row is a
// review-queue item an admin has to read, so the cap is small on purpose:
// enough to stage a couple of packs, far short of a flood. Override per
// Manager via MaxPendingPerUser.
defaultMaxPendingPerUser = 5
// defaultMaxStoredBytesPerUser caps the total bytes one user's stored
// contexts may occupy on the uploads store. The uploads PVC renders at a
// fixed 5Gi (platform/workloads.go); without a per-user budget one account
// could fill it and every other user's upload would start failing. Two GiB
// leaves room for a couple of full-size modpacks (a single blob may be 1 GiB)
// while keeping a small user base from exhausting the volume; size the PVC
// above users × this budget before raising it. Override per Manager via
// MaxStoredBytesPerUser.
defaultMaxStoredBytesPerUser = 2 << 30 // 2 GiB
)
// displayNameRE constrains the user-supplied label to a calm, single-line set:
@@ -149,6 +178,10 @@ type Submission struct {
type Store interface {
// CreateSubmission inserts a pending_review row.
CreateSubmission(ctx context.Context, s *Submission) error
// CountPendingSubmissionsBy reports how many of one user's submissions are
// still pending_review — the queue-length read behind the per-user pending
// cap in Create.
CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error)
// GetSubmission loads one submission, or ErrNotFound.
GetSubmission(ctx context.Context, id string) (*Submission, error)
// ListSubmissions returns every submission, newest first (admin queue).
@@ -196,6 +229,12 @@ type Builds interface {
type Blobs interface {
Put(ctx context.Context, id string, r io.Reader) (int64, error)
Exists(ctx context.Context, id string) (bool, error)
// Size returns the stored blob's size in bytes; ok=false means no blob is
// stored for id. UploadContext sums this over a user's submissions to enforce
// the per-user storage budget, so it must report what is actually on the
// store — never a recorded number that could drift from it (a re-upload
// supersedes the previous blob in place).
Size(ctx context.Context, id string) (int64, bool, error)
// Open returns the stored blob's bytes for the internal context-fetch route
// the build Pod's initContainer dials (cmd/felis fetch-context). It returns an
// error wrapping ErrBlobNotFound when no blob exists, so the route can answer
@@ -239,6 +278,12 @@ type Manager struct {
// MaxContextBytes overrides the uploaded-context size cap; 0 uses
// defaultMaxContextBytes.
MaxContextBytes int64
// MaxPendingPerUser overrides how many of one user's submissions may await
// review at once; 0 uses defaultMaxPendingPerUser.
MaxPendingPerUser int
// MaxStoredBytesPerUser overrides the per-user stored-context budget; 0 uses
// defaultMaxStoredBytesPerUser.
MaxStoredBytesPerUser int64
Now func() time.Time
IDGen func() string
@@ -251,6 +296,20 @@ func (m *Manager) maxContextBytes() int64 {
return defaultMaxContextBytes
}
func (m *Manager) maxPendingPerUser() int {
if m.MaxPendingPerUser > 0 {
return m.MaxPendingPerUser
}
return defaultMaxPendingPerUser
}
func (m *Manager) maxStoredBytesPerUser() int64 {
if m.MaxStoredBytesPerUser > 0 {
return m.MaxStoredBytesPerUser
}
return defaultMaxStoredBytesPerUser
}
func (m *Manager) now() time.Time {
if m.Now != nil {
return m.Now()
@@ -344,6 +403,18 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e
return nil, invalidf("submitter identity is required")
}
// Per-user pending cap: every pending row is a review-queue item an admin
// must read, so one account may not park an unbounded number of them. The
// check and the insert are not atomic (two concurrent creates may jointly
// land one row over the cap); that is a soft overshoot of a queue-length
// lever, not a resource bound, so it is deliberately not worth a lock.
if pending, err := m.Store.CountPendingSubmissionsBy(ctx, req.SubmittedBy); err != nil {
return nil, err
} else if pending >= m.maxPendingPerUser() {
return nil, fmt.Errorf("%w: %d submissions are already awaiting review (limit %d)",
ErrQuotaExceeded, pending, m.maxPendingPerUser())
}
id := m.newID()
s := &Submission{
ID: id,
@@ -373,7 +444,10 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e
// - 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.
// wrong-format or oversize upload is rejected as a 400 without persisting;
// - a user's stored contexts are budgeted (MaxStoredBytesPerUser): the write
// is capped at the remaining budget, so an upload that would exceed it is
// refused as a spent allowance (403) before the excess is persisted.
//
// A re-upload while still pending atomically supersedes the previous blob, so a
// user can fix their pack before an admin reviews it.
@@ -404,22 +478,72 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i
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 {
// Per-user storage budget: sum the bytes this user's OTHER submissions
// already hold (excluding this id, whose blob a re-upload supersedes) and cap
// the write at whatever remains. cappedReader trips on the first byte past
// the limit, so the store never persists a blob that would exceed the budget
// (it removes its temp file on the copy error) and the failure surfaces as a
// 403, not a 500. The read-then-write pair is not atomic in this package: a
// burst that reaches two api replicas (or any direct caller of the Manager)
// can overshoot by up to one blob per interleaved upload — each write still
// bounded by the single-blob cap — while the API's per-user upload
// reservation collapses the single-replica case.
used, err := m.storedBytes(ctx, submittedBy, id)
if err != nil {
return nil, err
}
remaining := m.maxStoredBytesPerUser() - used
if remaining <= 0 {
return nil, errStorageQuota
}
limit, over := m.maxContextBytes(), errContextTooLarge
if remaining < limit {
// The budget binds before the single-blob cap: an upload tripping here is
// refused as a spent allowance, never as a malformed request.
limit, over = remaining, errStorageQuota
}
if _, err := m.Blobs.Put(ctx, id, &cappedReader{r: br, left: limit, over: over}); err != nil {
return nil, err
}
return sub, nil
}
// storedBytes sums the stored-blob sizes of submittedBy's submissions, excluding
// excludeID — the submission a pending re-upload is about to replace, whose
// bytes must not be counted twice. Sizes are read from the blob store itself,
// the same source of truth uploads/approval consult, so the sum cannot drift
// from what is actually occupying the volume (including blobs uploaded before
// any budget existed).
func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string) (int64, error) {
subs, err := m.Store.ListSubmissionsBy(ctx, submittedBy)
if err != nil {
return 0, err
}
var total int64
for _, s := range subs {
if s.ID == excludeID {
continue
}
n, ok, err := m.Blobs.Size(ctx, s.ID)
if err != nil {
return 0, err
}
if ok {
total += n
}
}
return total, 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.
// trips over — errContextTooLarge for the single-blob cap, errStorageQuota when
// the per-user budget binds first. 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
over error
}
func (c *cappedReader) Read(p []byte) (int, error) {
@@ -429,7 +553,7 @@ func (c *cappedReader) Read(p []byte) (int, error) {
var probe [1]byte
n, err := c.r.Read(probe[:])
if n > 0 {
return 0, errContextTooLarge
return 0, c.over
}
if err == nil {
return 0, io.EOF
+148
View File
@@ -24,6 +24,7 @@ type fakeBlobs struct {
stored map[string][]byte
putErr error
existsErr error
sizeErr error
forceExists *bool // overrides the stored-map lookup for the approve-gate tests
}
@@ -52,6 +53,18 @@ func (f *fakeBlobs) Exists(_ context.Context, id string) (bool, error) {
return ok, nil
}
// Size mirrors the real stores: a missing blob is (0, false, nil).
func (f *fakeBlobs) Size(_ context.Context, id string) (int64, bool, error) {
if f.sizeErr != nil {
return 0, false, f.sizeErr
}
b, ok := f.stored[id]
if !ok {
return 0, false, nil
}
return int64(len(b)), true, nil
}
func (f *fakeBlobs) Open(_ context.Context, id string) (io.ReadCloser, error) {
b, ok := f.stored[id]
if !ok {
@@ -73,6 +86,7 @@ type fakeStore struct {
subs map[string]*Submission
createErr error
countErr error
approveErr error
rejectErr error
linkErr error
@@ -91,6 +105,19 @@ func (f *fakeStore) CreateSubmission(_ context.Context, s *Submission) error {
return nil
}
func (f *fakeStore) CountPendingSubmissionsBy(_ context.Context, by string) (int, error) {
if f.countErr != nil {
return 0, f.countErr
}
var n int
for _, s := range f.subs {
if s.SubmittedBy == by && s.Status == StatusPendingReview {
n++
}
}
return n, nil
}
func (f *fakeStore) GetSubmission(_ context.Context, id string) (*Submission, error) {
s, ok := f.subs[id]
if !ok {
@@ -671,6 +698,127 @@ func TestUploadContextNoTransportUnavailable(t *testing.T) {
}
}
// The per-user pending cap bounds the review queue: at the limit a new create is
// refused with ErrQuotaExceeded (403), a reviewed row frees a slot, and another
// user's queue is unaffected.
func TestCreatePendingCapRejects(t *testing.T) {
m, _, _ := newManager()
m.MaxPendingPerUser = 2
ctx := context.Background()
for i := 0; i < 2; i++ {
if _, err := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}); err != nil {
t.Fatalf("create %d: %v", i+1, err)
}
}
_, err := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if !errors.Is(err, ErrQuotaExceeded) {
t.Fatalf("create past the cap = %v, want ErrQuotaExceeded", err)
}
// A verdict moves the row out of pending_review, so the slot frees up.
if _, err := m.Reject(ctx, "sub-1", "[email protected]", "out of scope"); err != nil {
t.Fatalf("Reject: %v", err)
}
if _, err := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}); err != nil {
t.Fatalf("create after review: %v", err)
}
// The cap is per user, never global.
if _, err := m.Create(ctx, CreateRequest{DisplayName: "Other", SubmittedBy: "user-2"}); err != nil {
t.Fatalf("other user's first create: %v", err)
}
}
// The per-user storage budget bounds the sum of a user's stored contexts. It is
// charged against the blob store's actual sizes, refuses at the boundary with
// ErrQuotaExceeded (403 — not the 400 a single oversize blob gets), and does not
// double-charge a re-upload of the same submission.
func TestUploadContextStorageBudget(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
m.MaxStoredBytesPerUser = 10 // tiny budget; gzBody is 4 magic bytes + payload
ctx := context.Background()
a, err := m.Create(ctx, CreateRequest{DisplayName: "A", SubmittedBy: "user-1"})
if err != nil {
t.Fatalf("create A: %v", err)
}
b, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-1"})
if err != nil {
t.Fatalf("create B: %v", err)
}
if _, err := m.UploadContext(ctx, a.ID, "user-1", strings.NewReader(gzBody("a"))); err != nil {
t.Fatalf("first upload: %v", err)
}
// A's 5 bytes leave 5. B's exactly-5-byte blob must be accepted — the budget
// binds only past the limit, never at it.
if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader(gzBody("b"))); err != nil {
t.Fatalf("upload exactly at the budget = %v, want accepted", err)
}
if _, err := m.UploadContext(ctx, a.ID, "user-1", strings.NewReader(gzBody("far too much"))); err != nil {
// A re-upload is charged only for its NEW bytes (its old blob is
// superseded), and 16 bytes exceed the 5 bytes left after B.
if !errors.Is(err, ErrQuotaExceeded) {
t.Fatalf("re-upload past the budget = %v, want ErrQuotaExceeded", err)
}
} else {
t.Fatal("re-upload past the budget was accepted")
}
// Nothing oversize persisted, and A's good blob was not clobbered.
if string(fb.stored[a.ID]) != gzBody("a") {
t.Fatalf("A's blob = %q, want the original (a failed re-upload must not replace it)", fb.stored[a.ID])
}
// B now holds 5 and the budget is full: a fresh submission's upload is
// refused up front, before reading any body.
c, err := m.Create(ctx, CreateRequest{DisplayName: "C", SubmittedBy: "user-1"})
if err != nil {
t.Fatalf("create C: %v", err)
}
_, err = m.UploadContext(ctx, c.ID, "user-1", strings.NewReader(gzBody("c")))
if !errors.Is(err, ErrQuotaExceeded) {
t.Fatalf("upload with no budget left = %v, want ErrQuotaExceeded", err)
}
if _, ok := fb.stored[c.ID]; ok {
t.Fatal("a budget-refused upload must persist nothing")
}
// A 5-byte replacement of A fits exactly (10 − B's 5), proving the
// replacement is not double-charged against its own old bytes.
if _, err := m.UploadContext(ctx, a.ID, "user-1", strings.NewReader(gzBody("z"))); err != nil {
t.Fatalf("budget-exact replacement = %v, want accepted", err)
}
}
// A single oversize blob stays a 400 (ErrInvalid), distinct from the 403 the
// per-user budget answers with — the two failure classes must not collapse.
func TestUploadContextOversizeIsNotQuotaError(t *testing.T) {
m, _, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
m.MaxContextBytes = 4
seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
_, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(gzBody("too big")))
if !errors.Is(err, ErrInvalid) || errors.Is(err, ErrQuotaExceeded) {
t.Fatalf("err = %v, want ErrInvalid and NOT ErrQuotaExceeded", err)
}
}
// A store failure while counting the pending queue must surface as-is, never as
// a quota verdict that blames the user.
func TestCreatePendingCountFailureSurfaces(t *testing.T) {
m, st, _ := newManager()
st.countErr = errors.New("db is down")
_, err := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if err == nil || errors.Is(err, ErrQuotaExceeded) {
t.Fatalf("err = %v, want the raw store failure", 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.